Repository navigation
Conversation
Every keyed write checked one 32 MiB ceiling on a batch's whole Arrow memory, which includes a managed Blob value's bytes, so a value of exactly 32 MiB never fit beside its row. KeyedBytes splits that count into payload (each managed value's logical length) and framing (the rest), each under its own 32 MiB ceiling, and every keyed check uses it: staging, keyed-write preparation, the update carry scan, branch merge buffering and proven-insert chunking, and the compatibility loader's pre-decode forecast. The halves sum to the old count. A payload refusal names its resource with 'Blob payload bytes' in place of 'bytes'. Two Lance guards pin what the Blob write plan relies on: a whole-row merge-update keeps the row's stable row id, and merge-insert refuses an external reference outside the dataset's bases.
Session::put_blob_at_as and clear_blob_at_as replace or clear one Blob value of an existing node or edge, addressed by exact id like a Blob read. A cell write is a Mutation-protocol write: it stages one upsert of the carried row and runs the Mutation tail, now split into commit_staged_mutation and publish_committed_mutation, which mutate uses too. The target's old payload is never read; the row's other cells are carried as an update carries them and share the 32 MiB payload allowance with the new value, which is bounded at 32 MiB inclusive before anything is captured. An optional precondition (an ETag set, or any existing value) is evaluated against the attempt's base and re-evaluated when a competing commit forces a re-prepare; a change of branch incarnation or accepted schema between attempts fails closed. The receipt's ETag is read from the detached commit before publication.
The compatibility loader's pre-decode forecast charged every `base64:` string as Blob payload, so a declared String whose text starts with `base64:` used the payload ceiling: a 24 MiB Blob beside such a 12 MiB String was refused at 33 MiB although the batch check admits it. The forecast now takes the type's Blob properties, as the Arrow check does: only a Blob property's `base64:` value is payload; every other value is framing.
# Conflicts: # crates/omnigraph-server/src/blob_transport.rs
…udget # Conflicts: # crates/omnigraph/src/table_store.rs
# Conflicts: # docs/dev/blob.md
ragnorc
left a comment
There was a problem hiding this comment.
Recommendation: hold the stack's merge until the minimal GQT regression requested on #903 is included. I found no new blocking code defect in step 3A. Reviewed cffa7918f3467a3a73a3f92ab487b2f0117ab501, including the changes above parent a4429d4508a64f80df1b1877dc251e724fbd5b86. This is a COMMENT review because the authenticated account is the author.
This PR lets an embedded caller replace or clear one Blob on an existing node or edge. The caller supplies an exact entity ID and can require a matching ETag. The engine carries the rest of the row, then publishes one graph commit. Missing IDs cannot become inserts. Empty bytes and null remain distinct. HTTP and CLI write commands remain separate work.
Correctness and design. The contract needs one accepted write snapshot, atomic graph visibility, preserved sibling values, and a receipt that identifies this write. The implementation addresses those obligations in their existing owners:
- The shared mutation tail validates and stages the row, commits it detached, and publishes through the existing manifest CAS. Ordinary mutations use the same tail. The write guards remain held through publication.
- Each retry captures fresh authority, checks the precondition again, and carries fresh sibling values. Changed branch or schema identity ends the retry. This covers both nodes and edges without a special storage path.
- The receipt comes from the detached version before publication. A later writer cannot change its ETag. Reads and writes share the same descriptor lookup and ETag calculation.
Tradeoffs and liability. This suits bounded replacement of an existing cell. It still rewrites the whole row. The new value and carried Blobs share 32 MiB, so a large sibling can prevent an otherwise small put. An external sibling needs policy admission and a readable source. The target's old payload is omitted from the carry, so replacing a stale external target does not read it.
Those costs follow the pinned substrate. In Lance 11.0.0, UpdateBuilder refuses Blob-v2 columns. The merge-insert strategy selector also excludes them from RewriteColumns. Its whole-row writer uses default external-reference restrictions. I checked the complete relevant upstream guides against this exact implementation.
The ETag also changes after another row in the same table changes. This is conservative and can reject a conditional write whose payload remains unchanged. Each put adds a descriptor read before publication. I did not measure throughput, contention, cloud request cost, or memory peaks. The 32 MiB payload allowance is not an RSS limit.
The PR adds public APIs, precondition/error/receipt contracts, and retry identity checks that need continued tests. It adds no durable format or second publication protocol. Sharing the mutation tail avoids duplicate validation and publication code. Five similar features can use that owner. The added liability is proportionate to the new interface, provided future surfaces retain this common path.
Tests and the existing merge prerequisite. GQT cannot call these new Session methods or assert their byte streams and ETags. The existing Rust owners are appropriate for that contract and its races. The parent commit has no cell-write API, so a compile failure there would not prove a regression.
The inherited payload/framing fix is different: GQT can express its boundary. The existing #903 review requests a compact permanent case, and the author's response proposes one after #927. This head still contains no such case. Keep that requirement before merging the stack. I did not duplicate it as a new inline defect here. No optional improvement is required by this review.
Local validation on the exact reviewed head passed:
- 13 focused tests for cell writes, payload limits, carried siblings, policy, retries, receipts, and shared Blob reads.
- All 28 put/clear cases in the default crash matrix, including fresh, read-only, and same-handle recovery where applicable.
- All 35 architecture guards, formatting, the 233-file documentation check, and the agent-link check.
The test commands used this environment and common prefix:
export RUSTUP_TOOLCHAIN=1.97.1
export CARGO_TARGET_DIR=/tmp/review-928/target
export RUSTFLAGS='--cfg tokio_unstable --cfg tokio_unstable'
# Append the owner and filter from the table below:
cargo test --locked -p omnigraph-engine --features failpoints -j 3 --test OWNER FILTER -- --nocapture| Owner | Filter |
|---|---|
end_to_end |
blob_put_and_clear_replace_one_cell_by_exact_id (--exact), then blob_read_ |
writes |
blob_put_, then blob_write_under_deny_names_the_carried_reference_and_whole_row_writes_recover (--exact) |
writes |
exact_limit_blob_payload_fits_beside_its_row_and_one_more_byte_is_refused (--exact) |
failpoints |
blob_put |
policy_engine_chassis |
blob_put_and_clear_enforce_change_for_the_actor (--exact) |
detached_commit_matrix |
rfc_0067_failure_window_matrix (--exact), with OMNIGRAPH_MATRIX_WRITERS=BlobPut,BlobClear |
forbidden_apis |
No filter |
I also reran the identical temporary GQT boundary probe from the #903 review. It extends blob_update_carries_unassigned_blobs.gqt with one insert and a 32 MiB Blob parameter filled with byte 0x07. The assertion expects one affected node.
| Same regression | Commit | Result |
|---|---|---|
| Saved before-fix execution from the #903 review | b8b631ee8f813e557805baf33efe869011f48275 |
Failed at the added insert: keyed entity bytes for node:Doc, actual 33,556,272, limit 33,554,432 |
| Execution in this review | cffa7918f3467a3a73a3f92ab487b2f0117ab501 |
One ordinary filesystem execution passed in 8.19 seconds |
Both reports identify the same fixture digest, fb499d12646ca7dd0f58dddc909d8fa8a0ab508da8036511908c0bdbc0457892. The before result is retained evidence, not a new run. The head result executed the test, with no setup failure or skipped selected test. DST environments were unselected. The temporary file uses a 60-second timeout and 44.7 MB of literal text. It is evidence for the compact permanent case requested above, not a proposed fixture to commit.
Exact GQT commands, with the toolchain and flags above:
# Saved before-fix run at b8b631ee8f813e557805baf33efe869011f48275:
CARGO_TARGET_DIR=/tmp/review-903-followup/target cargo run -p omnigraph-gqt --locked --features omnigraph/failpoints --bin omnigraph-gqt -- /tmp/review-903-followup/blob_limit.gqt --target omnigraph-engine --storage local-filesystem --artifacts /tmp/review-903-followup/gqt-before-artifacts
# Current run at cffa7918f3467a3a73a3f92ab487b2f0117ab501:
CARGO_TARGET_DIR=/tmp/review-928/target cargo run -p omnigraph-gqt --locked --features omnigraph/failpoints --bin omnigraph-gqt -j 3 -- /tmp/review-903-followup/blob_limit.gqt --target omnigraph-engine --storage local-filesystem --artifacts /tmp/review-908/gqt-artifactsExact-commit workspace CI passed. Its log confirms execution of the new tests and crash matrix. GQT and the pinned DST suite also passed. RustFS and Azure integration jobs passed in the workspace run. These are CI results, separate from the local checks above.
All local checks used Rust 1.97.1, locked dependencies, and the isolated checkout. The checkout is clean. No source or test edits were needed, and no PR changes were pushed. Local cloud storage, DST, and performance checks were not run.
RFC 0033 step 3A: the engine half of the Blob write API. An embedded
Sessioncan now replace or clear one Blob value of an existing node or edge, addressed by exact id like a Blob read.Stacked on #903 (the payload/framing split, step 3-pre). The plan this implements is the RFC 0033 amendment in #900. Until #903 merges, this diff includes its commit.
What it adds
Session::put_blob_at_as(branch, cell, bytes, precondition, actor)stores managed bytes and returnsBlobWriteOutcome::Managed { length, etag, commit }.Session::clear_blob_at_as(branch, cell, precondition, actor)sets a nullable cell to null. It returnsNull { commit: None }when the cell is already null and no precondition is given.BlobPrecondition::{Tags, AnyExisting}and the newOmniError::BlobWritePreconditionFailed { current_etag }. The server maps it to its existing 412 Blob precondition response, so 3B only needs routing.BlobEtag::from_tag, so a caller can pass back a tag it received.Design
Not a new writer kind. A cell write is a Mutation-protocol write. Each attempt:
WriteTxn;PendingMode::Upsertrow;The tail is now split into
commit_staged_mutation, which does validation, staging andcommit_all's detached commit, andpublish_committed_mutation, the protocol's one publisher call.mutateuses the same two halves, so the publication path is unchanged.forbidden_apis.rsregistersexec/blob_write.rsunderMUTATION_V9, and no new durable call site appears.Whole-row replacement. Lance 11 has no single-cell Blob write:
UpdateBuilderandRewriteColumnsrefuse Blob-v2 columns. So the row is rebuilt:LargeBinarychild adopts the caller'sBytesbuffer without a copy.scan_with_pending_materialized_blobsof the exact id, with the target omitted. The old payload of the target is never read or charged.updatecarries them. A stored external sibling needs the external Blob policy; otherwise the write fails withStoredExternalBlobDeniednaming it.Bound. A put over 32 MiB is refused before anything is captured, with resource
Blob write payload bytes. Exactly 32 MiB is accepted.Preconditions and retries.
read_blob_atcomputes it.ReadSetChangedre-prepares the whole attempt, up to the insert-only bound of 32. Each new attempt re-reads the cell, re-evaluates the precondition and re-carries the row from the fresh base. A stale ETag therefore fails instead of overwriting.Receipt. The returned ETag is read from the detached version
commit_allproduced, before the manifest CAS. Nothing after publication reads storage, so a later write cannot change it.Deviations from the RFC text
current_etagisOption<String>, notOption<BlobEtag>: the error lives inomnigraph-core, which cannot name the engine'sBlobEtag. The string is the same quoted tag.Tests
end_to_end.rs::blob_put_and_clear_replace_one_cell_by_exact_idAnyExisting. Clear with a tag. A no-op clear returns no commit.AnyExistingon null fails withcurrent_etag: None. A missing id isNotFound.writes.rsBlob write payload bytes. A 20 MiB put beside a 20 MiB sibling is refused per type. The put never probes or reads its stale external target. The managed sibling stays byte-identical, and an external sibling becomes managed. Under deny, a put or clear names the carried reference on a node and an edge; a merge load recovers both, and a.gqupdate recovers the node.policy_engine_chassis.rschangefor the actor; a write with no actor is refused.detached_commit_matrix.rsBlobPutandBlobClearwriters in every window × fault × recovery actor. The oracle adds the cell value: the new value exactly when acknowledged, read from a fresh handle and from the writer's own.failpoints.rsblob_write.rsNot in this PR
PUT/DELETEon/blob.blob put/blob clear.The docs say the CLI and server do not offer these writes yet.