From 7d4a6769cf3298ffdee2ec869b9b11f19638b662 Mon Sep 17 00:00:00 2001 From: dawn Date: Thu, 16 Jul 2026 20:24:31 +0300 Subject: [PATCH] spindle/{ingester,server},jetstream: ingest repo.pull Signed-off-by: dawn --- jetstream/jetstream.go | 19 ++++++++++++++++--- spindle/ingester.go | 4 +++- spindle/server.go | 3 +++ 3 files changed, 22 insertions(+), 4 deletions(-) diff --git a/jetstream/jetstream.go b/jetstream/jetstream.go index 14d2b546..10e8e139 100644 --- a/jetstream/jetstream.go +++ b/jetstream/jetstream.go @@ -30,8 +30,9 @@ type JetstreamClient struct { ident string l *slog.Logger - logDids bool - wantedDids Set[string] + logDids bool + wantedDids Set[string] + unfilteredNsids Set[string] db DB waitForDid bool mu sync.RWMutex @@ -55,6 +56,12 @@ func (j *JetstreamClient) AddDid(did string) { j.mu.Unlock() } +func (j *JetstreamClient) ExemptCollection(nsid string) { + j.mu.Lock() + j.unfilteredNsids[nsid] = struct{}{} + j.mu.Unlock() +} + func (j *JetstreamClient) RemoveDid(did string) { if did == "" { return @@ -82,6 +89,11 @@ func (j *JetstreamClient) withDidFilter(processFunc processor) processor { matches = true } } + if !matches && evt.Commit != nil { + if _, ok := j.unfilteredNsids[evt.Commit.Collection]; ok { + matches = true + } + } j.mu.RUnlock() var err error @@ -106,7 +118,8 @@ func NewJetstreamClient(endpoint, ident string, collections []string, cfg *clien ident: ident, db: db, l: logger, - wantedDids: make(map[string]struct{}), + wantedDids: make(map[string]struct{}), + unfilteredNsids: make(map[string]struct{}), logDids: logDids, diff --git a/spindle/ingester.go b/spindle/ingester.go index b7366e94..f1fa8ae1 100644 --- a/spindle/ingester.go +++ b/spindle/ingester.go @@ -27,7 +27,7 @@ func (s *Spindle) ingest() Ingester { switch e.Commit.Collection { case tangled.SpindleMemberNSID: err = s.ingestMember(ctx, e) - case tangled.RepoNSID, tangled.RepoCollaboratorNSID: + case tangled.RepoNSID, tangled.RepoCollaboratorNSID, tangled.RepoPullNSID: if evt, ok := jetstreamToTapEvent(e); ok { err = s.tap.processEvent(ctx, evt) } @@ -68,6 +68,8 @@ func jetstreamToTapEvent(e *models.Event) (tapc.Event, bool) { Collection: syntax.NSID(e.Commit.Collection), Action: action, Record: e.Commit.Record, + // jetstream is only used for live + Live: true, }, }, true } diff --git a/spindle/server.go b/spindle/server.go index 26cd3643..dd51d281 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -123,12 +123,15 @@ func New(ctx context.Context, cfg *config.Config, d *db.DB, engines map[string]m tangled.SpindleMemberNSID, tangled.RepoNSID, tangled.RepoCollaboratorNSID, + tangled.RepoPullNSID, } 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) + // pull records are created by arbitrary users too, same hack as in tap + jc.ExemptCollection(tangled.RepoPullNSID) // Check if the spindle knows about any Dids; dids, err := d.GetAllDids() -- 2.51.2