initAndStartHookQueues create all queues defined in hooks
()
| 1041 | |
| 1042 | // initAndStartHookQueues create all queues defined in hooks |
| 1043 | func (op *ShellOperator) initAndStartHookQueues() { |
| 1044 | schHooks, _ := op.HookManager.GetHooksInOrder(types.Schedule) |
| 1045 | for _, hookName := range schHooks { |
| 1046 | h := op.HookManager.GetHook(hookName) |
| 1047 | for _, hookBinding := range h.Config.Schedules { |
| 1048 | if op.TaskQueues.GetByName(hookBinding.Queue) == nil { |
| 1049 | op.TaskQueues.NewNamedQueue(hookBinding.Queue, |
| 1050 | op.taskHandler, |
| 1051 | queue.WithCompactableTypes(task_metadata.HookRun), |
| 1052 | queue.WithLogger(op.logger.With(pkg.LogKeyOperatorComponent, "hookQueue", pkg.LogKeyHook, hookName, pkg.LogKeyQueue, hookBinding.Queue)), |
| 1053 | ) |
| 1054 | op.TaskQueues.GetByName(hookBinding.Queue).Start(op.ctx) |
| 1055 | } |
| 1056 | } |
| 1057 | } |
| 1058 | |
| 1059 | kubeHooks, _ := op.HookManager.GetHooksInOrder(types.OnKubernetesEvent) |
| 1060 | for _, hookName := range kubeHooks { |
| 1061 | h := op.HookManager.GetHook(hookName) |
| 1062 | for _, hookBinding := range h.Config.OnKubernetesEvents { |
| 1063 | if op.TaskQueues.GetByName(hookBinding.Queue) == nil { |
| 1064 | op.TaskQueues.NewNamedQueue(hookBinding.Queue, |
| 1065 | op.taskHandler, |
| 1066 | queue.WithCompactableTypes(task_metadata.HookRun), |
| 1067 | queue.WithLogger(op.logger.With(pkg.LogKeyOperatorComponent, "hookQueue", pkg.LogKeyHook, hookName, pkg.LogKeyQueue, hookBinding.Queue)), |
| 1068 | ) |
| 1069 | op.TaskQueues.GetByName(hookBinding.Queue).Start(op.ctx) |
| 1070 | } |
| 1071 | } |
| 1072 | } |
| 1073 | } |
no test coverage detected