Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion docs/benchmarks.md
Original file line number Diff line number Diff line change
Expand Up @@ -152,7 +152,8 @@ reader. The script's `measure_snapshot_shapes()` read `max(index_version)` from
`{prefix}_capture_raw` rather than from `{prefix}_index_watermark`, so the
"watermark" row timed a descriptor-table aggregate; it used two separate
per-column `argMax` expressions where the reader resolves one `argMax` over a
tuple of every column ordered on `(index_version, store_id, pack_id)`; it
tuple of every column ordered on `(index_version, store_id, pack_id)` (an
order since changed to `(member_version, store_id, pack_id, index_version)`); it
carried no manifest membership, which is half of what a pinned read pays for;
and it pinned at the raw maximum, so the historical case -- a pin below a later
publish that re-indexed a capture -- never arose. Those numbers established a
Expand Down
37 changes: 29 additions & 8 deletions docs/capture-storage-design.md
Original file line number Diff line number Diff line change
Expand Up @@ -382,8 +382,9 @@ built, nothing compares them -- and integrity proves nothing here, because a
forged pack is perfectly well formed. Anyone able to PUT into the bucket could
therefore write a pack whose footer carried another tenant's `tenant_id` and
`capture_id`, have it indexed under the victim's tenant, and -- since the
reader resolves a capture with `argMax` over `(index_version, store_id,
pack_id)` -- become the pack that capture resolves to at every fresh watermark.
reader resolves a capture with `argMax` over `(member_version, store_id,
pack_id, index_version)`, newest pack first -- become the pack that capture
resolves to at every fresh watermark.

`_descriptors` now refuses a pack whose records name a tenant other than the
one its key belongs to, comparing against the same `key_component` encoding
Expand Down Expand Up @@ -446,7 +447,7 @@ capture described by two packs -- a pack mirrored to a second store, or a
producer retrying a `capture_id` after the first pack was sealed -- is two
published rows and appears twice. Choosing between them is supersession, which
belongs to the reader (one `argMax` grouped on capture identity, ordered on
`(index_version, store_id, pack_id)` -- see *Phase 5*); a second copy of those
`(member_version, store_id, pack_id, index_version)` -- see *Phase 5*); a second copy of those
semantics in the view's SQL could drift away from the reader's without either
side failing.

Expand Down Expand Up @@ -1551,7 +1552,7 @@ Reads are pinned to a watermark. `CaptureQuery.filter_hash` identifies a query
independently of its page, keyset cursors carry that hash and the pinned
watermark, and `ClickHouseCaptureCatalog` resolves a capture out of
`*_capture_raw` with **one** `argMax` over a tuple of every non-grouped column,
ordered on the tuple `(index_version, store_id, pack_id)`.
ordered on the tuple `(member_version, store_id, pack_id, index_version)`.

