(task *Task, status TaskStatus, err error)
| 206 | } |
| 207 | |
| 208 | func (tm *TaskManager) completeTask(task *Task, status TaskStatus, err error) { |
| 209 | tm.mu.Lock() |
| 210 | task.Status = status |
| 211 | task.Error = err |
| 212 | task.FinishedAt = time.Now() |
| 213 | |
| 214 | key := taskKey(task.TargetNode, task.TargetVMID) |
| 215 | |
| 216 | // Remove from active |
| 217 | delete(tm.activeTasks, key) |
| 218 | |
| 219 | // Check queue for next task, respecting max running task cap. |
| 220 | nextTask := tm.dequeueNextStartableLocked(key) |
| 221 | tm.mu.Unlock() |
| 222 | |
| 223 | if task.OnComplete != nil { |
| 224 | go task.OnComplete(err) |
| 225 | } |
| 226 | |
| 227 | if tm.updateNotify != nil { |
| 228 | go tm.updateNotify() |
| 229 | } |
| 230 | |
| 231 | if nextTask != nil { |
| 232 | go tm.runTask(nextTask) |
| 233 | } |
| 234 | } |
| 235 | |
| 236 | func (tm *TaskManager) dequeueNextStartableLocked(preferredKey string) *Task { |
| 237 | if len(tm.activeTasks) >= tm.maxRunning { |
no test coverage detected