This repository has no description
0

Configure Feed

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

core / spindle / mill / restore_test.go
12 kB 386 lines
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}