This repository has no description
0

Configure Feed

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

core / spindle / db / mill_state_test.go
11 kB 394 lines
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}