Skip to content
Merged
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
60 changes: 59 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 Expand Up @@ -186,6 +187,63 @@ Core summaries run at roughly 80–140 M elements/s depending on dtype
the insert sweep, these are loopback numbers: repeat them on representative
hardware and duplicate ratios before treating any as a capacity claim.

**Native search pages over the snapshot join (2026-09-28).** A pinned read
ranks a capture's packs by `member_version`, the version at which each pack was
first published (see *The ranking version* in `capture-storage-design.md`). To
get that version, the native reader now joins `*_capture_raw` to a members
subquery, which runs `min(index_version)` over the paired manifest rows, grouped
on `(store_id, pack_id)`. Before, it filtered with `(store_id, pack_id) IN
(<manifest>)`. Both queries of the two-phase page read the join, so every page
builds the members set twice. What that costs depends on how many packs the
manifest holds.

The corpus was synthetic and built by SQL: 2,000,000 captures in `P` packs
published over 200 versions. On top of that, 5% of the captures are
re-described by a newer pack published five versions later. 1% of the packs
were rewritten at a version that was never published, and another 1% were
rewritten and published again. After `OPTIMIZE FINAL` the table holds
2,100,000 `capture_raw` rows. Each build's `search` statement was captured from
`system.query_log` and re-run nine times, interleaved, with its own settings
plus `use_query_condition_cache=0`. The figures below are the median
`query_duration_ms` for a page of `limit=100`. A first page has no cursor, and
a mid page sets its cursor at capture 1,000,000. "hook" filters on the tenant
and one of four hooks, and "layer" filters on the tenant and one of 32 layers.
The host was the reference host above, running ClickHouse 25.12.2.

| Packs | first | first, hook | first, layer | mid | mid, hook | mid, layer | members set, one build |
|---:|---:|---:|---:|---:|---:|---:|---:|
| 20k | 113 → 109 | 56 → 70 | 46 → 59 | 103 → 105 | 58 → 69 | 52 → 66 | 2 → 11 |
| 100k | 168 → 166 | 110 → 100 | 100 → 88 | 155 → 167 | 113 → 104 | 103 → 96 | 4 → 26 |
| 1M | 952 → 346 | 870 → 266 | 881 → 241 | 938 → 332 | 877 → 250 | 874 → 244 | 6 → 53 |

Each cell is main (`8b7991d`) → the join, in ms. The last column times the
membership subquery on its own: main's `IN` set, then the join's grouped
members.

- **Rows read are identical** in every cell, for example 2,188,994 at 20k packs
and 6,128,594 at 1M. `EXPLAIN indexes = 1` shows the same granules on
`capture_raw` for both builds (1/258 on a first page, 2/258 on a selective
mid page). The only difference is that main's primary-key condition carries
the `(store_id, pack_id)` set, 40,000 elements at 20k packs and 2,000,000 at
1M, where the join reads the manifest on its own.
- **At 20k packs, selective pages are 19-28% slower** (+11 to +14 ms), which is
about two builds of the members set. Unfiltered pages are unchanged.
- **At 100k packs the result is mixed, from 0.88x to 1.08x.** The members build
costs more, but so does main's larger key set. An independent review on a
differently shaped corpus measured 1.06-1.34x at 100k packs, so treat
20k-100k packs as a latency cost of up to about a third on selective pages.
Nothing is refused: `max_rows_to_read` sees the same row counts.
- **At 1M packs the join is 2.7-3.7x faster,** because main spends most of each
page on its 2M-element key set.
- **Page contents differ from main only where main ranks wrongly,** which is the
case this change fixes. That is a capture re-described by a newer pack while
its older pack was replayed, which main resolved back to the older pack: 1-10
rows of a 101-row page at 20k and 100k packs.

Building the members set once per statement and sharing it between the two
queries (a CTE or a named set) might remove the second build. It has not been
measured.

Baselines:

- **HuggingFace Ideal** — vanilla HF `generate`, no observation (used as 1.0)
Expand Down
43 changes: 34 additions & 9 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,12 +1605,17 @@ 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
shapes directly. Each page still scans the rows past the cursor. Removing that
needs index-usable cursor bounds.
needs index-usable cursor bounds. Each query builds the join's members set, so
a page builds it twice. On a 2.1M-row corpus that costs 19-28% on selective
pages at 20k packs, is mixed at 100k packs, and is 2.7-3.7x faster than the
`IN` set it replaced at 1M packs; see *Native search pages over the snapshot
join* in `benchmarks.md`.

