diff --git a/spindle/db/db.go b/spindle/db/db.go index f4a05c69..9b9a15c1 100644 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -130,6 +130,22 @@ func Make(ctx context.Context, dbPath string) (*DB, error) { expires_at text ); + create table if not exists mill_leases ( + lease_id text primary key, + node_id text not null, + engine text not null, + knot text not null, + rkey text not null, + workflow text not null, + state text not null, + created_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')) + ); + + create table if not exists mill_executor_cursors ( + node_id text primary key, + acked_offset integer not null + ); + create table if not exists migrations ( id integer primary key autoincrement, name text unique diff --git a/spindle/db/events.go b/spindle/db/events.go index 1879d2e1..1d8826ed 100644 --- a/spindle/db/events.go +++ b/spindle/db/events.go @@ -94,6 +94,39 @@ func (d *DB) InsertRelayedStatus( return d.insertEvent(event, n) } +// CompleteOrphanMillLease makes the terminal event and lease deletion one +// durable transition; an event without its matching deletion would be replayed. +func (d *DB) CompleteOrphanMillLease( + leaseID string, + pipelineAtUri string, + workflow string, + status string, + workflowError *string, + exitCode *int64, + n *notifier.Notifier, +) error { + event, err := statusEvent(pipelineAtUri, workflow, status, workflowError, exitCode) + if err != nil { + return err + } + tx, err := d.Begin() + if err != nil { + return err + } + defer tx.Rollback() + if err := eventstream.Insert(tx, event, nil); err != nil { + return err + } + if _, err := tx.Exec(`delete from mill_leases where lease_id = ?`, leaseID); err != nil { + return err + } + if err := tx.Commit(); err != nil { + return err + } + n.NotifyAll() + return nil +} + func (d *DB) GetStatus(workflowId models.WorkflowId) (*tangled.PipelineStatus, error) { pipelineAtUri := workflowId.PipelineId.AtUri() diff --git a/spindle/db/mill_state.go b/spindle/db/mill_state.go new file mode 100644 index 00000000..74fbc46b --- /dev/null +++ b/spindle/db/mill_state.go @@ -0,0 +1,86 @@ +package db + +// MillLease is the persisted view of a mill-side remote lease, enough to +// rebuild the fencing token and its workflow identity after a mill restart. +type MillLease struct { + LeaseID string + NodeID string + Engine string + Knot string + Rkey string + Workflow string + State string +} + +func (d *DB) SaveMillLease(l MillLease) error { + _, err := d.Exec( + `insert into mill_leases (lease_id, node_id, engine, knot, rkey, workflow, state) + values (?, ?, ?, ?, ?, ?, ?) + on conflict(lease_id) do update set state = excluded.state`, + l.LeaseID, l.NodeID, l.Engine, l.Knot, l.Rkey, l.Workflow, l.State, + ) + return err +} + +func (d *DB) DeleteMillLease(leaseID string) error { + _, err := d.Exec(`delete from mill_leases where lease_id = ?`, leaseID) + return err +} + +func (d *DB) DeleteMillLeasesByNode(nodeID string) error { + _, err := d.Exec(`delete from mill_leases where node_id = ?`, nodeID) + return err +} + +func (d *DB) ListMillLeases() ([]MillLease, error) { + rows, err := d.Query(`select lease_id, node_id, engine, knot, rkey, workflow, state from mill_leases`) + if err != nil { + return nil, err + } + defer rows.Close() + + var leases []MillLease + for rows.Next() { + var l MillLease + if err := rows.Scan(&l.LeaseID, &l.NodeID, &l.Engine, &l.Knot, &l.Rkey, &l.Workflow, &l.State); err != nil { + return nil, err + } + leases = append(leases, l) + } + return leases, rows.Err() +} + +// SetExecutorCursor persists the highest relay offset the mill has applied for +// a node. +func (d *DB) SetExecutorCursor(nodeID string, offset uint64) error { + _, err := d.Exec( + `insert into mill_executor_cursors (node_id, acked_offset) values (?, ?) + on conflict(node_id) do update set acked_offset = excluded.acked_offset`, + nodeID, offset, + ) + return err +} + +func (d *DB) DeleteExecutorCursor(nodeID string) error { + _, err := d.Exec(`delete from mill_executor_cursors where node_id = ?`, nodeID) + return err +} + +func (d *DB) ListExecutorCursors() (map[string]uint64, error) { + rows, err := d.Query(`select node_id, acked_offset from mill_executor_cursors`) + if err != nil { + return nil, err + } + defer rows.Close() + + cursors := make(map[string]uint64) + for rows.Next() { + var node string + var offset uint64 + if err := rows.Scan(&node, &offset); err != nil { + return nil, err + } + cursors[node] = offset + } + return cursors, rows.Err() +} diff --git a/spindle/db/mill_state_test.go b/spindle/db/mill_state_test.go new file mode 100644 index 00000000..c35a09ea --- /dev/null +++ b/spindle/db/mill_state_test.go @@ -0,0 +1,181 @@ +package db + +import ( + "testing" + + "tangled.org/core/notifier" +) + +func TestMillLeaseRoundTrip(t *testing.T) { + d := newTestDB(t) + + lease := MillLease{ + LeaseID: "lease-1", + NodeID: "node-1", + Engine: "dummy", + Knot: "knot.example", + Rkey: "rkey1", + Workflow: "build", + State: "reserved", + } + if err := d.SaveMillLease(lease); err != nil { + t.Fatalf("SaveMillLease: %v", err) + } + + lease.State = "running" + if err := d.SaveMillLease(lease); err != nil { + t.Fatalf("SaveMillLease(transition): %v", err) + } + + leases, err := d.ListMillLeases() + if err != nil { + t.Fatalf("ListMillLeases: %v", err) + } + if len(leases) != 1 { + t.Fatalf("ListMillLeases returned %d leases, want 1 (state transition must replace, not duplicate)", len(leases)) + } + if leases[0] != lease { + t.Fatalf("ListMillLeases[0] = %+v, want %+v", leases[0], lease) + } + + if err := d.DeleteMillLease("lease-1"); err != nil { + t.Fatalf("DeleteMillLease: %v", err) + } + if leases, _ = d.ListMillLeases(); len(leases) != 0 { + t.Fatalf("lease survived deletion: %+v", leases) + } +} + +func TestDeleteMillLeasesByNode(t *testing.T) { + d := newTestDB(t) + + for _, l := range []MillLease{ + {LeaseID: "a", NodeID: "node-1", Engine: "dummy", Knot: "k", Rkey: "r1", Workflow: "w", State: "running"}, + {LeaseID: "b", NodeID: "node-1", Engine: "dummy", Knot: "k", Rkey: "r2", Workflow: "w", State: "reserved"}, + {LeaseID: "c", NodeID: "node-2", Engine: "dummy", Knot: "k", Rkey: "r3", Workflow: "w", State: "running"}, + } { + if err := d.SaveMillLease(l); err != nil { + t.Fatalf("SaveMillLease(%s): %v", l.LeaseID, err) + } + } + + if err := d.DeleteMillLeasesByNode("node-1"); err != nil { + t.Fatalf("DeleteMillLeasesByNode: %v", err) + } + leases, err := d.ListMillLeases() + if err != nil { + t.Fatalf("ListMillLeases: %v", err) + } + if len(leases) != 1 || leases[0].LeaseID != "c" { + t.Fatalf("ListMillLeases = %+v, want only node-2's lease c", leases) + } +} + +func TestExecutorCursors(t *testing.T) { + d := newTestDB(t) + + if err := d.SetExecutorCursor("node-1", 5); err != nil { + t.Fatalf("SetExecutorCursor: %v", err) + } + if err := d.SetExecutorCursor("node-1", 9); err != nil { + t.Fatalf("SetExecutorCursor(advance): %v", err) + } + if err := d.SetExecutorCursor("node-2", 1); err != nil { + t.Fatalf("SetExecutorCursor(node-2): %v", err) + } + + cursors, err := d.ListExecutorCursors() + if err != nil { + t.Fatalf("ListExecutorCursors: %v", err) + } + if len(cursors) != 2 || cursors["node-1"] != 9 || cursors["node-2"] != 1 { + t.Fatalf("ListExecutorCursors = %v, want node-1:9 node-2:1", cursors) + } + + if err := d.DeleteExecutorCursor("node-1"); err != nil { + t.Fatalf("DeleteExecutorCursor: %v", err) + } + if cursors, _ = d.ListExecutorCursors(); len(cursors) != 1 || cursors["node-2"] != 1 { + t.Fatalf("ListExecutorCursors after delete = %v, want only node-2:1", cursors) + } +} + +func TestCompleteOrphanMillLeaseIsAtomic(t *testing.T) { + d := newTestDB(t) + lease := MillLease{ + LeaseID: "lease-1", NodeID: "node-1", Engine: "dummy", + Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: "running", + } + if err := d.SaveMillLease(lease); err != nil { + t.Fatalf("SaveMillLease: %v", err) + } + if _, err := d.Exec(` + create trigger reject_mill_lease_delete + before delete on mill_leases + begin + select raise(abort, 'forced delete failure'); + end + `); err != nil { + t.Fatalf("create failure trigger: %v", err) + } + + n := notifier.New() + notifications := n.Subscribe() + defer n.Unsubscribe(notifications) + err := d.CompleteOrphanMillLease( + "lease-1", + "at://knot.example/sh.tangled.pipeline/rkey1", + "build", + "failed", + nil, + nil, + &n, + ) + if err == nil { + t.Fatal("CompleteOrphanMillLease succeeded despite forced lease deletion failure") + } + var eventCount int + if err := d.QueryRow(`select count(*) from events`).Scan(&eventCount); err != nil { + t.Fatalf("count events after rollback: %v", err) + } + if eventCount != 0 { + t.Fatalf("terminal event count after rollback = %d, want 0", eventCount) + } + if leases, listErr := d.ListMillLeases(); listErr != nil || len(leases) != 1 { + t.Fatalf("leases after rollback = %+v, err = %v; want original lease", leases, listErr) + } + select { + case <-notifications: + t.Fatal("rollback notified event subscribers") + default: + } + + if _, err := d.Exec(`drop trigger reject_mill_lease_delete`); err != nil { + t.Fatalf("drop failure trigger: %v", err) + } + if err := d.CompleteOrphanMillLease( + "lease-1", + "at://knot.example/sh.tangled.pipeline/rkey1", + "build", + "failed", + nil, + nil, + &n, + ); err != nil { + t.Fatalf("CompleteOrphanMillLease retry: %v", err) + } + if err := d.QueryRow(`select count(*) from events`).Scan(&eventCount); err != nil { + t.Fatalf("count committed events: %v", err) + } + if eventCount != 1 { + t.Fatalf("terminal event count after commit = %d, want 1", eventCount) + } + if leases, listErr := d.ListMillLeases(); listErr != nil || len(leases) != 0 { + t.Fatalf("leases after commit = %+v, err = %v; want none", leases, listErr) + } + select { + case <-notifications: + default: + t.Fatal("committed terminal event did not notify subscribers") + } +} diff --git a/spindle/mill/executor/executor.go b/spindle/mill/executor/executor.go index f3b2eb0e..583ecec2 100644 --- a/spindle/mill/executor/executor.go +++ b/spindle/mill/executor/executor.go @@ -504,6 +504,10 @@ func (e *Executor) pushSnapshot() { e.mu.Lock() draining := e.draining active := len(e.active) + leaseIDs := make([]string, 0, len(e.active)) + for id := range e.active { + leaseIDs = append(leaseIDs, id) + } e.mu.Unlock() load := 0.0 @@ -529,8 +533,9 @@ func (e *Executor) pushSnapshot() { } } e.send(&millproto.Message{NodeSnapshot: &millv1.NodeSnapshot{ - NodeId: e.nodeID, - Engines: engines, + NodeId: e.nodeID, + Engines: engines, + ActiveLeaseIds: leaseIDs, }}) } diff --git a/spindle/mill/handler.go b/spindle/mill/handler.go index f9c6f149..c49005a5 100644 --- a/spindle/mill/handler.go +++ b/spindle/mill/handler.go @@ -71,7 +71,9 @@ func (m *Mill) HandleExecutorConn(w http.ResponseWriter, r *http.Request) { } m.sessionReady(sess) - sess.readLoop(m, dec) + if err := sess.readLoop(m, dec); err != nil { + m.l.Debug("session read ended", "node", sess.nodeID, "err", err) + } m.detachSession(sess) } diff --git a/spindle/mill/lease.go b/spindle/mill/lease.go index 11f67711..3f971f6d 100644 --- a/spindle/mill/lease.go +++ b/spindle/mill/lease.go @@ -40,15 +40,23 @@ type RemoteLease struct { nodeID string engine string wid models.WorkflowId // the job this lease carries; set once placed + // restored from persistence after a mill restart: no RunStep waits on it, + // so terminals and death are authored directly. Set before publication, + // never mutated. + orphaned bool mu sync.Mutex state leaseState cancel bool reason string + released bool terminal chan *millv1.AttemptResult // buffered(1); RunStep waits here dead chan struct{} // closed when the executor is lost past grace deadOnce sync.Once + finishMu sync.Mutex + cleanupOnce sync.Once + // log file for relayed lines, opened lazily and held open for the lease's // life so we don't reopen per line. Guarded by logMu, closed on cleanup. logMu sync.Mutex @@ -99,6 +107,17 @@ func (l *RemoteLease) getState() leaseState { return l.state } +// markDone CASes to done; exactly one caller wins. +func (l *RemoteLease) markDone() bool { + l.mu.Lock() + defer l.mu.Unlock() + if l.state == leaseDone { + return false + } + l.state = leaseDone + return true +} + func (l *RemoteLease) requestCancel(reason string) cancelAction { l.mu.Lock() defer l.mu.Unlock() @@ -119,6 +138,19 @@ func (l *RemoteLease) cancelRequested() (bool, string) { return l.cancel, l.reason } +func (l *RemoteLease) releaseState() (leaseState, bool) { + l.mu.Lock() + defer l.mu.Unlock() + l.released = true + return l.state, l.cancel +} + +func (l *RemoteLease) cleanupReady() bool { + l.mu.Lock() + defer l.mu.Unlock() + return l.released && l.state == leaseDone +} + func (l *RemoteLease) deliverCancelled(reason string) { l.deliverTerminal(&millv1.AttemptResult{ LeaseId: l.id, diff --git a/spindle/mill/mill.go b/spindle/mill/mill.go index c9a98001..52de7d1d 100644 --- a/spindle/mill/mill.go +++ b/spindle/mill/mill.go @@ -223,12 +223,29 @@ func (m *Mill) onGraceExpired(sess *millSession) { delete(m.nodeOffset, sess.nodeID) m.mu.Unlock() + if m.db != nil { + if err := m.db.DeleteExecutorCursor(sess.nodeID); err != nil { + m.l.Error("delete executor cursor", "node", sess.nodeID, "err", err) + } + } + m.l.Warn("executor declared dead; failing its in-flight jobs", "node", sess.nodeID, "jobs", len(dead)) + deadReason := "executor lost" for _, lease := range dead { - if cancelled, reason := lease.cancelRequested(); cancelled { - lease.deliverCancelled(reason) - } else { - lease.markDead() + switch { + case lease.orphaned: + if err := m.finishOrphan(lease, string(models.StatusKindFailed), &deadReason, nil); err != nil { + m.l.Error("finish orphan after executor loss", "lease", lease.id, "err", err) + } + default: + if cancelled, reason := lease.cancelRequested(); cancelled { + lease.deliverCancelled(reason) + if lease.cleanupReady() { + m.cleanupLease(lease) + } + } else { + lease.markDead() + } } } m.notifyChange() @@ -269,6 +286,10 @@ func (m *Mill) place(ctx context.Context, engineName string, wid models.Workflow } if lease != nil { lease.wid = wid + if err := m.persistLease(lease, leaseRowReserved); err != nil { + m.releaseRemote(lease) + return nil, fmt.Errorf("persist reserved mill lease: %w", err) + } m.mu.Lock() m.leases[lease.id] = lease if st, ok := wf.Data.(*millWorkflowState); ok && st != nil { @@ -554,6 +575,9 @@ func (m *Mill) commitAndWait(ctx context.Context, wf *models.Workflow, unlocked return engine.ErrWorkflowFailed } lease.markRunning() + if err := m.persistLease(lease, leaseRowRunning); err != nil { + m.l.Error("persist running mill lease", "lease", lease.id, "err", err) + } if cancelled, reason := lease.cancelRequested(); cancelled { m.sendCancel(sess, lease, reason) } @@ -651,16 +675,23 @@ func (m *Mill) destroy(wid models.WorkflowId) { func (m *Mill) releaseSlot(s *millSlot) { lease := s.lease - if lease.getState() == leaseReserved { + state, cancelled := lease.releaseState() + if state == leaseReserved { m.releaseRemote(lease) + } else if cancelled && state != leaseDone { + return } - lease.closeLog() - - m.mu.Lock() - delete(m.leases, lease.id) - m.mu.Unlock() - - m.notifyChange() + m.cleanupLease(lease) +} +func (m *Mill) cleanupLease(lease *RemoteLease) { + lease.cleanupOnce.Do(func() { + lease.closeLog() + m.mu.Lock() + delete(m.leases, lease.id) + m.mu.Unlock() + m.deleteLeaseRow(lease.id) + m.notifyChange() + }) } func (m *Mill) releaseRemote(lease *RemoteLease) { @@ -696,38 +727,51 @@ func (m *Mill) shouldAcceptOffset(sess *millSession, offset uint64) bool { return offset == current+1 } -func (m *Mill) ackOffset(sess *millSession, offset uint64) { +func (m *Mill) ackOffset(sess *millSession, offset uint64) error { + if m.db != nil { + if err := m.db.SetExecutorCursor(sess.nodeID, offset); err != nil { + return fmt.Errorf("persist executor cursor: %w", err) + } + } + m.mu.Lock() if offset > m.nodeOffset[sess.nodeID] { m.nodeOffset[sess.nodeID] = offset } m.mu.Unlock() - _ = sess.send(&millproto.Message{Ack: &millv1.Ack{UpToOffset: offset}}) + if err := sess.send(&millproto.Message{Ack: &millv1.Ack{UpToOffset: offset}}); err != nil { + return fmt.Errorf("send relay ack: %w", err) + } + return nil } -func (m *Mill) processRelay(sess *millSession, offset uint64, kind string, apply func() error) { +func (m *Mill) processRelay(sess *millSession, offset uint64, kind string, apply func() error) error { if !m.shouldAcceptOffset(sess, offset) { - return + return nil } - // ack only after the side effect lands, or a transient write error makes - // the message unreplayable. if err := apply(); err != nil { - m.l.Error("process relayed message failed", "kind", kind, "node", sess.nodeID, "offset", offset, "err", err) - return + return fmt.Errorf("process relayed %s at offset %d: %w", kind, offset, err) } - m.ackOffset(sess, offset) + if err := m.ackOffset(sess, offset); err != nil { + return fmt.Errorf("ack relayed %s at offset %d: %w", kind, offset, err) + } + return nil } -func (m *Mill) onSnapshot(sess *millSession, snap *millv1.NodeSnapshot) { +func (m *Mill) onSnapshot(sess *millSession, snap *millv1.NodeSnapshot) error { m.mu.Lock() sess.snapshot = snap m.mu.Unlock() + if err := m.reconcileOrphans(sess.nodeID, snap.GetActiveLeaseIds()); err != nil { + return err + } m.notifyChange() + return nil } -func (m *Mill) onStatusRelay(sess *millSession, ev *millv1.StatusEvent) { - m.processRelay(sess, ev.GetOffset(), "status", func() error { +func (m *Mill) onStatusRelay(sess *millSession, ev *millv1.StatusEvent) error { + return m.processRelay(sess, ev.GetOffset(), "status", func() error { if m.db == nil { return fmt.Errorf("mill db not attached") } @@ -753,8 +797,8 @@ func (m *Mill) onStatusRelay(sess *millSession, ev *millv1.StatusEvent) { }) } -func (m *Mill) onLogRelay(sess *millSession, ll *millv1.LogLine) { - m.processRelay(sess, ll.GetOffset(), "log", func() error { +func (m *Mill) onLogRelay(sess *millSession, ll *millv1.LogLine) error { + return m.processRelay(sess, ll.GetOffset(), "log", func() error { m.mu.Lock() lease := m.leases[ll.GetLeaseId()] m.mu.Unlock() @@ -767,15 +811,29 @@ func (m *Mill) onLogRelay(sess *millSession, ll *millv1.LogLine) { }) } -func (m *Mill) onAttemptResult(sess *millSession, ar *millv1.AttemptResult) { - m.processRelay(sess, ar.GetOffset(), "attempt-result", func() error { +func (m *Mill) onAttemptResult(sess *millSession, ar *millv1.AttemptResult) error { + return m.processRelay(sess, ar.GetOffset(), "attempt-result", func() error { m.mu.Lock() lease := m.leases[ar.GetLeaseId()] m.mu.Unlock() if lease == nil || lease.nodeID != sess.nodeID { return nil } + if lease.orphaned { + var errMsg *string + if e := ar.GetError(); e != "" { + errMsg = &e + } + var exit *int64 + if c := ar.GetExitCode(); c != 0 { + exit = &c + } + return m.finishOrphan(lease, ar.GetTerminalStatus(), errMsg, exit) + } lease.deliverTerminal(ar) + if lease.cleanupReady() { + m.cleanupLease(lease) + } return nil }) } diff --git a/spindle/mill/mill_test.go b/spindle/mill/mill_test.go index 7ce23b1e..cef8f1bf 100644 --- a/spindle/mill/mill_test.go +++ b/spindle/mill/mill_test.go @@ -1,11 +1,14 @@ package mill import ( + "bytes" "context" "errors" "io" "log/slog" + "path/filepath" "strings" + "sync" "testing" "time" @@ -445,26 +448,113 @@ func TestMaxPendingRejects(t *testing.T) { } } -func TestSessionRequestCancelledBeforeRegistrationOrSend(t *testing.T) { - sent := 0 - sess := newSession("node-1", scriptedEncoder(func(*millproto.Message) error { - sent++ +func TestCancelledRunningLeaseSurvivesReleaseForReconnectReplay(t *testing.T) { + m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) + wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + lease := newLease("lease-1", "node-1", "dummy") + lease.wid = wid + lease.setState(leaseRunning) + if err := m.persistLease(lease, leaseRowRunning); err != nil { + t.Fatalf("persistLease: %v", err) + } + m.mu.Lock() + m.leases[lease.id] = lease + m.mu.Unlock() + + m.destroy(wid) + slot := &millSlot{fleet: m, lease: lease} + slot.Release() + + m.mu.Lock() + _, retained := m.leases[lease.id] + m.mu.Unlock() + if !retained { + t.Fatal("slot release removed a cancellation-requested running lease before its terminal result") + } + if rows, err := bdb.ListMillLeases(); err != nil || len(rows) != 1 { + t.Fatalf("durable leases after slot release = %+v, err = %v; want retained lease", rows, err) + } + + var sentMu sync.Mutex + var sent []*millproto.Message + sess := newSession("node-1", scriptedEncoder(func(msg *millproto.Message) error { + sentMu.Lock() + sent = append(sent, msg) + sentMu.Unlock() return nil }), discardLogger()) - ctx, cancel := context.WithCancel(context.Background()) - cancel() - _, err := sess.request(ctx, "lease-1", &millproto.Message{ReleaseLease: &millv1.ReleaseLease{LeaseId: "lease-1"}}) - if !errors.Is(err, context.Canceled) { - t.Fatalf("request error = %v, want context.Canceled", err) + if _, ok := m.attachSession(sess); !ok { + t.Fatal("attachSession rejected reconnect") + } + m.sessionReady(sess) + sentMu.Lock() + var replayed bool + for _, msg := range sent { + if cancel := msg.GetCancelAttempt(); cancel != nil && cancel.GetLeaseId() == lease.id { + replayed = true + } } - if sent != 0 { - t.Fatalf("request sent %d messages for an already-cancelled context, want 0", sent) + sentMu.Unlock() + if !replayed { + t.Fatal("reconnect did not replay CancelAttempt for retained lease") } - sess.mu.Lock() - pending := len(sess.pending) - sess.mu.Unlock() - if pending != 0 { - t.Fatalf("request left %d pending waiters, want 0", pending) + + if err := m.onAttemptResult(sess, &millv1.AttemptResult{ + Offset: 1, + LeaseId: lease.id, + TerminalStatus: string(models.StatusKindCancelled), + }); err != nil { + t.Fatalf("onAttemptResult: %v", err) + } + m.mu.Lock() + _, retained = m.leases[lease.id] + m.mu.Unlock() + if retained { + t.Fatal("terminal result did not clean retained cancelled lease") + } + if rows, err := bdb.ListMillLeases(); err != nil || len(rows) != 0 { + t.Fatalf("durable leases after terminal = %+v, err = %v; want none", rows, err) + } + slot.Release() +} + +func TestRelaySideEffectFailureStopsReadLoopAtFailedOffset(t *testing.T) { + m := New(discardLogger(), Config{LogDir: filepath.Join(t.TempDir(), "missing")}) + lease := newLease("lease-1", "node-1", "dummy") + lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + m.mu.Lock() + m.leases[lease.id] = lease + m.mu.Unlock() + sess := newSession("node-1", nopEncoder(), discardLogger()) + if _, ok := m.attachSession(sess); !ok { + t.Fatal("attachSession rejected first session") + } + + var frames bytes.Buffer + enc := millproto.NewEncoder(&frames) + if err := enc.Encode(&millproto.Message{LogLine: &millv1.LogLine{ + Offset: 1, LeaseId: lease.id, RawJson: []byte(`{"line":"first"}`), + }}); err != nil { + t.Fatalf("encode failing log relay: %v", err) + } + if err := enc.Encode(&millproto.Message{AttemptResult: &millv1.AttemptResult{ + Offset: 2, LeaseId: lease.id, TerminalStatus: string(models.StatusKindSuccess), + }}); err != nil { + t.Fatalf("encode later relay: %v", err) + } + + err := sess.readLoop(m, millproto.NewDecoder(&frames)) + if err == nil { + t.Fatal("readLoop continued after relayed log side effect failed") + } + m.mu.Lock() + offset := m.nodeOffset[sess.nodeID] + m.mu.Unlock() + if offset != 0 { + t.Fatalf("in-memory relay offset = %d, want 0 so reconnect retries offset 1", offset) + } + if lease.getState() == leaseDone { + t.Fatal("readLoop applied offset 2 after offset 1 failed") } } @@ -496,6 +586,29 @@ func TestFallbackBidGetsFullTimeout(t *testing.T) { } } +func TestSessionRequestCancelledBeforeRegistrationOrSend(t *testing.T) { + sent := 0 + sess := newSession("node-1", scriptedEncoder(func(*millproto.Message) error { + sent++ + return nil + }), discardLogger()) + ctx, cancel := context.WithCancel(context.Background()) + cancel() + _, err := sess.request(ctx, "lease-1", &millproto.Message{ReleaseLease: &millv1.ReleaseLease{LeaseId: "lease-1"}}) + if !errors.Is(err, context.Canceled) { + t.Fatalf("request error = %v, want context.Canceled", err) + } + if sent != 0 { + t.Fatalf("request sent %d messages for an already-cancelled context, want 0", sent) + } + sess.mu.Lock() + pending := len(sess.pending) + sess.mu.Unlock() + if pending != 0 { + t.Fatalf("request left %d pending waiters, want 0", pending) + } +} + func TestAttachSessionReplacesSilentIncumbentButRejectsActiveDuplicate(t *testing.T) { m := New(discardLogger(), Config{ReconnectGrace: time.Minute}) old := newSession("node-1", nopEncoder(), discardLogger()) @@ -523,8 +636,61 @@ func TestAttachSessionReplacesSilentIncumbentButRejectsActiveDuplicate(t *testin m.mu.Lock() replacement.lastSeen = time.Now().Add(-2 * m.cfg.ReconnectGrace) m.mu.Unlock() - replacement.dispatch(m, &millproto.Message{NodeSnapshot: &millv1.NodeSnapshot{NodeId: "node-1"}}) + if err := replacement.dispatch(m, &millproto.Message{NodeSnapshot: &millv1.NodeSnapshot{NodeId: "node-1"}}); err != nil { + t.Fatalf("periodic snapshot dispatch: %v", err) + } if _, ok := m.attachSession(newSession("node-1", nopEncoder(), discardLogger())); ok { t.Fatal("active replacement did not reject a duplicate session") } } + +func TestPlaceReleasesRemoteReservationWhenInitialPersistenceFails(t *testing.T) { + m, bdb := restoreTestMill(t, Config{BidTimeout: time.Second}) + if err := bdb.Close(); err != nil { + t.Fatalf("close db: %v", err) + } + released := make(chan string, 1) + var sess *millSession + sess = addCandidateSession(t, m, "node-1", nil, 0, scriptedEncoder(func(msg *millproto.Message) error { + switch { + case msg.GetReserveSeat() != nil: + leaseID := msg.GetReserveSeat().GetLeaseId() + sess.deliver(leaseID, &millproto.Message{ReserveResult: &millv1.ReserveResult{ + LeaseId: leaseID, + Accepted: true, + }}) + case msg.GetReleaseLease() != nil: + released <- msg.GetReleaseLease().GetLeaseId() + } + return nil + })) + + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + slot, err := m.place( + ctx, + "dummy", + models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}, + testWorkflow("build"), + ) + if err == nil { + t.Fatal("place succeeded after reserved lease persistence failed") + } + if slot != nil { + t.Fatalf("place returned slot %T after persistence failure", slot) + } + select { + case leaseID := <-released: + if leaseID == "" { + t.Fatal("ReleaseLease had empty lease id") + } + case <-time.After(time.Second): + t.Fatal("persistence failure did not compensate with ReleaseLease") + } + m.mu.Lock() + leases := len(m.leases) + m.mu.Unlock() + if leases != 0 { + t.Fatalf("mill published %d leases after initial persistence failure, want 0", leases) + } +} diff --git a/spindle/mill/proto/gen/mill.pb.go b/spindle/mill/proto/gen/mill.pb.go index a393ee13..c9121c20 100644 --- a/spindle/mill/proto/gen/mill.pb.go +++ b/spindle/mill/proto/gen/mill.pb.go @@ -252,12 +252,15 @@ func (x *EngineAvailability) GetCapabilities() []string { // NodeSnapshot is pushed on connect, periodically, and right after any state // change (reserve, commit, terminal). type NodeSnapshot struct { - state protoimpl.MessageState `protogen:"open.v1"` - NodeId string `protobuf:"bytes,1,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"` - Seq uint64 `protobuf:"varint,2,opt,name=seq,proto3" json:"seq,omitempty"` - Engines map[string]*EngineAvailability `protobuf:"bytes,3,rep,name=engines,proto3" json:"engines,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + state protoimpl.MessageState `protogen:"open.v1"` + NodeId string `protobuf:"bytes,1,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"` + Seq uint64 `protobuf:"varint,2,opt,name=seq,proto3" json:"seq,omitempty"` + Engines map[string]*EngineAvailability `protobuf:"bytes,3,rep,name=engines,proto3" json:"engines,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"` + // every lease the executor currently holds (reserved or running), so the + // mill can reconcile restored leases against executor reality. + ActiveLeaseIds []string `protobuf:"bytes,4,rep,name=active_lease_ids,json=activeLeaseIds,proto3" json:"active_lease_ids,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *NodeSnapshot) Reset() { @@ -311,6 +314,13 @@ func (x *NodeSnapshot) GetEngines() map[string]*EngineAvailability { return nil } +func (x *NodeSnapshot) GetActiveLeaseIds() []string { + if x != nil { + return x.ActiveLeaseIds + } + return nil +} + // ReserveSeat asks an executor to hold a seat for a job. Zero secrets ride this // message; the raw pipeline/workflow are carried as JSON since processPipeline // already round-trips them through JSON. @@ -1171,11 +1181,12 @@ const file_spindle_mill_v1_mill_proto_rawDesc = "" + "\fcapabilities\x18\x03 \x03(\tR\fcapabilities\x1a7\n" + "\tLoadEntry\x12\x10\n" + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + - "\x05value\x18\x02 \x01(\x01R\x05value:\x028\x01\"\xe0\x01\n" + + "\x05value\x18\x02 \x01(\x01R\x05value:\x028\x01\"\x8a\x02\n" + "\fNodeSnapshot\x12\x17\n" + "\anode_id\x18\x01 \x01(\tR\x06nodeId\x12\x10\n" + "\x03seq\x18\x02 \x01(\x04R\x03seq\x12D\n" + - "\aengines\x18\x03 \x03(\v2*.spindle.mill.v1.NodeSnapshot.EnginesEntryR\aengines\x1a_\n" + + "\aengines\x18\x03 \x03(\v2*.spindle.mill.v1.NodeSnapshot.EnginesEntryR\aengines\x12(\n" + + "\x10active_lease_ids\x18\x04 \x03(\tR\x0eactiveLeaseIds\x1a_\n" + "\fEnginesEntry\x12\x10\n" + "\x03key\x18\x01 \x01(\tR\x03key\x129\n" + "\x05value\x18\x02 \x01(\v2#.spindle.mill.v1.EngineAvailabilityR\x05value:\x028\x01\"\x80\x02\n" + diff --git a/spindle/mill/proto/spindle/mill/v1/mill.proto b/spindle/mill/proto/spindle/mill/v1/mill.proto index f842d3ce..e7e4c908 100644 --- a/spindle/mill/proto/spindle/mill/v1/mill.proto +++ b/spindle/mill/proto/spindle/mill/v1/mill.proto @@ -42,6 +42,9 @@ message NodeSnapshot { string node_id = 1; uint64 seq = 2; map engines = 3; + // every lease the executor currently holds (reserved or running), so the + // mill can reconcile restored leases against executor reality. + repeated string active_lease_ids = 4; } // ReserveSeat asks an executor to hold a seat for a job. Zero secrets ride this diff --git a/spindle/mill/restore.go b/spindle/mill/restore.go new file mode 100644 index 00000000..205460d6 --- /dev/null +++ b/spindle/mill/restore.go @@ -0,0 +1,159 @@ +package mill + +import ( + "time" + + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/models" +) + +// mill_leases.state values. Committing persists as reserved: a crash mid-commit +// is indistinguishable from one before it, and reconciliation treats both the +// same (the executor either holds the lease or it doesn't). +const ( + leaseRowReserved = "reserved" + leaseRowRunning = "running" +) + +func (m *Mill) persistLease(lease *RemoteLease, state string) error { + if m.db == nil { + return nil + } + return m.db.SaveMillLease(db.MillLease{ + LeaseID: lease.id, + NodeID: lease.nodeID, + Engine: lease.engine, + Knot: lease.wid.Knot, + Rkey: lease.wid.Rkey, + Workflow: lease.wid.Name, + State: state, + }) +} + +func (m *Mill) deleteLeaseRow(leaseID string) { + if m.db == nil { + return + } + if err := m.db.DeleteMillLease(leaseID); err != nil { + m.l.Error("delete mill lease row", "lease", leaseID, "err", err) + } +} + +// RestoreState rebuilds leases and relay cursors persisted by a previous run. +// Restored leases are orphaned: no RunStep waits on them, so their terminal +// (or their executor's death) is authored directly as a status row. +func (m *Mill) RestoreState() error { + if m.db == nil { + return nil + } + cursors, err := m.db.ListExecutorCursors() + if err != nil { + return err + } + rows, err := m.db.ListMillLeases() + if err != nil { + return err + } + + m.mu.Lock() + for node, offset := range cursors { + m.nodeOffset[node] = offset + } + for _, r := range rows { + lease := newLease(r.LeaseID, r.NodeID, r.Engine) + lease.wid = models.WorkflowId{ + PipelineId: models.PipelineId{Knot: r.Knot, Rkey: r.Rkey}, + Name: r.Workflow, + } + lease.orphaned = true + if r.State == leaseRowRunning { + lease.state = leaseRunning + } + m.leases[r.LeaseID] = lease + } + restored := len(rows) + m.mu.Unlock() + + if restored > 0 { + m.l.Info("restored mill leases from previous run", "leases", restored, "cursors", len(cursors)) + time.AfterFunc(m.cfg.ReconnectGrace, m.sweepUnclaimedOrphans) + } + return nil +} + +func (m *Mill) sweepUnclaimedOrphans() { + m.mu.Lock() + var unclaimed []*RemoteLease + for _, lease := range m.leases { + if !lease.orphaned { + continue + } + if sess := m.sessions[lease.nodeID]; sess == nil || sess.disconnected { + unclaimed = append(unclaimed, lease) + } + } + m.mu.Unlock() + + reason := "executor did not reconnect after mill restart" + for _, lease := range unclaimed { + m.l.Warn("failing unclaimed restored lease", "lease", lease.id, "node", lease.nodeID) + if err := m.finishOrphan(lease, string(models.StatusKindFailed), &reason, nil); err != nil { + m.l.Error("finish unclaimed restored lease", "lease", lease.id, "err", err) + } + } +} + +// reconcileOrphans fails restored leases the executor no longer holds. Only +// orphaned leases are touched, so an in-flight bid can never race this. +func (m *Mill) reconcileOrphans(nodeID string, activeLeaseIDs []string) error { + active := make(map[string]struct{}, len(activeLeaseIDs)) + for _, id := range activeLeaseIDs { + active[id] = struct{}{} + } + + m.mu.Lock() + var gone []*RemoteLease + for _, lease := range m.leases { + if lease.orphaned && lease.nodeID == nodeID { + if _, ok := active[lease.id]; !ok { + gone = append(gone, lease) + } + } + } + m.mu.Unlock() + + reason := "executor no longer holds lease after mill restart" + for _, lease := range gone { + m.l.Warn("failing restored lease dropped by executor", "lease", lease.id, "node", nodeID) + if err := m.finishOrphan(lease, string(models.StatusKindFailed), &reason, nil); err != nil { + return err + } + } + return nil +} + +// finishOrphan authors the terminal status row (there is no RunStep waiter to +// do it) and unwinds the lease. markDone makes duplicates no-ops. +func (m *Mill) finishOrphan(lease *RemoteLease, status string, errMsg *string, exitCode *int64) error { + lease.finishMu.Lock() + defer lease.finishMu.Unlock() + if lease.getState() == leaseDone { + return nil + } + if m.db != nil { + if err := m.db.CompleteOrphanMillLease( + lease.id, + string(lease.wid.PipelineId.AtUri()), + lease.wid.Name, + status, + errMsg, + exitCode, + m.n, + ); err != nil { + return err + } + } + lease.markDone() + m.cleanupLease(lease) + return nil +} diff --git a/spindle/mill/restore_test.go b/spindle/mill/restore_test.go new file mode 100644 index 00000000..7b5f3615 --- /dev/null +++ b/spindle/mill/restore_test.go @@ -0,0 +1,287 @@ +package mill + +import ( + "context" + "path/filepath" + "testing" + "time" + + "tangled.org/core/notifier" + "tangled.org/core/spindle/db" + millv1 "tangled.org/core/spindle/mill/proto/gen" + "tangled.org/core/spindle/models" +) + +func restoreTestMill(t *testing.T, cfg Config) (*Mill, *db.DB) { + t.Helper() + bdb, err := db.Make(context.Background(), filepath.Join(t.TempDir(), "mill.db")) + if err != nil { + t.Fatalf("db.Make: %v", err) + } + t.Cleanup(func() { bdb.Close() }) + n := notifier.New() + m := New(discardLogger(), cfg) + m.Attach(bdb, &n) + return m, bdb +} + +func restoredMill(t *testing.T, bdb *db.DB, cfg Config) *Mill { + t.Helper() + n := notifier.New() + m := New(discardLogger(), cfg) + m.Attach(bdb, &n) + if err := m.RestoreState(); err != nil { + t.Fatalf("RestoreState: %v", err) + } + return m +} + +func TestRestoreStateRebuildsLeasesAndCursors(t *testing.T) { + _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) + + if err := bdb.SaveMillLease(db.MillLease{ + LeaseID: "lease-1", NodeID: "node-1", Engine: "dummy", + Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: leaseRowRunning, + }); err != nil { + t.Fatalf("SaveMillLease: %v", err) + } + if err := bdb.SetExecutorCursor("node-1", 7); err != nil { + t.Fatalf("SetExecutorCursor: %v", err) + } + + m2 := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute}) + + m2.mu.Lock() + lease := m2.leases["lease-1"] + offset := m2.nodeOffset["node-1"] + m2.mu.Unlock() + + if lease == nil { + t.Fatal("restored mill has no lease-1") + } + if !lease.orphaned { + t.Fatal("restored lease is not orphaned; a terminal would be delivered to a waiter that does not exist") + } + if lease.getState() != leaseRunning { + t.Fatalf("restored lease state = %v, want leaseRunning", lease.getState()) + } + wantWid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.example", Rkey: "rkey1"}, Name: "build"} + if lease.wid != wantWid { + t.Fatalf("restored lease wid = %+v, want %+v", lease.wid, wantWid) + } + if offset != 7 { + t.Fatalf("restored cursor = %d, want 7", offset) + } +} + +func TestOrphanTerminalAuthorsStatusRow(t *testing.T) { + _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) + if err := bdb.SaveMillLease(db.MillLease{ + LeaseID: "lease-1", NodeID: "node-1", Engine: "dummy", + Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: leaseRowRunning, + }); err != nil { + t.Fatalf("SaveMillLease: %v", err) + } + if err := bdb.SetExecutorCursor("node-1", 3); err != nil { + t.Fatalf("SetExecutorCursor: %v", err) + } + + m := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute}) + + sess := newSession("node-1", nopEncoder(), discardLogger()) + resume, ok := m.attachSession(sess) + if !ok { + t.Fatal("attachSession rejected the reconnecting executor") + } + if resume != 3 { + t.Fatalf("attachSession resume offset = %d, want restored cursor 3", resume) + } + + m.onAttemptResult(sess, &millv1.AttemptResult{ + Offset: 4, + LeaseId: "lease-1", + TerminalStatus: string(models.StatusKindSuccess), + }) + + wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.example", Rkey: "rkey1"}, Name: "build"} + st, err := bdb.GetStatus(wid) + if err != nil { + t.Fatalf("GetStatus after orphan terminal: %v", err) + } + if st.Status != string(models.StatusKindSuccess) { + t.Fatalf("orphan terminal authored status %q, want success", st.Status) + } + + m.mu.Lock() + _, still := m.leases["lease-1"] + m.mu.Unlock() + if still { + t.Fatal("finished orphan still in the lease map") + } + if rows, _ := bdb.ListMillLeases(); len(rows) != 0 { + t.Fatalf("finished orphan still persisted: %+v", rows) + } +} + +func TestSnapshotReconciliationFailsDroppedOrphans(t *testing.T) { + _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) + for _, l := range []db.MillLease{ + {LeaseID: "lease-kept", NodeID: "node-1", Engine: "dummy", Knot: "k", Rkey: "r1", Workflow: "w", State: leaseRowRunning}, + {LeaseID: "lease-gone", NodeID: "node-1", Engine: "dummy", Knot: "k", Rkey: "r2", Workflow: "w", State: leaseRowRunning}, + } { + if err := bdb.SaveMillLease(l); err != nil { + t.Fatalf("SaveMillLease(%s): %v", l.LeaseID, err) + } + } + + m := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute}) + sess := newSession("node-1", nopEncoder(), discardLogger()) + m.attachSession(sess) + + m.onSnapshot(sess, &millv1.NodeSnapshot{ + NodeId: "node-1", + ActiveLeaseIds: []string{"lease-kept"}, + }) + + m.mu.Lock() + _, kept := m.leases["lease-kept"] + _, gone := m.leases["lease-gone"] + m.mu.Unlock() + if !kept { + t.Fatal("reconciliation dropped a lease the executor still holds") + } + if gone { + t.Fatal("reconciliation kept a lease the executor no longer holds") + } + + st, err := bdb.GetStatus(models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r2"}, Name: "w"}) + if err != nil { + t.Fatalf("GetStatus for dropped orphan: %v", err) + } + if st.Status != string(models.StatusKindFailed) { + t.Fatalf("dropped orphan authored status %q, want failed", st.Status) + } +} + +func TestSweepFailsOrphansOfAbsentExecutors(t *testing.T) { + _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) + if err := bdb.SaveMillLease(db.MillLease{ + LeaseID: "lease-1", NodeID: "node-absent", Engine: "dummy", + Knot: "k", Rkey: "r1", Workflow: "w", State: leaseRowReserved, + }); err != nil { + t.Fatalf("SaveMillLease: %v", err) + } + + m := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute}) + m.sweepUnclaimedOrphans() + + m.mu.Lock() + _, still := m.leases["lease-1"] + m.mu.Unlock() + if still { + t.Fatal("sweep kept an orphan whose executor never reconnected") + } + st, err := bdb.GetStatus(models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r1"}, Name: "w"}) + if err != nil { + t.Fatalf("GetStatus after sweep: %v", err) + } + if st.Status != string(models.StatusKindFailed) { + t.Fatalf("sweep authored status %q, want failed", st.Status) + } + if rows, _ := bdb.ListMillLeases(); len(rows) != 0 { + t.Fatalf("swept orphan still persisted: %+v", rows) + } +} + +func TestAckOffsetPersistsCursor(t *testing.T) { + m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) + sess := newSession("node-1", nopEncoder(), discardLogger()) + m.attachSession(sess) + + m.ackOffset(sess, 12) + + cursors, err := bdb.ListExecutorCursors() + if err != nil { + t.Fatalf("ListExecutorCursors: %v", err) + } + if cursors["node-1"] != 12 { + t.Fatalf("persisted cursor = %d, want 12", cursors["node-1"]) + } +} + +func TestOrphanTerminalFailureKeepsLeaseAndOffsetRetryable(t *testing.T) { + _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) + if err := bdb.SaveMillLease(db.MillLease{ + LeaseID: "lease-1", NodeID: "node-1", Engine: "dummy", + Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: leaseRowRunning, + }); err != nil { + t.Fatalf("SaveMillLease: %v", err) + } + if err := bdb.SetExecutorCursor("node-1", 3); err != nil { + t.Fatalf("SetExecutorCursor: %v", err) + } + if _, err := bdb.Exec(` + create trigger reject_orphan_lease_delete + before delete on mill_leases + begin + select raise(abort, 'forced delete failure'); + end + `); err != nil { + t.Fatalf("create failure trigger: %v", err) + } + + m := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute}) + sess := newSession("node-1", nopEncoder(), discardLogger()) + if _, ok := m.attachSession(sess); !ok { + t.Fatal("attachSession rejected reconnect") + } + result := &millv1.AttemptResult{ + Offset: 4, + LeaseId: "lease-1", + TerminalStatus: string(models.StatusKindSuccess), + } + if err := m.onAttemptResult(sess, result); err == nil { + t.Fatal("orphan terminal relay succeeded despite forced transaction failure") + } + + m.mu.Lock() + lease := m.leases["lease-1"] + offset := m.nodeOffset["node-1"] + m.mu.Unlock() + if lease == nil || lease.getState() == leaseDone { + t.Fatal("failed orphan completion made the in-memory lease unretryable") + } + if offset != 3 { + t.Fatalf("in-memory relay offset = %d, want 3", offset) + } + if rows, err := bdb.ListMillLeases(); err != nil || len(rows) != 1 { + t.Fatalf("durable leases after transaction rollback = %+v, err = %v; want retained lease", rows, err) + } + var events int + if err := bdb.QueryRow(`select count(*) from events`).Scan(&events); err != nil { + t.Fatalf("count events: %v", err) + } + if events != 0 { + t.Fatalf("terminal events after transaction rollback = %d, want 0", events) + } + + if _, err := bdb.Exec(`drop trigger reject_orphan_lease_delete`); err != nil { + t.Fatalf("drop failure trigger: %v", err) + } + if err := m.onAttemptResult(sess, result); err != nil { + t.Fatalf("retry orphan terminal: %v", err) + } + m.mu.Lock() + _, still := m.leases["lease-1"] + offset = m.nodeOffset["node-1"] + m.mu.Unlock() + if still { + t.Fatal("successful orphan completion retained in-memory lease") + } + if offset != 4 { + t.Fatalf("in-memory relay offset after retry = %d, want 4", offset) + } + if rows, err := bdb.ListMillLeases(); err != nil || len(rows) != 0 { + t.Fatalf("durable leases after successful retry = %+v, err = %v; want none", rows, err) + } +} diff --git a/spindle/mill/session.go b/spindle/mill/session.go index 710575ed..503d5277 100644 --- a/spindle/mill/session.go +++ b/spindle/mill/session.go @@ -120,35 +120,37 @@ func (s *millSession) request(ctx context.Context, leaseID string, msg *millprot } // readLoop demuxes incoming frames until the decoder errors (connection gone). -func (s *millSession) readLoop(m *Mill, dec *millproto.Decoder) { +func (s *millSession) readLoop(m *Mill, dec *millproto.Decoder) error { for { msg, err := dec.Decode() if err != nil { - s.l.Debug("session read ended", "node", s.nodeID, "err", err) - return + return err + } + if err := s.dispatch(m, msg); err != nil { + return err } - s.dispatch(m, msg) } } -func (s *millSession) dispatch(m *Mill, msg *millproto.Message) { +func (s *millSession) dispatch(m *Mill, msg *millproto.Message) error { if !m.touchSession(s) { - return + return errSessionClosed } switch { case msg.GetNodeSnapshot() != nil: - m.onSnapshot(s, msg.GetNodeSnapshot()) + return m.onSnapshot(s, msg.GetNodeSnapshot()) case msg.GetReserveResult() != nil: s.deliver(msg.GetReserveResult().GetLeaseId(), msg) case msg.GetCommitted() != nil: s.deliver(msg.GetCommitted().GetLeaseId(), msg) case msg.GetStatusEvent() != nil: - m.onStatusRelay(s, msg.GetStatusEvent()) + return m.onStatusRelay(s, msg.GetStatusEvent()) case msg.GetLogLine() != nil: - m.onLogRelay(s, msg.GetLogLine()) + return m.onLogRelay(s, msg.GetLogLine()) case msg.GetAttemptResult() != nil: - m.onAttemptResult(s, msg.GetAttemptResult()) + return m.onAttemptResult(s, msg.GetAttemptResult()) default: s.l.Warn("session received unexpected message", "node", s.nodeID) } + return nil } diff --git a/spindle/server.go b/spindle/server.go index 486166ea..b28a3484 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -423,6 +423,9 @@ func Run(ctx context.Context) error { if err := m.SeedBootstrapToken(); err != nil { return fmt.Errorf("seeding mill bootstrap token: %w", err) } + if err := m.RestoreState(); err != nil { + return fmt.Errorf("restoring mill state: %w", err) + } } if cfg.Role == config.RoleExecutor { s.exec = executor.New(cfg, engines, s.DB(), s.Notifier(), log.SubLogger(logger, "executor"))