package db import ( "context" "encoding/json" "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("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", PipelineID: "r1", Workflow: "w1", State: "active", RepoDID: "did:plc:repo1"}, {LeaseID: "l2", NodeID: "n1", Epoch: "e1", Engine: "microvm", PipelineID: "r1", Workflow: "w2", State: "active", RepoDID: "did:plc:repo1"}, {LeaseID: "l3", NodeID: "n1", Epoch: "e1", Engine: "microvm", PipelineID: "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, pipeline_id, workflow, ref, hash) values (?, ?, ?, ?, ?, ?)`, l.LeaseID, l.RepoDID, l.PipelineID, 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 TestDeletePipelinesByRepo(t *testing.T) { d := newTestDB(t) repo := "did:plc:repo1" other := "did:plc:repo2" insertPipeline := func(rkey, repoDid string) { p := &tangled.CiPipeline{Id: rkey, Repo: repoDid} payload, err := json.Marshal(p) if err != nil { t.Fatal(err) } if _, err := d.Exec(`insert into pipelines (pipeline_id, repo_did, commit_sha, kind, payload) values (?, ?, '', '', ?)`, rkey, repoDid, string(payload)); err != nil { t.Fatal(err) } } insertPipeline("p1", repo) insertPipeline("p2", other) for _, rkey := range []string{"p1", "p2", "p9"} { if _, err := d.Exec(`insert into workflow_statuses (pipeline_id, workflow, status, created_at) values (?, 'build', 'pending', 'now')`, rkey); err != nil { t.Fatal(err) } } if err := d.DeletePipelinesByRepo(repo, []models.PipelineId{"p9"}); err != nil { t.Fatal(err) } var n int if err := d.QueryRow(`select count(*) from pipelines where pipeline_id = 'p1'`).Scan(&n); err != nil || n != 0 { t.Fatalf("expected repo pipeline gone, got %d err %v", n, err) } if err := d.QueryRow(`select count(*) from workflow_statuses where pipeline_id in ('p1', 'p9')`).Scan(&n); err != nil || n != 0 { t.Fatalf("expected wiped statuses gone, got %d err %v", n, err) } if err := d.QueryRow(`select count(*) from pipelines where pipeline_id = 'p2'`).Scan(&n); err != nil || n != 1 { t.Fatalf("expected other pipeline to survive, got %d err %v", n, err) } }