Both halves of that shape are load-bearing, and the per-column
`argMax(<column>, index_version)` this document used to describe has been
Expand All @@ -1564,9 +1565,28 @@ removed:
watermark resolved to a different pack at `max_threads = 1` than above it, and
to a different one again once a merge had put both rows in one part -- a
pinned selection silently reading different bytes before and after a
background merge. `(index_version, store_id, pack_id)` is a total order over
the rows in a group, and `index_version` still leads, so supersession is
unchanged.
background merge. `(member_version, store_id, pack_id, index_version)` is a
total order over the rows in a group.
- **The ranking version.** `member_version` is the version at which a pack's
FIRST publish reached the watermark, at or below the pin: `min(index_version)`
over the manifest rows paired with the watermark log, joined in on
`(store_id, pack_id)`. It is NOT a
descriptor row's own `index_version`, which is only the version the row was
written at. A pass that re-indexes an already-published pack -- after a crash
between publishing and `commit_packs`, an outcome-unknown publish that
landed, or a rebuild beside the live indexer -- rewrites the pack's rows at a
higher version before publishing anything. Ranked on the rows' version, a
superseded pack then outranked the newer pack inside snapshots already
pinned, and kept doing so if that pass never published. Found by model
checking the publish protocol. It is the first publish rather than the
newest because a replay that DOES publish makes the pack a member again at a
fresh version: ranked on that, the superseded pack would win every head from
the replay on, on the ordinary crash-recovery path. A replay adds nothing to
the catalog, so it does not move a pack's rank; a genuine re-capture is a new
pack and a mirror is another store, so both still get a fresh first publish.
Pins are stable either way, since every later publish lands above the pin.
`index_version` is kept as the last component, so within one pack the row a
merge keeps is also the row a read resolves.
- **One aggregate, not one per column.** Twenty-seven separate `argMax` calls
leave nothing forbidding `store_id` from one row and `object_key` from
another -- a descriptor describing no pack that exists. It could not be
Expand All @@ -1585,7 +1605,8 @@ past the cursor before `LIMIT` kept `limit + 1` of them. The keyset comparison
`(tenant_id, experiment_id, run_id, captured_at_ns, capture_id) > (...)` is not
usable by the primary-key index, so each page cost about the whole catalog
whatever its size. Measured on 25.12 over a 198k-row catalog, a 38-row page
took ~150 ms, three quarters of it that tuple. Both queries carry every filter,
took ~150 ms, three quarters of it that tuple. Both queries read the same
snapshot join (the one that supplies `member_version`) and carry every filter,
so the groups and their resolution are unchanged: that is the same immutability
rule, below, that makes the pre-aggregation filters safe. The Python reference
reader keeps the single-phase shape, so the parity suite compares the two
Expand Down
27 changes: 21 additions & 6 deletions docs/catalog-descriptor-key.md
Original file line number Diff line number Diff line change
Expand Up @@ -87,11 +87,22 @@ as everything else".
> merge had put both rows in one part.
>
> The shipped projection is **one** `argMax` over a tuple of every resolved
> column, ordered on the tuple `(index_version, store_id, pack_id)`. The tuple
> key restores a total order (supersession is unchanged -- `index_version`
> still leads); the single aggregate makes a mixed descriptor -- `store_id`
> from one row and `object_key` from another -- structurally impossible rather
> than merely unobserved. Twenty-seven aggregates ordered on the tuple are also
> column, ordered on the tuple `(member_version, store_id, pack_id,
> index_version)`. The tuple key restores a total order; the single aggregate
> makes a mixed descriptor -- `store_id` from one row and `object_key` from
> another -- structurally impossible rather than merely unobserved.
>
> Supersession no longer leads with a row's `index_version`. `e93a2c8`'s key
> was `(index_version, store_id, pack_id)`, and a pass that replays an
> already-published pack (a crash before `commit_packs`, an outcome-unknown
> publish that landed, a rebuild) rewrites its rows at a fresh, higher version,
> which let a superseded pack outrank the newer one inside snapshots already
> pinned. `member_version` is instead the version at which the pack's FIRST
> publish reached the watermark, at or below the pin -- `min(index_version)`
> over the manifest rows paired with the watermark log. First rather than
> newest, so that a replay which does publish cannot re-promote the superseded
> pack at every later head. `index_version` stays last, so within one pack a
> read resolves the row a merge keeps. Twenty-seven aggregates ordered on the tuple are also
> correct and cost +291% at a 100-row page, because ClickHouse compares a tuple
> ordering argument through a generic `Field` once per row per aggregate. See
> `clickhouse_reader._projection`.
Expand Down Expand Up @@ -253,7 +264,11 @@ goes without a contract change:
locator, which is exactly the field that may differ. The rewrite is
byte-identical rows at the winning version; the superseded rows share their
full sort key with them (pack identity included), so the engine collapses
each pair and `argMax` resolves the new version in the meantime.
each pair and `argMax` resolves the new version in the meantime. (That was
the resolution order when this shipped. The reader now ranks a pack by its
first paired publish, `member_version`, read from the manifest -- see the
note under the projection above -- so supersession no longer depends on the
rewrite either.)
- **The publish verifies that it owns the version, not that the version is
occupied.** Each attempt mints a `publish_id`, writes it on its manifest rows
and on its watermark row, and reads that column back. The check it replaced --
Expand Down
2 changes: 2 additions & 0 deletions docs/catalog-differential-review-2026-09-01.md
Original file line number Diff line number Diff line change
Expand Up @@ -144,6 +144,8 @@ except SnapshotPublishConflictError:
```
If the inventory INSERT raises (transport error, Code 159 timeout), the bare `raise` is never reached; the driver exception propagates with the conflict demoted to `__context__`. No `except CaptureStorageError` supervisor sees the must-not-retry anomaly, and the packs are *visible but not in the inventory*, so the next pass re-indexes and re-publishes the batch at a higher version — the retry `SnapshotPublishConflictError`'s contract (`catalog.py:35-50`) exists to forbid, while the foreign-writer anomaly is buried. Reader-visible corruption: none (rows are byte-identical and collapse).

*Correction (2026-09-24):* "none" holds only while one pack describes the capture. The re-publish writes the batch's descriptor rows at a new, higher version before publishing, and the reader ranked a capture's packs by that version, so a batch superseded in the meantime by a second pack describing the same capture outranked the newer pack inside snapshots already pinned. If the re-publish never landed, that stayed true for good. Model checking the publish protocol found it. Every replay route reaches it: this one, a crash before `commit_packs`, an outcome-unknown publish that landed, and a rebuild beside the live indexer. The reader now ranks a pack by the version its first publish reached the watermark at (`clickhouse_reader._snapshot`), so neither the rewritten rows nor a re-publish that lands moves its rank, and `indexer.cpp` now has this guard too.

**Recommendation:**
```python
except SnapshotPublishConflictError as conflict:
Expand Down
27 changes: 24 additions & 3 deletions native/csrc/catalog/indexer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -287,7 +287,8 @@ IndexResultData NativeIndexer::index(const std::vector<PackRefData>& refs) {
// the inventory, so a pack recorded there but never made visible is
// skipped forever AND invisible. Only a lost VERSION race is retried
// here, repaired by allocating higher and rewriting the descriptors at
// the winning version (supersession ranks by index_version).
// the winning version (supersession ranks by the pack's membership
// version, not by the rows' index_version; see catalog.py _publish).
uint64_t attempts = 0;
for (; attempts < static_cast<uint64_t>(config_.max_publish_attempts);
++attempts) {
Expand All @@ -298,10 +299,30 @@ IndexResultData NativeIndexer::index(const std::vector<PackRefData>& refs) {
all_rows.size(), indexed.size());
} catch (const CatalogError& e) {
if (e.kind() != CatalogError::Kind::kPublishRace) {
if (e.kind() == CatalogError::Kind::kPublishConflict) {
if (e.kind() == CatalogError::Kind::kPublishConflict &&
!indexed.empty()) {
// Visible, so skippable: record the packs before propagating.
if (!indexed.empty()) {
try {
writer_->commit_packs(RenderPackRows(indexed), version);
} catch (const std::exception& commit_failure) {
// The conflict is the finding (catalog.py's `raise conflict
// from commit_failure`). Left to propagate, the commit's
// transport error replaced kPublishConflict, so a supervisor
// matching on it never saw the second writer, and the visible
// packs stayed out of the inventory for the next pass to
// re-publish.
throw CatalogError(
CatalogError::Kind::kPublishConflict,
std::string(e.what()) +
" (recording its packs in the inventory then failed "
"too: " +
commit_failure.what() + ")");
} catch (...) {
throw CatalogError(
CatalogError::Kind::kPublishConflict,
std::string(e.what()) +
" (recording its packs in the inventory then failed "
"too)");
}
}
throw;
Expand Down
Loading
Loading