| 586 | } |
| 587 | |
| 588 | func recoverTaskRunsOnBoot( |
| 589 | ctx context.Context, |
| 590 | manager *taskpkg.Service, |
| 591 | store taskStore, |
| 592 | sessions taskBridgeSessionManager, |
| 593 | actor taskpkg.ActorContext, |
| 594 | ) (taskRecoveryStats, error) { |
| 595 | expired, err := manager.RecoverExpiredRunLeases(ctx, taskpkg.ExpiredLeaseRecovery{ |
| 596 | Reason: taskRecoveryReasonBoot, |
| 597 | }, actor) |
| 598 | if err != nil { |
| 599 | return taskRecoveryStats{}, fmt.Errorf("daemon: recover expired task run leases on boot: %w", err) |
| 600 | } |
| 601 | |
| 602 | runs, err := store.ListTaskRunsByStatus(ctx, []taskpkg.RunStatus{ |
| 603 | taskpkg.TaskRunStatusClaimed, |
| 604 | taskpkg.TaskRunStatusStarting, |
| 605 | taskpkg.TaskRunStatusRunning, |
| 606 | }) |
| 607 | if err != nil { |
| 608 | return taskRecoveryStats{}, fmt.Errorf("daemon: list task runs for boot recovery: %w", err) |
| 609 | } |
| 610 | |
| 611 | stats := taskRecoveryStats{requeued: len(expired)} |
| 612 | for _, run := range runs { |
| 613 | recovery, err := planTaskRunRecovery(ctx, sessions, run) |
| 614 | if err != nil { |
| 615 | return taskRecoveryStats{}, fmt.Errorf("daemon: plan boot recovery for task run %q: %w", run.ID, err) |
| 616 | } |
| 617 | if recovery == nil { |
| 618 | continue |
| 619 | } |
| 620 | if _, err := manager.RecoverRunOnBoot(ctx, run.ID, *recovery, actor); err != nil { |
| 621 | return taskRecoveryStats{}, fmt.Errorf("daemon: recover task run %q on boot: %w", run.ID, err) |
| 622 | } |
| 623 | switch recovery.Action.Normalize() { |
| 624 | case taskpkg.RunBootRecoveryRequeue: |
| 625 | stats.requeued++ |
| 626 | case taskpkg.RunBootRecoveryMarkRunning: |
| 627 | stats.markedRunning++ |
| 628 | case taskpkg.RunBootRecoveryFail: |
| 629 | stats.failed++ |
| 630 | } |
| 631 | } |
| 632 | |
| 633 | return stats, nil |
| 634 | } |
| 635 | |
| 636 | func planTaskRunRecovery( |
| 637 | ctx context.Context, |