From b66498b90178c5d967c3e76d811d68ba1cb86ba7 Mon Sep 17 00:00:00 2001 From: Lewis Date: Mon, 11 May 2026 14:40:29 +0300 Subject: [PATCH] spindle: key repos and collabs on repoDID Lewis: May this revision serve well! --- spindle/db/collaborators.go | 94 ++++++++++++++++++++ spindle/db/db.go | 107 ++++++++++++++++++++--- spindle/db/repos.go | 105 +++++++++++++++++++--- spindle/engine/engine.go | 5 +- spindle/models/pipeline.go | 5 +- spindle/xrpc/add_secret.go | 14 +-- spindle/xrpc/list_secrets.go | 14 +-- spindle/xrpc/pipeline_cancel_pipeline.go | 12 +-- spindle/xrpc/remove_secret.go | 14 +-- 9 files changed, 314 insertions(+), 56 deletions(-) create mode 100644 spindle/db/collaborators.go diff --git a/spindle/db/collaborators.go b/spindle/db/collaborators.go new file mode 100644 index 00000000..3356f946 --- /dev/null +++ b/spindle/db/collaborators.go @@ -0,0 +1,94 @@ +package db + +import ( + "database/sql" + "fmt" + + "github.com/bluesky-social/indigo/atproto/syntax" +) + +type RepoCollaborator struct { + OwnerDid syntax.DID + Rkey syntax.RecordKey + Subject syntax.DID + RepoDid syntax.DID +} + +func (d *DB) AddRepoCollaborator(c RepoCollaborator) error { + _, err := d.Exec( + `insert into repo_collaborators (owner_did, rkey, subject, repo_did) + values (?, ?, ?, ?) + on conflict(owner_did, rkey) do update set + subject = excluded.subject, + repo_did = excluded.repo_did`, + c.OwnerDid.String(), c.Rkey.String(), c.Subject.String(), c.RepoDid.String(), + ) + return err +} + +func scanCollab(row interface{ Scan(...any) error }) (*RepoCollaborator, error) { + var owner, rkey, subject, repoDid string + if err := row.Scan(&owner, &rkey, &subject, &repoDid); err != nil { + return nil, err + } + return &RepoCollaborator{ + OwnerDid: syntax.DID(owner), + Rkey: syntax.RecordKey(rkey), + Subject: syntax.DID(subject), + RepoDid: syntax.DID(repoDid), + }, nil +} + +func (d *DB) GetRepoCollaborator(ownerDid syntax.DID, rkey syntax.RecordKey) (*RepoCollaborator, error) { + return scanCollab(d.QueryRow( + `select owner_did, rkey, subject, repo_did from repo_collaborators where owner_did = ? and rkey = ?`, + ownerDid.String(), rkey.String(), + )) +} + +func (d *DB) DeleteRepoCollaborator(ownerDid syntax.DID, rkey syntax.RecordKey) error { + res, err := d.Exec(`delete from repo_collaborators where owner_did = ? and rkey = ?`, ownerDid.String(), rkey.String()) + if err != nil { + return err + } + n, err := res.RowsAffected() + if err != nil { + return err + } + if n == 0 { + return sql.ErrNoRows + } + return nil +} + +func (d *DB) DeleteRepoCollaboratorsByRepoDid(repoDid syntax.DID) error { + _, err := d.Exec(`delete from repo_collaborators where repo_did = ?`, repoDid.String()) + if err != nil { + return fmt.Errorf("delete collaborators for %s: %w", repoDid, err) + } + return nil +} + +func (d *DB) ListCollaboratorsByRepoDid(repoDid syntax.DID) ([]RepoCollaborator, error) { + rows, err := d.Query( + `select owner_did, rkey, subject, repo_did from repo_collaborators where repo_did = ?`, + repoDid.String(), + ) + if err != nil { + return nil, fmt.Errorf("list collaborators for %s: %w", repoDid, err) + } + defer rows.Close() + + var out []RepoCollaborator + for rows.Next() { + c, err := scanCollab(rows) + if err != nil { + return nil, err + } + out = append(out, *c) + } + if err := rows.Err(); err != nil { + return nil, err + } + return out, nil +} diff --git a/spindle/db/db.go b/spindle/db/db.go index e59f9bf0..0d89f357 100644 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -1,17 +1,21 @@ package db import ( + "context" "database/sql" + "log/slog" "strings" _ "github.com/mattn/go-sqlite3" + "tangled.org/core/log" + "tangled.org/core/orm" ) type DB struct { *sql.DB } -func Make(dbPath string) (*DB, error) { +func Make(ctx context.Context, dbPath string) (*DB, error) { // https://github.com/mattn/go-sqlite3#connection-string opts := []string{ "_foreign_keys=1", @@ -21,16 +25,21 @@ func Make(dbPath string) (*DB, error) { "_busy_timeout=5000", } + logger := log.FromContext(ctx) + logger = log.SubLogger(logger, "db") + db, err := sql.Open("sqlite3", dbPath+"?"+strings.Join(opts, "&")) if err != nil { return nil, err } - // NOTE: If any other migration is added here, you MUST - // copy the pattern in appview: use a single sql.Conn - // for every migration. + conn, err := db.Conn(ctx) + if err != nil { + return nil, err + } + defer conn.Close() - _, err = db.Exec(` + _, err = conn.ExecContext(ctx, ` create table if not exists _jetstream ( id integer primary key autoincrement, last_time_us integer not null @@ -41,13 +50,26 @@ func Make(dbPath string) (*DB, error) { ); create table if not exists repos ( - id integer primary key autoincrement, - knot text not null, - owner text not null, - name text not null, - addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), + id integer primary key autoincrement, + knot text not null, + owner text not null, + rkey text not null, + repo_did text, + created_at text, + addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), + + unique(owner, rkey) + ); - unique(owner, name) + create table if not exists repo_collaborators ( + id integer primary key autoincrement, + owner_did text not null, + rkey text not null, + subject text not null, + repo_did text not null, + addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), + + unique(owner_did, rkey) ); create table if not exists spindle_members ( @@ -72,14 +94,77 @@ func Make(dbPath string) (*DB, error) { event text not null, -- json created integer not null -- unix nanos ); + + create table if not exists migrations ( + id integer primary key autoincrement, + name text unique + ); `) if err != nil { return nil, err } + if err := runMigrations(ctx, conn, logger); err != nil { + return nil, err + } + return &DB{db}, nil } +func runMigrations(_ context.Context, conn *sql.Conn, logger *slog.Logger) error { + return orm.RunMigration(conn, logger, "repos-to-repo-did", func(tx *sql.Tx) error { + var hasName int + if err := tx.QueryRow( + `select count(*) from pragma_table_info('repos') where name = 'name'`, + ).Scan(&hasName); err != nil { + return err + } + + if hasName > 0 { + var totalRows, copiedRows int + if err := tx.QueryRow(`select count(*) from repos`).Scan(&totalRows); err != nil { + return err + } + if err := tx.QueryRow(`select count(*) from repos where coalesce(name, '') <> ''`).Scan(&copiedRows); err != nil { + return err + } + if dropped := totalRows - copiedRows; dropped > 0 { + logger.Warn("dropping repo rows with empty name during migration", "dropped", dropped, "kept", copiedRows) + } + + if _, err := tx.Exec(` + create table if not exists repos_new ( + id integer primary key autoincrement, + knot text not null, + owner text not null, + rkey text not null, + repo_did text, + created_at text, + addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), + + unique(owner, rkey) + ); + + insert into repos_new (id, knot, owner, rkey, addedAt) + select id, knot, owner, name, addedAt from repos where coalesce(name, '') <> ''; + + drop table repos; + alter table repos_new rename to repos; + `); err != nil { + return err + } + } + + _, err := tx.Exec(` + create index if not exists idx_repos_repo_did on repos(repo_did); + create index if not exists idx_repos_owner_repo_did on repos(owner, repo_did); + create index if not exists idx_repo_collaborators_repo_did + on repo_collaborators(repo_did); + `) + return err + }) +} + func (d *DB) SaveLastTimeUs(lastTimeUs int64) error { _, err := d.Exec(` insert into _jetstream (id, last_time_us) diff --git a/spindle/db/repos.go b/spindle/db/repos.go index d6ecf49d..89b7cc1d 100644 --- a/spindle/db/repos.go +++ b/spindle/db/repos.go @@ -1,18 +1,56 @@ package db +import ( + "database/sql" + + "github.com/bluesky-social/indigo/atproto/syntax" +) + type Repo struct { - Knot string - Owner string - Name string + Knot string + Owner syntax.DID + Rkey syntax.RecordKey + RepoDid syntax.DID + CreatedAt string } -func (d *DB) AddRepo(knot, owner, name string) error { - _, err := d.Exec(`insert or ignore into repos (knot, owner, name) values (?, ?, ?)`, knot, owner, name) +func (d *DB) AddRepo(repo Repo) error { + var createdAt sql.NullString + if repo.CreatedAt != "" { + createdAt = sql.NullString{String: repo.CreatedAt, Valid: true} + } + _, err := d.Exec( + `insert into repos (knot, owner, rkey, repo_did, created_at) + values (?, ?, ?, ?, ?) + on conflict(owner, rkey) do update set + knot = excluded.knot, + repo_did = excluded.repo_did, + created_at = coalesce(excluded.created_at, repos.created_at)`, + repo.Knot, repo.Owner.String(), repo.Rkey.String(), repo.RepoDid.String(), createdAt, + ) return err } +func (d *DB) CollapseRepoSiblings(owner, repoDid syntax.DID) (int64, error) { + res, err := d.Exec( + `delete from repos + where owner = ? + and repo_did = ? + and created_at is not null + and created_at < ( + select max(created_at) from repos + where owner = ? and repo_did = ? and created_at is not null + )`, + owner.String(), repoDid.String(), owner.String(), repoDid.String(), + ) + if err != nil { + return 0, err + } + return res.RowsAffected() +} + func (d *DB) Knots() ([]string, error) { - rows, err := d.Query(`select knot from repos`) + rows, err := d.Query(`select distinct knot from repos`) if err != nil { return nil, err } @@ -27,23 +65,64 @@ func (d *DB) Knots() ([]string, error) { knots = append(knots, knot) } - if err = rows.Err(); err != nil { + if err := rows.Err(); err != nil { return nil, err } return knots, nil } -func (d *DB) GetRepo(knot, owner, name string) (*Repo, error) { - var repo Repo +func scanRepo(row interface{ Scan(...any) error }) (*Repo, error) { + var knot, owner, rkey, repoDid string + if err := row.Scan(&knot, &owner, &rkey, &repoDid); err != nil { + return nil, err + } + return &Repo{ + Knot: knot, + Owner: syntax.DID(owner), + Rkey: syntax.RecordKey(rkey), + RepoDid: syntax.DID(repoDid), + }, nil +} - query := "select knot, owner, name from repos where knot = ? and owner = ? and name = ?" - err := d.DB.QueryRow(query, knot, owner, name). - Scan(&repo.Knot, &repo.Owner, &repo.Name) +func (d *DB) GetRepoByDid(repoDid syntax.DID) (*Repo, error) { + return scanRepo(d.QueryRow( + `select knot, owner, rkey, coalesce(repo_did, '') from repos where repo_did = ?`, + repoDid.String(), + )) +} + +func (d *DB) GetRepoByOwnerRkey(owner syntax.DID, rkey syntax.RecordKey) (*Repo, error) { + return scanRepo(d.QueryRow( + `select knot, owner, rkey, coalesce(repo_did, '') from repos where owner = ? and rkey = ?`, + owner.String(), rkey.String(), + )) +} +func (d *DB) AllRepos() ([]Repo, error) { + rows, err := d.Query(`select knot, owner, rkey, coalesce(repo_did, '') from repos`) if err != nil { return nil, err } + defer rows.Close() + + var repos []Repo + for rows.Next() { + r, err := scanRepo(rows) + if err != nil { + return nil, err + } + repos = append(repos, *r) + } - return &repo, nil + if err := rows.Err(); err != nil { + return nil, err + } + + return repos, nil +} + +func (d *DB) DeleteRepoByOwnerRkey(owner syntax.DID, rkey syntax.RecordKey) error { + _, err := d.Exec(`delete from repos where owner = ? and rkey = ?`, owner.String(), rkey.String()) + return err } diff --git a/spindle/engine/engine.go b/spindle/engine/engine.go index d09cb77f..b7927d2e 100644 --- a/spindle/engine/engine.go +++ b/spindle/engine/engine.go @@ -8,7 +8,6 @@ import ( "path/filepath" "sync" - securejoin "github.com/cyphar/filepath-securejoin" "tangled.org/core/notifier" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" @@ -26,8 +25,8 @@ func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, d // extract secrets var allSecrets []secrets.UnlockedSecret - if didSlashRepo, err := securejoin.SecureJoin(pipeline.RepoOwner, pipeline.RepoName); err == nil { - if res, err := vault.GetSecretsUnlocked(ctx, secrets.RepoIdentifier(didSlashRepo)); err == nil { + if pipeline.RepoDid != "" { + if res, err := vault.GetSecretsUnlocked(ctx, secrets.RepoIdentifier(pipeline.RepoDid.String())); err == nil { allSecrets = res } } diff --git a/spindle/models/pipeline.go b/spindle/models/pipeline.go index 8f35d622..fec501dc 100644 --- a/spindle/models/pipeline.go +++ b/spindle/models/pipeline.go @@ -1,8 +1,9 @@ package models +import "github.com/bluesky-social/indigo/atproto/syntax" + type Pipeline struct { - RepoOwner string - RepoName string + RepoDid syntax.DID Workflows map[Engine][]Workflow } diff --git a/spindle/xrpc/add_secret.go b/spindle/xrpc/add_secret.go index 1812ac39..92047cc9 100644 --- a/spindle/xrpc/add_secret.go +++ b/spindle/xrpc/add_secret.go @@ -9,7 +9,6 @@ import ( "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/xrpc" - securejoin "github.com/cyphar/filepath-securejoin" "tangled.org/core/api/tangled" "tangled.org/core/rbac" "tangled.org/core/spindle/secrets" @@ -61,24 +60,25 @@ func (x *Xrpc) AddSecret(w http.ResponseWriter, r *http.Request) { return } - if _, ok := resp.Value.Val.(*tangled.Repo); !ok { + repoRec, ok := resp.Value.Val.(*tangled.Repo) + if !ok { fail(xrpcerr.RepoNotFoundError) return } - didPath, err := securejoin.SecureJoin(ident.DID.String(), repoAt.RecordKey().String()) - if err != nil { - fail(xrpcerr.GenericError(err)) + if repoRec.RepoDid == nil || *repoRec.RepoDid == "" { + fail(xrpcerr.GenericError(fmt.Errorf("repo record %s has no repoDid", repoAt))) return } + repoDid := *repoRec.RepoDid - if ok, err := x.Enforcer.IsSettingsAllowed(actorDid.String(), rbac.ThisServer, didPath); !ok || err != nil { + if ok, err := x.Enforcer.IsSettingsAllowed(actorDid.String(), rbac.ThisServer, repoDid); !ok || err != nil { l.Error("insufficient permissions", "did", actorDid.String()) writeError(w, xrpcerr.AccessControlError(actorDid.String()), http.StatusUnauthorized) return } secret := secrets.UnlockedSecret{ - Repo: secrets.RepoIdentifier(didPath), + Repo: secrets.RepoIdentifier(repoDid), Key: data.Key, Value: data.Value, CreatedAt: time.Now(), diff --git a/spindle/xrpc/list_secrets.go b/spindle/xrpc/list_secrets.go index 1d935ee1..1db58030 100644 --- a/spindle/xrpc/list_secrets.go +++ b/spindle/xrpc/list_secrets.go @@ -9,7 +9,6 @@ import ( "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/xrpc" - securejoin "github.com/cyphar/filepath-securejoin" "tangled.org/core/api/tangled" "tangled.org/core/rbac" "tangled.org/core/spindle/secrets" @@ -56,23 +55,24 @@ func (x *Xrpc) ListSecrets(w http.ResponseWriter, r *http.Request) { return } - if _, ok := resp.Value.Val.(*tangled.Repo); !ok { + repoRec, ok := resp.Value.Val.(*tangled.Repo) + if !ok { fail(xrpcerr.RepoNotFoundError) return } - didPath, err := securejoin.SecureJoin(ident.DID.String(), repoAt.RecordKey().String()) - if err != nil { - fail(xrpcerr.GenericError(err)) + if repoRec.RepoDid == nil || *repoRec.RepoDid == "" { + fail(xrpcerr.GenericError(fmt.Errorf("repo record %s has no repoDid", repoAt))) return } + repoDid := *repoRec.RepoDid - if ok, err := x.Enforcer.IsSettingsAllowed(actorDid.String(), rbac.ThisServer, didPath); !ok || err != nil { + if ok, err := x.Enforcer.IsSettingsAllowed(actorDid.String(), rbac.ThisServer, repoDid); !ok || err != nil { l.Error("insufficient permissions", "did", actorDid.String()) writeError(w, xrpcerr.AccessControlError(actorDid.String()), http.StatusUnauthorized) return } - ls, err := x.Vault.GetSecretsLocked(r.Context(), secrets.RepoIdentifier(didPath)) + ls, err := x.Vault.GetSecretsLocked(r.Context(), secrets.RepoIdentifier(repoDid)) if err != nil { l.Error("failed to get secret from vault", "did", actorDid.String(), "err", err) writeError(w, xrpcerr.GenericError(err), http.StatusInternalServerError) diff --git a/spindle/xrpc/pipeline_cancel_pipeline.go b/spindle/xrpc/pipeline_cancel_pipeline.go index fb7877bb..cae58316 100644 --- a/spindle/xrpc/pipeline_cancel_pipeline.go +++ b/spindle/xrpc/pipeline_cancel_pipeline.go @@ -9,7 +9,6 @@ import ( "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/xrpc" - securejoin "github.com/cyphar/filepath-securejoin" "tangled.org/core/api/tangled" "tangled.org/core/rbac" "tangled.org/core/spindle/models" @@ -66,18 +65,19 @@ func (x *Xrpc) CancelPipeline(w http.ResponseWriter, r *http.Request) { return } - if _, ok := resp.Value.Val.(*tangled.Repo); !ok { + repoRec, ok := resp.Value.Val.(*tangled.Repo) + if !ok { fail(xrpcerr.RepoNotFoundError) return } - didSlashRepo, err := securejoin.SecureJoin(ident.DID.String(), repoAt.RecordKey().String()) - if err != nil { - fail(xrpcerr.GenericError(err)) + if repoRec.RepoDid == nil || *repoRec.RepoDid == "" { + fail(xrpcerr.GenericError(fmt.Errorf("repo record %s has no repoDid", repoAt))) return } + repoDid := *repoRec.RepoDid // TODO: fine-grained role based control - isRepoOwner, err := x.Enforcer.IsRepoOwner(actorDid.String(), rbac.ThisServer, didSlashRepo) + isRepoOwner, err := x.Enforcer.IsRepoOwner(actorDid.String(), rbac.ThisServer, repoDid) if err != nil || !isRepoOwner { fail(xrpcerr.AccessControlError(actorDid.String())) return diff --git a/spindle/xrpc/remove_secret.go b/spindle/xrpc/remove_secret.go index 7371aa2e..12b22342 100644 --- a/spindle/xrpc/remove_secret.go +++ b/spindle/xrpc/remove_secret.go @@ -8,7 +8,6 @@ import ( "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/bluesky-social/indigo/xrpc" - securejoin "github.com/cyphar/filepath-securejoin" "tangled.org/core/api/tangled" "tangled.org/core/rbac" "tangled.org/core/spindle/secrets" @@ -55,24 +54,25 @@ func (x *Xrpc) RemoveSecret(w http.ResponseWriter, r *http.Request) { return } - if _, ok := resp.Value.Val.(*tangled.Repo); !ok { + repoRec, ok := resp.Value.Val.(*tangled.Repo) + if !ok { fail(xrpcerr.RepoNotFoundError) return } - didPath, err := securejoin.SecureJoin(ident.DID.String(), repoAt.RecordKey().String()) - if err != nil { - fail(xrpcerr.GenericError(err)) + if repoRec.RepoDid == nil || *repoRec.RepoDid == "" { + fail(xrpcerr.GenericError(fmt.Errorf("repo record %s has no repoDid", repoAt))) return } + repoDid := *repoRec.RepoDid - if ok, err := x.Enforcer.IsSettingsAllowed(actorDid.String(), rbac.ThisServer, didPath); !ok || err != nil { + if ok, err := x.Enforcer.IsSettingsAllowed(actorDid.String(), rbac.ThisServer, repoDid); !ok || err != nil { l.Error("insufficient permissions", "did", actorDid.String()) writeError(w, xrpcerr.AccessControlError(actorDid.String()), http.StatusUnauthorized) return } secret := secrets.Secret[any]{ - Repo: secrets.RepoIdentifier(didPath), + Repo: secrets.RepoIdentifier(repoDid), Key: data.Key, } err = x.Vault.RemoveSecret(r.Context(), secret) -- 2.51.2