fix(p2p/compression) - upload header synchronization - #2919
Conversation
❌ 3 Tests Failed:
View the full list of 3 ❄️ flaky test(s)
To view more test analytics, go to the Test Analytics Dashboard |
There was a problem hiding this comment.
Code Review
This pull request refactors the block streaming chunker and build file reading logic to support on-demand frame table refreshing and recovery from peer transitions or missing object errors. It replaces the StreamingReader interface with RangeOpener and updates StorageDiff to dynamically refresh the frame table from the ancestor header when needed. No review comments were provided, and there are no critical findings to report.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 4d2309386f
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 87365d6629
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
Add a FullFrameTable wrapper for the two narrow contexts where the FT covers a whole file: (a) fresh out of compressStream / CompressBytes / StoreFile / UploadFramed, and (b) the future StorageDiff fallback latch (#2919). Producers return *FullFrameTable; downcast via Table() to feed the hot read path or assign to BuildData.FrameData. FrameTable itself stays the authoritative type — what BuildData stores, what DeserializeFrameTable and TrimToRanges return, what ReadAt and OpenRangeReader accept. No on-disk format change, no behavior change.
| // peer says "go to storage" → refresh authoritative header → continue. | ||
| // No 404-driven recovery. | ||
| func initialSize(ctx context.Context, upstream storage.Seekable) (size int64, ok bool, err error) { | ||
| if _, peerRouted := upstream.(peerclient.PeerRouted); !peerRouted { |
There was a problem hiding this comment.
can we avoid the cast "upstream.(peerclient.PeerRouted)" or not really?
There was a problem hiding this comment.
This is the signal - the upstream object was just resolved, and it's either a peerRouted or not... how else would I figure it out? We could add a return value and call a method but it'd be redundant.
| timer := time.NewTimer(transErr.RetryAfter) | ||
| defer timer.Stop() | ||
|
|
||
| select { | ||
| case <-timer.C: | ||
| case <-ctx.Done(): | ||
| return false, ctx.Err() | ||
| return ctx.Err() | ||
| } |
There was a problem hiding this comment.
tiny syntactic improvement, but you can use time.After(duration) instead of creating and managing a timer yourself:
select {
case <-time.After(transErr.RetryAfter):
case <-ctx.Done():
return false, ctx.Err()
}| } | ||
|
|
||
| if size == 0 { | ||
| // (d) and degenerate (a) where bd.Size was zero. Ask storage directly. |
| // resolve picks the (upstream, ft) the next read should use, given the | ||
| // caller's per-mapping FT hint. The contract: if there is no authoritative FT | ||
| // latched AND no peer currently serving this build, we MUST refresh before | ||
| // reading. The latched upstream was opened at the bootstrap-guessed CT path |
| return b.reloadSource(ctx, refreshCausePeerTransitioned) | ||
| } | ||
|
|
||
| // reloadSource is the idempotent ensure-latched entry. Both RefreshSource |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 34524c33a7
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| func (b *StorageDiff) RefreshSource(ctx context.Context) error { | ||
| return b.reloadSource(ctx, refreshCausePeerTransitioned) | ||
| } | ||
|
|
||
| // reloadSource is the idempotent ensure-latched entry. Both RefreshSource | ||
| // (PeerTransitionedError) and resolve (read-time peer-left fallback) funnel | ||
| // through it; the cause attribute distinguishes them in telemetry. A concurrent | ||
| // caller that wins the mutex short-circuits when the latch is already | ||
| // populated, so parallel segment reads on a fresh StorageDiff pay only one | ||
| // header fetch. | ||
| func (b *StorageDiff) reloadSource(ctx context.Context, cause string) error { | ||
| b.refreshMu.Lock() | ||
| defer b.refreshMu.Unlock() | ||
| if b.source.Load().fullDiffFrameTable != nil { | ||
| return nil | ||
| } | ||
|
|
||
| return b.reloadSourceLocked(ctx, cause) | ||
| } |
There was a problem hiding this comment.
nit: do we need both these functions?
| // authoritative empty &{} FT at construction, so resolve short-circuits and | ||
| // reloadSource is never called; V3 builds aren't peer-routed so | ||
| // PeerTransitionedError never fires against them either. | ||
| func (b *StorageDiff) reloadSourceLocked(ctx context.Context, cause string) error { |
There was a problem hiding this comment.
this looks to be called only in reloadSource, making merge of the 3 functions to 1
There was a problem hiding this comment.
I'll send a separate PR for this; agreed.
Main is broken: #2919 moved the upstream `Seekable` out of `block.NewChunker` and into `Chunker.Slice` (as a `storage.RangeOpener` param), and #2985 merged alongside it without picking up the new signatures, so `cmd/inspect-build` fails to compile (see [this run](https://github.com/e2b-dev/infra/actions/runs/27380475452/job/80915690448)). Updates `validate.go` to the new API: `openChunker` now also returns the storage object, which is passed to each `Slice` call — mirroring the production read path in `storage_diff.go`.
…ng (#2994) Since #2919, `createDiff`/`reloadSource` handle an ancestor build missing from a V4+ header's `Builds` map by loading the ancestor's own header and requiring a self entry (`SelfBuildData`). Pre-V4 ancestor headers have no `Builds` map at all — `appendAncestorBuilds` deliberately skips them on upload — so the lookup fails, `planRead` errors, the guest gets EIO from NBD, and envd init returns 400. On staging this kills every resume of pre-V4-era templates: ~650 creates/h failing with INTERNAL since the `ed6f90a` deploy (previously 0), observable as the disappearance of the 1 GiB snapshot cohort and all pause-path rootfs uploads. Fix: when the refreshed ancestor header is pre-V4, latch it as authoritatively uncompressed (`UncompressedFullFrameTable`, data at the basic path, size from storage) in both refresh paths — same semantics the read path had before #2919. V4+ headers keep the strict self-entry requirement (a missing self entry there still indicates a peer's in-flight header and must fail loudly). Regression test: V4 header mapping to a V3 ancestor with no Builds entry now reads raw bytes from the uncompressed path.
…ev#2989) Main is broken: e2b-dev#2919 moved the upstream `Seekable` out of `block.NewChunker` and into `Chunker.Slice` (as a `storage.RangeOpener` param), and e2b-dev#2985 merged alongside it without picking up the new signatures, so `cmd/inspect-build` fails to compile (see [this run](https://github.com/e2b-dev/infra/actions/runs/27380475452/job/80915690448)). Updates `validate.go` to the new API: `openChunker` now also returns the storage object, which is passed to each `Slice` call — mirroring the production read path in `storage_diff.go`.
…ng (e2b-dev#2994) Since e2b-dev#2919, `createDiff`/`reloadSource` handle an ancestor build missing from a V4+ header's `Builds` map by loading the ancestor's own header and requiring a self entry (`SelfBuildData`). Pre-V4 ancestor headers have no `Builds` map at all — `appendAncestorBuilds` deliberately skips them on upload — so the lookup fails, `planRead` errors, the guest gets EIO from NBD, and envd init returns 400. On staging this kills every resume of pre-V4-era templates: ~650 creates/h failing with INTERNAL since the `ed6f90a` deploy (previously 0), observable as the disappearance of the 1 GiB snapshot cohort and all pause-path rootfs uploads. Fix: when the refreshed ancestor header is pre-V4, latch it as authoritatively uncompressed (`UncompressedFullFrameTable`, data at the basic path, size from storage) in both refresh paths — same semantics the read path had before e2b-dev#2919. V4+ headers keep the strict self-entry requirement (a missing self entry there still indicates a peer's in-flight header and must fail loudly). Regression test: V4 header mapping to a V3 ancestor with no Builds entry now reads raw bytes from the uncompressed path.
Problem
A descendant B's header doesn't reflect ancestor A's data until B finishes uploading. Resume B early — peer-served, mid-build — and reads of A through B's header can't pick the right path (uncompressed vs .zstd) or hit a peer that just transitioned. Result: silent wrong-bytes or ErrObjectNotExist against the wrong CT.
Mechanism
Refresh loads A's own {A}/{filetype}.header sidecar, extracts A's full-file self-FT, reopens upstream at the resolved CT, and atomically stores {upstream, fullDiffFrameTable} into StorageDiff.source as one unit. Idempotent latch: concurrent calls short-circuit.
Diff highlights