Monorepo for Tangled
0

Configure Feed

Select the types of activity you want to include in your feed.

spindle/mill: persist leases and executor cursors across restarts

Signed-off-by: dawn <dawn@tangled.org>

spindle/mill: harden placement and recovery

dawn (Jul 18, 2026, 9:59 AM +0300) cfd64c95 86dc2ccf

+1110 -66
+16
spindle/db/db.go
··· 130 130 expires_at text 131 131 ); 132 132 133 + create table if not exists mill_leases ( 134 + lease_id text primary key, 135 + node_id text not null, 136 + engine text not null, 137 + knot text not null, 138 + rkey text not null, 139 + workflow text not null, 140 + state text not null, 141 + created_at text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')) 142 + ); 143 + 144 + create table if not exists mill_executor_cursors ( 145 + node_id text primary key, 146 + acked_offset integer not null 147 + ); 148 + 133 149 create table if not exists migrations ( 134 150 id integer primary key autoincrement, 135 151 name text unique
+33
spindle/db/events.go
··· 94 94 return d.insertEvent(event, n) 95 95 } 96 96 97 + // CompleteOrphanMillLease makes the terminal event and lease deletion one 98 + // durable transition; an event without its matching deletion would be replayed. 99 + func (d *DB) CompleteOrphanMillLease( 100 + leaseID string, 101 + pipelineAtUri string, 102 + workflow string, 103 + status string, 104 + workflowError *string, 105 + exitCode *int64, 106 + n *notifier.Notifier, 107 + ) error { 108 + event, err := statusEvent(pipelineAtUri, workflow, status, workflowError, exitCode) 109 + if err != nil { 110 + return err 111 + } 112 + tx, err := d.Begin() 113 + if err != nil { 114 + return err 115 + } 116 + defer tx.Rollback() 117 + if err := eventstream.Insert(tx, event, nil); err != nil { 118 + return err 119 + } 120 + if _, err := tx.Exec(`delete from mill_leases where lease_id = ?`, leaseID); err != nil { 121 + return err 122 + } 123 + if err := tx.Commit(); err != nil { 124 + return err 125 + } 126 + n.NotifyAll() 127 + return nil 128 + } 129 + 97 130 func (d *DB) GetStatus(workflowId models.WorkflowId) (*tangled.PipelineStatus, error) { 98 131 pipelineAtUri := workflowId.PipelineId.AtUri() 99 132
+86
spindle/db/mill_state.go
··· 1 + package db 2 + 3 + // MillLease is the persisted view of a mill-side remote lease, enough to 4 + // rebuild the fencing token and its workflow identity after a mill restart. 5 + type MillLease struct { 6 + LeaseID string 7 + NodeID string 8 + Engine string 9 + Knot string 10 + Rkey string 11 + Workflow string 12 + State string 13 + } 14 + 15 + func (d *DB) SaveMillLease(l MillLease) error { 16 + _, err := d.Exec( 17 + `insert into mill_leases (lease_id, node_id, engine, knot, rkey, workflow, state) 18 + values (?, ?, ?, ?, ?, ?, ?) 19 + on conflict(lease_id) do update set state = excluded.state`, 20 + l.LeaseID, l.NodeID, l.Engine, l.Knot, l.Rkey, l.Workflow, l.State, 21 + ) 22 + return err 23 + } 24 + 25 + func (d *DB) DeleteMillLease(leaseID string) error { 26 + _, err := d.Exec(`delete from mill_leases where lease_id = ?`, leaseID) 27 + return err 28 + } 29 + 30 + func (d *DB) DeleteMillLeasesByNode(nodeID string) error { 31 + _, err := d.Exec(`delete from mill_leases where node_id = ?`, nodeID) 32 + return err 33 + } 34 + 35 + func (d *DB) ListMillLeases() ([]MillLease, error) { 36 + rows, err := d.Query(`select lease_id, node_id, engine, knot, rkey, workflow, state from mill_leases`) 37 + if err != nil { 38 + return nil, err 39 + } 40 + defer rows.Close() 41 + 42 + var leases []MillLease 43 + for rows.Next() { 44 + var l MillLease 45 + if err := rows.Scan(&l.LeaseID, &l.NodeID, &l.Engine, &l.Knot, &l.Rkey, &l.Workflow, &l.State); err != nil { 46 + return nil, err 47 + } 48 + leases = append(leases, l) 49 + } 50 + return leases, rows.Err() 51 + } 52 + 53 + // SetExecutorCursor persists the highest relay offset the mill has applied for 54 + // a node. 55 + func (d *DB) SetExecutorCursor(nodeID string, offset uint64) error { 56 + _, err := d.Exec( 57 + `insert into mill_executor_cursors (node_id, acked_offset) values (?, ?) 58 + on conflict(node_id) do update set acked_offset = excluded.acked_offset`, 59 + nodeID, offset, 60 + ) 61 + return err 62 + } 63 + 64 + func (d *DB) DeleteExecutorCursor(nodeID string) error { 65 + _, err := d.Exec(`delete from mill_executor_cursors where node_id = ?`, nodeID) 66 + return err 67 + } 68 + 69 + func (d *DB) ListExecutorCursors() (map[string]uint64, error) { 70 + rows, err := d.Query(`select node_id, acked_offset from mill_executor_cursors`) 71 + if err != nil { 72 + return nil, err 73 + } 74 + defer rows.Close() 75 + 76 + cursors := make(map[string]uint64) 77 + for rows.Next() { 78 + var node string 79 + var offset uint64 80 + if err := rows.Scan(&node, &offset); err != nil { 81 + return nil, err 82 + } 83 + cursors[node] = offset 84 + } 85 + return cursors, rows.Err() 86 + }
+181
spindle/db/mill_state_test.go
··· 1 + package db 2 + 3 + import ( 4 + "testing" 5 + 6 + "tangled.org/core/notifier" 7 + ) 8 + 9 + func TestMillLeaseRoundTrip(t *testing.T) { 10 + d := newTestDB(t) 11 + 12 + lease := MillLease{ 13 + LeaseID: "lease-1", 14 + NodeID: "node-1", 15 + Engine: "dummy", 16 + Knot: "knot.example", 17 + Rkey: "rkey1", 18 + Workflow: "build", 19 + State: "reserved", 20 + } 21 + if err := d.SaveMillLease(lease); err != nil { 22 + t.Fatalf("SaveMillLease: %v", err) 23 + } 24 + 25 + lease.State = "running" 26 + if err := d.SaveMillLease(lease); err != nil { 27 + t.Fatalf("SaveMillLease(transition): %v", err) 28 + } 29 + 30 + leases, err := d.ListMillLeases() 31 + if err != nil { 32 + t.Fatalf("ListMillLeases: %v", err) 33 + } 34 + if len(leases) != 1 { 35 + t.Fatalf("ListMillLeases returned %d leases, want 1 (state transition must replace, not duplicate)", len(leases)) 36 + } 37 + if leases[0] != lease { 38 + t.Fatalf("ListMillLeases[0] = %+v, want %+v", leases[0], lease) 39 + } 40 + 41 + if err := d.DeleteMillLease("lease-1"); err != nil { 42 + t.Fatalf("DeleteMillLease: %v", err) 43 + } 44 + if leases, _ = d.ListMillLeases(); len(leases) != 0 { 45 + t.Fatalf("lease survived deletion: %+v", leases) 46 + } 47 + } 48 + 49 + func TestDeleteMillLeasesByNode(t *testing.T) { 50 + d := newTestDB(t) 51 + 52 + for _, l := range []MillLease{ 53 + {LeaseID: "a", NodeID: "node-1", Engine: "dummy", Knot: "k", Rkey: "r1", Workflow: "w", State: "running"}, 54 + {LeaseID: "b", NodeID: "node-1", Engine: "dummy", Knot: "k", Rkey: "r2", Workflow: "w", State: "reserved"}, 55 + {LeaseID: "c", NodeID: "node-2", Engine: "dummy", Knot: "k", Rkey: "r3", Workflow: "w", State: "running"}, 56 + } { 57 + if err := d.SaveMillLease(l); err != nil { 58 + t.Fatalf("SaveMillLease(%s): %v", l.LeaseID, err) 59 + } 60 + } 61 + 62 + if err := d.DeleteMillLeasesByNode("node-1"); err != nil { 63 + t.Fatalf("DeleteMillLeasesByNode: %v", err) 64 + } 65 + leases, err := d.ListMillLeases() 66 + if err != nil { 67 + t.Fatalf("ListMillLeases: %v", err) 68 + } 69 + if len(leases) != 1 || leases[0].LeaseID != "c" { 70 + t.Fatalf("ListMillLeases = %+v, want only node-2's lease c", leases) 71 + } 72 + } 73 + 74 + func TestExecutorCursors(t *testing.T) { 75 + d := newTestDB(t) 76 + 77 + if err := d.SetExecutorCursor("node-1", 5); err != nil { 78 + t.Fatalf("SetExecutorCursor: %v", err) 79 + } 80 + if err := d.SetExecutorCursor("node-1", 9); err != nil { 81 + t.Fatalf("SetExecutorCursor(advance): %v", err) 82 + } 83 + if err := d.SetExecutorCursor("node-2", 1); err != nil { 84 + t.Fatalf("SetExecutorCursor(node-2): %v", err) 85 + } 86 + 87 + cursors, err := d.ListExecutorCursors() 88 + if err != nil { 89 + t.Fatalf("ListExecutorCursors: %v", err) 90 + } 91 + if len(cursors) != 2 || cursors["node-1"] != 9 || cursors["node-2"] != 1 { 92 + t.Fatalf("ListExecutorCursors = %v, want node-1:9 node-2:1", cursors) 93 + } 94 + 95 + if err := d.DeleteExecutorCursor("node-1"); err != nil { 96 + t.Fatalf("DeleteExecutorCursor: %v", err) 97 + } 98 + if cursors, _ = d.ListExecutorCursors(); len(cursors) != 1 || cursors["node-2"] != 1 { 99 + t.Fatalf("ListExecutorCursors after delete = %v, want only node-2:1", cursors) 100 + } 101 + } 102 + 103 + func TestCompleteOrphanMillLeaseIsAtomic(t *testing.T) { 104 + d := newTestDB(t) 105 + lease := MillLease{ 106 + LeaseID: "lease-1", NodeID: "node-1", Engine: "dummy", 107 + Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: "running", 108 + } 109 + if err := d.SaveMillLease(lease); err != nil { 110 + t.Fatalf("SaveMillLease: %v", err) 111 + } 112 + if _, err := d.Exec(` 113 + create trigger reject_mill_lease_delete 114 + before delete on mill_leases 115 + begin 116 + select raise(abort, 'forced delete failure'); 117 + end 118 + `); err != nil { 119 + t.Fatalf("create failure trigger: %v", err) 120 + } 121 + 122 + n := notifier.New() 123 + notifications := n.Subscribe() 124 + defer n.Unsubscribe(notifications) 125 + err := d.CompleteOrphanMillLease( 126 + "lease-1", 127 + "at://knot.example/sh.tangled.pipeline/rkey1", 128 + "build", 129 + "failed", 130 + nil, 131 + nil, 132 + &n, 133 + ) 134 + if err == nil { 135 + t.Fatal("CompleteOrphanMillLease succeeded despite forced lease deletion failure") 136 + } 137 + var eventCount int 138 + if err := d.QueryRow(`select count(*) from events`).Scan(&eventCount); err != nil { 139 + t.Fatalf("count events after rollback: %v", err) 140 + } 141 + if eventCount != 0 { 142 + t.Fatalf("terminal event count after rollback = %d, want 0", eventCount) 143 + } 144 + if leases, listErr := d.ListMillLeases(); listErr != nil || len(leases) != 1 { 145 + t.Fatalf("leases after rollback = %+v, err = %v; want original lease", leases, listErr) 146 + } 147 + select { 148 + case <-notifications: 149 + t.Fatal("rollback notified event subscribers") 150 + default: 151 + } 152 + 153 + if _, err := d.Exec(`drop trigger reject_mill_lease_delete`); err != nil { 154 + t.Fatalf("drop failure trigger: %v", err) 155 + } 156 + if err := d.CompleteOrphanMillLease( 157 + "lease-1", 158 + "at://knot.example/sh.tangled.pipeline/rkey1", 159 + "build", 160 + "failed", 161 + nil, 162 + nil, 163 + &n, 164 + ); err != nil { 165 + t.Fatalf("CompleteOrphanMillLease retry: %v", err) 166 + } 167 + if err := d.QueryRow(`select count(*) from events`).Scan(&eventCount); err != nil { 168 + t.Fatalf("count committed events: %v", err) 169 + } 170 + if eventCount != 1 { 171 + t.Fatalf("terminal event count after commit = %d, want 1", eventCount) 172 + } 173 + if leases, listErr := d.ListMillLeases(); listErr != nil || len(leases) != 0 { 174 + t.Fatalf("leases after commit = %+v, err = %v; want none", leases, listErr) 175 + } 176 + select { 177 + case <-notifications: 178 + default: 179 + t.Fatal("committed terminal event did not notify subscribers") 180 + } 181 + }
+7 -2
spindle/mill/executor/executor.go
··· 504 504 e.mu.Lock() 505 505 draining := e.draining 506 506 active := len(e.active) 507 + leaseIDs := make([]string, 0, len(e.active)) 508 + for id := range e.active { 509 + leaseIDs = append(leaseIDs, id) 510 + } 507 511 e.mu.Unlock() 508 512 509 513 load := 0.0 ··· 529 533 } 530 534 } 531 535 e.send(&millproto.Message{NodeSnapshot: &millv1.NodeSnapshot{ 532 - NodeId: e.nodeID, 533 - Engines: engines, 536 + NodeId: e.nodeID, 537 + Engines: engines, 538 + ActiveLeaseIds: leaseIDs, 534 539 }}) 535 540 } 536 541
+3 -1
spindle/mill/handler.go
··· 71 71 } 72 72 m.sessionReady(sess) 73 73 74 - sess.readLoop(m, dec) 74 + if err := sess.readLoop(m, dec); err != nil { 75 + m.l.Debug("session read ended", "node", sess.nodeID, "err", err) 76 + } 75 77 m.detachSession(sess) 76 78 } 77 79
+32
spindle/mill/lease.go
··· 40 40 nodeID string 41 41 engine string 42 42 wid models.WorkflowId // the job this lease carries; set once placed 43 + // restored from persistence after a mill restart: no RunStep waits on it, 44 + // so terminals and death are authored directly. Set before publication, 45 + // never mutated. 46 + orphaned bool 43 47 44 48 mu sync.Mutex 45 49 state leaseState 46 50 cancel bool 47 51 reason string 52 + released bool 48 53 terminal chan *millv1.AttemptResult // buffered(1); RunStep waits here 49 54 dead chan struct{} // closed when the executor is lost past grace 50 55 deadOnce sync.Once 56 + 57 + finishMu sync.Mutex 58 + cleanupOnce sync.Once 51 59 52 60 // log file for relayed lines, opened lazily and held open for the lease's 53 61 // life so we don't reopen per line. Guarded by logMu, closed on cleanup. ··· 99 107 return l.state 100 108 } 101 109 110 + // markDone CASes to done; exactly one caller wins. 111 + func (l *RemoteLease) markDone() bool { 112 + l.mu.Lock() 113 + defer l.mu.Unlock() 114 + if l.state == leaseDone { 115 + return false 116 + } 117 + l.state = leaseDone 118 + return true 119 + } 120 + 102 121 func (l *RemoteLease) requestCancel(reason string) cancelAction { 103 122 l.mu.Lock() 104 123 defer l.mu.Unlock() ··· 117 136 l.mu.Lock() 118 137 defer l.mu.Unlock() 119 138 return l.cancel, l.reason 139 + } 140 + 141 + func (l *RemoteLease) releaseState() (leaseState, bool) { 142 + l.mu.Lock() 143 + defer l.mu.Unlock() 144 + l.released = true 145 + return l.state, l.cancel 146 + } 147 + 148 + func (l *RemoteLease) cleanupReady() bool { 149 + l.mu.Lock() 150 + defer l.mu.Unlock() 151 + return l.released && l.state == leaseDone 120 152 } 121 153 122 154 func (l *RemoteLease) deliverCancelled(reason string) {
+86 -28
spindle/mill/mill.go
··· 223 223 delete(m.nodeOffset, sess.nodeID) 224 224 m.mu.Unlock() 225 225 226 + if m.db != nil { 227 + if err := m.db.DeleteExecutorCursor(sess.nodeID); err != nil { 228 + m.l.Error("delete executor cursor", "node", sess.nodeID, "err", err) 229 + } 230 + } 231 + 226 232 m.l.Warn("executor declared dead; failing its in-flight jobs", "node", sess.nodeID, "jobs", len(dead)) 233 + deadReason := "executor lost" 227 234 for _, lease := range dead { 228 - if cancelled, reason := lease.cancelRequested(); cancelled { 229 - lease.deliverCancelled(reason) 230 - } else { 231 - lease.markDead() 235 + switch { 236 + case lease.orphaned: 237 + if err := m.finishOrphan(lease, string(models.StatusKindFailed), &deadReason, nil); err != nil { 238 + m.l.Error("finish orphan after executor loss", "lease", lease.id, "err", err) 239 + } 240 + default: 241 + if cancelled, reason := lease.cancelRequested(); cancelled { 242 + lease.deliverCancelled(reason) 243 + if lease.cleanupReady() { 244 + m.cleanupLease(lease) 245 + } 246 + } else { 247 + lease.markDead() 248 + } 232 249 } 233 250 } 234 251 m.notifyChange() ··· 269 286 } 270 287 if lease != nil { 271 288 lease.wid = wid 289 + if err := m.persistLease(lease, leaseRowReserved); err != nil { 290 + m.releaseRemote(lease) 291 + return nil, fmt.Errorf("persist reserved mill lease: %w", err) 292 + } 272 293 m.mu.Lock() 273 294 m.leases[lease.id] = lease 274 295 if st, ok := wf.Data.(*millWorkflowState); ok && st != nil { ··· 554 575 return engine.ErrWorkflowFailed 555 576 } 556 577 lease.markRunning() 578 + if err := m.persistLease(lease, leaseRowRunning); err != nil { 579 + m.l.Error("persist running mill lease", "lease", lease.id, "err", err) 580 + } 557 581 if cancelled, reason := lease.cancelRequested(); cancelled { 558 582 m.sendCancel(sess, lease, reason) 559 583 } ··· 651 675 652 676 func (m *Mill) releaseSlot(s *millSlot) { 653 677 lease := s.lease 654 - if lease.getState() == leaseReserved { 678 + state, cancelled := lease.releaseState() 679 + if state == leaseReserved { 655 680 m.releaseRemote(lease) 681 + } else if cancelled && state != leaseDone { 682 + return 656 683 } 657 - lease.closeLog() 658 - 659 - m.mu.Lock() 660 - delete(m.leases, lease.id) 661 - m.mu.Unlock() 662 - 663 - m.notifyChange() 684 + m.cleanupLease(lease) 685 + } 686 + func (m *Mill) cleanupLease(lease *RemoteLease) { 687 + lease.cleanupOnce.Do(func() { 688 + lease.closeLog() 689 + m.mu.Lock() 690 + delete(m.leases, lease.id) 691 + m.mu.Unlock() 692 + m.deleteLeaseRow(lease.id) 693 + m.notifyChange() 694 + }) 664 695 } 665 696 666 697 func (m *Mill) releaseRemote(lease *RemoteLease) { ··· 696 727 return offset == current+1 697 728 } 698 729 699 - func (m *Mill) ackOffset(sess *millSession, offset uint64) { 730 + func (m *Mill) ackOffset(sess *millSession, offset uint64) error { 731 + if m.db != nil { 732 + if err := m.db.SetExecutorCursor(sess.nodeID, offset); err != nil { 733 + return fmt.Errorf("persist executor cursor: %w", err) 734 + } 735 + } 736 + 700 737 m.mu.Lock() 701 738 if offset > m.nodeOffset[sess.nodeID] { 702 739 m.nodeOffset[sess.nodeID] = offset 703 740 } 704 741 m.mu.Unlock() 705 742 706 - _ = sess.send(&millproto.Message{Ack: &millv1.Ack{UpToOffset: offset}}) 743 + if err := sess.send(&millproto.Message{Ack: &millv1.Ack{UpToOffset: offset}}); err != nil { 744 + return fmt.Errorf("send relay ack: %w", err) 745 + } 746 + return nil 707 747 } 708 748 709 - func (m *Mill) processRelay(sess *millSession, offset uint64, kind string, apply func() error) { 749 + func (m *Mill) processRelay(sess *millSession, offset uint64, kind string, apply func() error) error { 710 750 if !m.shouldAcceptOffset(sess, offset) { 711 - return 751 + return nil 712 752 } 713 - // ack only after the side effect lands, or a transient write error makes 714 - // the message unreplayable. 715 753 if err := apply(); err != nil { 716 - m.l.Error("process relayed message failed", "kind", kind, "node", sess.nodeID, "offset", offset, "err", err) 717 - return 754 + return fmt.Errorf("process relayed %s at offset %d: %w", kind, offset, err) 718 755 } 719 - m.ackOffset(sess, offset) 756 + if err := m.ackOffset(sess, offset); err != nil { 757 + return fmt.Errorf("ack relayed %s at offset %d: %w", kind, offset, err) 758 + } 759 + return nil 720 760 } 721 761 722 - func (m *Mill) onSnapshot(sess *millSession, snap *millv1.NodeSnapshot) { 762 + func (m *Mill) onSnapshot(sess *millSession, snap *millv1.NodeSnapshot) error { 723 763 m.mu.Lock() 724 764 sess.snapshot = snap 725 765 m.mu.Unlock() 766 + if err := m.reconcileOrphans(sess.nodeID, snap.GetActiveLeaseIds()); err != nil { 767 + return err 768 + } 726 769 m.notifyChange() 770 + return nil 727 771 } 728 772 729 - func (m *Mill) onStatusRelay(sess *millSession, ev *millv1.StatusEvent) { 730 - m.processRelay(sess, ev.GetOffset(), "status", func() error { 773 + func (m *Mill) onStatusRelay(sess *millSession, ev *millv1.StatusEvent) error { 774 + return m.processRelay(sess, ev.GetOffset(), "status", func() error { 731 775 if m.db == nil { 732 776 return fmt.Errorf("mill db not attached") 733 777 } ··· 753 797 }) 754 798 } 755 799 756 - func (m *Mill) onLogRelay(sess *millSession, ll *millv1.LogLine) { 757 - m.processRelay(sess, ll.GetOffset(), "log", func() error { 800 + func (m *Mill) onLogRelay(sess *millSession, ll *millv1.LogLine) error { 801 + return m.processRelay(sess, ll.GetOffset(), "log", func() error { 758 802 m.mu.Lock() 759 803 lease := m.leases[ll.GetLeaseId()] 760 804 m.mu.Unlock() ··· 767 811 }) 768 812 } 769 813 770 - func (m *Mill) onAttemptResult(sess *millSession, ar *millv1.AttemptResult) { 771 - m.processRelay(sess, ar.GetOffset(), "attempt-result", func() error { 814 + func (m *Mill) onAttemptResult(sess *millSession, ar *millv1.AttemptResult) error { 815 + return m.processRelay(sess, ar.GetOffset(), "attempt-result", func() error { 772 816 m.mu.Lock() 773 817 lease := m.leases[ar.GetLeaseId()] 774 818 m.mu.Unlock() 775 819 if lease == nil || lease.nodeID != sess.nodeID { 776 820 return nil 777 821 } 822 + if lease.orphaned { 823 + var errMsg *string 824 + if e := ar.GetError(); e != "" { 825 + errMsg = &e 826 + } 827 + var exit *int64 828 + if c := ar.GetExitCode(); c != 0 { 829 + exit = &c 830 + } 831 + return m.finishOrphan(lease, ar.GetTerminalStatus(), errMsg, exit) 832 + } 778 833 lease.deliverTerminal(ar) 834 + if lease.cleanupReady() { 835 + m.cleanupLease(lease) 836 + } 779 837 return nil 780 838 }) 781 839 }
+183 -17
spindle/mill/mill_test.go
··· 1 1 package mill 2 2 3 3 import ( 4 + "bytes" 4 5 "context" 5 6 "errors" 6 7 "io" 7 8 "log/slog" 9 + "path/filepath" 8 10 "strings" 11 + "sync" 9 12 "testing" 10 13 "time" 11 14 ··· 445 448 } 446 449 } 447 450 448 - func TestSessionRequestCancelledBeforeRegistrationOrSend(t *testing.T) { 449 - sent := 0 450 - sess := newSession("node-1", scriptedEncoder(func(*millproto.Message) error { 451 - sent++ 451 + func TestCancelledRunningLeaseSurvivesReleaseForReconnectReplay(t *testing.T) { 452 + m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) 453 + wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} 454 + lease := newLease("lease-1", "node-1", "dummy") 455 + lease.wid = wid 456 + lease.setState(leaseRunning) 457 + if err := m.persistLease(lease, leaseRowRunning); err != nil { 458 + t.Fatalf("persistLease: %v", err) 459 + } 460 + m.mu.Lock() 461 + m.leases[lease.id] = lease 462 + m.mu.Unlock() 463 + 464 + m.destroy(wid) 465 + slot := &millSlot{fleet: m, lease: lease} 466 + slot.Release() 467 + 468 + m.mu.Lock() 469 + _, retained := m.leases[lease.id] 470 + m.mu.Unlock() 471 + if !retained { 472 + t.Fatal("slot release removed a cancellation-requested running lease before its terminal result") 473 + } 474 + if rows, err := bdb.ListMillLeases(); err != nil || len(rows) != 1 { 475 + t.Fatalf("durable leases after slot release = %+v, err = %v; want retained lease", rows, err) 476 + } 477 + 478 + var sentMu sync.Mutex 479 + var sent []*millproto.Message 480 + sess := newSession("node-1", scriptedEncoder(func(msg *millproto.Message) error { 481 + sentMu.Lock() 482 + sent = append(sent, msg) 483 + sentMu.Unlock() 452 484 return nil 453 485 }), discardLogger()) 454 - ctx, cancel := context.WithCancel(context.Background()) 455 - cancel() 456 - _, err := sess.request(ctx, "lease-1", &millproto.Message{ReleaseLease: &millv1.ReleaseLease{LeaseId: "lease-1"}}) 457 - if !errors.Is(err, context.Canceled) { 458 - t.Fatalf("request error = %v, want context.Canceled", err) 486 + if _, ok := m.attachSession(sess); !ok { 487 + t.Fatal("attachSession rejected reconnect") 488 + } 489 + m.sessionReady(sess) 490 + sentMu.Lock() 491 + var replayed bool 492 + for _, msg := range sent { 493 + if cancel := msg.GetCancelAttempt(); cancel != nil && cancel.GetLeaseId() == lease.id { 494 + replayed = true 495 + } 496 + } 497 + sentMu.Unlock() 498 + if !replayed { 499 + t.Fatal("reconnect did not replay CancelAttempt for retained lease") 500 + } 501 + 502 + if err := m.onAttemptResult(sess, &millv1.AttemptResult{ 503 + Offset: 1, 504 + LeaseId: lease.id, 505 + TerminalStatus: string(models.StatusKindCancelled), 506 + }); err != nil { 507 + t.Fatalf("onAttemptResult: %v", err) 508 + } 509 + m.mu.Lock() 510 + _, retained = m.leases[lease.id] 511 + m.mu.Unlock() 512 + if retained { 513 + t.Fatal("terminal result did not clean retained cancelled lease") 514 + } 515 + if rows, err := bdb.ListMillLeases(); err != nil || len(rows) != 0 { 516 + t.Fatalf("durable leases after terminal = %+v, err = %v; want none", rows, err) 517 + } 518 + slot.Release() 519 + } 520 + 521 + func TestRelaySideEffectFailureStopsReadLoopAtFailedOffset(t *testing.T) { 522 + m := New(discardLogger(), Config{LogDir: filepath.Join(t.TempDir(), "missing")}) 523 + lease := newLease("lease-1", "node-1", "dummy") 524 + lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} 525 + m.mu.Lock() 526 + m.leases[lease.id] = lease 527 + m.mu.Unlock() 528 + sess := newSession("node-1", nopEncoder(), discardLogger()) 529 + if _, ok := m.attachSession(sess); !ok { 530 + t.Fatal("attachSession rejected first session") 531 + } 532 + 533 + var frames bytes.Buffer 534 + enc := millproto.NewEncoder(&frames) 535 + if err := enc.Encode(&millproto.Message{LogLine: &millv1.LogLine{ 536 + Offset: 1, LeaseId: lease.id, RawJson: []byte(`{"line":"first"}`), 537 + }}); err != nil { 538 + t.Fatalf("encode failing log relay: %v", err) 539 + } 540 + if err := enc.Encode(&millproto.Message{AttemptResult: &millv1.AttemptResult{ 541 + Offset: 2, LeaseId: lease.id, TerminalStatus: string(models.StatusKindSuccess), 542 + }}); err != nil { 543 + t.Fatalf("encode later relay: %v", err) 544 + } 545 + 546 + err := sess.readLoop(m, millproto.NewDecoder(&frames)) 547 + if err == nil { 548 + t.Fatal("readLoop continued after relayed log side effect failed") 459 549 } 460 - if sent != 0 { 461 - t.Fatalf("request sent %d messages for an already-cancelled context, want 0", sent) 550 + m.mu.Lock() 551 + offset := m.nodeOffset[sess.nodeID] 552 + m.mu.Unlock() 553 + if offset != 0 { 554 + t.Fatalf("in-memory relay offset = %d, want 0 so reconnect retries offset 1", offset) 462 555 } 463 - sess.mu.Lock() 464 - pending := len(sess.pending) 465 - sess.mu.Unlock() 466 - if pending != 0 { 467 - t.Fatalf("request left %d pending waiters, want 0", pending) 556 + if lease.getState() == leaseDone { 557 + t.Fatal("readLoop applied offset 2 after offset 1 failed") 468 558 } 469 559 } 470 560 ··· 496 586 } 497 587 } 498 588 589 + func TestSessionRequestCancelledBeforeRegistrationOrSend(t *testing.T) { 590 + sent := 0 591 + sess := newSession("node-1", scriptedEncoder(func(*millproto.Message) error { 592 + sent++ 593 + return nil 594 + }), discardLogger()) 595 + ctx, cancel := context.WithCancel(context.Background()) 596 + cancel() 597 + _, err := sess.request(ctx, "lease-1", &millproto.Message{ReleaseLease: &millv1.ReleaseLease{LeaseId: "lease-1"}}) 598 + if !errors.Is(err, context.Canceled) { 599 + t.Fatalf("request error = %v, want context.Canceled", err) 600 + } 601 + if sent != 0 { 602 + t.Fatalf("request sent %d messages for an already-cancelled context, want 0", sent) 603 + } 604 + sess.mu.Lock() 605 + pending := len(sess.pending) 606 + sess.mu.Unlock() 607 + if pending != 0 { 608 + t.Fatalf("request left %d pending waiters, want 0", pending) 609 + } 610 + } 611 + 499 612 func TestAttachSessionReplacesSilentIncumbentButRejectsActiveDuplicate(t *testing.T) { 500 613 m := New(discardLogger(), Config{ReconnectGrace: time.Minute}) 501 614 old := newSession("node-1", nopEncoder(), discardLogger()) ··· 523 636 m.mu.Lock() 524 637 replacement.lastSeen = time.Now().Add(-2 * m.cfg.ReconnectGrace) 525 638 m.mu.Unlock() 526 - replacement.dispatch(m, &millproto.Message{NodeSnapshot: &millv1.NodeSnapshot{NodeId: "node-1"}}) 639 + if err := replacement.dispatch(m, &millproto.Message{NodeSnapshot: &millv1.NodeSnapshot{NodeId: "node-1"}}); err != nil { 640 + t.Fatalf("periodic snapshot dispatch: %v", err) 641 + } 527 642 if _, ok := m.attachSession(newSession("node-1", nopEncoder(), discardLogger())); ok { 528 643 t.Fatal("active replacement did not reject a duplicate session") 529 644 } 530 645 } 646 + 647 + func TestPlaceReleasesRemoteReservationWhenInitialPersistenceFails(t *testing.T) { 648 + m, bdb := restoreTestMill(t, Config{BidTimeout: time.Second}) 649 + if err := bdb.Close(); err != nil { 650 + t.Fatalf("close db: %v", err) 651 + } 652 + released := make(chan string, 1) 653 + var sess *millSession 654 + sess = addCandidateSession(t, m, "node-1", nil, 0, scriptedEncoder(func(msg *millproto.Message) error { 655 + switch { 656 + case msg.GetReserveSeat() != nil: 657 + leaseID := msg.GetReserveSeat().GetLeaseId() 658 + sess.deliver(leaseID, &millproto.Message{ReserveResult: &millv1.ReserveResult{ 659 + LeaseId: leaseID, 660 + Accepted: true, 661 + }}) 662 + case msg.GetReleaseLease() != nil: 663 + released <- msg.GetReleaseLease().GetLeaseId() 664 + } 665 + return nil 666 + })) 667 + 668 + ctx, cancel := context.WithTimeout(context.Background(), time.Second) 669 + defer cancel() 670 + slot, err := m.place( 671 + ctx, 672 + "dummy", 673 + models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}, 674 + testWorkflow("build"), 675 + ) 676 + if err == nil { 677 + t.Fatal("place succeeded after reserved lease persistence failed") 678 + } 679 + if slot != nil { 680 + t.Fatalf("place returned slot %T after persistence failure", slot) 681 + } 682 + select { 683 + case leaseID := <-released: 684 + if leaseID == "" { 685 + t.Fatal("ReleaseLease had empty lease id") 686 + } 687 + case <-time.After(time.Second): 688 + t.Fatal("persistence failure did not compensate with ReleaseLease") 689 + } 690 + m.mu.Lock() 691 + leases := len(m.leases) 692 + m.mu.Unlock() 693 + if leases != 0 { 694 + t.Fatalf("mill published %d leases after initial persistence failure, want 0", leases) 695 + } 696 + }
+19 -8
spindle/mill/proto/gen/mill.pb.go
··· 252 252 // NodeSnapshot is pushed on connect, periodically, and right after any state 253 253 // change (reserve, commit, terminal). 254 254 type NodeSnapshot struct { 255 - state protoimpl.MessageState `protogen:"open.v1"` 256 - NodeId string `protobuf:"bytes,1,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"` 257 - Seq uint64 `protobuf:"varint,2,opt,name=seq,proto3" json:"seq,omitempty"` 258 - 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"` 259 - unknownFields protoimpl.UnknownFields 260 - sizeCache protoimpl.SizeCache 255 + state protoimpl.MessageState `protogen:"open.v1"` 256 + NodeId string `protobuf:"bytes,1,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"` 257 + Seq uint64 `protobuf:"varint,2,opt,name=seq,proto3" json:"seq,omitempty"` 258 + 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"` 259 + // every lease the executor currently holds (reserved or running), so the 260 + // mill can reconcile restored leases against executor reality. 261 + ActiveLeaseIds []string `protobuf:"bytes,4,rep,name=active_lease_ids,json=activeLeaseIds,proto3" json:"active_lease_ids,omitempty"` 262 + unknownFields protoimpl.UnknownFields 263 + sizeCache protoimpl.SizeCache 261 264 } 262 265 263 266 func (x *NodeSnapshot) Reset() { ··· 307 310 func (x *NodeSnapshot) GetEngines() map[string]*EngineAvailability { 308 311 if x != nil { 309 312 return x.Engines 313 + } 314 + return nil 315 + } 316 + 317 + func (x *NodeSnapshot) GetActiveLeaseIds() []string { 318 + if x != nil { 319 + return x.ActiveLeaseIds 310 320 } 311 321 return nil 312 322 } ··· 1171 1181 "\fcapabilities\x18\x03 \x03(\tR\fcapabilities\x1a7\n" + 1172 1182 "\tLoadEntry\x12\x10\n" + 1173 1183 "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + 1174 - "\x05value\x18\x02 \x01(\x01R\x05value:\x028\x01\"\xe0\x01\n" + 1184 + "\x05value\x18\x02 \x01(\x01R\x05value:\x028\x01\"\x8a\x02\n" + 1175 1185 "\fNodeSnapshot\x12\x17\n" + 1176 1186 "\anode_id\x18\x01 \x01(\tR\x06nodeId\x12\x10\n" + 1177 1187 "\x03seq\x18\x02 \x01(\x04R\x03seq\x12D\n" + 1178 - "\aengines\x18\x03 \x03(\v2*.spindle.mill.v1.NodeSnapshot.EnginesEntryR\aengines\x1a_\n" + 1188 + "\aengines\x18\x03 \x03(\v2*.spindle.mill.v1.NodeSnapshot.EnginesEntryR\aengines\x12(\n" + 1189 + "\x10active_lease_ids\x18\x04 \x03(\tR\x0eactiveLeaseIds\x1a_\n" + 1179 1190 "\fEnginesEntry\x12\x10\n" + 1180 1191 "\x03key\x18\x01 \x01(\tR\x03key\x129\n" + 1181 1192 "\x05value\x18\x02 \x01(\v2#.spindle.mill.v1.EngineAvailabilityR\x05value:\x028\x01\"\x80\x02\n" +
+3
spindle/mill/proto/spindle/mill/v1/mill.proto
··· 42 42 string node_id = 1; 43 43 uint64 seq = 2; 44 44 map<string, EngineAvailability> engines = 3; 45 + // every lease the executor currently holds (reserved or running), so the 46 + // mill can reconcile restored leases against executor reality. 47 + repeated string active_lease_ids = 4; 45 48 } 46 49 47 50 // ReserveSeat asks an executor to hold a seat for a job. Zero secrets ride this
+159
spindle/mill/restore.go
··· 1 + package mill 2 + 3 + import ( 4 + "time" 5 + 6 + "tangled.org/core/spindle/db" 7 + "tangled.org/core/spindle/models" 8 + ) 9 + 10 + // mill_leases.state values. Committing persists as reserved: a crash mid-commit 11 + // is indistinguishable from one before it, and reconciliation treats both the 12 + // same (the executor either holds the lease or it doesn't). 13 + const ( 14 + leaseRowReserved = "reserved" 15 + leaseRowRunning = "running" 16 + ) 17 + 18 + func (m *Mill) persistLease(lease *RemoteLease, state string) error { 19 + if m.db == nil { 20 + return nil 21 + } 22 + return m.db.SaveMillLease(db.MillLease{ 23 + LeaseID: lease.id, 24 + NodeID: lease.nodeID, 25 + Engine: lease.engine, 26 + Knot: lease.wid.Knot, 27 + Rkey: lease.wid.Rkey, 28 + Workflow: lease.wid.Name, 29 + State: state, 30 + }) 31 + } 32 + 33 + func (m *Mill) deleteLeaseRow(leaseID string) { 34 + if m.db == nil { 35 + return 36 + } 37 + if err := m.db.DeleteMillLease(leaseID); err != nil { 38 + m.l.Error("delete mill lease row", "lease", leaseID, "err", err) 39 + } 40 + } 41 + 42 + // RestoreState rebuilds leases and relay cursors persisted by a previous run. 43 + // Restored leases are orphaned: no RunStep waits on them, so their terminal 44 + // (or their executor's death) is authored directly as a status row. 45 + func (m *Mill) RestoreState() error { 46 + if m.db == nil { 47 + return nil 48 + } 49 + cursors, err := m.db.ListExecutorCursors() 50 + if err != nil { 51 + return err 52 + } 53 + rows, err := m.db.ListMillLeases() 54 + if err != nil { 55 + return err 56 + } 57 + 58 + m.mu.Lock() 59 + for node, offset := range cursors { 60 + m.nodeOffset[node] = offset 61 + } 62 + for _, r := range rows { 63 + lease := newLease(r.LeaseID, r.NodeID, r.Engine) 64 + lease.wid = models.WorkflowId{ 65 + PipelineId: models.PipelineId{Knot: r.Knot, Rkey: r.Rkey}, 66 + Name: r.Workflow, 67 + } 68 + lease.orphaned = true 69 + if r.State == leaseRowRunning { 70 + lease.state = leaseRunning 71 + } 72 + m.leases[r.LeaseID] = lease 73 + } 74 + restored := len(rows) 75 + m.mu.Unlock() 76 + 77 + if restored > 0 { 78 + m.l.Info("restored mill leases from previous run", "leases", restored, "cursors", len(cursors)) 79 + time.AfterFunc(m.cfg.ReconnectGrace, m.sweepUnclaimedOrphans) 80 + } 81 + return nil 82 + } 83 + 84 + func (m *Mill) sweepUnclaimedOrphans() { 85 + m.mu.Lock() 86 + var unclaimed []*RemoteLease 87 + for _, lease := range m.leases { 88 + if !lease.orphaned { 89 + continue 90 + } 91 + if sess := m.sessions[lease.nodeID]; sess == nil || sess.disconnected { 92 + unclaimed = append(unclaimed, lease) 93 + } 94 + } 95 + m.mu.Unlock() 96 + 97 + reason := "executor did not reconnect after mill restart" 98 + for _, lease := range unclaimed { 99 + m.l.Warn("failing unclaimed restored lease", "lease", lease.id, "node", lease.nodeID) 100 + if err := m.finishOrphan(lease, string(models.StatusKindFailed), &reason, nil); err != nil { 101 + m.l.Error("finish unclaimed restored lease", "lease", lease.id, "err", err) 102 + } 103 + } 104 + } 105 + 106 + // reconcileOrphans fails restored leases the executor no longer holds. Only 107 + // orphaned leases are touched, so an in-flight bid can never race this. 108 + func (m *Mill) reconcileOrphans(nodeID string, activeLeaseIDs []string) error { 109 + active := make(map[string]struct{}, len(activeLeaseIDs)) 110 + for _, id := range activeLeaseIDs { 111 + active[id] = struct{}{} 112 + } 113 + 114 + m.mu.Lock() 115 + var gone []*RemoteLease 116 + for _, lease := range m.leases { 117 + if lease.orphaned && lease.nodeID == nodeID { 118 + if _, ok := active[lease.id]; !ok { 119 + gone = append(gone, lease) 120 + } 121 + } 122 + } 123 + m.mu.Unlock() 124 + 125 + reason := "executor no longer holds lease after mill restart" 126 + for _, lease := range gone { 127 + m.l.Warn("failing restored lease dropped by executor", "lease", lease.id, "node", nodeID) 128 + if err := m.finishOrphan(lease, string(models.StatusKindFailed), &reason, nil); err != nil { 129 + return err 130 + } 131 + } 132 + return nil 133 + } 134 + 135 + // finishOrphan authors the terminal status row (there is no RunStep waiter to 136 + // do it) and unwinds the lease. markDone makes duplicates no-ops. 137 + func (m *Mill) finishOrphan(lease *RemoteLease, status string, errMsg *string, exitCode *int64) error { 138 + lease.finishMu.Lock() 139 + defer lease.finishMu.Unlock() 140 + if lease.getState() == leaseDone { 141 + return nil 142 + } 143 + if m.db != nil { 144 + if err := m.db.CompleteOrphanMillLease( 145 + lease.id, 146 + string(lease.wid.PipelineId.AtUri()), 147 + lease.wid.Name, 148 + status, 149 + errMsg, 150 + exitCode, 151 + m.n, 152 + ); err != nil { 153 + return err 154 + } 155 + } 156 + lease.markDone() 157 + m.cleanupLease(lease) 158 + return nil 159 + }
+287
spindle/mill/restore_test.go
··· 1 + package mill 2 + 3 + import ( 4 + "context" 5 + "path/filepath" 6 + "testing" 7 + "time" 8 + 9 + "tangled.org/core/notifier" 10 + "tangled.org/core/spindle/db" 11 + millv1 "tangled.org/core/spindle/mill/proto/gen" 12 + "tangled.org/core/spindle/models" 13 + ) 14 + 15 + func restoreTestMill(t *testing.T, cfg Config) (*Mill, *db.DB) { 16 + t.Helper() 17 + bdb, err := db.Make(context.Background(), filepath.Join(t.TempDir(), "mill.db")) 18 + if err != nil { 19 + t.Fatalf("db.Make: %v", err) 20 + } 21 + t.Cleanup(func() { bdb.Close() }) 22 + n := notifier.New() 23 + m := New(discardLogger(), cfg) 24 + m.Attach(bdb, &n) 25 + return m, bdb 26 + } 27 + 28 + func restoredMill(t *testing.T, bdb *db.DB, cfg Config) *Mill { 29 + t.Helper() 30 + n := notifier.New() 31 + m := New(discardLogger(), cfg) 32 + m.Attach(bdb, &n) 33 + if err := m.RestoreState(); err != nil { 34 + t.Fatalf("RestoreState: %v", err) 35 + } 36 + return m 37 + } 38 + 39 + func TestRestoreStateRebuildsLeasesAndCursors(t *testing.T) { 40 + _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) 41 + 42 + if err := bdb.SaveMillLease(db.MillLease{ 43 + LeaseID: "lease-1", NodeID: "node-1", Engine: "dummy", 44 + Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: leaseRowRunning, 45 + }); err != nil { 46 + t.Fatalf("SaveMillLease: %v", err) 47 + } 48 + if err := bdb.SetExecutorCursor("node-1", 7); err != nil { 49 + t.Fatalf("SetExecutorCursor: %v", err) 50 + } 51 + 52 + m2 := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute}) 53 + 54 + m2.mu.Lock() 55 + lease := m2.leases["lease-1"] 56 + offset := m2.nodeOffset["node-1"] 57 + m2.mu.Unlock() 58 + 59 + if lease == nil { 60 + t.Fatal("restored mill has no lease-1") 61 + } 62 + if !lease.orphaned { 63 + t.Fatal("restored lease is not orphaned; a terminal would be delivered to a waiter that does not exist") 64 + } 65 + if lease.getState() != leaseRunning { 66 + t.Fatalf("restored lease state = %v, want leaseRunning", lease.getState()) 67 + } 68 + wantWid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.example", Rkey: "rkey1"}, Name: "build"} 69 + if lease.wid != wantWid { 70 + t.Fatalf("restored lease wid = %+v, want %+v", lease.wid, wantWid) 71 + } 72 + if offset != 7 { 73 + t.Fatalf("restored cursor = %d, want 7", offset) 74 + } 75 + } 76 + 77 + func TestOrphanTerminalAuthorsStatusRow(t *testing.T) { 78 + _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) 79 + if err := bdb.SaveMillLease(db.MillLease{ 80 + LeaseID: "lease-1", NodeID: "node-1", Engine: "dummy", 81 + Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: leaseRowRunning, 82 + }); err != nil { 83 + t.Fatalf("SaveMillLease: %v", err) 84 + } 85 + if err := bdb.SetExecutorCursor("node-1", 3); err != nil { 86 + t.Fatalf("SetExecutorCursor: %v", err) 87 + } 88 + 89 + m := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute}) 90 + 91 + sess := newSession("node-1", nopEncoder(), discardLogger()) 92 + resume, ok := m.attachSession(sess) 93 + if !ok { 94 + t.Fatal("attachSession rejected the reconnecting executor") 95 + } 96 + if resume != 3 { 97 + t.Fatalf("attachSession resume offset = %d, want restored cursor 3", resume) 98 + } 99 + 100 + m.onAttemptResult(sess, &millv1.AttemptResult{ 101 + Offset: 4, 102 + LeaseId: "lease-1", 103 + TerminalStatus: string(models.StatusKindSuccess), 104 + }) 105 + 106 + wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.example", Rkey: "rkey1"}, Name: "build"} 107 + st, err := bdb.GetStatus(wid) 108 + if err != nil { 109 + t.Fatalf("GetStatus after orphan terminal: %v", err) 110 + } 111 + if st.Status != string(models.StatusKindSuccess) { 112 + t.Fatalf("orphan terminal authored status %q, want success", st.Status) 113 + } 114 + 115 + m.mu.Lock() 116 + _, still := m.leases["lease-1"] 117 + m.mu.Unlock() 118 + if still { 119 + t.Fatal("finished orphan still in the lease map") 120 + } 121 + if rows, _ := bdb.ListMillLeases(); len(rows) != 0 { 122 + t.Fatalf("finished orphan still persisted: %+v", rows) 123 + } 124 + } 125 + 126 + func TestSnapshotReconciliationFailsDroppedOrphans(t *testing.T) { 127 + _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) 128 + for _, l := range []db.MillLease{ 129 + {LeaseID: "lease-kept", NodeID: "node-1", Engine: "dummy", Knot: "k", Rkey: "r1", Workflow: "w", State: leaseRowRunning}, 130 + {LeaseID: "lease-gone", NodeID: "node-1", Engine: "dummy", Knot: "k", Rkey: "r2", Workflow: "w", State: leaseRowRunning}, 131 + } { 132 + if err := bdb.SaveMillLease(l); err != nil { 133 + t.Fatalf("SaveMillLease(%s): %v", l.LeaseID, err) 134 + } 135 + } 136 + 137 + m := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute}) 138 + sess := newSession("node-1", nopEncoder(), discardLogger()) 139 + m.attachSession(sess) 140 + 141 + m.onSnapshot(sess, &millv1.NodeSnapshot{ 142 + NodeId: "node-1", 143 + ActiveLeaseIds: []string{"lease-kept"}, 144 + }) 145 + 146 + m.mu.Lock() 147 + _, kept := m.leases["lease-kept"] 148 + _, gone := m.leases["lease-gone"] 149 + m.mu.Unlock() 150 + if !kept { 151 + t.Fatal("reconciliation dropped a lease the executor still holds") 152 + } 153 + if gone { 154 + t.Fatal("reconciliation kept a lease the executor no longer holds") 155 + } 156 + 157 + st, err := bdb.GetStatus(models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r2"}, Name: "w"}) 158 + if err != nil { 159 + t.Fatalf("GetStatus for dropped orphan: %v", err) 160 + } 161 + if st.Status != string(models.StatusKindFailed) { 162 + t.Fatalf("dropped orphan authored status %q, want failed", st.Status) 163 + } 164 + } 165 + 166 + func TestSweepFailsOrphansOfAbsentExecutors(t *testing.T) { 167 + _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) 168 + if err := bdb.SaveMillLease(db.MillLease{ 169 + LeaseID: "lease-1", NodeID: "node-absent", Engine: "dummy", 170 + Knot: "k", Rkey: "r1", Workflow: "w", State: leaseRowReserved, 171 + }); err != nil { 172 + t.Fatalf("SaveMillLease: %v", err) 173 + } 174 + 175 + m := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute}) 176 + m.sweepUnclaimedOrphans() 177 + 178 + m.mu.Lock() 179 + _, still := m.leases["lease-1"] 180 + m.mu.Unlock() 181 + if still { 182 + t.Fatal("sweep kept an orphan whose executor never reconnected") 183 + } 184 + st, err := bdb.GetStatus(models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r1"}, Name: "w"}) 185 + if err != nil { 186 + t.Fatalf("GetStatus after sweep: %v", err) 187 + } 188 + if st.Status != string(models.StatusKindFailed) { 189 + t.Fatalf("sweep authored status %q, want failed", st.Status) 190 + } 191 + if rows, _ := bdb.ListMillLeases(); len(rows) != 0 { 192 + t.Fatalf("swept orphan still persisted: %+v", rows) 193 + } 194 + } 195 + 196 + func TestAckOffsetPersistsCursor(t *testing.T) { 197 + m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) 198 + sess := newSession("node-1", nopEncoder(), discardLogger()) 199 + m.attachSession(sess) 200 + 201 + m.ackOffset(sess, 12) 202 + 203 + cursors, err := bdb.ListExecutorCursors() 204 + if err != nil { 205 + t.Fatalf("ListExecutorCursors: %v", err) 206 + } 207 + if cursors["node-1"] != 12 { 208 + t.Fatalf("persisted cursor = %d, want 12", cursors["node-1"]) 209 + } 210 + } 211 + 212 + func TestOrphanTerminalFailureKeepsLeaseAndOffsetRetryable(t *testing.T) { 213 + _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) 214 + if err := bdb.SaveMillLease(db.MillLease{ 215 + LeaseID: "lease-1", NodeID: "node-1", Engine: "dummy", 216 + Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: leaseRowRunning, 217 + }); err != nil { 218 + t.Fatalf("SaveMillLease: %v", err) 219 + } 220 + if err := bdb.SetExecutorCursor("node-1", 3); err != nil { 221 + t.Fatalf("SetExecutorCursor: %v", err) 222 + } 223 + if _, err := bdb.Exec(` 224 + create trigger reject_orphan_lease_delete 225 + before delete on mill_leases 226 + begin 227 + select raise(abort, 'forced delete failure'); 228 + end 229 + `); err != nil { 230 + t.Fatalf("create failure trigger: %v", err) 231 + } 232 + 233 + m := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute}) 234 + sess := newSession("node-1", nopEncoder(), discardLogger()) 235 + if _, ok := m.attachSession(sess); !ok { 236 + t.Fatal("attachSession rejected reconnect") 237 + } 238 + result := &millv1.AttemptResult{ 239 + Offset: 4, 240 + LeaseId: "lease-1", 241 + TerminalStatus: string(models.StatusKindSuccess), 242 + } 243 + if err := m.onAttemptResult(sess, result); err == nil { 244 + t.Fatal("orphan terminal relay succeeded despite forced transaction failure") 245 + } 246 + 247 + m.mu.Lock() 248 + lease := m.leases["lease-1"] 249 + offset := m.nodeOffset["node-1"] 250 + m.mu.Unlock() 251 + if lease == nil || lease.getState() == leaseDone { 252 + t.Fatal("failed orphan completion made the in-memory lease unretryable") 253 + } 254 + if offset != 3 { 255 + t.Fatalf("in-memory relay offset = %d, want 3", offset) 256 + } 257 + if rows, err := bdb.ListMillLeases(); err != nil || len(rows) != 1 { 258 + t.Fatalf("durable leases after transaction rollback = %+v, err = %v; want retained lease", rows, err) 259 + } 260 + var events int 261 + if err := bdb.QueryRow(`select count(*) from events`).Scan(&events); err != nil { 262 + t.Fatalf("count events: %v", err) 263 + } 264 + if events != 0 { 265 + t.Fatalf("terminal events after transaction rollback = %d, want 0", events) 266 + } 267 + 268 + if _, err := bdb.Exec(`drop trigger reject_orphan_lease_delete`); err != nil { 269 + t.Fatalf("drop failure trigger: %v", err) 270 + } 271 + if err := m.onAttemptResult(sess, result); err != nil { 272 + t.Fatalf("retry orphan terminal: %v", err) 273 + } 274 + m.mu.Lock() 275 + _, still := m.leases["lease-1"] 276 + offset = m.nodeOffset["node-1"] 277 + m.mu.Unlock() 278 + if still { 279 + t.Fatal("successful orphan completion retained in-memory lease") 280 + } 281 + if offset != 4 { 282 + t.Fatalf("in-memory relay offset after retry = %d, want 4", offset) 283 + } 284 + if rows, err := bdb.ListMillLeases(); err != nil || len(rows) != 0 { 285 + t.Fatalf("durable leases after successful retry = %+v, err = %v; want none", rows, err) 286 + } 287 + }
+12 -10
spindle/mill/session.go
··· 120 120 } 121 121 122 122 // readLoop demuxes incoming frames until the decoder errors (connection gone). 123 - func (s *millSession) readLoop(m *Mill, dec *millproto.Decoder) { 123 + func (s *millSession) readLoop(m *Mill, dec *millproto.Decoder) error { 124 124 for { 125 125 msg, err := dec.Decode() 126 126 if err != nil { 127 - s.l.Debug("session read ended", "node", s.nodeID, "err", err) 128 - return 127 + return err 128 + } 129 + if err := s.dispatch(m, msg); err != nil { 130 + return err 129 131 } 130 - s.dispatch(m, msg) 131 132 } 132 133 } 133 134 134 - func (s *millSession) dispatch(m *Mill, msg *millproto.Message) { 135 + func (s *millSession) dispatch(m *Mill, msg *millproto.Message) error { 135 136 if !m.touchSession(s) { 136 - return 137 + return errSessionClosed 137 138 } 138 139 switch { 139 140 case msg.GetNodeSnapshot() != nil: 140 - m.onSnapshot(s, msg.GetNodeSnapshot()) 141 + return m.onSnapshot(s, msg.GetNodeSnapshot()) 141 142 case msg.GetReserveResult() != nil: 142 143 s.deliver(msg.GetReserveResult().GetLeaseId(), msg) 143 144 case msg.GetCommitted() != nil: 144 145 s.deliver(msg.GetCommitted().GetLeaseId(), msg) 145 146 case msg.GetStatusEvent() != nil: 146 - m.onStatusRelay(s, msg.GetStatusEvent()) 147 + return m.onStatusRelay(s, msg.GetStatusEvent()) 147 148 case msg.GetLogLine() != nil: 148 - m.onLogRelay(s, msg.GetLogLine()) 149 + return m.onLogRelay(s, msg.GetLogLine()) 149 150 case msg.GetAttemptResult() != nil: 150 - m.onAttemptResult(s, msg.GetAttemptResult()) 151 + return m.onAttemptResult(s, msg.GetAttemptResult()) 151 152 default: 152 153 s.l.Warn("session received unexpected message", "node", s.nodeID) 153 154 } 155 + return nil 154 156 }
+3
spindle/server.go
··· 423 423 if err := m.SeedBootstrapToken(); err != nil { 424 424 return fmt.Errorf("seeding mill bootstrap token: %w", err) 425 425 } 426 + if err := m.RestoreState(); err != nil { 427 + return fmt.Errorf("restoring mill state: %w", err) 428 + } 426 429 } 427 430 if cfg.Role == config.RoleExecutor { 428 431 s.exec = executor.New(cfg, engines, s.DB(), s.Notifier(), log.SubLogger(logger, "executor"))