Skip to content

Report one real command tag per DML statement on every engine, protocol, and path - #356

Merged
farhan-syah merged 29 commits into
mainfrom
fix/statement-command-tag
Sep 20, 2026
Merged

farhan-syah merged 29 commits into
mainfrom
fix/statement-command-tag

Conversation

@farhan-syah

@farhan-syah farhan-syah commented Sep 20, 2026 •

Copy link
Copy Markdown
Member

Closes #297.

What changes

A DML statement answers one command tag with the real row count, on every engine, every protocol, and every path a statement can take: autocommit, inside a transaction, forwarded to another node, or through a Calvin batch.

Command tag

  • StatementTag folds every task of a statement into one DmlOutcome { verb, affected }. Same-verb tasks sum. INSERT plus UPDATE folds to INSERT. A verb mismatch is an error, never a guess.
  • describe_plan classifies every client write as count-bearing. KV, CRDT, vector, graph, timeseries, and array writes report affected from their handler. UPSERT is its own verb. ON CONFLICT DO UPDATE reports INSERT or UPDATE by what it did.
  • TRUNCATE is the one verb with no count, expressed once as DmlOutcome::carries_count.
  • Native gains NativeResponse.command / QueryResult.command and reads counts from the payload instead of counting tasks.

Transaction overlay

Every engine stages a write at statement time, so the statement reports its real count, the transaction reads its own write, ROLLBACK and ROLLBACK TO SAVEPOINT undo it, and COMMIT applies it once:

  • Array Put/Delete, single-node and cluster fan-out.
  • Vector-primary DirectInsert/DirectUpsert/DirectDelete/DirectUpdate.
  • CRDT DocUpsert/DocDelete.
  • KV predicate UPDATE/DELETE.
  • TRUNCATE on every engine, with an overlay marker every read merge consults. RESTART IDENTITY applies at COMMIT.

TRUNCATE

TRUNCATE lowered to DocumentOp::Truncate for every collection and walked the document store only. It left kv, columnar, timeseries, spatial, and vector-primary rows in place while answering TRUNCATE. It routes per engine through EngineRules::plan_truncate now. An array refuses with a typed error naming DROP ARRAY.

Cluster

The multi-shard tag fold is pinned by 3-node tests. Writing them surfaced four cluster-only defects, all fixed here:

  • A write inside a transaction on a coordinator that does not own the vShard forwarded through the gateway before the transaction state was read. It applied at once and survived ROLLBACK. The gateway also shaped every payload as rows, so an INSERT answered SELECT 1.
  • A Calvin-applied row never got its pk -> surrogate binding on participants. WHERE id = ... missed on every node but the coordinator. bind_plan_identities is the one binder the replication decoder and the Calvin scheduler share.
  • MERGE, UPDATE ... FROM, and INSERT ... SELECT applied through dispatch_local on the coordinator. On a cluster the write landed on a node that did not own the target. They propose through Raft now.
  • Native folded only the first count-bearing task of a batch. pgwire and native share calvin_fold.

Fixes found on the way

  • A bare-value KV row (single value column, RESP SET) was decoded as msgpack by every ON CONFLICT DO UPDATE path. A long value failed to decode, a one-byte value was silently discarded.
  • KV writes never invalidated the aggregate cache, so COUNT(*) on a KV collection went stale.
  • Columnar/spatial DELETE and UPDATE left R-tree entries behind, so ST_DWithin returned deleted rows.
  • A columnar flush reassigned the memtable's segment id and dropped its tombstones.
  • ALTER SEQUENCE ... RESTART WITH n handed out n + 1.
  • An empty vector search answered a JSON literal the response collector read as one row.

Tests

New suites: command_complete_tag_conformance, truncate_engine_conformance*, sql_transactions_{array,vector_primary,crdt,truncate,kv_predicate}_overlay*, sql_insert_conflict_kv_bare_value, native_dml_outcome_conformance, calvin_multishard_*.

Companion

nodedb-lite main carries the QueryResult.command field and per-engine TRUNCATE parity.

KV insert/update, upsert, ON CONFLICT resolution (both the direct and
staged-transaction paths), CRDT document writes, and vector-primary
upserts answered with a bare OK response instead of a Postgres command
tag carrying an affected-row count. A bare tag parses to 0 in
tokio_postgres, so a driver reads a successful write as having
affected no rows.

