api_queue.go 2.4 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485
  1. package api
  2. import (
  3. ctx "context"
  4. "sort"
  5. "connectrpc.com/connect"
  6. apiv1 "github.com/OliveTin/OliveTin/gen/olivetin/api/v1"
  7. "github.com/OliveTin/OliveTin/internal/auth"
  8. authpublic "github.com/OliveTin/OliveTin/internal/auth/authpublic"
  9. "github.com/OliveTin/OliveTin/internal/executor"
  10. )
  11. func (api *oliveTinAPI) GetExecutionQueue(ctx ctx.Context, req *connect.Request[apiv1.GetExecutionQueueRequest]) (*connect.Response[apiv1.GetExecutionQueueResponse], error) {
  12. user := auth.UserFromApiCall(ctx, req, api.cfg)
  13. if err := api.checkDashboardAccess(user); err != nil {
  14. return nil, err
  15. }
  16. active := api.executor.GetActiveExecutionsACL(api.cfg, user)
  17. groups := buildExecutionQueueGroups(active, user, api)
  18. return connect.NewResponse(&apiv1.GetExecutionQueueResponse{
  19. Groups: groups,
  20. TotalActive: int32(len(active)),
  21. }), nil
  22. }
  23. func buildExecutionQueueGroups(active []*executor.InternalLogEntry, user *authpublic.AuthenticatedUser, api *oliveTinAPI) []*apiv1.ExecutionQueueGroup {
  24. grouped := make(map[string]*apiv1.ExecutionQueueGroup)
  25. for _, entry := range active {
  26. bindingID := entry.GetBindingId()
  27. group := grouped[bindingID]
  28. if group == nil {
  29. group = newExecutionQueueGroup(entry)
  30. grouped[bindingID] = group
  31. }
  32. group.Entries = append(group.Entries, api.internalLogEntryToPb(entry, user))
  33. }
  34. groups := make([]*apiv1.ExecutionQueueGroup, 0, len(grouped))
  35. for _, group := range grouped {
  36. sortQueueEntries(group.Entries)
  37. group.ActiveCount = int32(len(group.Entries))
  38. groups = append(groups, group)
  39. }
  40. sortExecutionQueueGroups(groups)
  41. return groups
  42. }
  43. func newExecutionQueueGroup(entry *executor.InternalLogEntry) *apiv1.ExecutionQueueGroup {
  44. group := &apiv1.ExecutionQueueGroup{
  45. BindingId: entry.GetBindingId(),
  46. ActionTitle: entry.ActionTitle,
  47. ActionIcon: entry.ActionIcon,
  48. EntityPrefix: entry.EntityPrefix,
  49. }
  50. if entry.Binding != nil && entry.Binding.Action != nil {
  51. group.MaxConcurrent = int32(entry.Binding.Action.MaxConcurrent)
  52. }
  53. return group
  54. }
  55. func sortQueueEntries(entries []*apiv1.LogEntry) {
  56. sort.Slice(entries, func(i, j int) bool {
  57. return entries[i].DatetimeStarted < entries[j].DatetimeStarted
  58. })
  59. }
  60. func sortExecutionQueueGroups(groups []*apiv1.ExecutionQueueGroup) {
  61. sort.Slice(groups, func(i, j int) bool {
  62. left := groups[i].ActionTitle
  63. right := groups[j].ActionTitle
  64. if left == right {
  65. return groups[i].EntityPrefix < groups[j].EntityPrefix
  66. }
  67. return left < right
  68. })
  69. }