Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352package 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) }}