From bc7d5dfc217a85e0962870f3c193f25d637ade1f Mon Sep 17 00:00:00 2001 From: Aly Raffauf Date: Wed, 16 Sep 2026 02:01:47 -0400 Subject: [PATCH] firehose: invalidate did directory cache on identity events --- internal/firehose/identity.go | 22 ++++++++++++++++++++++ internal/firehose/stream.go | 3 +++ 2 files changed, 25 insertions(+) create mode 100644 internal/firehose/identity.go diff --git a/internal/firehose/identity.go b/internal/firehose/identity.go new file mode 100644 index 0000000..347f893 --- /dev/null +++ b/internal/firehose/identity.go @@ -0,0 +1,22 @@ +package firehose + +import ( + "context" + "fmt" + + "github.com/bluesky-social/indigo/api/atproto" + "github.com/bluesky-social/indigo/atproto/syntax" +) + +func (c *Consumer) HandleIdentity(ctx context.Context, event *atproto.SyncSubscribeRepos_Identity) error { + did, err := syntax.ParseDID(event.Did) + if err != nil { + return fmt.Errorf("parse identity DID: %w", err) + } + + if err := c.Directory.Purge(ctx, did.AtIdentifier()); err != nil { + return fmt.Errorf("purge identity cache: %w", err) + } + + return nil +} diff --git a/internal/firehose/stream.go b/internal/firehose/stream.go index 616c193..1e4f8b2 100644 --- a/internal/firehose/stream.go +++ b/internal/firehose/stream.go @@ -57,6 +57,9 @@ func (c *Consumer) stream(ctx context.Context, conn *websocket.Conn) { RepoAccount: func(event *atproto.SyncSubscribeRepos_Account) error { return c.HandleAccount(ctx, event) }, + RepoIdentity: func(event *atproto.SyncSubscribeRepos_Identity) error { + return c.HandleIdentity(ctx, event) + }, } scheduler := sequential.NewScheduler("asterism", callbacks.EventHandler) -- 2.51.2