start queue task
(task *QueueTask)
| 95 | |
| 96 | //start queue task |
| 97 | func startQueueTask(task *QueueTask) { |
| 98 | taskCtx := task.getTaskContext() |
| 99 | handler := func() { |
| 100 | defer func() { |
| 101 | task.putTaskContext(taskCtx) |
| 102 | if err := recover(); err != nil { |
| 103 | task.CounterInfo().ErrorCounter.Inc(1) |
| 104 | if task.taskService.ExceptionHandler != nil { |
| 105 | task.taskService.ExceptionHandler(taskCtx, fmt.Errorf("%v", err)) |
| 106 | } |
| 107 | } |
| 108 | }() |
| 109 | |
| 110 | task.CounterInfo().RunCounter.Inc(1) |
| 111 | //get value from message chan |
| 112 | message := <-task.MessageChan |
| 113 | taskCtx.Message = message |
| 114 | |
| 115 | if task.taskService != nil && task.taskService.OnBeforeHandler != nil { |
| 116 | task.taskService.OnBeforeHandler(taskCtx) |
| 117 | } |
| 118 | |
| 119 | var err error |
| 120 | if !taskCtx.IsEnd { |
| 121 | err = task.handler(taskCtx) |
| 122 | } |
| 123 | |
| 124 | if err != nil { |
| 125 | taskCtx.Error = err |
| 126 | task.CounterInfo().ErrorCounter.Inc(1) |
| 127 | if task.taskService != nil && task.taskService.ExceptionHandler != nil { |
| 128 | task.taskService.ExceptionHandler(taskCtx, err) |
| 129 | } |
| 130 | } |
| 131 | |
| 132 | if task.taskService != nil && task.taskService.OnEndHandler != nil { |
| 133 | task.taskService.OnEndHandler(taskCtx) |
| 134 | } |
| 135 | } |
| 136 | dofunc := func() { |
| 137 | task.TimeTicker = time.NewTicker(time.Duration(task.Interval) * time.Millisecond) |
| 138 | handler() |
| 139 | for { |
| 140 | select { |
| 141 | case <-task.TimeTicker.C: |
| 142 | handler() |
| 143 | } |
| 144 | } |
| 145 | } |
| 146 | go dofunc() |
| 147 | } |
no test coverage detected