Add response_affected_with_op alongside the existing response_affected
helper so upsert paths can report INSERT 0 n or UPDATE n depending on
whether the write inserted or overwrote a key, and route the KV
conflict-resolution handlers through it instead of hand-encoding the
payload with response_codec.

Extend the CommandComplete tag conformance tests to cover multi-row
inserts, ON CONFLICT DO NOTHING, staged transactions, and every engine
(document, strict document, KV, timeseries, spatial, vector-primary,
CRDT, array) so each write path is pinned to the tag contract.
Writes on the KV, timeseries, vector-primary, CRDT, and cluster-array
paths fell through describe_plan's engine wildcards to
PlanKind::Execution, which pgwire renders as a bare OK tag: a driver
parses that as zero rows affected. A CRDT INSERT/UPSERT/UPDATE was
classified SingleDocument and answered with an empty row set instead
of a command tag.

Replace the wildcards with explicit per-variant arms so a new op forces
a classification. KV Put answers UPSERT n (matching DocumentOp::Upsert
on the autocommit, staged, and Calvin-folded paths); KV Insert/BatchPut,
timeseries Ingest, vector-primary DirectUpsert, and array Put answer
INSERT 0 n; array Delete answers DELETE n.

Add PlanKind::DmlResultByOp for KV INSERT ... ON CONFLICT DO UPDATE,
whose verb is decided per row: the handler reports {affected, op} and
the response layer renders INSERT 0 n or UPDATE n from it. The
RLS-gated resolve path for that op reports the same payload.

Add CrdtWriteVerb to CrdtOp::DocUpsert, set by the planner from the
statement (INSERT, UPSERT, UPDATE) and carried through WAL
replication, so the tag follows the SQL verb.

Emit the cluster-array coordinator's Put/Delete counts as a keyed map
so the shared affected-count reader accepts them.

Split response_shape/types/plan_kind.rs into one file per engine
family plus a describe.rs dispatcher.
Break the static/active transaction dispatch module into
active_dispatch, primary_write, and static_dispatch, with primary_write
holding the write/RETURNING/change-set classification shared by both
dispatch paths.

Decide the primary-write participant at transaction level: a slice
holding only derived side-effect writes (implicit graph edge, balance
delta) deposits an applied response only when the whole transaction has
no non-derived write. That keeps a lone balance-delta participant from
racing the user's own write for the statement's response, while a
Graph-only transaction (the standalone GRAPH INSERT/DELETE EDGE DSL)
still deposits the response its caller reads the count from.
Edge INSERT/DELETE and node LABEL/UNLABEL previously tagged as
Execution/count-less DDL. They now classify as DmlResult (INSERT,
DELETE, UPDATE) and carry a real affected count end to end: the Data
Plane computes it (existence checks for delete/unlabel, deterministic
1 for put/label), staged writes resolve it against BASE ∪ overlay, the
neutral DDL layer threads it through single-home, cross-shard, and
in-transaction paths, and the native wire layer surfaces it instead of
the 'one command ran' sentinel.

Split graph_edge_write.rs into a directory by operation (put,
put_batch, delete, delete_batch, shared) to carry the added logic
within the file-size limit.
A statement with several write tasks (Calvin batches, staged in-
transaction writes, pre-dispatch trigger short-circuits, cluster array
ops, and the per-task dispatch loop) previously emitted one
CommandComplete tag per task, so a multi-row statement answered with
several tags instead of the one a client expects.

Introduce DmlOutcome (a count-bearing verb) and StatementTag, which
folds every task's outcome into a single tag per statement: same-verb
tasks sum their counts, INSERT mixed with UPDATE folds to INSERT (the
ON CONFLICT DO UPDATE shape), any other verb mix is a DmlFoldError
surfaced as a pgwire internal error, and opaque (non-count-bearing)
tasks fold to a bare OK unless a DML outcome also folded. Every call
site that used to build its own Tag or push a Response per task now
folds through StatementTag and emits the result once, after any
RETURNING rows.
Array engine writes (Put, Delete) previously bypassed in-transaction
staging entirely, so a same-txn read against an array collection could
not see an uncommitted write. Add ArrayTxnOverlay and StagedCellPut to
buffer staged cells per array/coordinate, an array_merge module to fold
the overlay against base storage on read, a lease module shared across
overlay types, and stage_array to route ArrayOp::Put/Delete through the
staging path alongside the existing KV/Columnar/Timeseries/Spatial/
Graph writes.

