This repository has no description
1package mill
2
3import (
4 "context"
5 "path/filepath"
6 "testing"
7 "time"
8
9 "tangled.org/core/notifier"
10 "tangled.org/core/spindle/db"
11 millproto "tangled.org/core/spindle/mill/proto"
12 millv1 "tangled.org/core/spindle/mill/proto/gen"
13 "tangled.org/core/spindle/models"
14)
15
16func restoreTestMill(t *testing.T, cfg Config) (*Mill, *db.DB) {
17 t.Helper()
18 bdb, err := db.Make(context.Background(), filepath.Join(t.TempDir(), "mill.db"))
19 if err != nil {
20 t.Fatalf("db.Make: %v", err)
21 }
22 t.Cleanup(func() { bdb.Close() })
23 n := notifier.New()
24 m := New(discardLogger(), cfg)
25 m.Attach(bdb, &n)
26 return m, bdb
27}
28
29func restoredMill(t *testing.T, bdb *db.DB, cfg Config) *Mill {
30 t.Helper()
31 n := notifier.New()
32 m := New(discardLogger(), cfg)
33 m.Attach(bdb, &n)
34 if err := m.RestoreState(); err != nil {
35 t.Fatalf("RestoreState: %v", err)
36 }
37 return m
38}
39
40func TestRestoreStateRebuildsLeasesAndCursors(t *testing.T) {
41 _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
42
43 if err := bdb.SaveMillLease(db.MillLease{
44 LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy",
45 Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: leaseRowRunning,
46 }); err != nil {
47 t.Fatalf("SaveMillLease: %v", err)
48 }
49 if err := bdb.ApplyEventBatch(nil, func(tx *db.EventBatchTx) error { return tx.AdvanceCursor("node-1", "inc-1", 7) }); err != nil {
50 t.Fatalf("AdvanceCursor: %v", err)
51 }
52
53 m2 := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute})
54
55 m2.mu.Lock()
56 lease := m2.leases["lease-1"]
57 seqno := m2.nodeSeqno["node-1/inc-1"]
58 m2.mu.Unlock()
59
60 if lease == nil {
61 t.Fatal("restored mill has no lease-1")
62 }
63 if !lease.orphaned {
64 t.Fatal("restored lease is not orphaned; a terminal would be delivered to a waiter that does not exist")
65 }
66 if lease.getState() != leaseRunning {
67 t.Fatalf("restored lease state = %v, want leaseRunning", lease.getState())
68 }
69 wantWid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.example", Rkey: "rkey1"}, Name: "build"}
70 if lease.wid != wantWid {
71 t.Fatalf("restored lease wid = %+v, want %+v", lease.wid, wantWid)
72 }
73 if seqno != 7 {
74 t.Fatalf("restored cursor = %d, want 7", seqno)
75 }
76}
77
78func TestOrphanTerminalAuthorsStatusRow(t *testing.T) {
79 _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
80 if err := bdb.SaveMillLease(db.MillLease{
81 LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy",
82 Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: leaseRowRunning,
83 }); err != nil {
84 t.Fatalf("SaveMillLease: %v", err)
85 }
86 if err := bdb.ApplyEventBatch(nil, func(tx *db.EventBatchTx) error { return tx.AdvanceCursor("node-1", "inc-1", 3) }); err != nil {
87 t.Fatalf("AdvanceCursor: %v", err)
88 }
89
90 m := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute})
91
92 sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger())
93 resume, ok := m.attachSession(sess)
94 if !ok {
95 t.Fatal("attachSession rejected the reconnecting executor")
96 }
97 if resume != 3 {
98 t.Fatalf("attachSession resume seqno = %d, want restored cursor 3", resume)
99 }
100
101 _ = m.onEventBatch(sess, &millv1.EventBatch{
102 Epoch: sess.epoch,
103 Events: []*millv1.Event{
104 {
105 Seqno: 4,
106 LeaseId: "lease-1",
107 Payload: &millv1.Event_AttemptResult{
108 AttemptResult: &millv1.AttemptResult{
109 Status: millv1.TerminalStatus_SUCCESS,
110 },
111 },
112 },
113 },
114 })
115
116 wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.example", Rkey: "rkey1"}, Name: "build"}
117 st, err := bdb.GetStatus(wid)
118 if err != nil {
119 t.Fatalf("GetStatus after orphan terminal: %v", err)
120 }
121 if st.Status != string(models.StatusKindSuccess) {
122 t.Fatalf("orphan terminal authored status %q, want success", st.Status)
123 }
124
125 m.mu.Lock()
126 _, still := m.leases["lease-1"]
127 m.mu.Unlock()
128 if still {
129 t.Fatal("finished orphan still in the lease map")
130 }
131 if rows, _ := bdb.ListMillLeases(); len(rows) != 0 {
132 t.Fatalf("finished orphan still persisted: %+v", rows)
133 }
134}
135
136func TestSnapshotReconciliationFailsDroppedOrphans(t *testing.T) {
137 _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
138 for _, l := range []db.MillLease{
139 {LeaseID: "lease-kept", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", Knot: "k", Rkey: "r1", Workflow: "w", State: leaseRowRunning},
140 {LeaseID: "lease-gone", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", Knot: "k", Rkey: "r2", Workflow: "w", State: leaseRowRunning},
141 } {
142 if err := bdb.SaveMillLease(l); err != nil {
143 t.Fatalf("SaveMillLease(%s): %v", l.LeaseID, err)
144 }
145 }
146
147 m := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute})
148
149 sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger())
150 m.attachSession(sess)
151
152 m.onSnapshot(sess, &millv1.NodeSnapshot{
153 Seqno: 1,
154 ActiveLeaseIds: []string{"lease-kept"},
155 })
156
157 m.mu.Lock()
158 _, kept := m.leases["lease-kept"]
159 _, gone := m.leases["lease-gone"]
160 m.mu.Unlock()
161 if !kept {
162 t.Fatal("reconciliation dropped a lease the executor still holds")
163 }
164 if gone {
165 t.Fatal("reconciliation kept a lease the executor no longer holds")
166 }
167
168 st, err := bdb.GetStatus(models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r2"}, Name: "w"})
169 if err != nil {
170 t.Fatalf("GetStatus for dropped orphan: %v", err)
171 }
172 if st.Status != string(models.StatusKindFailed) {
173 t.Fatalf("dropped orphan authored status %q, want failed", st.Status)
174 }
175}
176func TestSnapshotReconciliationPreservesRequestedCancellation(t *testing.T) {
177 m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
178 lease := newLease("lease-1", "node-1", "inc-old", "dummy")
179 lease.wid = models.WorkflowId{
180 PipelineId: models.PipelineId{Knot: "k", Rkey: "r1"},
181 Name: "w",
182 }
183 lease.setState(leaseRunning)
184 lease.requestCancel("workflow destroyed")
185 if err := m.persistLease(lease, leaseRowRunning); err != nil {
186 t.Fatalf("persistLease: %v", err)
187 }
188 m.mu.Lock()
189 m.leases[lease.id] = lease
190 m.mu.Unlock()
191
192 sess := newSession("node-1", "inc-new", nil, nopEncoder(), discardLogger())
193 m.attachSession(sess)
194 if err := m.onSnapshot(sess, &millv1.NodeSnapshot{Seqno: 1}); err != nil {
195 t.Fatalf("onSnapshot: %v", err)
196 }
197
198 st, err := bdb.GetStatus(lease.wid)
199 if err != nil {
200 t.Fatalf("GetStatus: %v", err)
201 }
202 if st.Status != string(models.StatusKindCancelled) {
203 t.Fatalf("reconciled status = %q, want cancelled", st.Status)
204 }
205}
206
207func TestSnapshotReconciliationCancelsUnknownExecutorLease(t *testing.T) {
208 m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
209 cancelled := make(chan string, 1)
210 sess := newSession("node-1", "inc-1", nil, scriptedEncoder(func(msg *millproto.Message) error {
211 if cancel := msg.GetCancelAttempt(); cancel != nil {
212 cancelled <- cancel.GetLeaseId()
213 }
214 return nil
215 }), discardLogger())
216 m.attachSession(sess)
217
218 if err := m.onSnapshot(sess, &millv1.NodeSnapshot{
219 Seqno: 1,
220 ActiveLeaseIds: []string{"executor-only"},
221 }); err != nil {
222 t.Fatalf("onSnapshot: %v", err)
223 }
224 select {
225 case id := <-cancelled:
226 if id != "executor-only" {
227 t.Fatalf("cancelled lease = %q, want executor-only", id)
228 }
229 default:
230 t.Fatal("snapshot reconciliation left an executor-only lease running")
231 }
232}
233
234func TestSweepFailsOrphansOfAbsentExecutors(t *testing.T) {
235 _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
236 if err := bdb.SaveMillLease(db.MillLease{
237 LeaseID: "lease-1", NodeID: "node-absent", Epoch: "inc-absent", Engine: "dummy",
238 Knot: "k", Rkey: "r1", Workflow: "w", State: leaseRowReserved,
239 }); err != nil {
240 t.Fatalf("SaveMillLease: %v", err)
241 }
242
243 m := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute})
244 m.sweepUnclaimedOrphans()
245
246 m.mu.Lock()
247 _, still := m.leases["lease-1"]
248 m.mu.Unlock()
249 if still {
250 t.Fatal("sweep kept an orphan whose executor never reconnected")
251 }
252 st, err := bdb.GetStatus(models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r1"}, Name: "w"})
253 if err != nil {
254 t.Fatalf("GetStatus after sweep: %v", err)
255 }
256 if st.Status != string(models.StatusKindFailed) {
257 t.Fatalf("sweep authored status %q, want failed", st.Status)
258 }
259 if rows, _ := bdb.ListMillLeases(); len(rows) != 0 {
260 t.Fatalf("swept orphan still persisted: %+v", rows)
261 }
262}
263
264func TestAckSeqnoPersistsCursor(t *testing.T) {
265 m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
266 sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger())
267 m.attachSession(sess)
268
269 owned := newLease("lease-1", "node-1", "inc-1", "dummy")
270 owned.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
271 m.mu.Lock()
272 m.leases[owned.id] = owned
273 m.mu.Unlock()
274
275 err := m.onEventBatch(sess, &millv1.EventBatch{
276 Epoch: sess.epoch,
277 Events: []*millv1.Event{
278 {
279 Seqno: 1,
280 LeaseId: owned.id,
281 Payload: &millv1.Event_StatusEvent{
282 StatusEvent: &millv1.StatusEvent{
283 Status: millv1.NonterminalStatus_RUNNING,
284 },
285 },
286 },
287 },
288 })
289 if err != nil {
290 t.Fatalf("onEventBatch: %v", err)
291 }
292
293 cursors, err := bdb.ListExecutorCursors()
294 if err != nil {
295 t.Fatalf("ListExecutorCursors: %v", err)
296 }
297 if len(cursors) != 1 || cursors[0].AckedSeqno != 1 {
298 t.Fatalf("persisted cursor = %+v, want seqno 1", cursors)
299 }
300}
301
302func TestOrphanTerminalFailureKeepsLeaseAndSeqnoRetryable(t *testing.T) {
303 _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
304 if err := bdb.SaveMillLease(db.MillLease{
305 LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy",
306 Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: leaseRowRunning,
307 }); err != nil {
308 t.Fatalf("SaveMillLease: %v", err)
309 }
310 if err := bdb.ApplyEventBatch(nil, func(tx *db.EventBatchTx) error { return tx.AdvanceCursor("node-1", "inc-1", 3) }); err != nil {
311 t.Fatalf("AdvanceCursor: %v", err)
312 }
313 if _, err := bdb.Exec(`
314 create trigger reject_orphan_lease_delete
315 before delete on mill_leases
316 begin
317 select raise(abort, 'forced delete failure');
318 end
319 `); err != nil {
320 t.Fatalf("create failure trigger: %v", err)
321 }
322
323 m := restoredMill(t, bdb, Config{ReconnectGrace: time.Minute})
324 sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger())
325 if _, ok := m.attachSession(sess); !ok {
326 t.Fatal("attachSession rejected reconnect")
327 }
328 batch := &millv1.EventBatch{
329 Epoch: sess.epoch,
330 Events: []*millv1.Event{
331 {
332 Seqno: 4,
333 LeaseId: "lease-1",
334 Payload: &millv1.Event_AttemptResult{
335 AttemptResult: &millv1.AttemptResult{
336 Status: millv1.TerminalStatus_SUCCESS,
337 },
338 },
339 },
340 },
341 }
342 if err := m.onEventBatch(sess, batch); err == nil {
343 t.Fatal("orphan terminal stream succeeded despite forced transaction failure")
344 }
345
346 m.mu.Lock()
347 lease := m.leases["lease-1"]
348 seqno := m.nodeSeqno["node-1/inc-1"]
349 m.mu.Unlock()
350 if lease == nil || lease.getState() == leaseDone {
351 t.Fatal("failed orphan completion made the in-memory lease unretryable")
352 }
353 if seqno != 3 {
354 t.Fatalf("in-memory stream seqno = %d, want 3", seqno)
355 }
356 if rows, err := bdb.ListMillLeases(); err != nil || len(rows) != 1 {
357 t.Fatalf("durable leases after transaction rollback = %+v, err = %v; want retained lease", rows, err)
358 }
359 var events int
360 if err := bdb.QueryRow(`select count(*) from events`).Scan(&events); err != nil {
361 t.Fatalf("count events: %v", err)
362 }
363 if events != 0 {
364 t.Fatalf("terminal events after transaction rollback = %d, want 0", events)
365 }
366
367 if _, err := bdb.Exec(`drop trigger reject_orphan_lease_delete`); err != nil {
368 t.Fatalf("drop failure trigger: %v", err)
369 }
370 if err := m.onEventBatch(sess, batch); err != nil {
371 t.Fatalf("retry orphan terminal: %v", err)
372 }
373 m.mu.Lock()
374 _, still := m.leases["lease-1"]
375 seqno = m.nodeSeqno["node-1/inc-1"]
376 m.mu.Unlock()
377 if still {
378 t.Fatal("successful orphan completion retained in-memory lease")
379 }
380 if seqno != 4 {
381 t.Fatalf("in-memory stream seqno after retry = %d, want 4", seqno)
382 }
383 if rows, err := bdb.ListMillLeases(); err != nil || len(rows) != 0 {
384 t.Fatalf("durable leases after successful retry = %+v, err = %v; want none", rows, err)
385 }
386}