From f9ea3744b7795fcf072dd93946ba604719f712ab Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Wed, 2 Apr 2025 22:09:37 -0700 Subject: [PATCH 1/2] repo: clarify commit loading name, and return CID --- atproto/repo/car.go | 31 +++++++++++++++--------------- atproto/repo/cmd/repo-tool/main.go | 4 ++-- atproto/repo/sync.go | 4 ++-- atproto/repo/sync_test.go | 2 +- cmd/goat/firehose.go | 6 +++--- cmd/goat/repo.go | 6 +++--- 6 files changed, 27 insertions(+), 26 deletions(-) diff --git a/atproto/repo/car.go b/atproto/repo/car.go index d18f284f..b67f1c83 100644 --- a/atproto/repo/car.go +++ b/atproto/repo/car.go @@ -11,10 +11,14 @@ import ( "github.com/bluesky-social/indigo/atproto/syntax" blocks "github.com/ipfs/go-block-format" + "github.com/ipfs/go-cid" "github.com/ipld/go-car" ) -func LoadFromCAR(ctx context.Context, r io.Reader) (*Commit, *Repo, error) { +var ErrNoRoot = errors.New("CAR file missing root CID") +var ErrNoCommit = errors.New("no commit") + +func LoadRepoFromCAR(ctx context.Context, r io.Reader) (*Commit, *Repo, error) { //bs := blockstore.NewBlockstore(datastore.NewMapDatastore()) bs := NewTinyBlockstore() @@ -73,21 +77,18 @@ func LoadFromCAR(ctx context.Context, r io.Reader) (*Commit, *Repo, error) { return &commit, &repo, nil } -var ErrNoRoot = errors.New("CAR file missing root CID") -var ErrNoCommit = errors.New("no commit") - -// LoadCARCommit is like LoadFromCAR() but filters to only return the commit object. -// useful for subscribeRepos/firehose `#sync` message -func LoadCARCommit(ctx context.Context, r io.Reader) (*Commit, error) { +// LoadCommitFromCAR is like LoadRepoFromCAR() but filters to only return the commit object. +// Also returns the commit CID. +func LoadCommitFromCAR(ctx context.Context, r io.Reader) (*Commit, *cid.Cid, error) { cr, err := car.NewCarReader(r) if err != nil { - return nil, err + return nil, nil, err } if cr.Header.Version != 1 { - return nil, fmt.Errorf("unsupported CAR file version: %d", cr.Header.Version) + return nil, nil, fmt.Errorf("unsupported CAR file version: %d", cr.Header.Version) } if len(cr.Header.Roots) < 1 { - return nil, ErrNoRoot + return nil, nil, ErrNoRoot } commitCID := cr.Header.Roots[0] var commitBlock blocks.Block @@ -97,7 +98,7 @@ func LoadCARCommit(ctx context.Context, r io.Reader) (*Commit, error) { if err == io.EOF { break } - return nil, err + return nil, nil, err } if blk.Cid().Equals(commitCID) { @@ -106,14 +107,14 @@ func LoadCARCommit(ctx context.Context, r io.Reader) (*Commit, error) { } } if commitBlock == nil { - return nil, ErrNoCommit + return nil, nil, ErrNoCommit } var commit Commit if err := commit.UnmarshalCBOR(bytes.NewReader(commitBlock.RawData())); err != nil { - return nil, fmt.Errorf("parsing commit block from CAR file: %w", err) + return nil, nil, fmt.Errorf("parsing commit block from CAR file: %w", err) } if err := commit.VerifyStructure(); err != nil { - return nil, fmt.Errorf("parsing commit block from CAR file: %w", err) + return nil, nil, fmt.Errorf("parsing commit block from CAR file: %w", err) } - return &commit, nil + return &commit, &commitCID, nil } diff --git a/atproto/repo/cmd/repo-tool/main.go b/atproto/repo/cmd/repo-tool/main.go index 5c3848bc..a1137133 100644 --- a/atproto/repo/cmd/repo-tool/main.go +++ b/atproto/repo/cmd/repo-tool/main.go @@ -93,7 +93,7 @@ func runVerifyCarMst(cctx *cli.Context) error { } defer f.Close() - commit, repo, err := repo.LoadFromCAR(ctx, f) + commit, repo, err := repo.LoadRepoFromCAR(ctx, f) if err != nil { return err } @@ -125,7 +125,7 @@ func runVerifyCarSignature(cctx *cli.Context) error { } defer f.Close() - commit, _, err := repo.LoadFromCAR(ctx, f) + commit, _, err := repo.LoadRepoFromCAR(ctx, f) if err != nil { return err } diff --git a/atproto/repo/sync.go b/atproto/repo/sync.go index 6f94ed22..46eb5d65 100644 --- a/atproto/repo/sync.go +++ b/atproto/repo/sync.go @@ -40,7 +40,7 @@ func VerifyCommitMessage(ctx context.Context, msg *comatproto.SyncSubscribeRepos logger.Warn("event with rebase flag set") } - commit, repo, err := LoadFromCAR(ctx, bytes.NewReader([]byte(msg.Blocks))) + commit, repo, err := LoadRepoFromCAR(ctx, bytes.NewReader([]byte(msg.Blocks))) if err != nil { return nil, err } @@ -168,7 +168,7 @@ func parseCommitOps(ops []*comatproto.SyncSubscribeRepos_RepoOp) ([]Operation, e // // TODO: in real implementation, will want to merge this code with `VerifyCommitMessage` above, and have it hanging off some service struct with a configured `identity.Directory` func VerifyCommitSignature(ctx context.Context, dir identity.Directory, msg *comatproto.SyncSubscribeRepos_Commit) error { - commit, _, err := LoadFromCAR(ctx, bytes.NewReader([]byte(msg.Blocks))) + commit, _, err := LoadRepoFromCAR(ctx, bytes.NewReader([]byte(msg.Blocks))) if err != nil { return err } diff --git a/atproto/repo/sync_test.go b/atproto/repo/sync_test.go index 4ea75f98..313cf4ed 100644 --- a/atproto/repo/sync_test.go +++ b/atproto/repo/sync_test.go @@ -47,7 +47,7 @@ func testCommitFile(t *testing.T, p string) { _, err = VerifyCommitMessage(ctx, &msg) assert.NoError(err) if err != nil { - _, repo, err := LoadFromCAR(ctx, bytes.NewReader([]byte(msg.Blocks))) + _, repo, err := LoadRepoFromCAR(ctx, bytes.NewReader([]byte(msg.Blocks))) if err != nil { t.Fail() } diff --git a/cmd/goat/firehose.go b/cmd/goat/firehose.go index c2fd9508..9b331ba4 100644 --- a/cmd/goat/firehose.go +++ b/cmd/goat/firehose.go @@ -256,7 +256,7 @@ func (gfc *GoatFirehoseConsumer) handleAccountEvent(ctx context.Context, evt *co } func (gfc *GoatFirehoseConsumer) handleSyncEvent(ctx context.Context, evt *comatproto.SyncSubscribeRepos_Sync) error { - commit, err := repo.LoadCARCommit(ctx, bytes.NewReader(evt.Blocks)) + commit, _, err := repo.LoadCommitFromCAR(ctx, bytes.NewReader(evt.Blocks)) if err != nil { return err } @@ -296,7 +296,7 @@ func (gfc *GoatFirehoseConsumer) handleCommitEvent(ctx context.Context, evt *com return err } - commit, err := repo.LoadCARCommit(ctx, bytes.NewReader(evt.Blocks)) + commit, _, err := repo.LoadCommitFromCAR(ctx, bytes.NewReader(evt.Blocks)) if err != nil { return err } @@ -409,7 +409,7 @@ func (gfc *GoatFirehoseConsumer) handleCommitEventOps(ctx context.Context, evt * return nil } - _, rr, err := repo.LoadFromCAR(ctx, bytes.NewReader(evt.Blocks)) + _, rr, err := repo.LoadRepoFromCAR(ctx, bytes.NewReader(evt.Blocks)) if err != nil { logger.Error("failed to read repo from car", "err", err) return nil diff --git a/cmd/goat/repo.go b/cmd/goat/repo.go index d869fa95..b532d1b7 100644 --- a/cmd/goat/repo.go +++ b/cmd/goat/repo.go @@ -184,7 +184,7 @@ func runRepoList(cctx *cli.Context) error { } // read repository tree in to memory - _, r, err := repo.LoadFromCAR(ctx, fi) + _, r, err := repo.LoadRepoFromCAR(ctx, fi) if err != nil { return fmt.Errorf("failed to parse repo CAR file: %w", err) } @@ -211,7 +211,7 @@ func runRepoInspect(cctx *cli.Context) error { } // read repository tree in to memory - c, _, err := repo.LoadFromCAR(ctx, fi) + c, _, err := repo.LoadRepoFromCAR(ctx, fi) if err != nil { return err } @@ -255,7 +255,7 @@ func runRepoUnpack(cctx *cli.Context) error { return err } - c, r, err := repo.LoadFromCAR(ctx, fi) + c, r, err := repo.LoadRepoFromCAR(ctx, fi) if err != nil { return err } -- 2.51.2 From 874660fc43d6bc6e88fafde8da458fb0a90df120 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Sun, 6 Apr 2025 21:28:27 -0700 Subject: [PATCH 2/2] patch existing relay for repo rename --- cmd/relay/bgs/validator.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/cmd/relay/bgs/validator.go b/cmd/relay/bgs/validator.go index 11a7db3c..1df8d8ff 100644 --- a/cmd/relay/bgs/validator.go +++ b/cmd/relay/bgs/validator.go @@ -174,7 +174,7 @@ func (val *Validator) VerifyCommitMessage(ctx context.Context, host *models.PDS, hasWarning = true } - commit, repoFragment, err := atrepo.LoadFromCAR(ctx, bytes.NewReader([]byte(msg.Blocks))) + commit, repoFragment, err := atrepo.LoadRepoFromCAR(ctx, bytes.NewReader([]byte(msg.Blocks))) if err != nil { commitVerifyErrors.WithLabelValues(hostname, "car").Inc() return nil, err @@ -326,7 +326,7 @@ func (val *Validator) HandleSync(ctx context.Context, host *models.PDS, msg *atp return nil, err } - commit, err := atrepo.LoadCARCommit(ctx, bytes.NewReader([]byte(msg.Blocks))) + commit, _, err := atrepo.LoadCommitFromCAR(ctx, bytes.NewReader([]byte(msg.Blocks))) if err != nil { commitVerifyErrors.WithLabelValues(hostname, "car").Inc() return nil, err -- 2.51.2