Wire the overlay into the dispatch, aggregate, elementwise, mutate,
and read paths for array queries, into savepoint rollback/state
tracking, into overlay gauges and reaping, and into the cluster-array
physical plan and catalog scan so a staged write is visible to reads,
rolled back on savepoint abort, and reflected in overlay memory
accounting. Extend the buffered-atomicity and savepoint-overlay wire
tests and add a dedicated array-overlay transaction test.
Extend cluster array Slice/Agg wire requests with the reading
transaction's id so each shard folds its own staged overlay cells into
the result, matching the single-node ArrayOp read-your-own-writes
contract.

Fan a Put/Delete ClusterArrayOp out into one ArrayOp per owning vShard
at statement time and stage each into that shard's overlay on its own
leader (session::array_fanout_stage), instead of only buffering for
COMMIT. Every per-shard task is buffered before any shard stages, so
ROLLBACK still tears down every recorded vShard on error.

Fold the distributed aggregate partial path into the same overlay
gating as the single-node finalizing read, since each shard now applies
its own overlay before merging at the coordinator.
Distributed array Put/Delete answered with the number of cells/coords
named in the request instead of the number the Data Plane actually
wrote or removed, so a DELETE of an absent coordinate over-reported
its count. Shard responses now carry a real `affected` count read from
the local `{"inserted"|"deleted": n}` handler payload, and the
coordinator sums per-shard counts for the client-facing tally.

Extracts `ClusterArray` dispatch into a protocol-neutral
`shared::cluster_array_dispatch` module and wires it into the native
protocol, which previously had no `ClusterArray` handling at all;
pgwire's routing now delegates to the same core. Adds an
`Error::Shaping` variant so response-shaping failures classify with
their own code across both protocols instead of collapsing to a
generic internal error, and switches native/HTTP session setup to
`QueryContext::for_state_with_lease` so top-level connections plan
with the same lease/catalog/WAL context as pgwire.

Also adds `NODEDB_TEST_RUST_LOG` and `NODEDB_TEST_DUMP_SERVER_LOG` to
the wire test harness for debugging a failing spawned server.
A vector-primary collection keyed its surrogate on the first four
floats of the vector instead of the declared PRIMARY KEY, so rows with
the same key never collided and rows with the same vector did. An
INSERT onto an existing key silently overwrote the row and leaked the
old HNSW node and its payload bitmap entries. UPDATE and DELETE fell
through to the document engine and touched only the sidecar, leaving
the index stale.

Key the surrogate on the declared PRIMARY KEY value, exactly as the
document insert does. Lower INSERT / INSERT ... ON CONFLICT DO NOTHING /
UPSERT to DirectInsert / DirectInsertIfAbsent / DirectUpsert: a
duplicate key raises unique_violation, DO NOTHING reports 0, UPSERT
removes the old node before writing (ON CONFLICT DO UPDATE SET patches
the payload). Add DirectDelete and DirectUpdate so DELETE and UPDATE on
a vector-primary collection remove or re-embed the HNSW node, update
payload bitmaps and the sidecar, and report real counts; predicate
targets resolve on the Data Plane. UPDATE ... FROM and MERGE on a
vector-primary target are refused instead of silently degrading.
insert_with_surrogate tombstones a node already bound to the surrogate.

Add a resolve-before-propose path (ResolveDirectWrite /
ResolvedDirectWrite) so a write governed by a row-level write policy is
decided once on the leader and replicated as concrete mutations.

Wire the new ops through the SQL planner, Calvin classification, WAL
encode/decode/replay (three new record types), replication, response
shaping, and RLS injection. Split nodedb-physical's vector physical-plan
module and the gateway version_set module into per-concern files.
…rlay

Vector-primary DirectInsert/DirectInsertIfAbsent/DirectUpsert/
DirectDelete/DirectUpdate previously bypassed in-transaction staging
entirely, so a same-transaction read or search against a vector-primary
collection could not see an uncommitted write and every such write ran
autocommit-only.

