Recover worker leases after empty event replay
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user