From 8a07e7947b3c92552bb6b349dfe4fb63888eadbc Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Thu, 18 Dec 2025 12:14:09 +0000 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/vm.nix | 6 +++++- spindle/ingester.go | 276 ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ spindle/server.go | 49 +++++++++++++++++++------------------------------ spindle/tap.go | 294 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ nix/modules/spindle.nix | 10 ++-------- spindle/config/config.go | 22 +++++++++++----------- spindle/db/db.go | 42 ++++++++++++++++++++++++++++-------------- spindle/db/known_dids.go | 44 -------------------------------------------- spindle/db/repos.go | 132 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------ 9 file(s) changed, 479 insertion(s)(+), 396 deletion(s)(-) diff --git a/nix/vm.nix b/nix/vm.nix --- a/nix/vm.nix +++ b/nix/vm.nix @@ -58,6 +58,11 @@ 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 @@ 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/ingester.go b/spindle/ingester.go deleted file mode 100644 --- 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 --- a/spindle/server.go +++ b/spindle/server.go @@ -10,12 +10,12 @@ "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 @@ "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 @@ 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 @@ 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, @@ -138,11 +121,6 @@ cursorStore, err := cursor.NewSQLiteStore(cfg.Server.DBPath) if err != nil { 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 @@ -226,6 +204,18 @@ 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 @@ } // 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 --- /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 +} diff --git a/nix/modules/spindle.nix b/nix/modules/spindle.nix --- a/nix/modules/spindle.nix +++ b/nix/modules/spindle.nix @@ -53,12 +53,6 @@ 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 @@ 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 @@ "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 @@ "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/spindle/config/config.go b/spindle/config/config.go --- a/spindle/config/config.go +++ b/spindle/config/config.go @@ -9,17 +9,17 @@ ) 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 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -5,6 +5,7 @@ "database/sql" "strings" + "github.com/bluesky-social/indigo/atproto/syntax" _ "github.com/mattn/go-sqlite3" "tangled.org/core/log" "tangled.org/core/orm" @@ -55,6 +56,20 @@ addedAt text not null default (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')), 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 ( @@ -120,18 +135,17 @@ 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 --- 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 --- 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 } -func (d *DB) AddRepo(knot, owner, name string) error { - _, err := d.Exec(`insert or ignore into repos (knot, owner, name) values (?, ?, ?)`, knot, owner, name) +type RepoCollaborator struct { + Did syntax.DID + Rkey syntax.RecordKey + Repo syntax.ATURI + Subject syntax.DID +} + +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 @@ 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 } -- tangled.sh