Add StagedVectorRow to buffer a staged vector plus its exact sidecar
bytes, a vector_primary_merge module to fold the overlay against base
storage on read, and stage_vector/stage_vector_targets to route
DirectInsert-family, DirectDelete, and DirectUpdate through the staging
path alongside the existing KV/Columnar/Timeseries/Spatial/Graph/Array
writes. Each staged op decides its outcome against base plus overlay
using the live handler's own rules: a duplicate key refuses, ON
CONFLICT DO NOTHING skips, and a conflict patch merges through the
live handler's merge; COMMIT replays the buffered plan through the
live handler as the durable apply.

Wire the overlay into point lookup, document scan, and search-exec
result merging so a staged vector-primary row is visible to reads and
ranks in the same transaction. Replace staged_tag_kind with
stageable_write_shape so the staging gate can validate and classify a
write from the same computed shape instead of two separate lookups.
Add vector_for_id/vector_for_surrogate to VectorCollection to read the
live FP32 vector for a row, needed to seed the staged sidecar merge.
Extend staging_predicates for the new stageable ops and add a
dedicated vector-primary overlay transaction test.
CRDT DocUpsert/DocDelete previously bypassed in-transaction staging
entirely, so a same-transaction read against a CRDT collection could
not see an uncommitted write. Add stage_crdt to route these ops
through the staging path alongside the existing Document/KV/Columnar/
Timeseries/Spatial/Graph/Vector/Array writes: a staged write lives
only in the per-transaction overlay until COMMIT replays the buffered
plan through the live handler, which stays the sole durable apply.

Classify DocUpsert by its verb (Insert/Upsert/Update) and DocDelete as
Delete in staging_predicates so the staging gate can validate and tag
them. Extract stage_current_body as the one place every staged
read-modify-write (document update/upsert/delete, CRDT upsert/delete)
resolves a row's current body under BASE ∪ OVERLAY, replacing the
duplicated logic in stage_point_document and stage_upsert. Extract
crdt_row_body so the staged path and the live materialization build a
CRDT row's stored MessagePack body through the same encoding.

Add a dedicated CRDT overlay transaction test and lifecycle test.
Extract a shared RawPgConn wire harness for tests that need raw
pgwire message framing, and port pgwire_ddl_result_types onto it.
TRUNCATE inside BEGIN..COMMIT now marks the collection truncated on
the transaction's overlay instead of touching base storage. Every row
staged earlier in the transaction is tombstoned through the undo
journal, every base row with no newer overlay entry is hidden from
the transaction's own reads, and the buffered plan replays through
the live truncate at COMMIT in statement order. ROLLBACK and ROLLBACK
TO SAVEPOINT drop the marker.

RESTART IDENTITY restarts the collection's sequences once COMMIT has
made the truncate durable; ROLLBACK leaves them untouched. Autocommit
restarts them right after dispatch. SequenceHandle gains restart_at,
distinct from setval: the restart value counts as not yet called, so
the next nextval returns it directly.

Adds a KV write-version helper split out of write_version.rs to stay
under the file size limit, and makes every write-version match over
ColumnarOp/SpatialOp/TextOp/GraphOp/DocumentOp exhaustive instead of
falling back to a wildcard arm.
The Data Plane aggregate/facet result cache had no invalidation hook
for KV writes, so a cached COUNT/GROUP BY over a WITH (engine='kv')
collection could serve a stale result after a subsequent put, delete,
atomic op, rename, purge, or truncate.

Add a per-table write epoch to KvEngine, bumped at the single
table_mut_for_write / table_for_write_or_create chokepoint plus the
remove-shaped mutations that bypass it. Stamp each aggregate cache
entry with the epoch at compute time and treat a mismatch on lookup
as a miss, evicting the stale entry instead of the previous
insert-only cache.
An empty or exhausted vector index short-circuited the whole GraphRAG
query with an empty [] response, discarding results the other fusion
legs (graph, text) still had to offer. Return an empty score set for
that leg instead so fusion proceeds and the envelope keeps its shape.
TRUNCATE lowered to DocumentOp::Truncate for every collection, so it
walked the document store only and left kv, columnar, timeseries,
spatial, and vector-primary rows in place while answering TRUNCATE.

