From 5a2fe5d25f6f516b53537c224dca5885a1d5d2c4 Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Thu, 18 Dec 2025 21:14:09 +0900 Subject: [PATCH] nix,spindle: switch to tap from jetstream spindle-tap will collect/stream record events from: - users dynamically added by spindle (spindle members | collaborators of repos using spindle) - any users with `sh.tangled.repo.pull` collection It might be bit inefficient considering it will also stream repo creation events from PR authors due to second rule, but at least we now have backfill logic and Sync 1.1 based syncing. This inefficiency can be fixed later by modifying upstream tap cli or embedding tap into spindle. ``` +--------- all tangled users --------+ | | | +-- users known to spindle-tap --+ | | | (PR author / manually added) | | | | | | | | +----------------------------+ | | | | | users known to spindle | | | | | | (members / collaborators) | | | | | +----------------------------+ | | | +--------------------------------+ | +------------------------------------+ ``` Close: Signed-off-by: Seongmin Lee --- nix/modules/spindle.nix | 10 +- nix/vm.nix | 6 +- spindle/config/config.go | 22 +-- spindle/db/db.go | 42 ++++-- spindle/db/known_dids.go | 44 ------ spindle/db/repos.go | 130 +++++++++++++++-- spindle/ingester.go | 276 ------------------------------------ spindle/server.go | 49 +++---- spindle/tap.go | 294 +++++++++++++++++++++++++++++++++++++++ 9 files changed, 478 insertions(+), 395 deletions(-) delete mode 100644 spindle/db/known_dids.go delete mode 100644 spindle/ingester.go create mode 100644 spindle/tap.go diff --git a/nix/modules/spindle.nix b/nix/modules/spindle.nix index 591cb661..fad035fb 100644 --- a/nix/modules/spindle.nix +++ b/nix/modules/spindle.nix @@ -53,12 +53,6 @@ in description = "atproto PLC directory"; }; - jetstreamEndpoint = mkOption { - type = types.str; - default = "wss://jetstream1.us-west.bsky.network/subscribe"; - description = "Jetstream endpoint to subscribe to"; - }; - dev = mkOption { type = types.bool; default = false; @@ -151,7 +145,7 @@ in systemd.services.spindle = { description = "spindle service"; - after = ["network.target" "docker.service"]; + after = ["network.target" "docker.service" "spindle-tap.service"]; wantedBy = ["multi-user.target"]; serviceConfig = { LogsDirectory = "spindle"; @@ -161,7 +155,6 @@ in "SPINDLE_SERVER_DB_PATH=${cfg.server.dbPath}" "SPINDLE_SERVER_HOSTNAME=${cfg.server.hostname}" "SPINDLE_SERVER_PLC_URL=${cfg.server.plcUrl}" - "SPINDLE_SERVER_JETSTREAM_ENDPOINT=${cfg.server.jetstreamEndpoint}" "SPINDLE_SERVER_DEV=${lib.boolToString cfg.server.dev}" "SPINDLE_SERVER_OWNER=${cfg.server.owner}" "SPINDLE_SERVER_MAX_JOB_COUNT=${toString cfg.server.maxJobCount}" @@ -169,6 +162,7 @@ in "SPINDLE_SERVER_SECRETS_PROVIDER=${cfg.server.secrets.provider}" "SPINDLE_SERVER_SECRETS_OPENBAO_PROXY_ADDR=${cfg.server.secrets.openbao.proxyAddr}" "SPINDLE_SERVER_SECRETS_OPENBAO_MOUNT=${cfg.server.secrets.openbao.mount}" + "SPINDLE_SERVER_TAP_URL=http://localhost:2480" "SPINDLE_NIXERY_PIPELINES_NIXERY=${cfg.pipelines.nixery}" "SPINDLE_NIXERY_PIPELINES_WORKFLOW_TIMEOUT=${cfg.pipelines.workflowTimeout}" ]; diff --git a/nix/vm.nix b/nix/vm.nix index 24b42bac..14cb0290 100644 --- a/nix/vm.nix +++ b/nix/vm.nix @@ -58,6 +58,11 @@ in host.port = 6555; guest.port = 6555; } + { + from = "host"; + host.port = 6556; + guest.port = 2480; + } ]; sharedDirectories = { # We can't use the 9p mounts directly for most of these @@ -101,7 +106,6 @@ in owner = envVar "TANGLED_VM_SPINDLE_OWNER"; hostname = envVarOr "TANGLED_VM_SPINDLE_HOST" "localhost:6555"; plcUrl = plcUrl; - jetstreamEndpoint = jetstream; listenAddr = "0.0.0.0:6555"; dev = true; queueSize = 100; diff --git a/spindle/config/config.go b/spindle/config/config.go index 8be5357d..f7410519 100644 --- a/spindle/config/config.go +++ b/spindle/config/config.go @@ -9,17 +9,17 @@ import ( ) type Server struct { - ListenAddr string `env:"LISTEN_ADDR, default=0.0.0.0:6555"` - DBPath string `env:"DB_PATH, default=spindle.db"` - Hostname string `env:"HOSTNAME, required"` - JetstreamEndpoint string `env:"JETSTREAM_ENDPOINT, default=wss://jetstream1.us-west.bsky.network/subscribe"` - PlcUrl string `env:"PLC_URL, default=https://plc.directory"` - Dev bool `env:"DEV, default=false"` - Owner syntax.DID `env:"OWNER, required"` - Secrets Secrets `env:",prefix=SECRETS_"` - LogDir string `env:"LOG_DIR, default=/var/log/spindle"` - QueueSize int `env:"QUEUE_SIZE, default=100"` - MaxJobCount int `env:"MAX_JOB_COUNT, default=2"` // max number of jobs that run at a time + ListenAddr string `env:"LISTEN_ADDR, default=0.0.0.0:6555"` + DBPath string `env:"DB_PATH, default=spindle.db"` + Hostname string `env:"HOSTNAME, required"` + TapUrl string `env:"TAP_URL, required"` + PlcUrl string `env:"PLC_URL, default=https://plc.directory"` + Dev bool `env:"DEV, default=false"` + Owner syntax.DID `env:"OWNER, required"` + Secrets Secrets `env:",prefix=SECRETS_"` + LogDir string `env:"LOG_DIR, default=/var/log/spindle"` + QueueSize int `env:"QUEUE_SIZE, default=100"` + MaxJobCount int `env:"MAX_JOB_COUNT, default=2"` // max number of jobs that run at a time } func (s Server) Did() syntax.DID { diff --git a/spindle/db/db.go b/spindle/db/db.go index 942e5c92..1406ead6 100644 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -5,6 +5,7 @@ import ( "database/sql" "strings" + "github.com/bluesky-social/indigo/atproto/syntax" _ "github.com/mattn/go-sqlite3" "tangled.org/core/log" "tangled.org/core/orm" @@ -57,6 +58,20 @@ func Make(ctx context.Context, dbPath string) (*DB, error) { unique(owner, name) ); + create table if not exists repo_collaborators ( + -- identifiers + id integer primary key autoincrement, + did text not null, + rkey text not null, + at_uri text generated always as ('at://' || did || '/' || 'sh.tangled.repo.collaborator' || '/' || rkey) stored, + + repo text not null, + subject text not null, + + addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), + unique(did, rkey) + ); + create table if not exists spindle_members ( -- identifiers for the record id integer primary key autoincrement, @@ -120,18 +135,17 @@ func Make(ctx context.Context, dbPath string) (*DB, error) { return &DB{db}, nil } -func (d *DB) SaveLastTimeUs(lastTimeUs int64) error { - _, err := d.Exec(` - insert into _jetstream (id, last_time_us) - values (1, ?) - on conflict(id) do update set last_time_us = excluded.last_time_us - `, lastTimeUs) - return err -} - -func (d *DB) GetLastTimeUs() (int64, error) { - var lastTimeUs int64 - row := d.QueryRow(`select last_time_us from _jetstream where id = 1;`) - err := row.Scan(&lastTimeUs) - return lastTimeUs, err +func (d *DB) IsKnownDid(did syntax.DID) (bool, error) { + // is spindle member / repo collaborator + var exists bool + err := d.QueryRow( + `select exists ( + select 1 from repo_collaborators where subject = ? + union all + select 1 from spindle_members where did = ? + )`, + did, + did, + ).Scan(&exists) + return exists, err } diff --git a/spindle/db/known_dids.go b/spindle/db/known_dids.go deleted file mode 100644 index 8acdb695..00000000 --- a/spindle/db/known_dids.go +++ /dev/null @@ -1,44 +0,0 @@ -package db - -func (d *DB) AddDid(did string) error { - _, err := d.Exec(`insert or ignore into known_dids (did) values (?)`, did) - return err -} - -func (d *DB) RemoveDid(did string) error { - _, err := d.Exec(`delete from known_dids where did = ?`, did) - return err -} - -func (d *DB) GetAllDids() ([]string, error) { - var dids []string - - rows, err := d.Query(`select did from known_dids`) - if err != nil { - return nil, err - } - defer rows.Close() - - for rows.Next() { - var did string - if err := rows.Scan(&did); err != nil { - return nil, err - } - dids = append(dids, did) - } - - if err := rows.Err(); err != nil { - return nil, err - } - - return dids, nil -} - -func (d *DB) HasKnownDids() bool { - var count int - err := d.QueryRow(`select count(*) from known_dids`).Scan(&count) - if err != nil { - return false - } - return count > 0 -} diff --git a/spindle/db/repos.go b/spindle/db/repos.go index d6ecf49d..44190316 100644 --- a/spindle/db/repos.go +++ b/spindle/db/repos.go @@ -1,13 +1,42 @@ package db +import "github.com/bluesky-social/indigo/atproto/syntax" + type Repo struct { - Knot string - Owner string - Name string + Did syntax.DID + Rkey syntax.RecordKey + Name string + Knot string +} + +type RepoCollaborator struct { + Did syntax.DID + Rkey syntax.RecordKey + Repo syntax.ATURI + Subject syntax.DID } -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) PutRepo(repo *Repo) error { + _, err := d.Exec( + `insert or ignore into repos (did, rkey, name, knot) + values (?, ?, ?, ?) + on conflict(did, rkey) do update set + name = excluded.name, + knot = excluded.knot`, + repo.Did, + repo.Rkey, + repo.Name, + repo.Knot, + ) + return err +} + +func (d *DB) DeleteRepo(did syntax.DID, rkey syntax.RecordKey) error { + _, err := d.Exec( + `delete from repos where did = ? and rkey = ?`, + did, + rkey, + ) return err } @@ -34,16 +63,95 @@ func (d *DB) Knots() ([]string, error) { return knots, nil } -func (d *DB) GetRepo(knot, owner, name string) (*Repo, error) { +func (d *DB) GetRepo(repoAt syntax.ATURI) (*Repo, error) { var repo Repo - - 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) - + err := d.DB.QueryRow( + `select + did, + rkey, + name, + knot + from repos where at_uri = ?`, + repoAt, + ).Scan( + &repo.Did, + &repo.Rkey, + &repo.Name, + &repo.Knot, + ) if err != nil { return nil, err } + return &repo, nil +} +func (d *DB) GetRepoWithName(did syntax.DID, name string) (*Repo, error) { + var repo Repo + err := d.DB.QueryRow( + `select + did, + rkey, + name, + knot + from repos where did = ? and name = ?`, + did, + name, + ).Scan( + &repo.Did, + &repo.Rkey, + &repo.Name, + &repo.Knot, + ) + if err != nil { + return nil, err + } return &repo, nil } + +func (d *DB) PutRepoCollaborator(collaborator *RepoCollaborator) error { + _, err := d.Exec( + `insert into repo_collaborators (did, rkey, repo, subject) + values (?, ?, ?, ?) + on conflict(did, rkey) do update set + repo = excluded.repo, + subject = excluded.subject`, + collaborator.Did, + collaborator.Rkey, + collaborator.Repo, + collaborator.Subject, + ) + return err +} + +func (d *DB) RemoveRepoCollaborator(did syntax.DID, rkey syntax.RecordKey) error { + _, err := d.Exec( + `delete from repo_collaborators where did = ? and rkey = ?`, + did, + rkey, + ) + return err +} + +func (d *DB) GetRepoCollaborator(did syntax.DID, rkey syntax.RecordKey) (*RepoCollaborator, error) { + var collaborator RepoCollaborator + err := d.DB.QueryRow( + `select + did, + rkey, + repo, + subject + from repo_collaborators + where did = ? and rkey = ?`, + did, + rkey, + ).Scan( + &collaborator.Did, + &collaborator.Rkey, + &collaborator.Repo, + &collaborator.Subject, + ) + if err != nil { + return nil, err + } + return &collaborator, nil +} diff --git a/spindle/ingester.go b/spindle/ingester.go deleted file mode 100644 index 3cb56037..00000000 --- a/spindle/ingester.go +++ /dev/null @@ -1,276 +0,0 @@ -package spindle - -import ( - "context" - "encoding/json" - "errors" - "fmt" - "time" - - "tangled.org/core/api/tangled" - "tangled.org/core/eventconsumer" - "tangled.org/core/spindle/db" - - comatproto "github.com/bluesky-social/indigo/api/atproto" - "github.com/bluesky-social/indigo/atproto/syntax" - "github.com/bluesky-social/indigo/xrpc" - "github.com/bluesky-social/jetstream/pkg/models" -) - -type Ingester func(ctx context.Context, e *models.Event) error - -func (s *Spindle) ingest() Ingester { - return func(ctx context.Context, e *models.Event) error { - var err error - defer func() { - eventTime := e.TimeUS - lastTimeUs := eventTime + 1 - if err := s.db.SaveLastTimeUs(lastTimeUs); err != nil { - err = fmt.Errorf("(deferred) failed to save last time us: %w", err) - } - }() - - if e.Kind != models.EventKindCommit { - return nil - } - - switch e.Commit.Collection { - case tangled.SpindleMemberNSID: - err = s.ingestMember(ctx, e) - case tangled.RepoNSID: - err = s.ingestRepo(ctx, e) - case tangled.RepoCollaboratorNSID: - err = s.ingestCollaborator(ctx, e) - } - - if err != nil { - s.l.Debug("failed to process message", "nsid", e.Commit.Collection, "err", err) - } - - return nil - } -} - -func (s *Spindle) ingestMember(_ context.Context, e *models.Event) error { - var err error - did := e.Did - rkey := e.Commit.RKey - - l := s.l.With("component", "ingester", "record", tangled.SpindleMemberNSID) - - switch e.Commit.Operation { - case models.CommitOperationCreate, models.CommitOperationUpdate: - raw := e.Commit.Record - record := tangled.SpindleMember{} - err = json.Unmarshal(raw, &record) - if err != nil { - l.Error("invalid record", "error", err) - return err - } - - domain := s.cfg.Server.Hostname - recordInstance := record.Instance - - if recordInstance != domain { - l.Error("domain mismatch", "domain", recordInstance, "expected", domain) - return fmt.Errorf("domain mismatch: %s != %s", record.Instance, domain) - } - - ok, err := s.e.IsSpindleMemberInviteAllowed(syntax.DID(did), s.cfg.Server.Did()) - if err != nil || !ok { - l.Error("failed to add member", "did", did, "error", err) - return fmt.Errorf("failed to enforce permissions: %w", err) - } - - if err := db.AddSpindleMember(s.db, db.SpindleMember{ - Did: syntax.DID(did), - Rkey: rkey, - Instance: recordInstance, - Subject: syntax.DID(record.Subject), - Created: time.Now(), - }); err != nil { - l.Error("failed to add member", "error", err) - return fmt.Errorf("failed to add member: %w", err) - } - - if err := s.e.AddSpindleMember(syntax.DID(record.Subject), s.cfg.Server.Did()); err != nil { - l.Error("failed to add member", "error", err) - return fmt.Errorf("failed to add member: %w", err) - } - l.Info("added member from firehose", "member", record.Subject) - - if err := s.db.AddDid(record.Subject); err != nil { - l.Error("failed to add did", "error", err) - return fmt.Errorf("failed to add did: %w", err) - } - s.jc.AddDid(record.Subject) - - return nil - - case models.CommitOperationDelete: - record, err := db.GetSpindleMember(s.db, did, rkey) - if err != nil { - l.Error("failed to find member", "error", err) - return fmt.Errorf("failed to find member: %w", err) - } - - if err := db.RemoveSpindleMember(s.db, did, rkey); err != nil { - l.Error("failed to remove member", "error", err) - return fmt.Errorf("failed to remove member: %w", err) - } - - if err := s.e.RemoveSpindleMember(record.Subject, s.cfg.Server.Did()); err != nil { - l.Error("failed to add member", "error", err) - return fmt.Errorf("failed to add member: %w", err) - } - l.Info("added member from firehose", "member", record.Subject) - - if err := s.db.RemoveDid(record.Subject.String()); err != nil { - l.Error("failed to add did", "error", err) - return fmt.Errorf("failed to add did: %w", err) - } - s.jc.RemoveDid(record.Subject.String()) - - } - return nil -} - -func (s *Spindle) ingestRepo(ctx context.Context, e *models.Event) error { - var err error - did := e.Did - - l := s.l.With("component", "ingester", "record", tangled.RepoNSID) - - l.Info("ingesting repo record", "did", did) - - 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 - - // no spindle configured for this repo - if record.Spindle == nil { - l.Info("no spindle configured", "name", record.Name) - return nil - } - - // this repo did not want this spindle - if *record.Spindle != domain { - l.Info("different spindle configured", "name", record.Name, "spindle", *record.Spindle, "domain", domain) - return nil - } - - // add this repo to the watch list - if err := s.db.AddRepo(record.Knot, did, record.Name); err != nil { - l.Error("failed to add repo", "error", err) - return fmt.Errorf("failed to add repo: %w", err) - } - - repoAt := syntax.ATURI(fmt.Sprintf("at://%s/%s/%s", did, e.Commit.Collection, e.Commit.RKey)) - - // add repo to rbac - if err := s.e.AddRepo(repoAt); err != nil { - l.Error("failed to add repo to enforcer", "error", err) - return fmt.Errorf("failed to add repo: %w", err) - } - - // add collaborators to rbac - if err := s.fetchAndAddCollaborators(ctx, repoAt); err != nil { - return err - } - - // add this knot to the event consumer - src := eventconsumer.NewKnotSource(record.Knot) - s.ks.AddSource(context.Background(), src) - - return nil - - } - return nil -} - -func (s *Spindle) ingestCollaborator(ctx context.Context, e *models.Event) error { - var err error - - l := s.l.With("component", "ingester", "record", tangled.RepoCollaboratorNSID, "did", e.Did) - - l.Info("ingesting collaborator record") - - switch e.Commit.Operation { - case models.CommitOperationCreate, models.CommitOperationUpdate: - raw := e.Commit.Record - record := tangled.RepoCollaborator{} - err = json.Unmarshal(raw, &record) - if err != nil { - l.Error("invalid record", "error", err) - return err - } - - subjectId, err := s.res.ResolveIdent(ctx, record.Subject) - if err != nil || subjectId.Handle.IsInvalidHandle() { - return err - } - - repoAt, err := syntax.ParseATURI(record.Repo) - if err != nil { - l.Info("rejecting record, invalid repoAt", "repoAt", record.Repo) - return nil - } - - // check perms for this user - if ok, err := s.e.IsRepoCollaboratorInviteAllowed(syntax.DID(e.Did), repoAt); !ok || err != nil { - return fmt.Errorf("insufficient permissions: %w", err) - } - - // add collaborator to rbac - if err := s.e.AddRepoCollaborator(syntax.DID(record.Subject), repoAt); err != nil { - l.Error("failed to add repo to enforcer", "error", err) - return fmt.Errorf("failed to add repo: %w", err) - } - - return nil - } - return nil -} - -func (s *Spindle) fetchAndAddCollaborators(ctx context.Context, repo syntax.ATURI) error { - l := s.l.With("component", "ingester", "handler", "fetchAndAddCollaborators") - - l.Info("fetching and adding existing collaborators") - - ident, err := s.res.ResolveIdent(ctx, repo.Authority().String()) - if err != nil || ident.Handle.IsInvalidHandle() { - return fmt.Errorf("failed to resolve handle: %w", err) - } - - xrpcc := xrpc.Client{ - Host: ident.PDSEndpoint(), - } - - resp, err := comatproto.RepoListRecords(ctx, &xrpcc, tangled.RepoCollaboratorNSID, "", 50, ident.DID.String(), false) - if err != nil { - return err - } - - var errs error - for _, r := range resp.Records { - if r == nil { - continue - } - record := r.Value.Val.(*tangled.RepoCollaborator) - - if err := s.e.AddRepoCollaborator(syntax.DID(record.Subject), syntax.ATURI(record.Repo)); err != nil { - l.Error("failed to add repo to enforcer", "error", err) - errors.Join(errs, fmt.Errorf("failed to add repo: %w", err)) - } - } - - return errs -} diff --git a/spindle/server.go b/spindle/server.go index e9863194..2343c83b 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -10,12 +10,12 @@ import ( "net/http" "sync" + "github.com/bluesky-social/indigo/atproto/syntax" "github.com/go-chi/chi/v5" "tangled.org/core/api/tangled" "tangled.org/core/eventconsumer" "tangled.org/core/eventconsumer/cursor" "tangled.org/core/idresolver" - "tangled.org/core/jetstream" "tangled.org/core/log" "tangled.org/core/notifier" "tangled.org/core/rbac2" @@ -27,6 +27,7 @@ import ( "tangled.org/core/spindle/queue" "tangled.org/core/spindle/secrets" "tangled.org/core/spindle/xrpc" + "tangled.org/core/tapc" "tangled.org/core/xrpc/serviceauth" ) @@ -34,7 +35,7 @@ import ( var defaultMotd []byte type Spindle struct { - jc *jetstream.JetstreamClient + tap *tapc.Client db *db.DB e *rbac2.Enforcer l *slog.Logger @@ -93,30 +94,12 @@ func New(ctx context.Context, cfg *config.Config, engines map[string]models.Engi jq := queue.NewQueue(cfg.Server.QueueSize, cfg.Server.MaxJobCount) logger.Info("initialized queue", "queueSize", cfg.Server.QueueSize, "numWorkers", cfg.Server.MaxJobCount) - collections := []string{ - tangled.SpindleMemberNSID, - tangled.RepoNSID, - tangled.RepoCollaboratorNSID, - } - jc, err := jetstream.NewJetstreamClient(cfg.Server.JetstreamEndpoint, "spindle", collections, nil, log.SubLogger(logger, "jetstream"), d, true, true) - if err != nil { - return nil, fmt.Errorf("failed to setup jetstream client: %w", err) - } - jc.AddDid(cfg.Server.Owner.String()) - - // Check if the spindle knows about any Dids; - dids, err := d.GetAllDids() - if err != nil { - return nil, fmt.Errorf("failed to get all dids: %w", err) - } - for _, d := range dids { - jc.AddDid(d) - } + tap := tapc.NewClient(cfg.Server.TapUrl, "") resolver := idresolver.DefaultResolver(cfg.Server.PlcUrl) spindle := &Spindle{ - jc: jc, + tap: &tap, e: e, db: d, l: logger, @@ -140,11 +123,6 @@ func New(ctx context.Context, cfg *config.Config, engines map[string]models.Engi return nil, fmt.Errorf("failed to setup sqlite3 cursor store: %w", err) } - err = jc.StartJetstream(ctx, spindle.ingest()) - if err != nil { - return nil, fmt.Errorf("failed to start jetstream consumer: %w", err) - } - // for each incoming sh.tangled.pipeline, we execute // spindle.processPipeline, which in turn enqueues the pipeline // job in the above registered queue. @@ -226,6 +204,18 @@ func (s *Spindle) Start(ctx context.Context) error { s.ks.Start(ctx) }() + // ensure server owner is tracked + if err := s.tap.AddRepos(ctx, []syntax.DID{s.cfg.Server.Owner}); err != nil { + return err + } + + go func() { + s.l.Info("starting tap stream consumer") + s.tap.Connect(ctx, &tapc.SimpleIndexer{ + EventHandler: s.processEvent, + }) + }() + s.l.Info("starting spindle server", "address", s.cfg.Server.ListenAddr) return http.ListenAndServe(s.cfg.Server.ListenAddr, s.Router()) } @@ -306,9 +296,8 @@ func (s *Spindle) processPipeline(ctx context.Context, src eventconsumer.Source, } // filter by repos - _, err = s.db.GetRepo( - tpl.TriggerMetadata.Repo.Knot, - tpl.TriggerMetadata.Repo.Did, + _, err = s.db.GetRepoWithName( + syntax.DID(tpl.TriggerMetadata.Repo.Did), tpl.TriggerMetadata.Repo.Repo, ) if err != nil { diff --git a/spindle/tap.go b/spindle/tap.go new file mode 100644 index 00000000..31e62fb0 --- /dev/null +++ b/spindle/tap.go @@ -0,0 +1,294 @@ +package spindle + +import ( + "context" + "encoding/json" + "fmt" + "time" + + "github.com/bluesky-social/indigo/atproto/syntax" + "tangled.org/core/api/tangled" + "tangled.org/core/eventconsumer" + "tangled.org/core/spindle/db" + "tangled.org/core/tapc" +) + +func (s *Spindle) processEvent(ctx context.Context, evt tapc.Event) error { + l := s.l.With("component", "tapIndexer") + + var err error + switch evt.Type { + case tapc.EvtRecord: + switch evt.Record.Collection.String() { + case tangled.SpindleMemberNSID: + err = s.processMember(ctx, evt) + case tangled.RepoNSID: + err = s.processRepo(ctx, evt) + case tangled.RepoCollaboratorNSID: + err = s.processCollaborator(ctx, evt) + case tangled.RepoPullNSID: + err = s.processPull(ctx, evt) + } + case tapc.EvtIdentity: + // no-op + } + + if err != nil { + l.Error("failed to process message. will retry later", "event.ID", evt.ID, "err", err) + return err + } + return nil +} + +// NOTE: make sure to return nil if we don't need to retry (e.g. forbidden, unrelated) + +func (s *Spindle) processMember(ctx context.Context, evt tapc.Event) error { + l := s.l.With("component", "tapIndexer", "record", evt.Record.AtUri()) + + l.Info("processing spindle.member record") + + // only listen to members + if ok, err := s.e.IsSpindleMemberInviteAllowed(evt.Record.Did, s.cfg.Server.Did()); !ok || err != nil { + l.Warn("forbidden request: member invite not allowed", "did", evt.Record.Did, "error", err) + return nil + } + + switch evt.Record.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + record := tangled.SpindleMember{} + if err := json.Unmarshal(evt.Record.Record, &record); err != nil { + return fmt.Errorf("parsing record: %w", err) + } + + domain := s.cfg.Server.Hostname + if record.Instance != domain { + l.Info("domain mismatch", "domain", record.Instance, "expected", domain) + return nil + } + + created, err := time.Parse(record.CreatedAt, time.RFC3339) + if err != nil { + created = time.Now() + } + if err := db.AddSpindleMember(s.db, db.SpindleMember{ + Did: evt.Record.Did, + Rkey: evt.Record.Rkey.String(), + Instance: record.Instance, + Subject: syntax.DID(record.Subject), + Created: created, + }); err != nil { + l.Error("failed to add member", "error", err) + return fmt.Errorf("adding member to db: %w", err) + } + if err := s.e.AddSpindleMember(syntax.DID(record.Subject), s.cfg.Server.Did()); err != nil { + return fmt.Errorf("adding member to rbac: %w", err) + } + if err := s.tap.AddRepos(ctx, []syntax.DID{syntax.DID(record.Subject)}); err != nil { + return fmt.Errorf("adding did to tapc: %w", err) + } + + l.Info("added member", "member", record.Subject) + return nil + + case tapc.RecordDeleteAction: + var ( + did = evt.Record.Did.String() + rkey = evt.Record.Rkey.String() + ) + member, err := db.GetSpindleMember(s.db, did, rkey) + if err != nil { + return fmt.Errorf("finding member: %w", err) + } + + if err := db.RemoveSpindleMember(s.db, did, rkey); err != nil { + return fmt.Errorf("removing member from db: %w", err) + } + if err := s.e.RemoveSpindleMember(member.Subject, s.cfg.Server.Did()); err != nil { + return fmt.Errorf("removing member from rbac: %w", err) + } + if err := s.tapSafeRemoveDid(ctx, member.Subject); err != nil { + return fmt.Errorf("removing did from tapc: %w", err) + } + + l.Info("removed member", "member", member.Subject) + return nil + } + return nil +} + +func (s *Spindle) processCollaborator(ctx context.Context, evt tapc.Event) error { + l := s.l.With("component", "tapIndexer", "record", evt.Record.AtUri()) + + l.Info("processing repo.collaborator record") + + // only listen to members + if ok, err := s.e.IsSpindleMember(evt.Record.Did, s.cfg.Server.Did()); !ok || err != nil { + l.Warn("forbidden request: not spindle member", "did", evt.Record.Did, "err", err) + return nil + } + + switch evt.Record.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + record := tangled.RepoCollaborator{} + if err := json.Unmarshal(evt.Record.Record, &record); err != nil { + l.Error("invalid record", "err", err) + return fmt.Errorf("parsing record: %w", err) + } + + // retry later if target repo is not ingested yet + if _, err := s.db.GetRepo(syntax.ATURI(record.Repo)); err != nil { + l.Warn("target repo is not ingested yet", "repo", record.Repo, "err", err) + return fmt.Errorf("target repo is unknown") + } + + // check perms for this user + if ok, err := s.e.IsRepoCollaboratorInviteAllowed(evt.Record.Did, syntax.ATURI(record.Repo)); !ok || err != nil { + l.Warn("forbidden request collaborator invite not allowed", "did", evt.Record.Did, "err", err) + return nil + } + + if err := s.db.PutRepoCollaborator(&db.RepoCollaborator{ + Did: evt.Record.Did, + Rkey: evt.Record.Rkey, + Repo: syntax.ATURI(record.Repo), + Subject: syntax.DID(record.Subject), + }); err != nil { + return fmt.Errorf("adding collaborator to db: %w", err) + } + if err := s.e.AddRepoCollaborator(syntax.DID(record.Subject), syntax.ATURI(record.Repo)); err != nil { + return fmt.Errorf("adding collaborator to rbac: %w", err) + } + if err := s.tap.AddRepos(ctx, []syntax.DID{syntax.DID(record.Subject)}); err != nil { + return fmt.Errorf("adding did to tapc: %w", err) + } + + l.Info("add repo collaborator", "subejct", record.Subject, "repo", record.Repo) + return nil + + case tapc.RecordDeleteAction: + // get existing collaborator + collaborator, err := s.db.GetRepoCollaborator(evt.Record.Did, evt.Record.Rkey) + if err != nil { + return fmt.Errorf("failed to get existing collaborator info: %w", err) + } + + // check perms for this user + if ok, err := s.e.IsRepoCollaboratorInviteAllowed(evt.Record.Did, collaborator.Repo); !ok || err != nil { + l.Warn("forbidden request collaborator invite not allowed", "did", evt.Record.Did, "err", err) + return nil + } + + if err := s.db.RemoveRepoCollaborator(collaborator.Subject, collaborator.Rkey); err != nil { + return fmt.Errorf("removing collaborator from db: %w", err) + } + if err := s.e.RemoveRepoCollaborator(collaborator.Subject, collaborator.Repo); err != nil { + return fmt.Errorf("removing collaborator from rbac: %w", err) + } + if err := s.tapSafeRemoveDid(ctx, collaborator.Subject); err != nil { + return fmt.Errorf("removing did from tapc: %w", err) + } + + l.Info("removed repo collaborator", "subejct", collaborator.Subject, "repo", collaborator.Repo) + return nil + } + return nil +} + +func (s *Spindle) processRepo(ctx context.Context, evt tapc.Event) error { + l := s.l.With("component", "tapIndexer", "record", evt.Record.AtUri()) + + l.Info("processing repo record") + + // only listen to members + if ok, err := s.e.IsSpindleMember(evt.Record.Did, s.cfg.Server.Did()); !ok || err != nil { + l.Warn("forbidden request: not spindle member", "did", evt.Record.Did, "err", err) + return nil + } + + switch evt.Record.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + record := tangled.Repo{} + if err := json.Unmarshal(evt.Record.Record, &record); err != nil { + return fmt.Errorf("parsing record: %w", err) + } + + domain := s.cfg.Server.Hostname + if record.Spindle == nil || *record.Spindle != domain { + if record.Spindle == nil { + l.Info("spindle isn't configured", "name", record.Name) + } else { + l.Info("different spindle configured", "name", record.Name, "spindle", *record.Spindle, "domain", domain) + } + if err := s.db.DeleteRepo(evt.Record.Did, evt.Record.Rkey); err != nil { + return fmt.Errorf("deleting repo from db: %w", err) + } + return nil + } + + if err := s.db.PutRepo(&db.Repo{ + Did: evt.Record.Did, + Rkey: evt.Record.Rkey, + Name: record.Name, + Knot: record.Knot, + }); err != nil { + return fmt.Errorf("adding repo to db: %w", err) + } + + if err := s.e.AddRepo(evt.Record.AtUri()); err != nil { + return fmt.Errorf("adding repo to rbac") + } + + // add this knot to the event consumer + src := eventconsumer.NewKnotSource(record.Knot) + s.ks.AddSource(context.Background(), src) + + l.Info("added repo", "repo", evt.Record.AtUri()) + return nil + + case tapc.RecordDeleteAction: + // check perms for this user + if ok, err := s.e.IsRepoOwner(evt.Record.Did, evt.Record.AtUri()); !ok || err != nil { + l.Warn("forbidden request: not repo owner", "did", evt.Record.Did, "err", err) + return nil + } + + if err := s.db.DeleteRepo(evt.Record.Did, evt.Record.Rkey); err != nil { + return fmt.Errorf("deleting repo from db: %w", err) + } + + if err := s.e.DeleteRepo(evt.Record.AtUri()); err != nil { + return fmt.Errorf("deleting repo from rbac: %w", err) + } + + l.Info("deleted repo", "repo", evt.Record.AtUri()) + return nil + } + return nil +} + +func (s *Spindle) processPull(ctx context.Context, evt tapc.Event) error { + l := s.l.With("component", "tapIndexer", "record", evt.Record.AtUri()) + + l.Info("processing pull record") + + switch evt.Record.Action { + case tapc.RecordCreateAction, tapc.RecordUpdateAction: + // TODO + case tapc.RecordDeleteAction: + // TODO + } + return nil +} + +func (s *Spindle) tapSafeRemoveDid(ctx context.Context, did syntax.DID) error { + known, err := s.db.IsKnownDid(syntax.DID(did)) + if err != nil { + return fmt.Errorf("ensuring did known state: %w", err) + } + if !known { + if err := s.tap.RemoveRepos(ctx, []syntax.DID{did}); err != nil { + return fmt.Errorf("removing did from tapc: %w", err) + } + } + return nil +} -- 2.51.2