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)