Plan TRUNCATE through a new EngineRules::plan_truncate. SqlPlan::Truncate
carries its engine, plan_truncate_stmt resolves the target through the
catalog, a vector-primary collection lowers to SqlPlan::VectorPrimaryTruncate,
and an array refuses with a typed error naming DROP ARRAY and
DELETE FROM ARRAY. The convert layer emits KvOp::Truncate for kv, the
document op for the document family, and a typed "not yet routed" error
for the columnar family until those engines gain their own truncate op.

Add VectorOp::DirectTruncate: every live row of the primary index is
removed through the DirectDelete row path, then the collection resets
its HNSW segments and payload bitmaps. Wire its WAL record, replication
encode/decode, replay, transaction staging, resolve, and Calvin routing.

KvOp::Truncate carries restart_identity. PhysicalPlan::truncate_target
lets RESTART IDENTITY apply engine-neutrally on autocommit (pgwire and
native) and at COMMIT.

An empty vector search answers an empty msgpack hit array instead of a
JSON literal, so the response collector decodes it as zero rows.

Add per-engine TRUNCATE conformance tests and planner routing tests.
append.rs accreted every WAL append method across vector, batch,
index, metadata, and transaction records in one file. Split each
family into its own module (append_batch, append_index,
append_metadata, append_transaction, append_vector), keeping
append.rs to the base Put/Delete appends and the shared
append_record helper.
on_memtable_flushed reset the memtable's virtual segment id to
next_segment_id instead of leaving it fixed, so a second flush could
remap PK-index entries that already belonged to the first flushed
segment, and left the memtable's delete-bitmap tombstones behind under
the old id instead of moving them with the rows they mark.

Keep the virtual segment id constant: it names the memtable, not a
flushed segment, so a later remap can never sweep flushed rows along.
Move the memtable's delete bitmap to the new segment id at flush time,
so a row tombstoned while still in the memtable stays tombstoned once
flushed. Expose flush_threshold() as a read accessor.

Adds regression coverage for both a second flush leaving the first
segment's rows untouched and a memtable tombstone surviving flush.
Columnar/spatial UPDATE and DELETE removed the mutation-engine row but
left its R-tree entries in place, so a spatial predicate kept matching
rows the collection no longer held; the checkpoint-restore rebuild
carried the same stale-entry risk across a restart. Neither path fed
its R-tree changes to the transaction undo log, so a rollback left a
committed index change behind.

Add remove_columnar_row_spatial_entries and the shared
apply_columnar_update_rows / apply_columnar_delete_pks row-apply loop
(columnar_mutation_apply.rs) that both the predicate and
resolved-row-set UPDATE/DELETE handlers now call, cascading the R-tree
removal for a spatial collection's geometry columns per affected row.
Reading the row's current field values (read_columnar_row_by_pk) is
what lets the cascade derive the entry to remove.

index_columnar_geometry_columns now returns a GeometryIndexDelta of
what it removed and inserted; push_geometry_index_undo records it on
the transaction undo log so a rollback reverses an insert's R-tree
maintenance too. Wire spatial_undo through ColumnarInsertParams and
its callers. The checkpoint-restore rebuild clears a collection's
spatial index and doc-map entries before rebuilding from the restored
rows, so a stale checkpoint entry can never survive.
TRUNCATE against a columnar, spatial, or timeseries collection refused
with FeatureNotSupported: those engines had no truncate op, so the SQL
layer could plan neither.

