From c859f5b887038386ad0de2d9d448103d4b8357ba Mon Sep 17 00:00:00 2001 From: Aly Raffauf Date: Sat, 4 Jul 2026 11:25:50 -0400 Subject: [PATCH] firehose: verify commit sigs against repo signing keys --- README.md | 1 + cmd/asterism/main.go | 1 + internal/firehose/commit.go | 46 ++++++++++++++++++++++++++++++++++- internal/firehose/consumer.go | 2 ++ 4 files changed, 49 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index 6743547..fce99ec 100644 --- a/README.md +++ b/README.md @@ -179,6 +179,7 @@ Note the DID filter parameter is `did` here, not `linkDid` like `getManyToMany` - [x] Graceful shutdown and Firehose reconnect - [x] CI + Dockerfile - [x] Add health endpoint +- [x] Verify commits against the repo's signing key (not just CID/hash consistency) **Medium term** diff --git a/cmd/asterism/main.go b/cmd/asterism/main.go index 13beb5a..ff8bc1f 100644 --- a/cmd/asterism/main.go +++ b/cmd/asterism/main.go @@ -103,6 +103,7 @@ func main() { consumer := &firehose.Consumer{ WantedCollections: wantedCollections, Store: linkStore, + Directory: bf.Directory, Backfill: bf, } diff --git a/internal/firehose/commit.go b/internal/firehose/commit.go index 4bd0009..f968cdd 100644 --- a/internal/firehose/commit.go +++ b/internal/firehose/commit.go @@ -5,16 +5,18 @@ import ( "context" "fmt" "strings" + "time" "github.com/alyraffauf/asterism/internal/index" "github.com/bluesky-social/indigo/api/atproto" + "github.com/bluesky-social/indigo/atproto/syntax" indigorepo "github.com/bluesky-social/indigo/repo" "github.com/ipfs/go-cid" ) func (c *Consumer) HandleCommit(ctx context.Context, event *atproto.SyncSubscribeRepos_Commit) error { // fmt.Println("repo:", event.Repo, "commit:", event.Rev) - // + if err := c.Store.SaveCursor(ctx, event.Seq); err != nil { fmt.Println("could not save cursor:", err) } @@ -33,6 +35,15 @@ func (c *Consumer) HandleCommit(ctx context.Context, event *atproto.SyncSubscrib return err } + if !c.hasWantedOps(event.Ops) { + return nil + } + + if err := c.verifyCommit(ctx, repo); err != nil { + fmt.Println("commit verification failed:", event.Repo, err) + return nil + } + for _, operation := range event.Ops { if err := c.handleOperation(ctx, event, repo, operation); err != nil { fmt.Println("could not handle operation:", err) @@ -73,3 +84,36 @@ func (c *Consumer) handleOperation(ctx context.Context, event *atproto.SyncSubsc return index.Record(ctx, c.Store, event.Repo, collection, recordKey, recordCid.String(), event.Rev, *recordBytes) } + +func (c *Consumer) hasWantedOps(ops []*atproto.SyncSubscribeRepos_RepoOp) bool { + for _, op := range ops { + collection, _, ok := strings.Cut(op.Path, "/") + if ok && c.wants(collection) { + return true + } + } + return false +} + +func (c *Consumer) verifyCommit(ctx context.Context, repo *indigorepo.Repo) error { + sc := repo.SignedCommit() + + resolveCtx, cancel := context.WithTimeout(ctx, 10*time.Second) + identity, err := c.Directory.LookupDID(resolveCtx, syntax.DID(sc.Did)) + cancel() + if err != nil { + return fmt.Errorf("resolve did: %w", err) + } + + pubKey, err := identity.PublicKey() + if err != nil { + return fmt.Errorf("get public key: %w", err) + } + + unsignedBytes, err := sc.Unsigned().BytesForSigning() + if err != nil { + return fmt.Errorf("marshal unsigned commit: %w", err) + } + + return pubKey.HashAndVerify(unsignedBytes, sc.Sig) +} diff --git a/internal/firehose/consumer.go b/internal/firehose/consumer.go index 7bfb54b..8f5048e 100644 --- a/internal/firehose/consumer.go +++ b/internal/firehose/consumer.go @@ -3,11 +3,13 @@ package firehose import ( "github.com/alyraffauf/asterism/internal/backfill" "github.com/alyraffauf/asterism/internal/store" + "github.com/bluesky-social/indigo/atproto/identity" ) type Consumer struct { WantedCollections map[string]struct{} Store *store.Store + Directory identity.Directory Backfill *backfill.Backfill } -- 2.51.2