diff --git a/spindle/models/clone.go b/spindle/models/clone.go --- a/spindle/models/clone.go +++ b/spindle/models/clone.go @@ -103,10 +103,6 @@ } // BuildRepoURL constructs the repository URL from repo metadata. func BuildRepoURL(repo *tangled.Pipeline_TriggerRepo, devMode bool) string { - if repo == nil { - return "" - } - scheme := "https://" if devMode { scheme = "http://" @@ -120,14 +116,7 @@ if devMode && strings.Contains(host, "localhost") { host = strings.ReplaceAll(host, "localhost", "host.docker.internal") } - switch { - case repo.RepoDid != nil: - return fmt.Sprintf("%s%s/%s", scheme, host, *repo.RepoDid) - case repo.Repo != nil: - return fmt.Sprintf("%s%s/%s/%s", scheme, host, repo.Did, *repo.Repo) - default: - return "" - } + return fmt.Sprintf("%s%s/%s", scheme, host, *repo.RepoDid) } // buildFetchArgs constructs the arguments for git fetch based on clone options diff --git a/spindle/models/clone_test.go b/spindle/models/clone_test.go --- a/spindle/models/clone_test.go +++ b/spindle/models/clone_test.go @@ -26,9 +26,10 @@ OldSha: "def456", Ref: "refs/heads/main", }, Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "example.com", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "example.com", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } @@ -64,7 +65,7 @@ } if !strings.Contains(allCmds, "git checkout FETCH_HEAD") { t.Error("Commands should contain 'git checkout FETCH_HEAD'") } - if !strings.Contains(allCmds, "https://example.com/did:plc:user123/my-repo") { + if !strings.Contains(allCmds, "https://example.com/did:plc:boltless") { t.Error("Commands should contain expected repo URL") } } @@ -85,9 +86,10 @@ TargetBranch: "main", Action: "opened", }, Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "example.com", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "example.com", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } @@ -112,9 +114,10 @@ Manual: &tangled.Pipeline_ManualTriggerData{ Inputs: nil, }, Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "example.com", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "example.com", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } @@ -143,9 +146,10 @@ Push: &tangled.Pipeline_PushTriggerData{ NewSha: "abc123", }, Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "example.com", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "example.com", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } @@ -173,9 +177,10 @@ Push: &tangled.Pipeline_PushTriggerData{ NewSha: "abc123", }, Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "localhost:3000", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "localhost:3000", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } @@ -183,7 +188,7 @@ step := BuildCloneStep(twf, tr, true) // In dev mode, should use http:// and replace localhost with host.docker.internal allCmds := strings.Join(step.Commands(), " ") - expectedURL := "http://host.docker.internal:3000/did:plc:user123/my-repo" + expectedURL := "http://host.docker.internal:3000/did:plc:boltless" if !strings.Contains(allCmds, expectedURL) { t.Errorf("Expected dev mode URL '%s' in commands", expectedURL) } @@ -203,9 +208,10 @@ Push: &tangled.Pipeline_PushTriggerData{ NewSha: "abc123", }, Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "example.com", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "example.com", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } @@ -234,9 +240,10 @@ Push: &tangled.Pipeline_PushTriggerData{ NewSha: "abc123", }, Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "example.com", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "example.com", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } @@ -259,9 +266,10 @@ tr := tangled.Pipeline_TriggerMetadata{ Kind: string(workflow.TriggerKindPush), Push: nil, // Nil push data should create error step Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "example.com", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "example.com", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } @@ -292,9 +300,10 @@ tr := tangled.Pipeline_TriggerMetadata{ Kind: string(workflow.TriggerKindPullRequest), PullRequest: nil, // Nil PR data should create error step Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "example.com", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "example.com", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } @@ -321,9 +330,10 @@ } tr := tangled.Pipeline_TriggerMetadata{ Kind: "unknown_trigger", Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "example.com", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "example.com", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } @@ -350,9 +360,10 @@ Push: &tangled.Pipeline_PushTriggerData{ NewSha: "abc123", }, Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "example.com", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "example.com", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } diff --git a/spindle/models/pipeline_env_test.go b/spindle/models/pipeline_env_test.go --- a/spindle/models/pipeline_env_test.go +++ b/spindle/models/pipeline_env_test.go @@ -19,6 +19,7 @@ Repo: &tangled.Pipeline_TriggerRepo{ Knot: "example.com", Did: "did:plc:user123", Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), DefaultBranch: "main", }, } @@ -65,8 +66,8 @@ } if env["TANGLED_REPO_DEFAULT_BRANCH"] != "main" { t.Errorf("Expected TANGLED_REPO_DEFAULT_BRANCH='main', got '%s'", env["TANGLED_REPO_DEFAULT_BRANCH"]) } - if env["TANGLED_REPO_URL"] != "https://example.com/did:plc:user123/my-repo" { - t.Errorf("Expected TANGLED_REPO_URL='https://example.com/did:plc:user123/my-repo', got '%s'", env["TANGLED_REPO_URL"]) + if env["TANGLED_REPO_URL"] != "https://example.com/did:plc:boltless" { + t.Errorf("Expected TANGLED_REPO_URL='https://example.com/did:plc:boltless', got '%s'", env["TANGLED_REPO_URL"]) } } @@ -79,9 +80,10 @@ OldSha: "000000000000", Ref: "refs/tags/v1.2.3", }, Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "example.com", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "example.com", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } id := PipelineId{ @@ -111,9 +113,10 @@ SourceSha: "pr-sha-789", Action: "opened", }, Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "example.com", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "example.com", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } id := PipelineId{ @@ -166,9 +169,10 @@ {Key: "environment", Value: "production"}, }, }, Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "example.com", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "example.com", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } id := PipelineId{ @@ -202,9 +206,10 @@ NewSha: "abc123", Ref: "refs/heads/main", }, Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "localhost:3000", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "localhost:3000", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } id := PipelineId{ @@ -214,7 +219,7 @@ } env := PipelineEnvVars(tr, id, true) // Dev mode should use http:// and replace localhost with host.docker.internal - expectedURL := "http://host.docker.internal:3000/did:plc:user123/my-repo" + expectedURL := "http://host.docker.internal:3000/did:plc:boltless" if env["TANGLED_REPO_URL"] != expectedURL { t.Errorf("Expected TANGLED_REPO_URL='%s', got '%s'", expectedURL, env["TANGLED_REPO_URL"]) } @@ -237,9 +242,10 @@ tr := &tangled.Pipeline_TriggerMetadata{ Kind: string(workflow.TriggerKindPush), Push: nil, Repo: &tangled.Pipeline_TriggerRepo{ - Knot: "example.com", - Did: "did:plc:user123", - Repo: sp("my-repo"), + Knot: "example.com", + Did: "did:plc:user123", + Repo: sp("my-repo"), + RepoDid: sp("did:plc:boltless"), }, } id := PipelineId{ diff --git a/spindle/secret_copy.go b/spindle/secret_copy.go new file mode 100644 --- /dev/null +++ b/spindle/secret_copy.go @@ -0,0 +1,39 @@ +package spindle + +import ( + "context" + "errors" + "fmt" + + "tangled.org/core/spindle/secrets" +) + +func copyRepoSecrets(ctx context.Context, mgr secrets.Manager, src, dst secrets.RepoIdentifier) (int, error) { + cur, err := mgr.GetSecretsUnlocked(ctx, src) + if err != nil { + return 0, fmt.Errorf("get %s: %w", src, err) + } + var step func(remaining []secrets.UnlockedSecret, copied int) (int, error) + step = func(remaining []secrets.UnlockedSecret, copied int) (int, error) { + if len(remaining) == 0 { + return copied, nil + } + s := remaining[0] + addErr := mgr.AddSecret(ctx, secrets.UnlockedSecret{ + Repo: dst, + Key: s.Key, + Value: s.Value, + CreatedAt: s.CreatedAt, + CreatedBy: s.CreatedBy, + }) + switch { + case addErr == nil: + return step(remaining[1:], copied+1) + case errors.Is(addErr, secrets.ErrKeyAlreadyPresent): + return step(remaining[1:], copied) + default: + return copied, fmt.Errorf("add %s/%s: %w", dst, s.Key, addErr) + } + } + return step(cur, 0) +} diff --git a/spindle/secrets/openbao.go b/spindle/secrets/openbao.go --- a/spindle/secrets/openbao.go +++ b/spindle/secrets/openbao.go @@ -98,11 +98,15 @@ v.logger.Debug("secret already exists", "path", secretPath) return ErrKeyAlreadyPresent } + createdAt := secret.CreatedAt + if createdAt.IsZero() { + createdAt = time.Now() + } secretData := map[string]interface{}{ "value": secret.Value, "repo": string(secret.Repo), "key": secret.Key, - "created_at": secret.CreatedAt.Format(time.RFC3339), + "created_at": createdAt.UTC().Format(time.RFC3339), "created_by": secret.CreatedBy.String(), } diff --git a/spindle/secrets/sqlite.go b/spindle/secrets/sqlite.go --- a/spindle/secrets/sqlite.go +++ b/spindle/secrets/sqlite.go @@ -64,11 +64,15 @@ } func (s *SqliteManager) AddSecret(ctx context.Context, secret UnlockedSecret) error { query := fmt.Sprintf(` - insert or ignore into %s (repo, key, value, created_by) - values (?, ?, ?, ?); + insert or ignore into %s (repo, key, value, created_at, created_by) + values (?, ?, ?, ?, ?); `, s.tableName) - res, err := s.db.ExecContext(ctx, query, secret.Repo, secret.Key, secret.Value, secret.CreatedBy) + createdAt := secret.CreatedAt + if createdAt.IsZero() { + createdAt = time.Now() + } + res, err := s.db.ExecContext(ctx, query, secret.Repo, secret.Key, secret.Value, createdAt.UTC().Format(time.RFC3339), secret.CreatedBy) if err != nil { return err } diff --git a/spindle/server.go b/spindle/server.go --- a/spindle/server.go +++ b/spindle/server.go @@ -100,6 +100,10 @@ default: return nil, fmt.Errorf("unknown secrets provider: %s", cfg.Server.Secrets.Provider) } + if err := runStartupMigrations(ctx, d, vault, logger); err != nil { + return nil, fmt.Errorf("failed to run startup migrations: %w", err) + } + jq := queue.NewQueue(cfg.Server.QueueSize, cfg.Server.MaxJobCount) logger.Info("initialized queue", "queueSize", cfg.Server.QueueSize, "numWorkers", cfg.Server.MaxJobCount) diff --git a/spindle/startup_migrations.go b/spindle/startup_migrations.go new file mode 100644 --- /dev/null +++ b/spindle/startup_migrations.go @@ -0,0 +1,80 @@ +package spindle + +import ( + "context" + "database/sql" + "fmt" + "log/slog" + + "tangled.org/core/orm" + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/secrets" +) + +func runStartupMigrations(ctx context.Context, d *db.DB, vault secrets.Manager, logger *slog.Logger) error { + conn, err := d.DB.Conn(ctx) + if err != nil { + return fmt.Errorf("acquire spindle conn: %w", err) + } + defer conn.Close() + + return orm.RunMigration(conn, logger, "copy-owner-rkey-secrets-to-repo-did", func(tx *sql.Tx) error { + return copyOwnerRkeySecretsToRepoDid(ctx, tx, vault, logger) + }) +} + +type repoSecretPair struct { + oldID, newID secrets.RepoIdentifier +} + +func loadRepoSecretPairs(ctx context.Context, tx *sql.Tx) ([]repoSecretPair, error) { + rows, err := tx.QueryContext(ctx, + `select owner, rkey, repo_did from repos + where repo_did is not null and repo_did <> ''`, + ) + if err != nil { + return nil, fmt.Errorf("select repos: %w", err) + } + defer rows.Close() + + var collect func(acc []repoSecretPair) ([]repoSecretPair, error) + collect = func(acc []repoSecretPair) ([]repoSecretPair, error) { + if !rows.Next() { + return acc, rows.Err() + } + var owner, rkey, repoDid string + if err := rows.Scan(&owner, &rkey, &repoDid); err != nil { + return acc, fmt.Errorf("scan repos row: %w", err) + } + return collect(append(acc, repoSecretPair{ + oldID: secrets.RepoIdentifier(owner + "/" + rkey), + newID: secrets.RepoIdentifier(repoDid), + })) + } + return collect(nil) +} + +func copyOwnerRkeySecretsToRepoDid(ctx context.Context, tx *sql.Tx, vault secrets.Manager, logger *slog.Logger) error { + pairs, err := loadRepoSecretPairs(ctx, tx) + if err != nil { + return err + } + + var step func(remaining []repoSecretPair, totalCopied int) error + step = func(remaining []repoSecretPair, totalCopied int) error { + if len(remaining) == 0 { + logger.Info("secret copy migration complete", "rows", len(pairs), "copied", totalCopied) + return nil + } + p := remaining[0] + n, err := copyRepoSecrets(ctx, vault, p.oldID, p.newID) + if err != nil { + return fmt.Errorf("copy %s -> %s: %w", p.oldID, p.newID, err) + } + if n > 0 { + logger.Info("secrets copied", "old", p.oldID, "new", p.newID, "count", n) + } + return step(remaining[1:], totalCopied+n) + } + return step(pairs, 0) +} diff --git a/spindle/startup_migrations_test.go b/spindle/startup_migrations_test.go new file mode 100644 --- /dev/null +++ b/spindle/startup_migrations_test.go @@ -0,0 +1,209 @@ +package spindle + +import ( + "context" + "io" + "log/slog" + "path/filepath" + "testing" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/secrets" +) + +func newTestSpindleDB(t *testing.T) *db.DB { + t.Helper() + d, err := db.Make(context.Background(), filepath.Join(t.TempDir(), "spindle.db")) + if err != nil { + t.Fatalf("db.Make: %v", err) + } + t.Cleanup(func() { d.Close() }) + return d +} + +func newTestVault(t *testing.T) *secrets.SqliteManager { + t.Helper() + vault, err := secrets.NewSQLiteManager(filepath.Join(t.TempDir(), "vault.db")) + if err != nil { + t.Fatalf("vault.New: %v", err) + } + return vault +} + +func mustAddRepo(t *testing.T, d *db.DB, knot, owner, rkey, repoDid string) { + t.Helper() + if err := d.AddRepo(db.Repo{ + Knot: knot, + Owner: syntax.DID(owner), + Rkey: syntax.RecordKey(rkey), + RepoDid: syntax.DID(repoDid), + }); err != nil { + t.Fatalf("AddRepo(%s): %v", rkey, err) + } +} + +func mustAddSecret(t *testing.T, vault secrets.Manager, repo, key, value string, createdAt time.Time, by string) { + t.Helper() + err := vault.AddSecret(context.Background(), secrets.UnlockedSecret{ + Repo: secrets.RepoIdentifier(repo), + Key: key, + Value: value, + CreatedAt: createdAt, + CreatedBy: syntax.DID(by), + }) + if err != nil { + t.Fatalf("AddSecret(%s/%s): %v", repo, key, err) + } +} + +func TestStartupMigrations_CopyOwnerRkeySecretsToRepoDid(t *testing.T) { + ctx := context.Background() + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + + d := newTestSpindleDB(t) + vault := newTestVault(t) + + owner := "did:plc:akshay" + migratedRepoDid := "did:plc:boltless" + skippedRkey := "3kspindlerkey00b" + migratedRkey := "3kspindlerkey00a" + + mustAddRepo(t, d, "knot.test", owner, migratedRkey, migratedRepoDid) + mustAddRepo(t, d, "knot.test", owner, skippedRkey, "") + + created := time.Date(2024, 6, 1, 12, 0, 0, 0, time.UTC) + oldRepoKey := owner + "/" + migratedRkey + skippedKey := owner + "/" + skippedRkey + + mustAddSecret(t, vault, oldRepoKey, "API_KEY", "alpha", created, owner) + mustAddSecret(t, vault, oldRepoKey, "DB_PASSWORD", "bravo", created.Add(1*time.Hour), owner) + mustAddSecret(t, vault, skippedKey, "STRAY", "delta", created, owner) + + if err := runStartupMigrations(ctx, d, vault, logger); err != nil { + t.Fatalf("first migration run: %v", err) + } + + copied, err := vault.GetSecretsUnlocked(ctx, secrets.RepoIdentifier(migratedRepoDid)) + if err != nil { + t.Fatalf("GetSecretsUnlocked(new): %v", err) + } + if len(copied) != 2 { + t.Fatalf("expected 2 secrets under new repo_did key, got %d", len(copied)) + } + + want := map[string]struct { + value string + createdAt time.Time + }{ + "API_KEY": {"alpha", created}, + "DB_PASSWORD": {"bravo", created.Add(1 * time.Hour)}, + } + for _, s := range copied { + w, ok := want[s.Key] + if !ok { + t.Errorf("unexpected key %q under %s", s.Key, migratedRepoDid) + continue + } + if s.Value != w.value { + t.Errorf("%s: value got %q, want %q", s.Key, s.Value, w.value) + } + if !s.CreatedAt.Equal(w.createdAt) { + t.Errorf("%s: CreatedAt got %s, want %s", s.Key, s.CreatedAt, w.createdAt) + } + if string(s.Repo) != migratedRepoDid { + t.Errorf("%s: Repo got %s, want %s", s.Key, s.Repo, migratedRepoDid) + } + } + + orig, err := vault.GetSecretsUnlocked(ctx, secrets.RepoIdentifier(oldRepoKey)) + if err != nil { + t.Fatalf("GetSecretsUnlocked(old): %v", err) + } + if len(orig) != 2 { + t.Errorf("expected old-key secrets preserved, got %d", len(orig)) + } + + stray, err := vault.GetSecretsUnlocked(ctx, secrets.RepoIdentifier(skippedKey)) + if err != nil { + t.Fatalf("GetSecretsUnlocked(skipped): %v", err) + } + if len(stray) != 1 { + t.Errorf("expected skipped repo's old-key secret untouched, got %d", len(stray)) + } + + if err := runStartupMigrations(ctx, d, vault, logger); err != nil { + t.Fatalf("second migration run: %v", err) + } + + again, err := vault.GetSecretsUnlocked(ctx, secrets.RepoIdentifier(migratedRepoDid)) + if err != nil { + t.Fatalf("GetSecretsUnlocked(new) after re-run: %v", err) + } + if len(again) != 2 { + t.Errorf("re-run should not duplicate or drop secrets, got %d", len(again)) + } + + var marked int + if err := d.QueryRow( + `select count(*) from migrations where name = ?`, + "copy-owner-rkey-secrets-to-repo-did", + ).Scan(&marked); err != nil { + t.Fatalf("query migrations: %v", err) + } + if marked != 1 { + t.Errorf("expected migration recorded exactly once, got %d", marked) + } +} + +func TestStartupMigrations_NoRepos(t *testing.T) { + ctx := context.Background() + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + d := newTestSpindleDB(t) + vault := newTestVault(t) + + if err := runStartupMigrations(ctx, d, vault, logger); err != nil { + t.Fatalf("migration on empty db: %v", err) + } +} + +func TestStartupMigrations_PartialPreExisting(t *testing.T) { + ctx := context.Background() + logger := slog.New(slog.NewTextHandler(io.Discard, nil)) + d := newTestSpindleDB(t) + vault := newTestVault(t) + + owner := "did:plc:akshay" + repoDid := "did:plc:boltless" + rkey := "3kspindlerkey00a" + mustAddRepo(t, d, "knot.test", owner, rkey, repoDid) + + created := time.Date(2024, 6, 1, 12, 0, 0, 0, time.UTC) + oldKey := owner + "/" + rkey + mustAddSecret(t, vault, oldKey, "API_KEY", "alpha", created, owner) + mustAddSecret(t, vault, oldKey, "DB_PASSWORD", "bravo", created, owner) + + mustAddSecret(t, vault, repoDid, "API_KEY", "pre-existing", created.Add(-24*time.Hour), owner) + + if err := runStartupMigrations(ctx, d, vault, logger); err != nil { + t.Fatalf("migration: %v", err) + } + + got, err := vault.GetSecretsUnlocked(ctx, secrets.RepoIdentifier(repoDid)) + if err != nil { + t.Fatalf("GetSecretsUnlocked: %v", err) + } + if len(got) != 2 { + t.Fatalf("expected 2 secrets under new key, got %d", len(got)) + } + for _, s := range got { + if s.Key == "API_KEY" && s.Value != "pre-existing" { + t.Errorf("API_KEY should preserve pre-existing value, got %q", s.Value) + } + if s.Key == "DB_PASSWORD" && s.Value != "bravo" { + t.Errorf("DB_PASSWORD should be copied, got %q", s.Value) + } + } +}