package db import ( "context" "fmt" "testing" "github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/api/tangled" "tangled.org/core/spindle/models" ) func TestDeleteReposByDid(t *testing.T) { d := newTestDB(t) repoDid := syntax.DID("did:plc:repo1") other := syntax.DID("did:plc:repo2") for _, r := range []Repo{ {Knot: "k1", Owner: syntax.DID("did:plc:o1"), Rkey: syntax.RecordKey("a"), RepoDid: repoDid}, {Knot: "k2", Owner: syntax.DID("did:plc:o1"), Rkey: syntax.RecordKey("b"), RepoDid: repoDid}, {Knot: "k1", Owner: syntax.DID("did:plc:o2"), Rkey: syntax.RecordKey("c"), RepoDid: other}, } { if err := d.AddRepo(r); err != nil { t.Fatal(err) } } repos, err := d.ReposByDid(repoDid) if err != nil { t.Fatal(err) } if len(repos) != 2 { t.Fatalf("expected 2 sibling rows, got %d", len(repos)) } if err := d.DeleteReposByDid(repoDid); err != nil { t.Fatal(err) } if repos, _ := d.ReposByDid(repoDid); len(repos) != 0 { t.Fatalf("expected rows gone, got %d", len(repos)) } if repos, _ := d.ReposByDid(other); len(repos) != 1 { t.Fatalf("expected other repo untouched, got %d", len(repos)) } } func TestDeleteJobsByRepo(t *testing.T) { d := newTestDB(t) ctx := context.Background() pid := models.PipelineId{Knot: "k1", Rkey: "p1"} if err := d.EnqueueJob(ctx, "did:plc:repo1", pid, nil, tangled.Pipeline{}, "", ""); err != nil { t.Fatal(err) } if err := d.EnqueueJob(ctx, "did:plc:repo2", pid, nil, tangled.Pipeline{}, "", ""); err != nil { t.Fatal(err) } if err := d.DeleteJobsByRepo(ctx, "did:plc:repo1"); err != nil { t.Fatal(err) } var n int if err := d.QueryRow(`select count(*) from jobs`).Scan(&n); err != nil { t.Fatal(err) } if n != 1 { t.Fatalf("expected 1 job left, got %d", n) } if err := d.QueryRow(`select count(*) from jobs where repo_did = 'did:plc:repo2'`).Scan(&n); err != nil || n != 1 { t.Fatalf("expected surviving job to be repo2's, count=%d err=%v", n, err) } } func TestDeleteMillLeasesByRepo(t *testing.T) { d := newTestDB(t) for _, l := range []MillLease{ {LeaseID: "l1", NodeID: "n1", Epoch: "e1", Engine: "microvm", Knot: "k1", Rkey: "r1", Workflow: "w1", State: "active", RepoDID: "did:plc:repo1"}, {LeaseID: "l2", NodeID: "n1", Epoch: "e1", Engine: "microvm", Knot: "k1", Rkey: "r1", Workflow: "w2", State: "active", RepoDID: "did:plc:repo1"}, {LeaseID: "l3", NodeID: "n1", Epoch: "e1", Engine: "microvm", Knot: "k1", Rkey: "r9", Workflow: "w1", State: "active", RepoDID: "did:plc:repo2"}, } { if err := d.SaveMillLease(l); err != nil { t.Fatal(err) } if _, err := d.Exec(`insert into mill_artifacts (lease_id, repo_did, knot, rkey, workflow, ref, hash) values (?, ?, ?, ?, ?, ?, ?)`, l.LeaseID, l.RepoDID, l.Knot, l.Rkey, l.Workflow, "logs/"+l.LeaseID+".log", "h"); err != nil { t.Fatal(err) } if _, err := d.Exec(`insert into executor_pending_artifacts (lease_id, workflow, status, ref, hash) values (?, ?, 'done', 'r', 'h')`, l.LeaseID, l.Workflow); err != nil { t.Fatal(err) } } removed, err := d.DeleteMillLeasesByRepo("did:plc:repo1") if err != nil { t.Fatal(err) } if len(removed) != 2 { t.Fatalf("expected 2 lease ids, got %v", removed) } leases, err := d.ListMillLeases() if err != nil { t.Fatal(err) } if len(leases) != 1 || leases[0].LeaseID != "l3" { t.Fatalf("expected only l3 left, got %+v", leases) } if err := d.DeleteArtifactRefsByRepo("did:plc:repo1"); err != nil { t.Fatal(err) } var n int if err := d.QueryRow(`select count(*) from mill_artifacts`).Scan(&n); err != nil || n != 1 { t.Fatalf("expected 1 artifact row left, got %d err %v", n, err) } if err := d.QueryRow(`select count(*) from executor_pending_artifacts`).Scan(&n); err != nil || n != 1 { t.Fatalf("expected 1 pending artifact row left, got %d err %v", n, err) } } func TestDeleteQuotaStateForRepo(t *testing.T) { d := newTestDB(t) ctx := context.Background() repo := "did:plc:repo1" other := "did:plc:repo2" ins := []string{ fmt.Sprintf(`insert into quota_reservations (id, resource, kind, key, amount, repo_did, owner_did, phase, created_at) values ('r1', 'compute', 'generic', 'k', 1, '%s', 'did:plc:o1', 'active', 1)`, repo), fmt.Sprintf(`insert into quota_reservations (id, resource, kind, key, amount, repo_did, owner_did, phase, created_at) values ('r2', 'compute', 'generic', 'k', 1, '%s', 'did:plc:o1', 'active', 1)`, other), fmt.Sprintf(`insert into quota_allocations (repo_did, resource, kind, key, amount) values ('%s', 'compute', 'generic', 'k', 2)`, repo), fmt.Sprintf(`insert into quota_allocations (repo_did, resource, kind, key, amount) values ('%s', 'compute', 'generic', 'k', 2)`, other), fmt.Sprintf(`insert into quota_limits (did, resource, max_amount) values ('%s', 'compute', 5)`, repo), `insert into quota_limits (did, resource, max_amount) values ('did:plc:o1', 'compute', 5)`, fmt.Sprintf(`insert into quota_repo_owners (repo_did, owner_did) values ('%s', 'did:plc:o1')`, repo), } for _, q := range ins { if _, err := d.Exec(q); err != nil { t.Fatalf("fixture %q: %v", q, err) } } if err := d.DeleteQuotaStateForRepo(ctx, repo); err != nil { t.Fatal(err) } checks := []struct { q, want string }{ {`select id from quota_reservations`, "r2"}, {`select repo_did from quota_allocations`, other}, {`select did from quota_limits`, "did:plc:o1"}, } for _, c := range checks { var got string if err := d.QueryRow(c.q).Scan(&got); err != nil { t.Fatalf("%s: %v", c.q, err) } if got != c.want { t.Fatalf("%s: surviving row = %s, want %s", c.q, got, c.want) } } var n int if err := d.QueryRow(`select count(*) from quota_repo_owners`).Scan(&n); err != nil || n != 0 { t.Fatalf("expected repo owner rows gone, got %d err %v", n, err) } } func TestDeleteEventsByRepo(t *testing.T) { d := newTestDB(t) repo := "did:plc:repo1" other := "did:plc:repo2" pipelineEvent := func(did, knot string) string { return fmt.Sprintf(`{"triggerMetadata":{"repo":{"repoDid":"%s","knot":"%s"}}}`, did, knot) } statusEvent := func(knot, rkey string) string { return fmt.Sprintf(`{"pipeline":"at://did:web:%s/sh.tangled.pipeline/%s"}`, knot, rkey) } ins := []struct{ nsid, rkey, event string }{ {"sh.tangled.pipeline", "p1", pipelineEvent(repo, "k1")}, {"sh.tangled.pipeline", "p2", pipelineEvent(other, "k1")}, {"sh.tangled.pipeline.status", "s1", statusEvent("k1", "p1")}, {"sh.tangled.pipeline.status", "s2", statusEvent("k1", "p2")}, {"sh.tangled.pipeline.status", "s3", statusEvent("k9", "p1")}, {"sh.tangled.pipeline.status", "s4", statusEvent("k7", "p9")}, {"sh.tangled.repo", "x1", `{"whatever":1}`}, } for _, e := range ins { if _, err := d.Exec(`insert into events (nsid, rkey, event, created) values (?, ?, ?, 1)`, e.nsid, e.rkey, e.event); err != nil { t.Fatal(err) } } if err := d.DeleteEventsByRepo(repo, []PipelineKey{{Knot: "k7", Rkey: "p9"}}); err != nil { t.Fatal(err) } var n int if err := d.QueryRow(`select count(*) from events`).Scan(&n); err != nil { t.Fatal(err) } if n != 4 { t.Fatalf("expected 4 events left, got %d", n) } if err := d.QueryRow(`select count(*) from events where rkey in ('p1', 's1', 's4')`).Scan(&n); err != nil || n != 0 { t.Fatalf("expected wiped events gone, got %d err %v", n, err) } if err := d.QueryRow(`select count(*) from events where rkey = 's3'`).Scan(&n); err != nil || n != 1 { t.Fatalf("expected status on unrecorded knot to survive, got %d err %v", n, err) } }