The two queries must filter alike. If the inner query drops a filter that the
outer one keeps, such as snapshot membership or a hook filter, its `LIMIT` fills
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
81 changes: 81 additions & 0 deletions native/csrc/catalog/conformance_catalog.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -11,9 +11,11 @@

#include <chrono>
#include <cstdlib>
#include <functional>
#include <iostream>
#include <limits>
#include <memory>
#include <optional>
#include <set>
#include <string>
#include <thread>
Expand Down Expand Up @@ -1044,6 +1046,78 @@ std::string respond(const std::string& line, Session* session) {
{{"version", version + 1}});
};
}
// A second writer publishing the SAME version: its watermark row lands
// after this pass's own and before the pass reads the version's owners
// back, so publish_snapshot finds two publishes there and reports
// kPublishConflict -- visible, and never retried. That window is one
// round trip wide, so the seam sits in the request hook the storage
// service's lease scope uses (RequestDeadline's before_request), which
// runs before every request this thread sends: once the pass's row
// stands at its version, the next request, the owners read, finds the
// foreign row beside it.
//
// after_conflict then decides what the conflict path's inventory
// INSERT, the request after that read, meets:
// "transport" its connection fails;
// "lease_refused" a rival has claimed the lease's term, and the
// renewal the storage service runs before a request
// once one is due (keep_lease_in_pass) is refused.
struct ConflictSeam {
uint64_t version = 0;
enum { kWaiting, kConflicted, kDone } phase = kWaiting;
};
const auto seam = std::make_shared<ConflictSeam>();
std::function<void()> before_request;
if (jc::FindBool(line, "conflict_at_publish")) {
const std::function<void(uint64_t)> earlier =
index_config.after_allocate;
index_config.after_allocate = [seam, earlier](uint64_t version) {
if (earlier) earlier(version);
seam->version = version;
};
const auto client = session->client;
const std::string qualified = "`" + session->database + "`.`" +
session->table_prefix;
const std::string after_conflict =
jc::FindString(line, "after_conflict");
CatalogWriter* publisher = &writer;
before_request = [seam, client, qualified, after_conflict,
publisher] {
if (seam->version == 0 || seam->phase == ConflictSeam::kDone) return;
if (seam->phase == ConflictSeam::kConflicted) {
seam->phase = ConflictSeam::kDone;
if (after_conflict == "transport") {
throw ClickHouseError(
"simulated connection reset while recording the packs");
}
if (after_conflict == "lease_refused") {
const PublisherLease* held = publisher->held_lease();
if (held == nullptr) return;
client->execute(
"INSERT INTO " + qualified + "_publisher_lease` "
"(term, lease_id, holder, acquired_at_ns, expires_at_ns) "
"SELECT toUInt64(%(term)s), generateUUIDv4(), 'rival', "
"now_ns, now_ns + 60000000000 FROM (SELECT "
"toUnixTimestamp64Nano(now64(9)) AS now_ns)",
{{"term", held->term}});
publisher->renew_lease();
}
return;
}
const std::vector<dmi_catalog::Row> own = client->execute(
"SELECT count() FROM " + qualified +
"_index_watermark` WHERE index_version = %(version)s",
{{"version", seam->version}});
if (own.empty() || own[0].empty() || own[0][0] != "1") return;
client->execute(
"INSERT INTO " + qualified +
"_index_watermark` (index_version, publish_id, "
"published_at_ns, indexed_rows, indexed_packs) VALUES "
"(%(version)s, generateUUIDv4(), 1, 0, 0)",
{{"version", seam->version}});
seam->phase = ConflictSeam::kConflicted;
};
}
dmi_catalog::NativeIndexer indexer(&s3, &writer, index_config);
dmi_catalog::IndexPlan plan = indexer.plan(refs);
indexer.read(&plan);
Expand All @@ -1064,7 +1138,14 @@ std::string respond(const std::string& line, Session* session) {
}
writer.acquire_lease(holder);
}
std::optional<dmi_catalog::RequestDeadline> seam_scope;
if (before_request) {
seam_scope.emplace([] { return uint64_t{0}; },
"the conformance driver's conflict seam",
before_request);
}
const dmi_catalog::IndexResultData result = indexer.commit(&plan);
seam_scope.reset();
out = ",\"result\":{\"requested_packs\":" +
std::to_string(result.requested_packs) +
",\"skipped_packs\":" + std::to_string(result.skipped_packs) +
Expand Down
Loading
Loading