From dcde726528b4f5d179e7ddf139e200b8fd4a4fd5 Mon Sep 17 00:00:00 2001 From: Lewis Date: Tue, 02 Jun 2026 11:09:18 +0000 Subject: [PATCH] knotserver/db: knot-owned member & collaborator tables, pagination, and ACL update events Lewis: May this revision serve well! --- knotserver/db/aclupdate.go | 25 +++++++++++++++++++++++++ knotserver/db/collaborator.go | 94 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotserver/db/db.go | 84 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotserver/db/events.go | 38 ++++++++++++++++++++++++++++++++++++++ knotserver/db/member.go | 84 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------- knotserver/db/page.go | 74 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ notifier/notifier.go | 3 +++ 7 file(s) changed, 389 insertion(s)(+), 13 deletion(s)(-) diff --git a/knotserver/db/aclupdate.go b/knotserver/db/aclupdate.go new file mode 100644 --- /dev/null +++ b/knotserver/db/aclupdate.go @@ -0,0 +1,25 @@ +package db + +const ( + KnotMemberUpdateNSID = "sh.tangled.knot.memberUpdate" + RepoCollaboratorUpdateNSID = "sh.tangled.repo.collaboratorUpdate" +) + +type AclOp string + +const ( + AclOpAdd AclOp = "add" + AclOpRemove AclOp = "remove" +) + +type KnotMemberUpdate struct { + Op AclOp `json:"op"` + Subject string `json:"subject"` +} + +// NOTE: no "addedBy" so for now we can't deduce who to suggest a vouch to, about having added a collaborator. +type RepoCollaboratorUpdate struct { + Op AclOp `json:"op"` + Subject string `json:"subject"` + Repo string `json:"repo"` +} diff --git a/knotserver/db/collaborator.go b/knotserver/db/collaborator.go new file mode 100644 --- /dev/null +++ b/knotserver/db/collaborator.go @@ -0,0 +1,94 @@ +package db + +import ( + "context" + "database/sql" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/orm" +) + +type Collaborator struct { + Id int + RepoDid syntax.DID + Subject syntax.DID + AddedBy syntax.DID + Created string +} + +func AddCollaborator(q DBTX, c Collaborator) error { + _, err := q.Exec( + `insert or ignore into collaborators (repo_did, subject_did, added_by_did) values (?, ?, ?)`, + c.RepoDid, + c.Subject, + c.AddedBy, + ) + return err +} + +func IsCollaborator(q DBTX, repoDid, subject syntax.DID) (bool, error) { + var exists bool + err := q.QueryRow( + `select exists (select 1 from collaborators where repo_did = ? and subject_did = ?)`, + repoDid, + subject, + ).Scan(&exists) + return exists, err +} + +func RemoveCollaborator(q DBTX, repoDid, subject syntax.DID) error { + _, err := q.Exec( + `delete from collaborators where repo_did = ? and subject_did = ?`, + repoDid, + subject, + ) + return err +} + +func (d *DB) ApplyCollaboratorBackfill(ctx context.Context, rows []Collaborator, migrationName string, markApplied bool) error { + insert := func(tx *sql.Tx) error { + for _, c := range rows { + if err := AddDid(tx, c.Subject.String()); err != nil { + return err + } + if err := AddCollaborator(tx, c); err != nil { + return err + } + } + return nil + } + + if markApplied { + conn, err := d.db.Conn(ctx) + if err != nil { + return err + } + defer conn.Close() + return orm.RunMigration(conn, d.logger, migrationName, insert) + } + + tx, err := d.db.BeginTx(ctx, nil) + if err != nil { + return err + } + defer tx.Rollback() + if err := insert(tx); err != nil { + return err + } + return tx.Commit() +} + +func ListCollaborators(q DBTX, repoDid syntax.DID, p ListPage) ([]Collaborator, *int, error) { + return listPaged(q, + `select id, repo_did, subject_did, added_by_did, created + from collaborators + where repo_did = ?`, + []any{repoDid}, p, + func(r *sql.Rows) (Collaborator, error) { + var c Collaborator + err := r.Scan(&c.Id, &c.RepoDid, &c.Subject, &c.AddedBy, &c.Created) + return c, err + }, + func(c Collaborator) int { return c.Id }, + ) +} diff --git a/knotserver/db/db.go b/knotserver/db/db.go --- a/knotserver/db/db.go +++ b/knotserver/db/db.go @@ -26,6 +26,7 @@ } type DBTX interface { QueryRow(query string, args ...any) *sql.Row + Query(query string, args ...any) (*sql.Rows, error) Exec(query string, args ...any) (sql.Result, error) } @@ -39,6 +40,10 @@ } func (d *DB) QueryRow(query string, args ...any) *sql.Row { return d.db.QueryRow(query, args...) +} + +func (d *DB) Query(query string, args ...any) (*sql.Rows, error) { + return d.db.Query(query, args...) } func Setup(ctx context.Context, dbPath string) (*DB, error) { @@ -244,6 +249,63 @@ }); err != nil { return nil, err } + if err := orm.RunMigration(conn, logger, "knot-members-nullable-rkey", func(tx *sql.Tx) error { + _, mErr := tx.ExecContext(ctx, ` + create table knot_members_new ( + id integer primary key autoincrement, + did text not null, + rkey text, + subject text not null, + created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), + unique (did, rkey) + ); + insert into knot_members_new (id, did, rkey, subject, created) + select id, did, rkey, subject, created from knot_members; + drop table knot_members; + alter table knot_members_new rename to knot_members; + `) + return mErr + }); err != nil { + return nil, err + } + + if err := orm.RunMigration(conn, logger, "create-collaborators", func(tx *sql.Tx) error { + _, mErr := tx.ExecContext(ctx, ` + create table if not exists collaborators ( + id integer primary key autoincrement, + repo_did text not null, + subject_did text not null, + added_by_did text not null, + created text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), + unique (repo_did, subject_did) + ); + create index if not exists idx_collaborators_repo_id + on collaborators(repo_did, id); + `) + return mErr + }); err != nil { + return nil, err + } + + if err := orm.RunMigration(conn, logger, "knot-members-direct-subject-unique", func(tx *sql.Tx) error { + _, mErr := tx.ExecContext(ctx, ` + create unique index if not exists idx_knot_members_direct_subject + on knot_members(subject) where rkey is null; + create index if not exists idx_knot_members_subject + on knot_members(subject); + `) + return mErr + }); err != nil { + return nil, err + } + + if err := orm.RunMigration(conn, logger, "add-events-created-index", func(tx *sql.Tx) error { + _, mErr := tx.ExecContext(ctx, `create index if not exists idx_events_created on events(created)`) + return mErr + }); err != nil { + return nil, err + } + return &DB{ db: db, logger: logger, @@ -299,6 +361,10 @@ if _, err := tx.Exec(`DELETE FROM repo_keys WHERE repo_did = ?`, repoDid); err != nil { return err } + if _, err := tx.Exec(`DELETE FROM collaborators WHERE repo_did = ?`, repoDid); err != nil { + return err + } + return tx.Commit() } @@ -306,6 +372,24 @@ func (d *DB) RepoDidExists(repoDid string) (bool, error) { var count int err := d.db.QueryRow(`SELECT count(1) FROM repo_keys WHERE repo_did = ?`, repoDid).Scan(&count) return count > 0, err +} + +func (d *DB) ListRepoDids() ([]string, error) { + rows, err := d.db.Query(`SELECT repo_did FROM repo_keys`) + if err != nil { + return nil, err + } + defer rows.Close() + + dids := []string{} + for rows.Next() { + var did string + if err := rows.Scan(&did); err != nil { + return nil, err + } + dids = append(dids, did) + } + return dids, rows.Err() } func (d *DB) GetRepoDid(ownerDid, rkey string) (string, error) { diff --git a/knotserver/db/events.go b/knotserver/db/events.go --- a/knotserver/db/events.go +++ b/knotserver/db/events.go @@ -4,6 +4,7 @@ import ( "encoding/json" "fmt" + "github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/eventstream" "tangled.org/core/notifier" "tangled.org/core/tid" @@ -28,6 +29,43 @@ return d.InsertEvent(eventstream.Event{ Rkey: tid.TID(), Nsid: RepoDIDAssignNSID, + EventJson: eventJson, + }, n) +} + +func (d *DB) EmitKnotMemberUpdate(n *notifier.Notifier, op AclOp, subject syntax.DID) error { + payload := KnotMemberUpdate{ + Op: op, + Subject: subject.String(), + } + + eventJson, err := json.Marshal(payload) + if err != nil { + return fmt.Errorf("marshal memberUpdate event: %w", err) + } + + return d.InsertEvent(eventstream.Event{ + Rkey: tid.TID(), + Nsid: KnotMemberUpdateNSID, + EventJson: eventJson, + }, n) +} + +func (d *DB) EmitCollaboratorUpdate(n *notifier.Notifier, op AclOp, subject, repoDid syntax.DID) error { + payload := RepoCollaboratorUpdate{ + Op: op, + Subject: subject.String(), + Repo: repoDid.String(), + } + + eventJson, err := json.Marshal(payload) + if err != nil { + return fmt.Errorf("marshal collaboratorUpdate event: %w", err) + } + + return d.InsertEvent(eventstream.Event{ + Rkey: tid.TID(), + Nsid: RepoCollaboratorUpdateNSID, EventJson: eventJson, }, n) } diff --git a/knotserver/db/member.go b/knotserver/db/member.go --- a/knotserver/db/member.go +++ b/knotserver/db/member.go @@ -13,6 +13,7 @@ Id int Did syntax.DID Rkey string Subject syntax.DID + Created string } func (d *DB) IsMigrationApplied(name string) (bool, error) { @@ -24,6 +25,75 @@ ).Scan(&exists) return exists, err } +func (d *DB) ApplyKnotMemberBackfill(ctx context.Context, rows []KnotMember, migrationName string) error { + conn, err := d.db.Conn(ctx) + if err != nil { + return err + } + defer conn.Close() + + return orm.RunMigration(conn, d.logger, migrationName, func(tx *sql.Tx) error { + for _, m := range rows { + if err := AddDid(tx, m.Subject.String()); err != nil { + return err + } + if err := AddKnotMemberDirect(tx, m.Did, m.Subject); err != nil { + return err + } + } + return nil + }) +} + +func AddKnotMemberDirect(q DBTX, addedBy, subject syntax.DID) error { + _, err := q.Exec( + `insert or ignore into knot_members (did, rkey, subject) values (?, NULL, ?)`, + addedBy, + subject, + ) + return err +} + +func RemoveKnotMemberBySubject(q DBTX, subject syntax.DID) error { + _, err := q.Exec( + "delete from knot_members where subject = ?", + subject, + ) + return err +} + +func RemoveKnotMemberDirect(q DBTX, subject syntax.DID) error { + _, err := q.Exec( + "delete from knot_members where subject = ? and rkey is null", + subject, + ) + return err +} + +func CountKnotMembersBySubject(q DBTX, subject string) (int, error) { + var count int + err := q.QueryRow( + `select count(*) from knot_members where subject = ?`, + subject, + ).Scan(&count) + return count, err +} + +func ListKnotMembers(q DBTX, p ListPage) ([]KnotMember, *int, error) { + return listPaged(q, + `select id, did, subject, created + from knot_members + where id in (select min(id) from knot_members group by subject)`, + nil, p, + func(r *sql.Rows) (KnotMember, error) { + var m KnotMember + err := r.Scan(&m.Id, &m.Did, &m.Subject, &m.Created) + return m, err + }, + func(m KnotMember) int { return m.Id }, + ) +} + func (d *DB) ApplyKnotMembersBackfill(ctx context.Context, rows []KnotMember, migrationName string) error { conn, err := d.db.Conn(ctx) if err != nil { @@ -33,10 +103,7 @@ defer conn.Close() return orm.RunMigration(conn, d.logger, migrationName, func(tx *sql.Tx) error { for _, m := range rows { - if _, err := tx.ExecContext(ctx, - `insert or ignore into known_dids (did) values (?)`, - m.Subject, - ); err != nil { + if err := AddDid(tx, m.Subject.String()); err != nil { return err } if _, err := tx.ExecContext(ctx, @@ -67,15 +134,6 @@ ownerDid, rkey, ) return err -} - -func CountKnotMembersBySubject(q DBTX, subject string) (int, error) { - var count int - err := q.QueryRow( - `select count(*) from knot_members where subject = ?`, - subject, - ).Scan(&count) - return count, err } func GetKnotMember(q DBTX, did, rkey string) (*KnotMember, error) { diff --git a/knotserver/db/page.go b/knotserver/db/page.go new file mode 100644 --- /dev/null +++ b/knotserver/db/page.go @@ -0,0 +1,74 @@ +package db + +import ( + "database/sql" + "fmt" +) + +const ( + ListDefaultLimit = 50 + ListMaxLimit = 1000 +) + +type ListPage struct { + Limit int + Cursor *int + Desc bool +} + +func (p ListPage) limit() int { + switch { + case p.Limit <= 0: + return ListDefaultLimit + case p.Limit > ListMaxLimit: + return ListMaxLimit + default: + return p.Limit + } +} + +func (p ListPage) clause() (string, []any) { + dir, cmp := "asc", ">" + if p.Desc { + dir, cmp = "desc", "<" + } + if p.Cursor != nil { + return fmt.Sprintf("where id %s ? order by id %s limit ?", cmp, dir), []any{*p.Cursor, p.limit() + 1} + } + return fmt.Sprintf("order by id %s limit ?", dir), []any{p.limit() + 1} +} + +func listPaged[T any]( + q DBTX, + query string, + args []any, + p ListPage, + scan func(*sql.Rows) (T, error), + idOf func(T) int, +) ([]T, *int, error) { + clause, pArgs := p.clause() + rows, err := q.Query("select * from ("+query+") "+clause, append(args, pArgs...)...) + if err != nil { + return nil, nil, err + } + defer rows.Close() + + out := []T{} + for rows.Next() { + v, err := scan(rows) + if err != nil { + return nil, nil, err + } + out = append(out, v) + } + if err := rows.Err(); err != nil { + return nil, nil, err + } + + if limit := p.limit(); len(out) > limit { + out = out[:limit] + last := idOf(out[len(out)-1]) + return out, &last, nil + } + return out, nil, nil +} diff --git a/notifier/notifier.go b/notifier/notifier.go --- a/notifier/notifier.go +++ b/notifier/notifier.go @@ -31,6 +31,9 @@ n.mu.Unlock() } func (n *Notifier) NotifyAll() { + if n == nil { + return + } n.mu.Lock() for ch := range n.subscribers { select { -- tangled.sh