This repository has no description
1package db
2
3import (
4 "context"
5 "fmt"
6 "path/filepath"
7 "testing"
8
9 "tangled.org/core/notifier"
10)
11
12func TestMillLeaseRoundTrip(t *testing.T) {
13 d := newTestDB(t)
14
15 lease := MillLease{
16 LeaseID: "lease-1",
17 NodeID: "node-1",
18 Epoch: "inc-1",
19 Engine: "dummy",
20 Knot: "knot.example",
21 Rkey: "rkey1",
22 Workflow: "build",
23 State: "reserved",
24 }
25 if err := d.SaveMillLease(lease); err != nil {
26 t.Fatalf("SaveMillLease: %v", err)
27 }
28
29 lease.State = "running"
30 if err := d.SaveMillLease(lease); err != nil {
31 t.Fatalf("SaveMillLease(transition): %v", err)
32 }
33
34 leases, err := d.ListMillLeases()
35 if err != nil {
36 t.Fatalf("ListMillLeases: %v", err)
37 }
38 if len(leases) != 1 {
39 t.Fatalf("ListMillLeases returned %d leases, want 1 (state transition must replace, not duplicate)", len(leases))
40 }
41 if leases[0].LeaseID != lease.LeaseID || leases[0].State != "running" {
42 t.Fatalf("ListMillLeases[0] = %+v, want %+v", leases[0], lease)
43 }
44
45 if err := d.DeleteMillLease("lease-1"); err != nil {
46 t.Fatalf("DeleteMillLease: %v", err)
47 }
48 if leases, _ = d.ListMillLeases(); len(leases) != 0 {
49 t.Fatalf("lease survived deletion: %+v", leases)
50 }
51}
52
53func TestExecutorCursors(t *testing.T) {
54 d := newTestDB(t)
55
56 if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-1", "inc-1", 5) }); err != nil {
57 t.Fatalf("AdvanceCursor: %v", err)
58 }
59 if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-1", "inc-1", 9) }); err != nil {
60 t.Fatalf("AdvanceCursor(advance): %v", err)
61 }
62 if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-2", "inc-1", 1) }); err != nil {
63 t.Fatalf("AdvanceCursor(node-2): %v", err)
64 }
65
66 cursors, err := d.ListExecutorCursors()
67 if err != nil {
68 t.Fatalf("ListExecutorCursors: %v", err)
69 }
70 if len(cursors) != 2 {
71 t.Fatalf("ListExecutorCursors = %v, want 2 cursors", cursors)
72 }
73
74 cursorMap := make(map[string]uint64)
75 for _, c := range cursors {
76 cursorMap[c.NodeID+"/"+c.Epoch] = c.AckedSeqno
77 }
78
79 if cursorMap["node-1/inc-1"] != 9 || cursorMap["node-2/inc-1"] != 1 {
80 t.Fatalf("unexpected cursors: %v", cursorMap)
81 }
82
83}
84
85func TestCompleteMillLeaseIsAtomic(t *testing.T) {
86 d := newTestDB(t)
87 lease := MillLease{
88 LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy",
89 Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: "running",
90 }
91 if err := d.SaveMillLease(lease); err != nil {
92 t.Fatalf("SaveMillLease: %v", err)
93 }
94 if _, err := d.Exec(`
95 create trigger reject_mill_lease_delete
96 before delete on mill_leases
97 begin
98 select raise(abort, 'forced delete failure');
99 end
100 `); err != nil {
101 t.Fatalf("create failure trigger: %v", err)
102 }
103
104 n := notifier.New()
105 notifications := n.Subscribe()
106 defer n.Unsubscribe(notifications)
107 err := d.CompleteMillLease(
108 "lease-1",
109 "at://knot.example/sh.tangled.pipeline/rkey1",
110 "build",
111 "failed",
112 nil,
113 nil,
114 &n,
115 )
116 if err == nil {
117 t.Fatal("CompleteMillLease succeeded despite forced lease deletion failure")
118 }
119 var eventCount int
120 if err := d.QueryRow(`select count(*) from events`).Scan(&eventCount); err != nil {
121 t.Fatalf("count events after rollback: %v", err)
122 }
123 if eventCount != 0 {
124 t.Fatalf("terminal event count after rollback = %d, want 0", eventCount)
125 }
126 if leases, listErr := d.ListMillLeases(); listErr != nil || len(leases) != 1 {
127 t.Fatalf("leases after rollback = %+v, err = %v; want original lease", leases, listErr)
128 }
129 select {
130 case <-notifications:
131 t.Fatal("rollback notified event subscribers")
132 default:
133 }
134
135 if _, err := d.Exec(`drop trigger reject_mill_lease_delete`); err != nil {
136 t.Fatalf("drop failure trigger: %v", err)
137 }
138 if err := d.CompleteMillLease(
139 "lease-1",
140 "at://knot.example/sh.tangled.pipeline/rkey1",
141 "build",
142 "failed",
143 nil,
144 nil,
145 &n,
146 ); err != nil {
147 t.Fatalf("CompleteMillLease retry: %v", err)
148 }
149 if err := d.QueryRow(`select count(*) from events`).Scan(&eventCount); err != nil {
150 t.Fatalf("count committed events: %v", err)
151 }
152 if eventCount != 1 {
153 t.Fatalf("terminal event count after commit = %d, want 1", eventCount)
154 }
155 if leases, listErr := d.ListMillLeases(); listErr != nil || len(leases) != 0 {
156 t.Fatalf("leases after commit = %+v, err = %v; want none", leases, listErr)
157 }
158 select {
159 case <-notifications:
160 default:
161 t.Fatal("committed terminal event did not notify subscribers")
162 }
163}
164
165func TestRestartPersistence(t *testing.T) {
166 dbPath := filepath.Join(t.TempDir(), "persist.db")
167 ctx := context.Background()
168 d, err := Make(ctx, dbPath)
169 if err != nil {
170 t.Fatalf("Make: %v", err)
171 }
172
173 lease := MillLease{
174 LeaseID: "lease-p", NodeID: "node-p", Epoch: "inc-p", Engine: "dummy",
175 Knot: "k", Rkey: "r", Workflow: "w", State: "running",
176 }
177 if err := d.SaveMillLease(lease); err != nil {
178 t.Fatalf("SaveMillLease: %v", err)
179 }
180
181 if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-p", "inc-p", 42) }); err != nil {
182 t.Fatalf("AdvanceCursor: %v", err)
183 }
184
185 if err := d.SetOutboxEpoch("inc-p"); err != nil {
186 t.Fatalf("SetOutboxEpoch: %v", err)
187 }
188 if _, err := d.AppendOutboxRow([]byte("hello world"), true); err != nil {
189 t.Fatalf("AppendOutboxRow: %v", err)
190 }
191
192 if err := d.Close(); err != nil {
193 t.Fatalf("Close: %v", err)
194 }
195
196 d2, err := Make(ctx, dbPath)
197 if err != nil {
198 t.Fatalf("Make reopen: %v", err)
199 }
200 defer d2.Close()
201
202 leases, err := d2.ListMillLeases()
203 if err != nil {
204 t.Fatalf("ListMillLeases: %v", err)
205 }
206 if len(leases) != 1 || leases[0].LeaseID != "lease-p" || leases[0].Epoch != "inc-p" {
207 t.Fatalf("unexpected leases: %+v", leases)
208 }
209
210 cursors, err := d2.ListExecutorCursors()
211 if err != nil {
212 t.Fatalf("ListExecutorCursors: %v", err)
213 }
214 if len(cursors) != 1 || cursors[0].NodeID != "node-p" || cursors[0].Epoch != "inc-p" || cursors[0].AckedSeqno != 42 {
215 t.Fatalf("unexpected cursors: %+v", cursors)
216 }
217
218 inc, nextSeqno, err := d2.GetOutboxState()
219 if err != nil {
220 t.Fatalf("GetOutboxState: %v", err)
221 }
222 if inc != "inc-p" || nextSeqno != 2 {
223 t.Fatalf("unexpected outbox state: inc=%q, next=%d", inc, nextSeqno)
224 }
225 rows, err := d2.ListOutboxRows()
226 if err != nil {
227 t.Fatalf("ListOutboxRows: %v", err)
228 }
229 if len(rows) != 1 || string(rows[0].Payload) != "hello world" || !rows[0].Control {
230 t.Fatalf("unexpected outbox rows: %+v", rows)
231 }
232}
233
234func TestOutboxPrefixAck(t *testing.T) {
235 d := newTestDB(t)
236
237 if err := d.SetOutboxEpoch("inc-1"); err != nil {
238 t.Fatalf("SetOutboxEpoch: %v", err)
239 }
240
241 o1, err := d.AppendOutboxRow([]byte("msg1"), false)
242 if err != nil || o1 != 1 {
243 t.Fatalf("AppendOutboxRow 1: %v, seqno=%d", err, o1)
244 }
245 o2, err := d.AppendOutboxRow([]byte("msg2"), false)
246 if err != nil || o2 != 2 {
247 t.Fatalf("AppendOutboxRow 2: %v, seqno=%d", err, o2)
248 }
249 o3, err := d.AppendOutboxRow([]byte("msg3"), true)
250 if err != nil || o3 != 3 {
251 t.Fatalf("AppendOutboxRow 3: %v, seqno=%d", err, o3)
252 }
253
254 n, err := d.DeleteOutboxPrefix(2)
255 if err != nil {
256 t.Fatalf("DeleteOutboxPrefix: %v", err)
257 }
258 if n.Rows != 2 || n.Bytes != 8 {
259 t.Fatalf("deleted prefix = %+v, want 2 rows and 8 bytes", n)
260 }
261
262 rows, err := d.ListOutboxRows()
263 if err != nil {
264 t.Fatalf("ListOutboxRows: %v", err)
265 }
266 if len(rows) != 1 || rows[0].Seqno != 3 || string(rows[0].Payload) != "msg3" {
267 t.Fatalf("expected only msg3 (seqno 3) to remain, got: %+v", rows)
268 }
269}
270
271func TestCompositeCursors(t *testing.T) {
272 d := newTestDB(t)
273
274 if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-1", "inc-1", 10) }); err != nil {
275 t.Fatalf("AdvanceCursor: %v", err)
276 }
277 if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-1", "inc-2", 20) }); err != nil {
278 t.Fatalf("AdvanceCursor: %v", err)
279 }
280 if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-2", "inc-1", 5) }); err != nil {
281 t.Fatalf("AdvanceCursor: %v", err)
282 }
283
284 cursors, err := d.ListExecutorCursors()
285 if err != nil {
286 t.Fatalf("ListExecutorCursors: %v", err)
287 }
288 if len(cursors) != 3 {
289 t.Fatalf("expected 3 cursors, got %d", len(cursors))
290 }
291
292 cursorMap := make(map[string]uint64)
293 for _, c := range cursors {
294 key := c.NodeID + "/" + c.Epoch
295 cursorMap[key] = c.AckedSeqno
296 }
297
298 if cursorMap["node-1/inc-1"] != 10 || cursorMap["node-1/inc-2"] != 20 || cursorMap["node-2/inc-1"] != 5 {
299 t.Fatalf("unexpected cursor values: %v", cursorMap)
300 }
301
302}
303
304func TestBatchRollback(t *testing.T) {
305 d := newTestDB(t)
306
307 lease := MillLease{
308 LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy",
309 Knot: "k", Rkey: "r", Workflow: "w", State: "running",
310 }
311 if err := d.SaveMillLease(lease); err != nil {
312 t.Fatalf("SaveMillLease: %v", err)
313 }
314
315 n := notifier.New()
316 err := d.ApplyEventBatch(&n, func(tx *EventBatchTx) error {
317 if err := tx.DeleteLease("lease-1"); err != nil {
318 return err
319 }
320 if err := tx.AdvanceCursor("node-1", "inc-1", 100); err != nil {
321 return err
322 }
323 return fmt.Errorf("forced batch failure")
324 })
325
326 if err == nil {
327 t.Fatal("expected ApplyEventBatch to return error")
328 }
329
330 leases, err := d.ListMillLeases()
331 if err != nil {
332 t.Fatalf("ListMillLeases: %v", err)
333 }
334 if len(leases) != 1 {
335 t.Fatalf("lease was deleted despite rollback: %+v", leases)
336 }
337
338 cursors, err := d.ListExecutorCursors()
339 if err != nil {
340 t.Fatalf("ListExecutorCursors: %v", err)
341 }
342 if len(cursors) != 0 {
343 t.Fatalf("cursor was advanced despite rollback: %+v", cursors)
344 }
345}
346
347func TestTerminalCursorAtomicity(t *testing.T) {
348 d := newTestDB(t)
349
350 lease := MillLease{
351 LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy",
352 Knot: "k", Rkey: "r", Workflow: "w", State: "running",
353 }
354 if err := d.SaveMillLease(lease); err != nil {
355 t.Fatalf("SaveMillLease: %v", err)
356 }
357
358 n := notifier.New()
359 notifications := n.Subscribe()
360 defer n.Unsubscribe(notifications)
361
362 err := d.ApplyEventBatch(&n, func(tx *EventBatchTx) error {
363 if err := tx.DeleteLease("lease-1"); err != nil {
364 return err
365 }
366 return tx.AdvanceCursor("node-1", "inc-1", 100)
367 })
368
369 if err != nil {
370 t.Fatalf("ApplyEventBatch: %v", err)
371 }
372
373 leases, err := d.ListMillLeases()
374 if err != nil {
375 t.Fatalf("ListMillLeases: %v", err)
376 }
377 if len(leases) != 0 {
378 t.Fatalf("lease not deleted: %+v", leases)
379 }
380
381 cursors, err := d.ListExecutorCursors()
382 if err != nil {
383 t.Fatalf("ListExecutorCursors: %v", err)
384 }
385 if len(cursors) != 1 || cursors[0].AckedSeqno != 100 {
386 t.Fatalf("cursor not advanced correctly: %+v", cursors)
387 }
388
389 select {
390 case <-notifications:
391 default:
392 t.Fatal("notifier was not fired after batch commit")
393 }
394}