| 133 | } |
| 134 | |
| 135 | func (r *taskRoleRuntime) Recover(ctx context.Context) { |
| 136 | if r == nil { |
| 137 | return |
| 138 | } |
| 139 | if ctx == nil { |
| 140 | ctx = context.Background() |
| 141 | } |
| 142 | runs, err := r.store.ListTaskRunsByStatus(ctx, []taskpkg.RunStatus{taskpkg.TaskRunStatusQueued}) |
| 143 | if err != nil { |
| 144 | r.logTaskRoleError("daemon: list queued task runs for role recovery", err, hookspkg.TaskRunEnqueuedPayload{}) |
| 145 | return |
| 146 | } |
| 147 | for _, run := range runs { |
| 148 | taskRecord, err := r.store.GetTask(ctx, run.TaskID) |
| 149 | if err != nil { |
| 150 | r.logTaskRoleError("daemon: load task for role recovery", err, hookspkg.TaskRunEnqueuedPayload{ |
| 151 | TaskRunContext: hookspkg.TaskRunContext{RunID: run.ID, TaskID: run.TaskID}, |
| 152 | }) |
| 153 | continue |
| 154 | } |
| 155 | if err := r.activateRun(ctx, taskRecord, run, taskRoleActivationReasonRecovery); err != nil { |
| 156 | r.logTaskRoleError( |
| 157 | "daemon: recover task role session for queued run", |
| 158 | err, |
| 159 | hookspkg.TaskRunEnqueuedPayload{ |
| 160 | TaskRunContext: hookspkg.TaskRunContext{ |
| 161 | RunID: run.ID, |
| 162 | TaskID: run.TaskID, |
| 163 | WorkspaceID: taskRecord.WorkspaceID, |
| 164 | CoordinationChannelID: run.CoordinationChannelID, |
| 165 | }, |
| 166 | }, |
| 167 | ) |
| 168 | } |
| 169 | } |
| 170 | } |
| 171 | |
| 172 | func (r *taskRoleRuntime) activateRun( |
| 173 | ctx context.Context, |