diff --git a/cmd/orchestra-worker/main.go b/cmd/orchestra-worker/main.go index 3130f74..79c966d 100644 --- a/cmd/orchestra-worker/main.go +++ b/cmd/orchestra-worker/main.go @@ -260,6 +260,16 @@ func (w *worker) once(ctx context.Context) error { w.cursor = e.Seq } } + // A coordinator restart can restore its task snapshot without retaining + // the in-memory event tail. In that state a worker with a persisted cursor + // receives an empty page even though a lease is currently assigned to it. + // Reconcile the authoritative projection before treating an empty page as + // "nothing to do"; otherwise the lease remains invisible until expiry. + if len(es) == 0 { + if err := w.reconcileLeases(ctx); err != nil { + return err + } + } // State is only a cache. If a lease survived but its TaskCreated event is // older than the worker's cursor (or the event has been compacted), hydrate // the authoritative task projection before deciding whether to start. @@ -294,6 +304,29 @@ func (w *worker) once(ctx context.Context) error { return w.api.Ack(ctx, w.cursor) } +func (w *worker) reconcileLeases(ctx context.Context) error { + tasks, err := w.api.Tasks(ctx) + if err != nil { + return fmt.Errorf("reconcile leased tasks: %w", err) + } + active := make(map[string]lease) + for _, task := range tasks { + w.tasks[task.ID] = task + if task.State == domain.StateLeased && task.Lease != nil && task.Lease.HarnessID == w.harnessID { + active[task.ID] = lease{HandoffRef: task.HandoffRef} + } + } + for taskID := range w.leases { + if _, ok := active[taskID]; !ok { + delete(w.leases, taskID) + } + } + for taskID, l := range active { + w.leases[taskID] = l + } + return nil +} + // reRegisterAfterCoordinatorRestart restores the coordinator's in-memory // worker registry. A worker must survive a server restart without operator // intervention; its persisted cursor and sessions remain valid. diff --git a/cmd/orchestra-worker/main_test.go b/cmd/orchestra-worker/main_test.go index 839839a..95730ee 100644 --- a/cmd/orchestra-worker/main_test.go +++ b/cmd/orchestra-worker/main_test.go @@ -83,6 +83,29 @@ func TestWorkerHydratesLeaseMissingTaskCache(t *testing.T) { } } +func TestWorkerReconcilesLeaseAfterEmptyEventReplay(t *testing.T) { + s := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/v1/federation/events": + _, _ = w.Write([]byte(`{"cursor":42,"events":[]}`)) + case "/v1/tasks": + _, _ = w.Write([]byte(`[{"id":"t","source":"s","external_id":"x","project":"p","state":"leased","lease":{"harness_id":"h"}}]`)) + case "/v1/federation/events/ack": + w.WriteHeader(http.StatusNoContent) + default: + t.Fatalf("unexpected %s", r.URL.Path) + } + })) + defer s.Close() + w := &worker{api: federation.Client{BaseURL: s.URL, WorkerID: "h", Token: "t"}, harnessID: "h", cursor: 42, tasks: map[string]domain.Task{}, sessions: map[string]herdr.Session{"t": {}}, leases: map[string]lease{}, statePath: t.TempDir() + "/state.json", hard: .75} + if err := w.once(context.Background()); err != nil { + t.Fatal(err) + } + if _, ok := w.leases["t"]; !ok { + t.Fatal("empty event replay did not recover current lease") + } +} + func TestWorkerStartsRouterIssuedLeaseInLocalGitWorktree(t *testing.T) { remote := filepath.Join(t.TempDir(), "remote.git") if out, err := exec.Command("git", "init", "--bare", remote).CombinedOutput(); err != nil {