Add ColumnarOp::Truncate and TimeseriesOp::Truncate, carrying the
target collection and RESTART IDENTITY. convert_truncate lowers a
columnar-family collection (spatial shares columnar's DML ops) to the
former and a timeseries collection to the latter. Wire both through
Calvin routing, write-class, RLS injection (refused under a write
policy: a truncate reads no row to evaluate a predicate against),
write-resolve, permission checks, CDC change-event extraction,
transaction buffering classification, response shaping, and the
dedicated ColumnarTruncate / TimeseriesTruncate WAL record types with
their encode/decode/replay/replication paths.

MutationEngine::truncate on the columnar engine returns every row it
held (TruncatedRows) so a transactional caller can restore the exact
pre-image on rollback. The timeseries truncate renames the
collection's partition directory aside rather than removing it
outright, so a rollback renames it back and a commit removes it in the
background (finalize_timeseries_truncates, with a maintenance-tick
retry backlog for a removal that raced a still-open reader). Boot
replay and checkpoint load skip an aside directory as a truncate's
leftover rather than a collection.

Stage TRUNCATE through the transaction overlay: ColumnarOp::Truncate
and TimeseriesOp::Truncate are staged as an overlay marker at
statement time (stage_columnar_family.rs, extracted from the
statement-time dispatch table alongside the resolved-row-set paths)
and replayed as the durable truncate at COMMIT, with a full pre-image
captured on the undo log (ColumnarTruncateUndo, TimeseriesTruncateUndo)
for atomic rollback. A committed truncate resets its collection's
continuous aggregates to their freshly-registered state and raises a
truncate floor so a WAL catch-up redelivery from before the truncate
can never resurrect a removed row.

Extract the shared {truncated: n} response encoder
(truncate_response.rs) used by every engine's truncate handler.

Add per-engine TRUNCATE conformance tests for the columnar family.
The native SQL loop counted one row per dispatched task, per buffered
in-transaction write, and per opaque response instead of reading the
Data Plane payload, and direct-op writes answered no count at all.
It also carried no command verb.

Move the payload-to-outcome readers out of pgwire into response_shape
so both protocols derive a DmlOutcome the same way. The native loop
folds every task through StatementTag exactly like pgwire: DML counts
sum, buffered and opaque tasks fold to nothing, a verb mismatch is an
error. Direct ops, staged writes, cluster-array writes, and Calvin
batches all report the payload's count.

NativeResponse and QueryResult gain an optional command field carrying
the verb, threaded through dispatch, session streaming, the client,
and the test harness. TRUNCATE stays the one verb with no row count,
expressed once as DmlOutcome::carries_count for both protocols.

Add a native DML outcome conformance test.
KV predicate UPDATE/DELETE previously required autocommit: transaction
resolve and the TransactionBatch dispatcher both rejected
`KvOp::PredicateUpdate`/`PredicateDelete` because the row set resolved
from committed state at apply time, which per-row KV redo couldn't
express.

Stage them like Document `BulkUpdate`/`BulkDelete` instead: resolve the
matched rows against BASE ∪ OVERLAY at statement time, stage each
row's post-image or tombstone, and classify it as a keyed write during
resolve/classification. COMMIT replays the live handler in statement
order rather than through redo.

Factor the staged-row surrogate lookup shared by keyed delete and the
new predicate paths into `resolve_kv_stage_surrogate`. Fix the overlay
scan-merge predicate to take both the raw key and the value: a
`WHERE key = ...` predicate needs the key, which `kv_row_to_doc` folds
in as a row field, and the old signature only ever saw the value.
Move resp_client.rs from crash_harness into tests/support so both the
crash harness and the wire harness can drive AUTH/SELECT/data commands
over RESP. Expose TestServer::resp_port and a resp_session() helper
that creates a readwrite user with a runtime-generated password.

Also drop the hardcoded RESP password in crash_resp_kv_write.rs for a
fresh, process-scoped one.
A KV row written through the single-`value` SQL column, or through
RESP SET, stores its scalar as raw bytes rather than a msgpack map.
Every merge path (autocommit ON CONFLICT DO UPDATE, field set, atomic
transfer, predicate UPDATE/DELETE, transaction staging, and WAL
replay) used to decode that body as msgpack unconditionally: a
multi-byte value failed to decode, and a one-byte value (a valid
msgpack fixint) was read as an integer, discarded, and silently
re-encoded as a map.

Add scalar_to_raw_bytes/NotScalar in nodedb-types and
kv_body_shape/kv_body_to_row/row_to_kv_body in nodedb-query to decode
a stored body into the {"value": ...} row every read already presents,
and re-encode a merged row back into the shape it came from. Route
every KV read-modify-write handler and WAL replayer through these
functions instead of ad hoc msgpack decode/encode.

HGET/HSET against a bare-value key now return WRONGTYPE (RESP) /
TYPE_MISMATCH (SQL) instead of silently treating the row as a hash.
Move RawPgConn and command_tags from nodedb/tests/wire/harness into
nodedb-test-support::pgwire_harness so both the wire test crate and the
cluster harness can read wire-level command tags and message bytes
tokio_postgres does not expose.

Add TestClusterNode::raw_pgwire, a connection pre-authenticated as the
harness superuser on the harness database, so a cluster test can
assert the exact tag verb a statement answers.
On a coordinator that does not own a statement's vShard, the pgwire
gateway fast path ran before the session's transaction state was read,
so an in-block write forwarded there applied durably at once: ROLLBACK
could not undo it and the transaction never saw its own row. The same
path shaped every forwarded payload as PlanKind::MultiRow and pushed
one response per task, so a plain INSERT answered `SELECT 1` and a
three-row insert answered six of them.

Skip the gateway fast path inside a transaction block so the dispatch
loop's staging gate handles the write (it already forwards to the
owner). Fold forwarded payloads through the statement tag by the
task's real plan kind (gateway_fold).

