From 5102fa4b679aeafd73af3fd78eee069c17b4b908 Mon Sep 17 00:00:00 2001 From: oppiliappan Date: Sat, 14 Jun 2025 18:32:38 +0100 Subject: [PATCH] spindle: ingest repo records in spindle allows the spindle to dynamically configure the knots it is listening to. Signed-off-by: oppiliappan --- spindle/db/db.go | 10 +++++++++ spindle/db/repos.go | 28 +++++++++++++++++++++++ spindle/ingester.go | 54 +++++++++++++++++++++++++++++++++++++++++++-- 3 files changed, 90 insertions(+), 2 deletions(-) create mode 100644 spindle/db/repos.go diff --git a/spindle/db/db.go b/spindle/db/db.go index df6eac32..a81edc85 100644 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -35,6 +35,16 @@ func Make(dbPath string) (*DB, error) { did text primary key ); + 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')), + + unique(owner, name) + ); + -- status event for a single workflow create table if not exists events ( rkey text not null, diff --git a/spindle/db/repos.go b/spindle/db/repos.go new file mode 100644 index 00000000..656ade55 --- /dev/null +++ b/spindle/db/repos.go @@ -0,0 +1,28 @@ +package db + +func (d *DB) AddRepo(knot, owner, name string) error { + _, err := d.Exec(`insert or ignore into repos (knot, owner, name) values (?, ?, ?)`, knot, owner, name) + return err +} + +func (d *DB) Knots() ([]string, error) { + rows, err := d.Query(`select knot from repos`) + if err != nil { + return nil, err + } + + var knots []string + for rows.Next() { + var knot string + if err := rows.Scan(&knot); err != nil { + return nil, err + } + knots = append(knots, knot) + } + + if err = rows.Err(); err != nil { + return nil, err + } + + return knots, nil +} diff --git a/spindle/ingester.go b/spindle/ingester.go index 59ecdf32..4a5da90f 100644 --- a/spindle/ingester.go +++ b/spindle/ingester.go @@ -5,8 +5,10 @@ import ( "encoding/json" "fmt" - "github.com/bluesky-social/jetstream/pkg/models" "tangled.sh/tangled.sh/core/api/tangled" + "tangled.sh/tangled.sh/core/knotclient" + + "github.com/bluesky-social/jetstream/pkg/models" ) type Ingester func(ctx context.Context, e *models.Event) error @@ -29,6 +31,8 @@ func (s *Spindle) ingest() Ingester { switch e.Commit.Collection { case tangled.SpindleMemberNSID: s.ingestMember(ctx, e) + case tangled.RepoNSID: + s.ingestRepo(ctx, e) } return err @@ -68,7 +72,7 @@ func (s *Spindle) ingestMember(_ context.Context, e *models.Event) error { return fmt.Errorf("failed to enforce permissions: %w", err) } - if err := s.e.AddMember(rbacDomain, record.Subject); err != nil { + if err := s.e.AddKnotMember(rbacDomain, record.Subject); err != nil { l.Error("failed to add member", "error", err) return fmt.Errorf("failed to add member: %w", err) } @@ -85,3 +89,49 @@ func (s *Spindle) ingestMember(_ context.Context, e *models.Event) error { } return nil } + +func (s *Spindle) ingestRepo(_ context.Context, e *models.Event) error { + var err error + + l := s.l.With("component", "ingester", "record", tangled.RepoNSID) + + switch e.Commit.Operation { + case models.CommitOperationCreate, models.CommitOperationUpdate: + raw := e.Commit.Record + record := tangled.Repo{} + err = json.Unmarshal(raw, &record) + if err != nil { + l.Error("invalid record", "error", err) + return err + } + + domain := s.cfg.Server.Hostname + if s.cfg.Server.Dev { + domain = s.cfg.Server.ListenAddr + } + + // no spindle configured for this repo + if record.Spindle == nil { + return nil + } + + // this repo did not want this spindle + if *record.Spindle != domain { + return nil + } + + // add this repo to the watch list + if err := s.db.AddRepo(record.Knot, record.Owner, record.Name); err != nil { + l.Error("failed to add repo", "error", err) + return fmt.Errorf("failed to add repo: %w", err) + } + + // add this knot to the event consumer + src := knotclient.NewEventSource(record.Knot) + s.ks.AddSource(context.Background(), src) + + return nil + + } + return nil +} -- 2.51.2