RunTasksStandalone run tasks in parallel
(parentLogger log.Logger, taskIds []uint64)
| 176 | |
| 177 | // RunTasksStandalone run tasks in parallel |
| 178 | func RunTasksStandalone(parentLogger log.Logger, taskIds []uint64) errors.Error { |
| 179 | if len(taskIds) == 0 { |
| 180 | return nil |
| 181 | } |
| 182 | results := make(chan error) |
| 183 | for _, taskId := range taskIds { |
| 184 | go func(id uint64) { |
| 185 | taskLog.Info("run task #%d in background ", id) |
| 186 | var err errors.Error |
| 187 | taskErr := runTaskStandalone(parentLogger, id) |
| 188 | if taskErr != nil { |
| 189 | err = errors.Default.Wrap(taskErr, fmt.Sprintf("Error running task %d.", id)) |
| 190 | } |
| 191 | results <- err |
| 192 | }(taskId) |
| 193 | } |
| 194 | errs := make([]error, 0) |
| 195 | var err error |
| 196 | finished := 0 |
| 197 | for err = range results { |
| 198 | if err != nil { |
| 199 | taskLog.Error(err, "task failed") |
| 200 | errs = append(errs, err) |
| 201 | } |
| 202 | finished++ |
| 203 | if finished == len(taskIds) { |
| 204 | close(results) |
| 205 | } |
| 206 | } |
| 207 | if len(errs) > 0 { |
| 208 | var sb strings.Builder |
| 209 | for _, e := range errs { |
| 210 | _, _ = sb.WriteString(e.Error()) |
| 211 | _, _ = sb.WriteString("\n") |
| 212 | if errors.Is(e, context.Canceled) { |
| 213 | parentLogger.Info("task canceled") |
| 214 | return errors.Convert(e) |
| 215 | } |
| 216 | } |
| 217 | err = errors.Default.New(sb.String()) |
| 218 | } |
| 219 | return errors.Convert(err) |
| 220 | } |
| 221 | |
| 222 | // RerunTask reruns specified task |
| 223 | func RerunTask(taskId uint64) (*models.Task, errors.Error) { |
no test coverage detected