Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639package db
import ( "context" "database/sql" "errors" "fmt" "path/filepath" "sync" "testing"
"github.com/mattn/go-sqlite3" "tangled.org/core/notifier" "tangled.org/core/spindle/models")
func TestMillLeaseRoundTrip(t *testing.T) { d := newTestDB(t)
lease := MillLease{ LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", PipelineID: "rkey1", Workflow: "build", State: "reserved", MillRecordsTerminalMetrics: true, } if err := d.SaveMillLease(lease); err != nil { t.Fatalf("SaveMillLease: %v", err) }
lease.State = "running" if err := d.SaveMillLease(lease); err != nil { t.Fatalf("SaveMillLease(transition): %v", err) }
leases, err := d.ListMillLeases() if err != nil { t.Fatalf("ListMillLeases: %v", err) } if len(leases) != 1 { t.Fatalf("ListMillLeases returned %d leases, want 1 (state transition must replace, not duplicate)", len(leases)) } if leases[0].LeaseID != lease.LeaseID || leases[0].State != "running" || !leases[0].MillRecordsTerminalMetrics { t.Fatalf("ListMillLeases[0] = %+v, want %+v", leases[0], lease) }
if err := d.DeleteMillLease("lease-1"); err != nil { t.Fatalf("DeleteMillLease: %v", err) } if leases, _ = d.ListMillLeases(); len(leases) != 0 { t.Fatalf("lease survived deletion: %+v", leases) }}
func TestExecutorCursors(t *testing.T) { d := newTestDB(t)
if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-1", "inc-1", 5) }); err != nil { t.Fatalf("AdvanceCursor: %v", err) } if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-1", "inc-1", 9) }); err != nil { t.Fatalf("AdvanceCursor(advance): %v", err) } if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-2", "inc-1", 1) }); err != nil { t.Fatalf("AdvanceCursor(node-2): %v", err) }
cursors, err := d.ListExecutorCursors() if err != nil { t.Fatalf("ListExecutorCursors: %v", err) } if len(cursors) != 2 { t.Fatalf("ListExecutorCursors = %v, want 2 cursors", cursors) }
cursorMap := make(map[string]uint64) for _, c := range cursors { cursorMap[c.NodeID+"/"+c.Epoch] = c.AckedSeqno }
if cursorMap["node-1/inc-1"] != 9 || cursorMap["node-2/inc-1"] != 1 { t.Fatalf("unexpected cursors: %v", cursorMap) }
if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-1", "inc-1", 4) }); err != nil { t.Fatalf("AdvanceCursor(stale): %v", err) } cur, err := d.GetExecutorCursor("node-1", "inc-1") if err != nil || cur != 9 { t.Fatalf("cursor regressed to %d (err: %v), want 9", cur, err) }}
func TestApplyEventBatchPublishesResultOnlyOnCommit(t *testing.T) { d := newTestDB(t) attempts := 0 res, err := ApplyEventBatch(d, nil, func(tx *EventBatchTx) ([]string, error) { attempts++ if attempts == 1 { // a result produced by a failed attempt must never escape return []string{"ghost"}, sqlite3.Error{Code: sqlite3.ErrBusy} } if err := tx.AdvanceCursor("node-1", "inc-1", 7); err != nil { return nil, err } return []string{"committed"}, nil }) if err != nil { t.Fatalf("ApplyEventBatch: %v", err) } if attempts != 2 { t.Fatalf("callback attempts = %d, want 2", attempts) } if len(res) != 1 || res[0] != "committed" { t.Fatalf("result = %v, want [committed]", res) } if got, err := d.GetExecutorCursor("node-1", "inc-1"); err != nil || got != 7 { t.Fatalf("cursor = %d, err = %v; want 7", got, err) }}
func TestApplyEventBatchRetriesBusy(t *testing.T) { d := newTestDB(t) attempts := 0 err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { attempts++ if attempts == 1 { return sqlite3.Error{Code: sqlite3.ErrBusy} } return tx.AdvanceCursor("node-1", "inc-1", 5) }) if err != nil { t.Fatalf("ApplyEventBatch: %v", err) } if attempts != 2 { t.Fatalf("callback attempts = %d, want 2", attempts) } if got, err := d.GetExecutorCursor("node-1", "inc-1"); err != nil || got != 5 { t.Fatalf("cursor = %d, err = %v; want 5", got, err) }}
func TestApplyEventBatchCtxAbortsOnCanceledContext(t *testing.T) { d := newTestDB(t) ctx, cancel := context.WithCancel(context.Background()) cancel() _, err := ApplyEventBatchCtx(ctx, d, nil, func(tx *EventBatchTx) (string, error) { return "committed", tx.AdvanceCursor("node-1", "inc-1", 10) }) if !errors.Is(err, context.Canceled) { t.Fatalf("ApplyEventBatchCtx err = %v, want context.Canceled", err) } cur, err := d.GetExecutorCursor("node-1", "inc-1") if err != nil { t.Fatalf("GetExecutorCursor: %v", err) } if cur != 0 { t.Fatalf("cursor advanced to %d on canceled context, want 0", cur) }}
func TestCompleteMillLeaseIsAtomic(t *testing.T) { d := newTestDB(t) lease := MillLease{ LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", PipelineID: "rkey1", Workflow: "build", State: "running", } if err := d.SaveMillLease(lease); err != nil { t.Fatalf("SaveMillLease: %v", err) } if _, err := d.Exec(` create trigger reject_mill_lease_delete before delete on mill_leases begin select raise(abort, 'forced delete failure'); end `); err != nil { t.Fatalf("create failure trigger: %v", err) }
n := notifier.New() notifications := n.Subscribe() defer n.Unsubscribe(notifications) err := d.CompleteMillLease( "lease-1", models.WorkflowId{PipelineId: models.PipelineId("rkey1"), Name: "build"}, "failed", nil, nil, &n, ) if err == nil { t.Fatal("CompleteMillLease succeeded despite forced lease deletion failure") } var eventCount int if err := d.QueryRow(`select count(*) from workflow_statuses`).Scan(&eventCount); err != nil { t.Fatalf("count events after rollback: %v", err) } if eventCount != 0 { t.Fatalf("terminal status count after rollback = %d, want 0", eventCount) } if leases, listErr := d.ListMillLeases(); listErr != nil || len(leases) != 1 { t.Fatalf("leases after rollback = %+v, err = %v; want original lease", leases, listErr) } select { case <-notifications: t.Fatal("rollback notified event subscribers") default: }
if _, err := d.Exec(`drop trigger reject_mill_lease_delete`); err != nil { t.Fatalf("drop failure trigger: %v", err) } if err := d.CompleteMillLease( "lease-1", models.WorkflowId{PipelineId: models.PipelineId("rkey1"), Name: "build"}, "failed", nil, nil, &n, ); err != nil { t.Fatalf("CompleteMillLease retry: %v", err) } if err := d.QueryRow(`select count(*) from workflow_statuses`).Scan(&eventCount); err != nil { t.Fatalf("count committed events: %v", err) } if eventCount != 1 { t.Fatalf("terminal status count after commit = %d, want 1", eventCount) } if leases, listErr := d.ListMillLeases(); listErr != nil || len(leases) != 0 { t.Fatalf("leases after commit = %+v, err = %v; want none", leases, listErr) } select { case <-notifications: default: t.Fatal("committed terminal event did not notify subscribers") }}
func TestRestartPersistence(t *testing.T) { dbPath := filepath.Join(t.TempDir(), "persist.db") ctx := context.Background() d, err := Make(ctx, dbPath) if err != nil { t.Fatalf("Make: %v", err) }
lease := MillLease{ LeaseID: "lease-p", NodeID: "node-p", Epoch: "inc-p", Engine: "dummy", PipelineID: "r", Workflow: "w", State: "running", } if err := d.SaveMillLease(lease); err != nil { t.Fatalf("SaveMillLease: %v", err) }
if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-p", "inc-p", 42) }); err != nil { t.Fatalf("AdvanceCursor: %v", err) }
if err := d.SetOutboxEpoch("inc-p"); err != nil { t.Fatalf("SetOutboxEpoch: %v", err) } if _, err := d.AppendOutboxRow([]byte("hello world"), true); err != nil { t.Fatalf("AppendOutboxRow: %v", err) }
if err := d.Close(); err != nil { t.Fatalf("Close: %v", err) }
d2, err := Make(ctx, dbPath) if err != nil { t.Fatalf("Make reopen: %v", err) } defer d2.Close()
leases, err := d2.ListMillLeases() if err != nil { t.Fatalf("ListMillLeases: %v", err) } if len(leases) != 1 || leases[0].LeaseID != "lease-p" || leases[0].Epoch != "inc-p" { t.Fatalf("unexpected leases: %+v", leases) }
cursors, err := d2.ListExecutorCursors() if err != nil { t.Fatalf("ListExecutorCursors: %v", err) } if len(cursors) != 1 || cursors[0].NodeID != "node-p" || cursors[0].Epoch != "inc-p" || cursors[0].AckedSeqno != 42 { t.Fatalf("unexpected cursors: %+v", cursors) }
inc, nextSeqno, err := d2.GetOutboxState() if err != nil { t.Fatalf("GetOutboxState: %v", err) } if inc != "inc-p" || nextSeqno != 2 { t.Fatalf("unexpected outbox state: inc=%q, next=%d", inc, nextSeqno) } rows, err := d2.ListOutboxRows() if err != nil { t.Fatalf("ListOutboxRows: %v", err) } if len(rows) != 1 || string(rows[0].Payload) != "hello world" || !rows[0].Control { t.Fatalf("unexpected outbox rows: %+v", rows) }}
func TestOutboxPrefixAck(t *testing.T) { d := newTestDB(t)
if err := d.SetOutboxEpoch("inc-1"); err != nil { t.Fatalf("SetOutboxEpoch: %v", err) }
o1, err := d.AppendOutboxRow([]byte("msg1"), false) if err != nil || o1 != 1 { t.Fatalf("AppendOutboxRow 1: %v, seqno=%d", err, o1) } o2, err := d.AppendOutboxRow([]byte("msg2"), false) if err != nil || o2 != 2 { t.Fatalf("AppendOutboxRow 2: %v, seqno=%d", err, o2) } o3, err := d.AppendOutboxRow([]byte("msg3"), true) if err != nil || o3 != 3 { t.Fatalf("AppendOutboxRow 3: %v, seqno=%d", err, o3) }
n, err := d.DeleteOutboxPrefix(2) if err != nil { t.Fatalf("DeleteOutboxPrefix: %v", err) } if n.Rows != 2 || n.Bytes != 8 { t.Fatalf("deleted prefix = %+v, want 2 rows and 8 bytes", n) }
rows, err := d.ListOutboxRows() if err != nil { t.Fatalf("ListOutboxRows: %v", err) } if len(rows) != 1 || rows[0].Seqno != 3 || string(rows[0].Payload) != "msg3" { t.Fatalf("expected only msg3 (seqno 3) to remain, got: %+v", rows) }}
func TestOutboxConcurrentAppend(t *testing.T) { d := newTestDB(t) if err := d.SetOutboxEpoch("inc-1"); err != nil { t.Fatalf("SetOutboxEpoch: %v", err) }
const count = 32 type result struct { seqno uint64 err error } start := make(chan struct{}) results := make(chan result, count) var wg sync.WaitGroup for i := range count { wg.Add(1) go func() { defer wg.Done() <-start seqno, err := d.AppendOutboxRow([]byte(fmt.Sprintf("msg%d", i)), true) results <- result{seqno: seqno, err: err} }() } close(start) wg.Wait() close(results)
seen := make(map[uint64]struct{}, count) for result := range results { if result.err != nil { t.Errorf("AppendOutboxRow: %v", result.err) continue } if _, ok := seen[result.seqno]; ok { t.Errorf("duplicate outbox seqno %d", result.seqno) } seen[result.seqno] = struct{}{} } if len(seen) != count { t.Fatalf("successful appends = %d, want %d", len(seen), count) }
rows, err := d.ListOutboxRows() if err != nil { t.Fatalf("ListOutboxRows: %v", err) } if len(rows) != count { t.Fatalf("persisted outbox rows = %d, want %d", len(rows), count) }}
func TestCompositeCursors(t *testing.T) { d := newTestDB(t)
if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-1", "inc-1", 10) }); err != nil { t.Fatalf("AdvanceCursor: %v", err) } if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-1", "inc-2", 20) }); err != nil { t.Fatalf("AdvanceCursor: %v", err) } if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor("node-2", "inc-1", 5) }); err != nil { t.Fatalf("AdvanceCursor: %v", err) }
cursors, err := d.ListExecutorCursors() if err != nil { t.Fatalf("ListExecutorCursors: %v", err) } if len(cursors) != 3 { t.Fatalf("expected 3 cursors, got %d", len(cursors)) }
cursorMap := make(map[string]uint64) for _, c := range cursors { key := c.NodeID + "/" + c.Epoch cursorMap[key] = c.AckedSeqno }
if cursorMap["node-1/inc-1"] != 10 || cursorMap["node-1/inc-2"] != 20 || cursorMap["node-2/inc-1"] != 5 { t.Fatalf("unexpected cursor values: %v", cursorMap) }
}
func TestBatchRollback(t *testing.T) { d := newTestDB(t)
lease := MillLease{ LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", PipelineID: "r", Workflow: "w", State: "running", } if err := d.SaveMillLease(lease); err != nil { t.Fatalf("SaveMillLease: %v", err) }
n := notifier.New() err := d.ApplyEventBatch(&n, func(tx *EventBatchTx) error { if err := tx.DeleteLease("lease-1"); err != nil { return err } if err := tx.AdvanceCursor("node-1", "inc-1", 100); err != nil { return err } return fmt.Errorf("forced batch failure") })
if err == nil { t.Fatal("expected ApplyEventBatch to return error") }
leases, err := d.ListMillLeases() if err != nil { t.Fatalf("ListMillLeases: %v", err) } if len(leases) != 1 { t.Fatalf("lease was deleted despite rollback: %+v", leases) }
cursors, err := d.ListExecutorCursors() if err != nil { t.Fatalf("ListExecutorCursors: %v", err) } if len(cursors) != 0 { t.Fatalf("cursor was advanced despite rollback: %+v", cursors) }}
func TestTerminalCursorAtomicity(t *testing.T) { d := newTestDB(t)
lease := MillLease{ LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", PipelineID: "r", Workflow: "w", State: "running", } if err := d.SaveMillLease(lease); err != nil { t.Fatalf("SaveMillLease: %v", err) }
n := notifier.New() notifications := n.Subscribe() defer n.Unsubscribe(notifications)
err := d.ApplyEventBatch(&n, func(tx *EventBatchTx) error { if err := tx.DeleteLease("lease-1"); err != nil { return err } return tx.AdvanceCursor("node-1", "inc-1", 100) })
if err != nil { t.Fatalf("ApplyEventBatch: %v", err) }
leases, err := d.ListMillLeases() if err != nil { t.Fatalf("ListMillLeases: %v", err) } if len(leases) != 0 { t.Fatalf("lease not deleted: %+v", leases) }
cursors, err := d.ListExecutorCursors() if err != nil { t.Fatalf("ListExecutorCursors: %v", err) } if len(cursors) != 1 || cursors[0].AckedSeqno != 100 { t.Fatalf("cursor not advanced correctly: %+v", cursors) }
select { case <-notifications: default: t.Fatal("notifier was not fired after batch commit") }}
func TestExecutorCursorResetHelpers(t *testing.T) { d := newTestDB(t)
seed := func(node, epoch string, seqno uint64) { if err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { return tx.AdvanceCursor(node, epoch, seqno) }); err != nil { t.Fatalf("seed cursor: %v", err) } } seed("node-1", "inc-1", 7) seed("node-1", "inc-2", 3) seed("node-2", "inc-1", 42)
if cur, err := d.GetExecutorCursor("node-1", "inc-1"); err != nil || cur != 7 { t.Fatalf("GetExecutorCursor = %d, %v; want 7", cur, err) } // an unknown stream looks the same as a fresh one if cur, err := d.GetExecutorCursor("node-1", "missing"); err != nil || cur != 0 { t.Fatalf("GetExecutorCursor(missing) = %d, %v; want 0", cur, err) }
// skip-forward touches every epoch of the node and nothing else if n, err := d.SetExecutorCursors("node-1", 28); err != nil || n != 2 { t.Fatalf("SetExecutorCursors = %d rows, %v; want 2", n, err) } if cur, _ := d.GetExecutorCursor("node-1", "inc-2"); cur != 28 { t.Fatalf("SetExecutorCursors left inc-2 at %d, want 28", cur) } if cur, _ := d.GetExecutorCursor("node-2", "inc-1"); cur != 42 { t.Fatalf("SetExecutorCursors clobbered node-2: %d, want 42", cur) }
if n, err := d.DeleteExecutorCursors("node-1"); err != nil || n != 2 { t.Fatalf("DeleteExecutorCursors = %d rows, %v; want 2", n, err) } if cur, _ := d.GetExecutorCursor("node-1", "inc-1"); cur != 0 { t.Fatalf("DeleteExecutorCursors left %d, want 0", cur) }}
func TestPendingArtifactWorkflowIdentityMigration(t *testing.T) { path := filepath.Join(t.TempDir(), "spindle.db") legacy, err := sql.Open("sqlite3", path) if err != nil { t.Fatal(err) } if _, err := legacy.Exec(` create table executor_pending_artifacts ( lease_id text primary key, workflow text not null, status text not null, error text not null default '', exit_code integer not null default 0, ref text not null, hash text not null, failure_class text not null default '', failure_reason text not null default '', mill_records_terminal_metrics integer not null default 0 ); insert into executor_pending_artifacts (lease_id, workflow, status, ref, hash) values ('lease-1', 'build', 'success', 'logs/lease-1.log', 'sha256:test'); `); err != nil { legacy.Close() t.Fatal(err) } if err := legacy.Close(); err != nil { t.Fatal(err) }
d, err := Make(context.Background(), path) if err != nil { t.Fatal(err) } defer d.Close() rows, err := d.ListPendingArtifacts() if err != nil { t.Fatal(err) } if len(rows) != 1 || rows[0].PipelineID != "" || rows[0].Workflow != "build" { t.Fatalf("migrated pending artifact = %+v", rows) } var identityColumns int if err := d.QueryRow(` select count(*) from pragma_table_info('executor_pending_artifacts') where name = 'pipeline_id' `).Scan(&identityColumns); err != nil { t.Fatal(err) } if identityColumns != 1 { t.Fatalf("pending artifact identity columns = %d, want 1", identityColumns) }}
func TestClearPendingArtifacts(t *testing.T) { d := newTestDB(t) wid := models.WorkflowId{PipelineId: models.PipelineId("rkey1"), Name: "build"} if err := d.SavePendingArtifact("lease-1", wid, "success", "", 0, "ref", "sha256:x", "none", "success", true); err != nil { t.Fatalf("SavePendingArtifact: %v", err) } rows, err := d.ListPendingArtifacts() if err != nil { t.Fatal(err) } if len(rows) != 1 || rows[0].PipelineID != string(wid.PipelineId) || rows[0].Workflow != wid.Name || rows[0].FailureClass != "none" || rows[0].FailureReason != "success" || !rows[0].MillRecordsTerminalMetrics { t.Fatalf("pending artifact attribution = %+v", rows) } if err := d.ClearPendingArtifacts(); err != nil { t.Fatalf("ClearPendingArtifacts: %v", err) } if rows, err := d.ListPendingArtifacts(); err != nil || len(rows) != 0 { t.Fatalf("ListPendingArtifacts after clear = %+v, %v; want empty", rows, err) }}