diff --git a/nix/modules/spindle.nix b/nix/modules/spindle.nix index 803c5e4c..c4fa6108 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; @@ -149,7 +143,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"; @@ -159,7 +153,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}" @@ -167,6 +160,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..6b9f9caf 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/tap" "tangled.org/core/xrpc/serviceauth" ) @@ -34,7 +35,7 @@ import ( var defaultMotd []byte type Spindle struct { - jc *jetstream.JetstreamClient + tap *tap.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 := tap.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, &tap.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..2df38724 --- /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/tap" +) + +func (s *Spindle) processEvent(ctx context.Context, evt tap.Event) error { + l := s.l.With("component", "tapIndexer") + + var err error + switch evt.Type { + case tap.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 tap.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 tap.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 tap.RecordCreateAction, tap.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 tap: %w", err) + } + + l.Info("added member", "member", record.Subject) + return nil + + case tap.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 tap: %w", err) + } + + l.Info("removed member", "member", member.Subject) + return nil + } + return nil +} + +func (s *Spindle) processCollaborator(ctx context.Context, evt tap.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 tap.RecordCreateAction, tap.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 tap: %w", err) + } + + l.Info("add repo collaborator", "subejct", record.Subject, "repo", record.Repo) + return nil + + case tap.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 tap: %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 tap.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 tap.RecordCreateAction, tap.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 tap.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 tap.Event) error { + l := s.l.With("component", "tapIndexer", "record", evt.Record.AtUri()) + + l.Info("processing pull record") + + switch evt.Record.Action { + case tap.RecordCreateAction, tap.RecordUpdateAction: + // TODO + case tap.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 tap: %w", err) + } + } + return nil +}