pgwire and native also each re-derived a Calvin batch's rows and tag by
hand, and native took only the first count-bearing task. Replace
calvin_response.rs with response_shape::calvin_fold, one
protocol-neutral fold both handlers call.

A derived side effect (an implicit graph edge beside the user's own
write) never answers the statement's own tag: plans_have_user_write
and plan_counts_toward_statement_tag in write_class.rs are the one
rule every fold and the Calvin scheduler apply.
…rites

These Control-Plane-orchestrated writes resolved their arms locally
and applied through dispatch_local, which never proposed through Raft:
a resolved MERGE or UPDATE ... FROM landed only on the node that
planned it, and a restart or a follower read never saw the write.

Add orchestrated_write::apply_orchestrated_write as the one apply seam
for these three statement kinds: on a cluster the resolved plan
proposes through Raft to the target vShard's owner and every replica,
standalone it dispatches locally and mints the redo record the write
funnel would have. Add ReplicatedWrite::MergeApply and
UpdateFromJoinApply, their WAL encode/decode (document_join.rs), and
route insert_select's per-page apply through the same seam.

DocumentOp::Merge gains resolved_insert_identities, the NOT-MATCHED
surrogates for its resolved inserts, index-aligned with
resolved_inserts, so every applying node installs the same
(collection, pk) -> surrogate binding.

Installing that binding was previously duplicated: the replicated-write
decoder bound each surrogate inline per op, and the Calvin scheduler
had no equivalent, so a Calvin-dispatched write never revalidated a
carried surrogate against another coordinator's earlier binding. Add
surrogate::bind_plan_identities, one walk over a PhysicalPlan shared by
both the decoder (wal_replication::decode::entry) and the scheduler
(dispatch::bind_identities), and simplify every decode helper to
rebuild the plan's surrogate verbatim from the record rather than
resolving it inline.
…k reads

Add calvin_multishard_dml_tag_fold, calvin_multishard_pk_read, and
calvin_multishard_txn_staging_cross_node, sharing fixtures pulled into
calvin_multishard_fixture. They cover the statement-level tag fold
across vShards, MERGE/UPDATE FROM replication to the target's owner
and every replica, and cross-node primary-key resolution after a
Calvin write installs its surrogate binding.

Extract distinct_vshard_collections and key_on_other_vshard into
vshard_names, replacing the copy inlined in four existing cases.

Add TestClusterNode::sequence_next_value: a restart stores a sequence
one increment below its RESTART WITH value with called = false, so
sequence_current_value could never observe the value a following
nextval actually returns. Fix ddl_objects.rs's assertion to check the
next value instead of the raw counter.
Compute PlanMeteringInfo before dispatch instead of cloning the
physical plan for later extraction, and pass the info through to
meter_gateway_task. Use the imported TransactionState alias in
execute.rs instead of the fully qualified path. Tidy stale comment
wording in the calvin dispatch and native conversion modules.
check_authorized_dispatch.py still referenced the old routing handler
path for into_physical_task; point it at the new shared module.
@farhan-syah farhan-syah added the run-ci Opt this PR into the full test suite; re-add to force a re-run label Sep 20, 2026
@farhan-syah
farhan-syah merged commit 4664c82 into main Sep 20, 2026
14 checks passed
@farhan-syah
farhan-syah deleted the fix/statement-command-tag branch September 23, 2026 01:41
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

run-ci Opt this PR into the full test suite; re-add to force a re-run

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Multi-row INSERT emits one CommandComplete per row — drivers report rowcount 1; kv and timeseries DML answer a bare OK tag

1 participant