package spindle import ( "context" "fmt" "os" "path/filepath" "testing" comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/syntax" jmodels "github.com/bluesky-social/jetstream/pkg/models" "tangled.org/core/api/tangled" "tangled.org/core/jetstream" "tangled.org/core/log" "tangled.org/core/notifier" "tangled.org/core/rbac" "tangled.org/core/spindle/artifactstore" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/models" "tangled.org/core/spindle/secrets" ) func newWipeTestSpindle(t *testing.T) *Spindle { t.Helper() tmp := t.TempDir() ctx := context.Background() d, err := db.Make(ctx, filepath.Join(tmp, "spindle.db")) if err != nil { t.Fatal(err) } t.Cleanup(func() { d.Close() }) e, err := rbac.NewEnforcer(filepath.Join(tmp, "rbac.db")) if err != nil { t.Fatal(err) } vault, err := secrets.NewSQLiteManager(filepath.Join(tmp, "secrets.db")) if err != nil { t.Fatal(err) } stores, err := artifactstore.NewStores(config.ArtifactStores{ Disk: config.ArtifactStoreDisk{Dir: filepath.Join(tmp, "artifacts")}, }, "", "") if err != nil { t.Fatal(err) } jc, err := jetstream.NewJetstreamClient("", "test", nil, nil, log.New("test"), d, false, false) if err != nil { t.Fatal(err) } return &Spindle{ l: log.New("test"), db: d, e: e, vault: vault, stores: stores, jc: jc, n: ptr(notifier.New()), engs: map[string]models.Engine{}, cfg: &config.Config{ Server: config.Server{ Hostname: "spindle.example.com", RepoDir: filepath.Join(tmp, "repos"), LogDir: filepath.Join(tmp, "logs"), }, }, } } func seedWipeRepo(t *testing.T, s *Spindle, repoDid, owner syntax.DID) { t.Helper() ctx := context.Background() sfx := repoDid.String()[len(repoDid.String())-1:] if err := s.db.AddRepo(db.Repo{Knot: "knot.example.com", Owner: owner, Rkey: syntax.RecordKey("r" + sfx), RepoDid: repoDid}); err != nil { t.Fatal(err) } if err := db.AddDid(s.db, owner.String()); err != nil { t.Fatal(err) } if err := s.db.AddRepoCollaborator(db.RepoCollaborator{OwnerDid: owner, Rkey: syntax.RecordKey("c" + sfx), Subject: "did:plc:collab", RepoDid: repoDid}); err != nil { t.Fatal(err) } if err := s.db.EnqueueJob(ctx, repoDid.String(), models.PipelineId{Knot: "knot.example.com", Rkey: "p" + sfx}, nil, tangled.Pipeline{}, "", ""); err != nil { t.Fatal(err) } if err := s.db.SaveMillLease(db.MillLease{LeaseID: "l" + sfx, NodeID: "n1", Epoch: "e1", Engine: "microvm", Knot: "knot.example.com", Rkey: "p" + sfx, Workflow: "w" + sfx, State: "active", RepoDID: repoDid.String()}); err != nil { t.Fatal(err) } if _, err := s.db.Exec(`insert into mill_artifacts (lease_id, repo_did, knot, rkey, workflow, ref, hash) values (?, ?, 'knot.example.com', ?, ?, ?, 'h')`, "l"+sfx, repoDid.String(), "p"+sfx, "w"+sfx, "out/l"+sfx+".bin"); err != nil { t.Fatal(err) } if _, err := s.db.Exec(`insert into quota_allocations (repo_did, resource, kind, key, amount) values (?, 'compute', 'generic', 'k', 1)`, repoDid.String()); err != nil { t.Fatal(err) } if _, err := s.db.Exec(`insert into webhooks (repo_did, url, events) values (?, 'https://example.com/hook', 'push')`, repoDid.String()); err != nil { t.Fatal(err) } if _, err := s.db.Exec(`insert into webhook_deliveries (webhook_id, event, delivery_id, url, success) values ((select max(id) from webhooks), 'push', ?, 'https://example.com/hook', 1)`, "d"+sfx); err != nil { t.Fatal(err) } if _, err := s.db.Exec(`insert into pull_rounds (repo_did, rkey, rounds) values (?, 'pull1', 2)`, repoDid.String()); err != nil { t.Fatal(err) } if err := s.vault.AddSecret(ctx, secrets.UnlockedSecret{Key: "api_key", Value: "v", Repo: secrets.RepoIdentifier(repoDid.String()), CreatedBy: owner}); err != nil { t.Fatal(err) } pipelineEvent := fmt.Sprintf(`{"triggerMetadata":{"repo":{"repoDid":"%s"}}}`, repoDid) if _, err := s.db.Exec(`insert into events (nsid, rkey, event, created) values ('sh.tangled.pipeline', ?, ?, 1)`, "p"+sfx, pipelineEvent); err != nil { t.Fatal(err) } statusEvent := fmt.Sprintf(`{"pipeline":"at://did:web:knot.example.com/sh.tangled.pipeline/%s"}`, "p"+sfx) if _, err := s.db.Exec(`insert into events (nsid, rkey, event, created) values ('sh.tangled.pipeline.status', ?, ?, 1)`, "s"+sfx, statusEvent); err != nil { t.Fatal(err) } if err := os.MkdirAll(s.newRepoPath(repoDid), 0755); err != nil { t.Fatal(err) } for _, ref := range []string{"out/l" + sfx + ".bin", "logs/l" + sfx + ".log"} { if err := s.stores.PutFile(ctx, ref, writeTempFile(t, "data")); err != nil && len(err) > 0 { t.Fatal(err) } } } func writeTempFile(t *testing.T, content string) string { t.Helper() p := filepath.Join(t.TempDir(), "src") if err := os.WriteFile(p, []byte(content), 0644); err != nil { t.Fatal(err) } return p } func TestWipeRepoRemovesAllState(t *testing.T) { s := newWipeTestSpindle(t) ctx := context.Background() repoDid := syntax.DID("did:plc:repo1") owner := syntax.DID("did:plc:owner1") seedWipeRepo(t, s, repoDid, owner) if err := s.WipeRepo(ctx, repoDid, "test"); err != nil { t.Fatalf("WipeRepo: %v", err) } var n int for _, table := range []string{"repos", "jobs", "mill_leases", "mill_artifacts", "quota_allocations", "repo_collaborators", "events", "webhooks", "webhook_deliveries", "pull_rounds"} { if err := s.db.QueryRow(`select count(*) from ` + table).Scan(&n); err != nil { t.Fatal(err) } if n != 0 { t.Fatalf("%s: expected empty, got %d rows", table, n) } } got, err := s.vault.GetSecretsUnlocked(ctx, secrets.RepoIdentifier(repoDid.String())) if err != nil { t.Fatal(err) } if len(got) != 0 { t.Fatalf("expected secrets gone, got %d", len(got)) } if _, err := os.Stat(s.newRepoPath(repoDid)); !os.IsNotExist(err) { t.Fatalf("expected clone dir gone, stat err: %v", err) } if _, err := s.stores.Open(ctx, "out/l1.bin"); err == nil { t.Fatal("expected artifact gone") } if _, err := s.stores.Open(ctx, "logs/l1.log"); err == nil { t.Fatal("expected log artifact gone") } dids, err := s.db.GetAllDids() if err != nil { t.Fatal(err) } if len(dids) != 0 { t.Fatalf("expected owner did released, got %v", dids) } if err := s.WipeRepo(ctx, repoDid, "test again"); err != nil { t.Fatalf("second WipeRepo: %v", err) } } func TestWipeRepoKeepsMemberInterest(t *testing.T) { s := newWipeTestSpindle(t) ctx := context.Background() repoDid := syntax.DID("did:plc:repo1") owner := syntax.DID("did:plc:owner1") seedWipeRepo(t, s, repoDid, owner) if err := db.AddSpindleMember(s.db, db.SpindleMember{ Did: owner, Rkey: "m1", Instance: "spindle.example.com", Subject: owner, }); err != nil { t.Fatal(err) } if err := s.WipeRepo(ctx, repoDid, "test"); err != nil { t.Fatalf("WipeRepo: %v", err) } dids, err := s.db.GetAllDids() if err != nil { t.Fatal(err) } if len(dids) != 1 || dids[0] != owner.String() { t.Fatalf("expected member did kept, got %v", dids) } } func TestWipeRepoKeepsCasbinMemberInterest(t *testing.T) { s := newWipeTestSpindle(t) repoDid := syntax.DID("did:plc:repo1") owner := syntax.DID("did:plc:owner1") seedWipeRepo(t, s, repoDid, owner) // the server owner inherits membership via casbin without a members row if err := s.e.AddSpindleMember(rbac.ThisServer, owner.String()); err != nil { t.Fatal(err) } if err := s.WipeRepo(context.Background(), repoDid, "test"); err != nil { t.Fatalf("WipeRepo: %v", err) } dids, err := s.db.GetAllDids() if err != nil { t.Fatal(err) } if len(dids) != 1 || dids[0] != owner.String() { t.Fatalf("expected casbin member did kept, got %v", dids) } } func TestWipeRepoRemovesLocalEngineArtifacts(t *testing.T) { s := newWipeTestSpindle(t) ctx := context.Background() repoDid := syntax.DID("did:plc:repo1") owner := syntax.DID("did:plc:owner1") seedWipeRepo(t, s, repoDid, owner) wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.example.com", Rkey: "p1"}, Name: "w1"} localRef := "logs/" + wid.String() + ".log" if err := s.db.SaveArtifactRef(wid.String(), repoDid.String(), wid, localRef, "h"); err != nil { t.Fatal(err) } if err := s.stores.PutFile(ctx, localRef, writeTempFile(t, "log")); err != nil && len(err) > 0 { t.Fatal(err) } if err := s.WipeRepo(ctx, repoDid, "test"); err != nil { t.Fatalf("WipeRepo: %v", err) } if _, err := s.stores.Open(ctx, localRef); err == nil { t.Fatal("expected local engine artifact gone") } var n int if err := s.db.QueryRow(`select count(*) from mill_artifacts`).Scan(&n); err != nil { t.Fatal(err) } if n != 0 { t.Fatalf("expected mill_artifacts empty, got %d rows", n) } } func TestWipeOwnerWipesAllOwnedRepos(t *testing.T) { s := newWipeTestSpindle(t) ctx := context.Background() owner := syntax.DID("did:plc:owner1") other := syntax.DID("did:plc:owner2") seedWipeRepo(t, s, "did:plc:repo1", owner) seedWipeRepo(t, s, "did:plc:repo2", owner) seedWipeRepo(t, s, "did:plc:repo3", other) if err := s.WipeOwner(ctx, owner, "test"); err != nil { t.Fatalf("WipeOwner: %v", err) } var n int if err := s.db.QueryRow(`select count(*) from repos`).Scan(&n); err != nil { t.Fatal(err) } if n != 1 { t.Fatalf("expected only the other owner's repo left, got %d", n) } repos, err := s.db.ReposByDid("did:plc:repo3") if err != nil || len(repos) != 1 { t.Fatalf("expected repo3 intact, got %v err %v", repos, err) } } func ptr[T any](v T) *T { return &v } func TestIngestAccountTakedownWipes(t *testing.T) { s := newWipeTestSpindle(t) owner := syntax.DID("did:plc:owner1") seedWipeRepo(t, s, "did:plc:repo1", owner) status := "takendown" err := s.ingest()(context.Background(), &jmodels.Event{ Kind: jmodels.EventKindAccount, Did: owner.String(), Account: &comatproto.SyncSubscribeRepos_Account{Did: owner.String(), Active: false, Status: &status}, }) if err != nil { t.Fatalf("ingest account event: %v", err) } var n int if err := s.db.QueryRow(`select count(*) from repos`).Scan(&n); err != nil { t.Fatal(err) } if n != 0 { t.Fatalf("expected repos wiped after takedown, got %d", n) } seedWipeRepo(t, s, "did:plc:repo1", owner) err = s.ingest()(context.Background(), &jmodels.Event{ Kind: jmodels.EventKindAccount, Did: owner.String(), Account: &comatproto.SyncSubscribeRepos_Account{Did: owner.String(), Active: true}, }) if err != nil { t.Fatal(err) } if err := s.db.QueryRow(`select count(*) from repos`).Scan(&n); err != nil { t.Fatal(err) } if n != 1 { t.Fatalf("expected active account untouched, got %d", n) } }