From 5da78291cfbf5b34ce774983fc2a1fa47e632c6d Mon Sep 17 00:00:00 2001 From: Aly Raffauf Date: Wed, 16 Sep 2026 02:17:22 -0400 Subject: [PATCH] firehose: handle sync events --- internal/firehose/stream.go | 3 +++ internal/firehose/sync.go | 33 +++++++++++++++++++++++++++++++++ 2 files changed, 36 insertions(+) create mode 100644 internal/firehose/sync.go diff --git a/internal/firehose/stream.go b/internal/firehose/stream.go index 1e4f8b2..6beb982 100644 --- a/internal/firehose/stream.go +++ b/internal/firehose/stream.go @@ -60,6 +60,9 @@ func (c *Consumer) stream(ctx context.Context, conn *websocket.Conn) { RepoIdentity: func(event *atproto.SyncSubscribeRepos_Identity) error { return c.HandleIdentity(ctx, event) }, + RepoSync: func(event *atproto.SyncSubscribeRepos_Sync) error { + return c.HandleSync(ctx, event) + }, } scheduler := sequential.NewScheduler("asterism", callbacks.EventHandler) diff --git a/internal/firehose/sync.go b/internal/firehose/sync.go new file mode 100644 index 0000000..243c838 --- /dev/null +++ b/internal/firehose/sync.go @@ -0,0 +1,33 @@ +package firehose + +import ( + "context" + "fmt" + + "github.com/bluesky-social/indigo/api/atproto" + "github.com/bluesky-social/indigo/atproto/syntax" +) + +func (c *Consumer) HandleSync(ctx context.Context, event *atproto.SyncSubscribeRepos_Sync) error { + if err := c.Store.SaveCursor(ctx, event.Seq); err != nil { + c.Logger.Error("could not save cursor", "err", err) + } + + did, err := syntax.ParseDID(event.Did) + if err != nil { + return fmt.Errorf("parse sync DID: %w", err) + } + + go func() { + if err := c.Store.DeleteAllLinks(ctx, did.String()); err != nil { + c.Logger.Error("could not clear synced repo links", "did", did, "err", err) + return + } + + if err := c.Backfill.Repo(ctx, did.String(), c.WantedCollections); err != nil { + c.Logger.Error("could not backfill synced repo", "did", did, "err", err) + } + }() + + return nil +} -- 2.51.2