diff --git a/docs/benchmarks.md b/docs/benchmarks.md index 024979c2f..53e2d5d87 100644 --- a/docs/benchmarks.md +++ b/docs/benchmarks.md @@ -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 @@ -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 +()`. 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) diff --git a/docs/capture-storage-design.md b/docs/capture-storage-design.md index 4ce5b1a31..8fdc5c4ed 100644 --- a/docs/capture-storage-design.md +++ b/docs/capture-storage-design.md @@ -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 @@ -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. @@ -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(, index_version)` this document used to describe has been @@ -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 @@ -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 diff --git a/docs/catalog-descriptor-key.md b/docs/catalog-descriptor-key.md index 34d3dbf7c..c910fc022 100644 --- a/docs/catalog-descriptor-key.md +++ b/docs/catalog-descriptor-key.md @@ -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`. @@ -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 -- diff --git a/docs/catalog-differential-review-2026-09-01.md b/docs/catalog-differential-review-2026-09-01.md index 03425346a..dec9e8565 100644 --- a/docs/catalog-differential-review-2026-09-01.md +++ b/docs/catalog-differential-review-2026-09-01.md @@ -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: diff --git a/native/csrc/catalog/conformance_catalog.cpp b/native/csrc/catalog/conformance_catalog.cpp index 842e88f35..e66c2f19e 100644 --- a/native/csrc/catalog/conformance_catalog.cpp +++ b/native/csrc/catalog/conformance_catalog.cpp @@ -11,9 +11,11 @@ #include #include +#include #include #include #include +#include #include #include #include @@ -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(); + std::function before_request; + if (jc::FindBool(line, "conflict_at_publish")) { + const std::function 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 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); @@ -1064,7 +1138,14 @@ std::string respond(const std::string& line, Session* session) { } writer.acquire_lease(holder); } + std::optional 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) + diff --git a/native/csrc/catalog/indexer.cpp b/native/csrc/catalog/indexer.cpp index 25793b897..516eabf61 100644 --- a/native/csrc/catalog/indexer.cpp +++ b/native/csrc/catalog/indexer.cpp @@ -337,7 +337,8 @@ IndexResultData NativeIndexer::commit(IndexPlan* planned) { // 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(config_.max_publish_attempts); ++attempts) { @@ -348,10 +349,51 @@ IndexResultData NativeIndexer::commit(IndexPlan* planned) { 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()) { + // + // If that fails, the conflict is still 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. + const auto conflict_then = [&e](const char* failure) { + return CatalogError( + CatalogError::Kind::kPublishConflict, + std::string(e.what()) + + " (recording its packs in the inventory then failed " + "too" + + (failure != nullptr ? std::string(": ") + failure + : std::string()) + + ")"); + }; + try { writer_->commit_packs(RenderPackRows(indexed), version); + } catch (const CatalogError& commit_failure) { + // Except a lease refusal, which keeps its kind. Under the + // storage service this INSERT runs behind the lease scope's + // hook, which renews first once a renewal is due, and a rival + // holding the lease -- as one may when a second writer is + // publishing -- refuses it. index_bounded rethrows a lost + // lease, by its kind, so that the pass it cut short is owed; + // relabelled a conflict, it was handled as an ordinary failed + // batch. The Python oracle's commit_packs never renews, so it + // has no such case. + if (is_lease_refusal(commit_failure)) { + throw CatalogError( + commit_failure.kind(), + std::string(commit_failure.what()) + + " (while recording the packs of a publish that " + "conflicted: " + + e.what() + ")"); + } + throw conflict_then(commit_failure.what()); + } catch (const std::exception& commit_failure) { + throw conflict_then(commit_failure.what()); + } catch (...) { + throw conflict_then(nullptr); } } throw; diff --git a/native/csrc/catalog/lease_coordinator.h b/native/csrc/catalog/lease_coordinator.h index 33b4ede68..e43e353d1 100644 --- a/native/csrc/catalog/lease_coordinator.h +++ b/native/csrc/catalog/lease_coordinator.h @@ -150,6 +150,15 @@ class CatalogError : public std::runtime_error { Kind kind_; }; +// The publisher lease is gone: a claim or renewal met another holder +// (kHeld), or this writer holds none or was fenced out (kLease). The storage +// service rethrows these so that the pass they cut short is owed, so code +// that wraps a failed request in another error has to keep these kinds. +inline bool is_lease_refusal(const CatalogError& exc) { + return exc.kind() == CatalogError::Kind::kHeld || + exc.kind() == CatalogError::Kind::kLease; +} + std::string new_uuid_v4(); class LeaseCoordinator { diff --git a/native/csrc/catalog/reader.cpp b/native/csrc/catalog/reader.cpp index 7f2d3342f..9b033428c 100644 --- a/native/csrc/catalog/reader.cpp +++ b/native/csrc/catalog/reader.cpp @@ -32,7 +32,13 @@ constexpr const char* kProjection[] = { "object_key", "object_bytes", "pack_checksum", "pack_record_count", "payload_offset", "stored_length", "decoded_length", "codec", "payload_checksum"}; -constexpr const char* kResolutionOrder = "(index_version, store_id, pack_id)"; +// clickhouse_reader._RESOLUTION_ORDER: a pack ranks by the version its +// FIRST publish reached the watermark at (member_version, from snapshot()), +// never by the version a descriptor row was written at -- a replayed pack's +// rows sit above that -- nor by its newest publish, which a replay also +// moves. index_version last picks, within one pack, the row a merge keeps. +constexpr const char* kResolutionOrder = + "(member_version, store_id, pack_id, index_version)"; std::string quoted(const std::string& name) { return "`" + name + "`"; } @@ -472,17 +478,22 @@ NativeCaptureCatalog::bounded_read_settings() const { return out; } -std::string NativeCaptureCatalog::membership() const { - // clickhouse_sql.membership_predicate, bounded: the snapshot is the set - // of packs whose publish reached the watermark at or before the bound. +std::string NativeCaptureCatalog::snapshot() const { + // clickhouse_reader._snapshot over clickhouse_sql.member_versions: the + // descriptor rows of the packs whose publish reached the watermark at or + // before the bound, each joined to the version its pack FIRST became a + // member at (min, so a replay's publish cannot re-promote a superseded + // pack), which is what kResolutionOrder ranks on. const std::string manifest = qualified("snapshot_manifest"); const std::string watermark = qualified("index_watermark"); return ( - "(store_id, pack_id) IN (" - "SELECT store_id, pack_id FROM " + manifest + " " + qualified("capture_raw") + + " INNER JOIN (SELECT store_id, pack_id, min(index_version) AS " + "member_version FROM " + manifest + " " "WHERE index_version <= %(watermark)s AND (index_version, publish_id) IN " "(SELECT index_version, publish_id FROM " + watermark + - " WHERE index_version <= %(watermark)s))"); + " WHERE index_version <= %(watermark)s) GROUP BY store_id, pack_id) " + "AS `members` USING (store_id, pack_id)"); } std::string NativeCaptureCatalog::projection() const { @@ -699,7 +710,12 @@ SearchPage NativeCaptureCatalog::search(const SearchFilters& filters) const { } Params params{{"watermark", watermark}}; - std::string clauses = membership(); + // The snapshot bound is the join in the FROM clause, so these are only + // the caller's filters, and there may be none. + std::string clauses; + auto add = [&clauses](const std::string& clause) { + clauses += (clauses.empty() ? "" : " AND ") + clause; + }; for (const auto& [value, name] : std::vector*, const char*>>{ {&filters.tenant_id, "tenant_id"}, @@ -708,7 +724,7 @@ SearchPage NativeCaptureCatalog::search(const SearchFilters& filters) const { {&filters.session_id, "session_id"}, {&filters.model_id, "model_id"}}) { if (value->has_value()) { - clauses += " AND " + quoted(name) + " = %(" + name + ")s"; + add(quoted(name) + " = %(" + name + ")s"); params.emplace(name, **value); } } @@ -719,7 +735,7 @@ SearchPage NativeCaptureCatalog::search(const SearchFilters& filters) const { rendered += sql_quote(filters.hook_names[i]); } rendered += ")"; - clauses += " AND hook_name IN " + rendered; + add("hook_name IN " + rendered); } if (!filters.layer_numbers.empty()) { std::string rendered = "("; @@ -728,14 +744,14 @@ SearchPage NativeCaptureCatalog::search(const SearchFilters& filters) const { rendered += std::to_string(filters.layer_numbers[i]); } rendered += ")"; - clauses += " AND layer_number IN " + rendered; + add("layer_number IN " + rendered); } if (filters.captured_after_ns.has_value()) { - clauses += " AND captured_at_ns >= %(captured_after_ns)s"; + add("captured_at_ns >= %(captured_after_ns)s"); params.emplace("captured_after_ns", *filters.captured_after_ns); } if (filters.captured_before_ns.has_value()) { - clauses += " AND captured_at_ns <= %(captured_before_ns)s"; + add("captured_at_ns <= %(captured_before_ns)s"); params.emplace("captured_before_ns", *filters.captured_before_ns); } if (after.has_value()) { @@ -751,7 +767,7 @@ SearchPage NativeCaptureCatalog::search(const SearchFilters& filters) const { placeholders += "%(after_" + std::string(names[i]) + ")s"; params.emplace("after_" + std::string(names[i]), (*after)[i]); } - clauses += " AND (" + columns + ") > (" + placeholders + ")"; + add("(" + columns + ") > (" + placeholders + ")"); } // One row beyond the page tells whether a cursor is owed, without a @@ -773,12 +789,14 @@ SearchPage NativeCaptureCatalog::search(const SearchFilters& filters) const { // about the whole catalog whatever its size. The Python reference reader // keeps that single-phase shape, and the parity suite compares the two. // - // Both queries carry every filter (the same `clauses`), so groups and - // resolution are unchanged. An inner query missing one -- snapshot - // membership, a hook filter -- fills its LIMIT with keys the outer query - // then drops: the page comes back short, owes no cursor, and a walk ends - // early. test_a_page_walk_skips_unpublished_keys_without_ending_early pins - // that. + // Both queries read the same snapshot() join and carry every filter (the + // same `clauses`), so groups and resolution are unchanged. An inner query + // missing one -- snapshot membership, a hook filter -- fills its LIMIT with + // keys the outer query then drops: the page comes back short, owes no + // cursor, and a walk ends early. + // test_a_page_walk_skips_unpublished_keys_without_ending_early pins that. + // The inner query needs only the membership the join applies, not its + // member_version; the outer one ranks on it. // // The second read counts against the read guard. max_rows_to_read limits // the whole statement, and both queries read capture_raw: the inner one @@ -790,11 +808,12 @@ SearchPage NativeCaptureCatalog::search(const SearchFilters& filters) const { // shape refuses this one with Code 158). Size max_rows_to_read for two // passes over the rows past the cursor. const std::vector rows = client_->execute( - "SELECT " + projection() + " FROM " + qualified("capture_raw") + - " WHERE " + clauses + " AND (" + grouped + ") IN (SELECT " + - grouped + " FROM " + qualified("capture_raw") + " WHERE " + - clauses + " GROUP BY " + grouped + " ORDER BY " + order + - " LIMIT %(row_limit)s) GROUP BY " + grouped + " ORDER BY " + order + + "SELECT " + projection() + " FROM " + snapshot() + " WHERE " + + (clauses.empty() ? "" : clauses + " AND ") + "(" + grouped + + ") IN (SELECT " + grouped + " FROM " + snapshot() + + (clauses.empty() ? "" : " WHERE " + clauses) + " GROUP BY " + + grouped + " ORDER BY " + order + " LIMIT %(row_limit)s) GROUP BY " + + grouped + " ORDER BY " + order + " LIMIT %(row_limit)s", params, bounded_read_settings()); @@ -925,7 +944,7 @@ std::vector> NativeCaptureCatalog::get_by_ids( // primary index narrows the read to one tenant's range and the bloom // filter prunes granules inside it. const std::string head = - "SELECT " + projection() + " FROM " + qualified("capture_raw") + + "SELECT " + projection() + " FROM " + snapshot() + " WHERE tenant_id = %(tenant_id)s AND capture_id IN "; // Chunked by rendered bytes: the ids land in the statement TEXT, and a // full-size lookup can breach max_query_size. Ids are sent once each. @@ -941,9 +960,9 @@ std::vector> NativeCaptureCatalog::get_by_ids( ids += sql_quote(chunk[i]); } // The ids land in the statement TEXT (chunked inline, like the - // writer's members); the membership + snapshot bound ride as params. + // writer's members); the snapshot bound rides as a param. const std::vector rows = client_->execute( - head + "(" + ids + ") AND " + membership() + + head + "(" + ids + ")" + " GROUP BY `tenant_id`,`experiment_id`,`run_id`," "`captured_at_ns`,`capture_id`", {{"tenant_id", tenant_id}, {"watermark", requested}}, diff --git a/native/csrc/catalog/reader.h b/native/csrc/catalog/reader.h index 256268a83..cbc105b9a 100644 --- a/native/csrc/catalog/reader.h +++ b/native/csrc/catalog/reader.h @@ -90,7 +90,7 @@ class NativeCaptureCatalog { const std::string& tenant_id, const std::string& watermark) const; private: - std::string membership() const; + std::string snapshot() const; std::string projection() const; std::string qualified(const std::string& table) const; std::map settings() const; diff --git a/native/csrc/catalog/storage_service.cpp b/native/csrc/catalog/storage_service.cpp index fd07c42f6..df57f79ca 100644 --- a/native/csrc/catalog/storage_service.cpp +++ b/native/csrc/catalog/storage_service.cpp @@ -54,11 +54,6 @@ bool is_pack_id(const std::string& value) { return dmi_pack::ParseUuid(value, &bytes, &canonical) && canonical == value; } -bool is_lease_refusal(const CatalogError& exc) { - return exc.kind() == CatalogError::Kind::kHeld || - exc.kind() == CatalogError::Kind::kLease; -} - // The lease thread's tick: the longest it sleeps, and the retry interval for // a claim another holder refused or one that wrote nothing -- a sixth of the // TTL. A renewal falls due a third of the TTL after the claim that stamped diff --git a/src/dmi/storage/capture/catalog.py b/src/dmi/storage/capture/catalog.py index d137b3539..a9bd0c30f 100644 --- a/src/dmi/storage/capture/catalog.py +++ b/src/dmi/storage/capture/catalog.py @@ -526,18 +526,23 @@ def _publish( A losing publish made nothing visible, so recovery is a fresh version and another attempt -- and the DESCRIPTORS are rewritten at that - version, not only the manifest rows. VISIBILITY does not need the - rewrite (membership decides it, and membership is rewritten at the - version that wins), but SUPERSESSION does: ``index_version`` leads - ``clickhouse_reader._RESOLUTION_ORDER``, the primary ordering between - two rows describing one capture in two DIFFERENT packs. Rows left at - the lost version would rank below another pack's rows written between - the lost and the winning version, so the reader would resolve a capture - to the OLDER publish's pack -- and to its locator, which is exactly - what may differ. The rewrite is byte-identical rows at the new version; - the superseded rows share their full sort key with them (pack identity - included), so the ReplacingMergeTree collapses each pair to the new - version and ``argMax`` resolves the same rows in the meantime. + version, not only the manifest rows. Neither visibility nor + supersession depends on that rewrite any more. Membership decides + visibility, and it is rewritten at the version that wins. Supersession + ranks a pack by the version its first publish reached the watermark at + (``member_version`` in ``clickhouse_reader._RESOLUTION_ORDER``), read + from the manifest, and not by the ``index_version`` its rows were + written at. It used to rank on the rows' version, and that is what + broke on a REPLAY: a pass re-indexing a pack that was already published + -- after a crash before ``commit_packs`` below, or an outcome-unknown + publish that landed -- writes its rows at a new, higher version before + publishing, so a pack that had already been superseded outranked the + newer pack inside snapshots it could no longer change, pinned ones + included. The rewrite still keeps a published pack's rows at the + version that published them. They are byte-identical rows that share + their full sort key (pack identity included) with the rows at the lost + version, so the ReplacingMergeTree collapses each pair to the new + version. A publish that CONFLICTED is the opposite of a loss: it is visible, by that error's own contract, so its packs enter the replay inventory diff --git a/src/dmi/storage/capture/clickhouse_reader.py b/src/dmi/storage/capture/clickhouse_reader.py index 68bb19d7b..028cf9ce7 100644 --- a/src/dmi/storage/capture/clickhouse_reader.py +++ b/src/dmi/storage/capture/clickhouse_reader.py @@ -34,11 +34,17 @@ Rows describing one capture in DIFFERENT packs survive side by side, and the ``argMax`` projection grouped on capture identity picks between them: -newest-wins. A reader pinned before the second pack was committed never sees -its rows at all, because the membership clause excludes that pack, so the pin -still resolves to the pack it was taken over. Two packs indexed in one batch -share an ``index_version``, so version alone does not order those rows; what -does is described at :meth:`ClickHouseCaptureCatalog._projection`. +newest-wins, where "newest" is the version at which each pack's publish +reached the watermark -- read from the manifest, at or below the pin -- and +never the version a descriptor row happens to carry. A reader pinned before +the second pack was committed never sees its rows at all, because the +membership clause excludes that pack. And a pass that re-indexes a pack which +is already a member rewrites its rows at a HIGHER version without changing +when that pack was published, so the pin still resolves to the pack it was +taken over; :meth:`ClickHouseCaptureCatalog._snapshot` has the routes that +do that. Two packs published in one batch share a version, so version alone +does not order their rows; what does is described at +:meth:`ClickHouseCaptureCatalog._projection`. The identity rule ----------------- @@ -91,11 +97,12 @@ from .clickhouse_schema import CAPTURE_COLUMNS from .clickhouse_sql import ( DECIDING_READ, + MEMBER_VERSION, ClickHouseClient, identifier, inline_chunks, inline_text_bytes, - membership_predicate, + member_versions, quoted, ) from .cursor import CursorKey, decode_cursor, encode_cursor @@ -142,7 +149,7 @@ # The ordering argument the projection's argMax resolves on. It is a tuple, not # ``index_version``, because it has to be a TOTAL order over the rows in one # group; ``_projection`` explains why, and what breaks without it. -_RESOLUTION_ORDER = "(index_version, store_id, pack_id)" +_RESOLUTION_ORDER = f"({MEMBER_VERSION}, store_id, pack_id, index_version)" @dataclass(frozen=True, slots=True) class ClickHouseReaderConfig: @@ -319,9 +326,9 @@ def search(self, query: CaptureQuery) -> CapturePage: # One row beyond the page tells us whether a cursor is owed, without a # second counting query. params["row_limit"] = query.limit + 1 + where = f"WHERE {' AND '.join(clauses)} " if clauses else "" sql = ( - f"SELECT {self._projection()} FROM {self._qualified()} " - f"WHERE {' AND '.join(clauses)} " + f"SELECT {self._projection()} FROM {self._snapshot()} {where}" f"GROUP BY {', '.join(quoted(name) for name in _SORT_KEY)} " f"ORDER BY {', '.join(quoted(name) for name in _SORT_KEY)} " "LIMIT %(row_limit)s" @@ -380,9 +387,8 @@ def get_by_ids( # the primary index narrows the read to one tenant's range, and the # capture_id bloom-filter skip index prunes granules inside it. sql = ( - f"SELECT {self._projection()} FROM {self._qualified()} " - "WHERE tenant_id = %(tenant_id)s AND " - f"capture_id IN %(capture_ids)s AND {self._membership()} " + f"SELECT {self._projection()} FROM {self._snapshot()} " + "WHERE tenant_id = %(tenant_id)s AND capture_id IN %(capture_ids)s " f"GROUP BY {', '.join(quoted(name) for name in _SORT_KEY)}" ) # Chunked by rendered bytes, because the ids land in the statement TEXT @@ -414,14 +420,33 @@ def _qualified(self, table: str | None = None) -> str: f"{quoted(table or self._capture_raw)}" ) - def _membership(self) -> str: - """The packs inside the snapshot, as a subquery on (store_id, pack_id). - - Two conditions, and the second is the whole point. A manifest row is - written before its watermark row, so requiring the publish to appear in - the watermark table is what stops a publish that lost the race -- which - never wrote one -- from leaking its packs into a snapshot that was - pinned before it ran. + def _snapshot(self) -> str: + """The descriptor rows of the packs inside the snapshot, each carrying + the version its pack first became a member at. + + An INNER JOIN on ``(store_id, pack_id)`` against + ``clickhouse_sql.member_versions`` rather than an ``IN``, because the + reader needs more than whether a pack is inside the snapshot: it ranks + a capture's packs by WHEN each was published (``_RESOLUTION_ORDER``), + and only the manifest knows that. A descriptor row's own + ``index_version`` does not. It is the version the row was WRITTEN at, + and a pass that re-indexes an already-published pack -- after a crash + between publishing and ``commit_packs``, after an outcome-unknown + publish that landed, or as a rebuild running beside the live indexer + -- writes that pack's rows again at a fresh, higher version before it + publishes anything, and may never publish it. Ranked on the row's + version, those rows outranked a newer pack's inside every snapshot the + old pack was already a member of, pinned ones included, and a merge + then makes the higher version the only one left. The rank is the + pack's FIRST publish for the same reason: a replay that does publish + makes the pack a member again at a fresh version, and ranked on its + newest publish a superseded pack would win every head from then on. + + The membership subquery has two conditions, and the second is the + whole point. A manifest row is written before its watermark row, so + requiring the publish to appear in the watermark table is what stops + a publish that lost the race -- which never wrote one -- from leaking + its packs into a snapshot that was pinned before it ran. That second test pairs ``(index_version, publish_id)`` rather than matching the version alone, so a manifest row counts only when the SAME @@ -435,15 +460,18 @@ def _membership(self) -> str: UUID published by a second store at a later version slip inside a pinned snapshot. - The predicate itself has ONE definition, - ``clickhouse_catalog.membership_predicate``, shared with the public - view's DDL so the two cannot drift apart about what exists; this - method only supplies the reader's snapshot bound. + Which manifest rows count has ONE definition in ``clickhouse_sql``, + shared by ``member_versions`` here and ``membership_predicate`` in the + public view's DDL, so the two cannot drift apart about what exists; + this method only supplies the reader's snapshot bound. """ - return membership_predicate( + members = member_versions( self._qualified(self._manifest), self._qualified(self._watermark_table), - bounded=True, + ) + return ( + f"{self._qualified()} INNER JOIN ({members}) AS `members` " + "USING (store_id, pack_id)" ) @staticmethod @@ -474,10 +502,10 @@ def _projection() -> str: is not a reason to unpick the tuple back into per-column aggregates. **A total ordering argument, so the row that wins cannot move.** This is - the failure that reproduces. Ordering on ``index_version`` alone ties + the failure that reproduces. Ordering on a version alone ties routinely: every pack indexed in one ``CatalogIndexer.index`` call is - written at one version, so two packs describing the same capture in one - batch produce rows whose ``index_version`` is equal. The engine breaks + published at one version, so two packs describing the same capture in + one batch produce rows whose version is equal. The engine breaks those ties consistently within a query but not across physical layouts: one pinned corpus resolved to a different pack at ``max_threads=1`` than it did above it, and to a different one again once a merge had put both @@ -485,19 +513,32 @@ def _projection() -> str: controls, so a selection resolved before one and hydrated after it resolves to different bytes with nothing reporting a change. - ``(index_version, store_id, pack_id)`` is a total order over the rows in - a group. They differ by pack identity -- that is exactly why it is in - the table's physical sort key -- so the tuple is unique per distinct row - and the maximum is one row. Rows that still tie on the whole tuple are - one pack re-indexed at one version, which rewrites byte-identical rows, - so which of those wins cannot be observed. - - Across versions this is unchanged newest-wins: ``index_version`` leads - the tuple, so a later pack still supersedes an earlier one. Within a - version the winner is the highest ``(store_id, pack_id)`` -- there is no - version ordering left to honour, and an arbitrary but FIXED choice is - what a reader needs, so that a selection resolved twice resolves to the - same bytes. + ``(member_version, store_id, pack_id, index_version)`` is a total order + over the rows in a group. Different packs differ by pack identity -- + that is exactly why it is in the table's physical sort key. Rows of ONE + pack share its ``member_version`` and differ only by the version they + were written at, which is last: the highest wins, which is the row a + merge keeps (``ReplacingMergeTree(index_version)``), so a read resolves + the same row before a merge and after it. Under the identity rule those + rows are byte identical anyway; the tiebreak is for a row that breaks + the rule, which then fails hydration against the pack footer rather + than winning or losing depending on the physical layout. Rows that tie + on the whole tuple are one pack re-indexed at one version, which + rewrites byte-identical rows, so which of those wins cannot be observed. + + Across versions this is newest-wins: ``member_version`` leads the + tuple, so a pack first published later supersedes an earlier one. It + is the version the pack's first publish reached the watermark at, at + or below the pin (``_snapshot``), and deliberately NOT the descriptor + row's own ``index_version``: a replayed pack's rows sit at a version + above that publish -- above the pin, or at a version never published at + all -- and leading with it let those rows outrank the pack that really + is newest, flipping a pinned read. Nor is it the pack's newest publish, + for the same reason one step later: a replay that publishes would + re-promote the superseded pack at every head after it. Within a version the winner is the highest + ``(store_id, pack_id)`` -- there is no version ordering left to honour, + and an arbitrary but FIXED choice is what a reader needs, so that a + selection resolved twice resolves to the same bytes. The shape is also what keeps determinism affordable, which is why the two halves arrived together. ClickHouse compares a tuple ordering @@ -509,6 +550,15 @@ def _projection() -> str: and 171.7 ms ordering one -- +22.6% for determinism where the per-column form cost +291%. Across page sizes, pagination depth and the selectivity cases this shape runs +17% to +43%. + + Ranking on ``member_version`` put a join where the membership ``IN`` + was. Both build one hash table over the snapshot's packs and probe it + once per row. Measured against the ``IN`` form on the same data, with + the two interleaved, on embedded ClickHouse 26.7 (chdb), from 100k + rows in 10 packs up to 1M rows in 100k packs: 100- and 1000-row pages + and a 100-id lookup ran from 39% faster to 9% slower, so no cost + stood out above the noise. Primary-key pruning survives the join + (``test_selection_resolve_prunes_to_the_tenant_range``). """ # Deliberately unaliased: naming an aggregate after a source column # shadows that column everywhere else in the statement, and ClickHouse @@ -529,8 +579,9 @@ def _filters( # descriptor rows is not durable, because ReplacingMergeTree deletes # rows sharing a sort key at a time nobody controls. Bounding on packs # is also what makes a pin resolve to the pack it was taken over when a - # later pack re-describes the same capture. - clauses = [self._membership()] + # later pack re-describes the same capture. That bound is the join + # `_snapshot` puts in the FROM clause, so it is not among these filters. + clauses: list[str] = [] params: dict[str, object] = {"watermark": watermark} # Equality and range filters apply to raw rows before grouping. That is diff --git a/src/dmi/storage/capture/clickhouse_sql.py b/src/dmi/storage/capture/clickhouse_sql.py index b0ed65549..2ad008839 100644 --- a/src/dmi/storage/capture/clickhouse_sql.py +++ b/src/dmi/storage/capture/clickhouse_sql.py @@ -128,12 +128,52 @@ def inline_chunks( yield chunk -def membership_predicate(manifest: str, watermark: str, *, bounded: bool) -> str: +def _published_manifest_rows(manifest: str, watermark: str, *, bounded: bool) -> str: + """The manifest rows whose own publish reached the watermark log.""" manifest_bound = "index_version <= %(watermark)s AND " if bounded else "" watermark_bound = " WHERE index_version <= %(watermark)s" if bounded else "" return ( - "(store_id, pack_id) IN (" - f"SELECT store_id, pack_id FROM {manifest} " + f"FROM {manifest} " f"WHERE {manifest_bound}(index_version, publish_id) IN " - f"(SELECT index_version, publish_id FROM {watermark}{watermark_bound}))" + f"(SELECT index_version, publish_id FROM {watermark}{watermark_bound})" + ) + + +def membership_predicate(manifest: str, watermark: str, *, bounded: bool) -> str: + return ( + "(store_id, pack_id) IN (" + "SELECT store_id, pack_id " + f"{_published_manifest_rows(manifest, watermark, bounded=bounded)})" + ) + + +# The column `member_versions` names a pack's membership version under. +MEMBER_VERSION = "member_version" + + +def member_versions(manifest: str, watermark: str) -> str: + """The packs inside the snapshot at ``%(watermark)s``, one row each, with + the FIRST version at which a publish that reached the watermark made the + pack a member. + + The same manifest rows ``membership_predicate`` admits, bounded at the + watermark, so the two cannot disagree about what is inside a snapshot; + this one also says WHEN each pack got there, which is what the reader + ranks a capture's packs by. + + The first publish, not the newest. A pass that replays an already-published + pack -- a crash before ``commit_packs``, an outcome-unknown publish that + landed, a rebuild -- publishes it again at a fresh version. Ranked on its + newest publish, a pack superseded in the meantime by a second pack + describing the same capture would win again at every head from the + replay on. A replay adds nothing to the catalog, so it must not move a + pack's rank; ``min`` fixes the rank once the pack is first published. + Either way a pin is stable: a later publish lands above it and the bound + excludes it. A genuine re-capture is a new pack and a mirror is another + store, so each still gets a fresh first publish. + """ + return ( + f"SELECT store_id, pack_id, min(index_version) AS {MEMBER_VERSION} " + f"{_published_manifest_rows(manifest, watermark, bounded=True)} " + "GROUP BY store_id, pack_id" ) diff --git a/src/dmi/storage/capture/pack.py b/src/dmi/storage/capture/pack.py index 9e2792ae1..3161c2801 100644 --- a/src/dmi/storage/capture/pack.py +++ b/src/dmi/storage/capture/pack.py @@ -627,9 +627,9 @@ def reject_a_foreign_tenant(ref: PackRef, tenants: Iterable[str]) -> None: them, so anyone able to PUT into the bucket could write a well-formed pack whose footer carried another tenant's ``tenant_id`` and ``capture_id``, have it indexed under the victim's tenant, and -- because the reader - resolves a capture with ``argMax`` over ``(index_version, store_id, - pack_id)`` -- become the pack that capture resolves to for every fresh - watermark. Integrity of the pack proves nothing here: the attacker's pack + resolves a capture with ``argMax`` over ``(member_version, store_id, + pack_id, index_version)``, newest pack first -- become the pack that + capture resolves to for every fresh watermark. Integrity of the pack proves nothing here: the attacker's pack is perfectly well-formed. Only its LOCATION is evidence, and this is where the two meet. diff --git a/tests/test_capture_review_findings.py b/tests/test_capture_review_findings.py index 0dba864df..97e20b2fe 100644 --- a/tests/test_capture_review_findings.py +++ b/tests/test_capture_review_findings.py @@ -402,8 +402,8 @@ def test_a_pack_whose_footer_names_another_tenant_is_refused(tmp_path: Path): the descriptors are built, nothing has ever compared what the pack CLAIMS to be against where it was found. Indexed, it would be admitted under the victim's tenant, and since the reader resolves a capture with `argMax` over - `(index_version, store_id, pack_id)` it can become the pack that capture - resolves to at every fresh watermark. + `(member_version, store_id, pack_id, index_version)`, newest pack first, + it can become the pack that capture resolves to at every fresh watermark. """ store, ref = _pack_at( tmp_path, diff --git a/tests/test_clickhouse_capture_reader.py b/tests/test_clickhouse_capture_reader.py index 5e55eaae8..d51b035fb 100644 --- a/tests/test_clickhouse_capture_reader.py +++ b/tests/test_clickhouse_capture_reader.py @@ -36,7 +36,12 @@ # The ordering argument the projection's argMax must carry. Spelled out here # rather than imported so that a change to it fails these tests instead of # silently travelling through them. -_ORDER = "(index_version, store_id, pack_id)" +_ORDER = "(member_version, store_id, pack_id, index_version)" + + +# The published-head read, told apart from the descriptor reads -- whose +# membership join also aggregates `index_version`, per pack -- by its shape. +_HEAD_READ = "SELECT max(index_version) FROM" def _source(descriptor: CaptureDescriptor) -> dict: @@ -89,7 +94,7 @@ def __init__(self, *, descriptors=(), watermark=_WATERMARK, pages=None): def execute(self, query, params=None, **kwargs): self.calls.append((" ".join(query.split()), params, kwargs)) - if "max(index_version)" in query: + if _HEAD_READ in query: # The watermark now comes from the published log, not the # descriptor table. assert "_index_watermark" in query, query @@ -99,7 +104,7 @@ def execute(self, query, params=None, **kwargs): @property def selects(self) -> list[str]: - return [call[0] for call in self.calls if "max(index_version)" not in call[0]] + return [call[0] for call in self.calls if _HEAD_READ not in call[0]] def _catalog(**kwargs) -> tuple[ClickHouseCaptureCatalog, _Client]: @@ -220,7 +225,8 @@ def test_pack_identity_is_resolved_by_argmax_not_grouped_on(): catalog.search(CaptureQuery(limit=10)) sql = client.selects[0] - group_by = sql.split("GROUP BY")[1].split("ORDER BY")[0] + # The outer GROUP BY: the membership join groups its own subquery by pack. + group_by = sql.rsplit("GROUP BY", 1)[1].split("ORDER BY")[0] resolved = _resolved_tuple(sql) for name in ("store_id", "pack_id"): assert f"`{name}`" in resolved @@ -261,19 +267,70 @@ def test_one_aggregate_on_a_total_order_resolves_both_query_sites(): # every resolved column, in _RESOLVED order so the row maps positionally. assert _resolved_tuple(sql) == ", ".join(f"`{n}`" for n in _RESOLVED) assert f"argMax(tuple({_resolved_tuple(sql)}), {_ORDER})" in sql - # And nothing is left resolving on the version alone. - assert ", index_version)" not in sql + # And nothing is left resolving on a version alone. + assert "`, index_version)" not in sql # The grouping columns still project directly, not through the tuple. for name in _SORT_KEY: assert f"`{name}`" in sql.split("argMax(")[0] - # The ordering key is exactly (index_version, store_id, pack_id): version - # first, so a later pack still supersedes an earlier one, then the columns - # the table is physically ordered on beyond capture identity -- the only - # ones a capture's rows can differ in, and therefore the only ones that can - # break the tie a shared version leaves. + # The ordering key is exactly (member_version, store_id, pack_id, + # index_version): the pack's publish version first, so a pack published + # later still supersedes an earlier one, then the columns the table is + # physically ordered on beyond capture identity -- the only ones two packs' + # rows can differ in, and therefore the only ones that can break the tie a + # shared version leaves -- and last the version a row was written at, which + # is all that separates one pack's rows from each other. assert _RESOLUTION_ORDER == _ORDER tail = _CAPTURE_TABLE_ORDER[len(_SORT_KEY) :] - assert _ORDER == "(" + ", ".join(("index_version",) + tail) + ")" + assert _ORDER == ( + "(" + ", ".join(("member_version",) + tail + ("index_version",)) + ")" + ) + + +def test_a_replayed_pack_ranks_by_its_publish_not_by_its_rewritten_rows(): + """A pass that re-indexes a published pack must not flip a pinned read. + + The scenario, found by model checking the publish protocol: pass A + publishes pack P1 at v1 and dies before ``commit_packs``; pass B publishes + P2, a second pack describing the same capture, at v2, and a reader pins + W=2 and resolves the capture to P2. A later pass does not find P1 in the + inventory, re-indexes it and writes P1's descriptor rows at v3 -- before it + publishes anything, and perhaps never. P1 is still a member at W=2, so if + the ranking led with the descriptor row's own ``index_version``, (3, P1) + would outrank (2, P2) and the pinned read would now resolve to P1. + + So the rank has to be the version the pack's publish reached the + watermark at, at or below the pin -- which only the manifest paired with + the watermark log knows -- and it has to lead the ordering at both query + sites. It is the pack's FIRST such publish, ``min(index_version)``: if the + replay does publish P1 at v3, the newest publish would rank P1 at 3 and + resolve every head from v3 on back to the superseded pack, while the + first keeps P1 at 1, below P2. The live suite runs the scenario itself + (``test_a_replayed_pack_does_not_flip_a_pinned_read``). + """ + expected = synthetic_descriptors(1) + catalog, client = _catalog(pages=[expected, expected]) + + catalog.search(CaptureQuery(limit=10)) + catalog.get_by_ids( + [expected[0].capture_id], tenant_id="tenant-a", watermark=str(_WATERMARK) + ) + + manifest = "`default`.`dmi_snapshot_manifest`" + watermark = "`default`.`dmi_index_watermark`" + members = ( + "INNER JOIN (SELECT store_id, pack_id, min(index_version) AS member_version " + f"FROM {manifest} WHERE index_version <= %(watermark)s AND " + "(index_version, publish_id) IN (SELECT index_version, publish_id " + f"FROM {watermark} WHERE index_version <= %(watermark)s) " + "GROUP BY store_id, pack_id) AS `members` USING (store_id, pack_id)" + ) + assert len(client.selects) == 2 + for sql in client.selects: + assert f"FROM `default`.`dmi_capture_raw` {members}" in sql + # The pack's publish version leads, not the row's written version. + order = sql.split(f"argMax(tuple({_resolved_tuple(sql)}), ")[1] + assert order.startswith("(member_version, ") + assert not order.startswith("(index_version") def test_only_the_locator_may_differ_between_a_captures_rows(): @@ -440,7 +497,7 @@ def test_only_a_cursor_bearing_search_reads_the_head_as_deciding(): heads = [ kwargs["settings"] for sql, _, kwargs in client.calls - if "max(index_version)" in sql + if _HEAD_READ in sql ] assert heads == [config.settings, {**config.settings, **_DECIDING_READ}] @@ -459,7 +516,7 @@ def test_a_cursor_at_a_published_watermark_survives_replica_lag(): class _LaggingClient(_Client): def execute(self, query, params=None, **kwargs): - if "max(index_version)" in query: + if _HEAD_READ in query: self.calls.append((" ".join(query.split()), params, kwargs)) settings = kwargs.get("settings") or {} if settings.get("select_sequential_consistency"): @@ -642,7 +699,7 @@ def test_get_by_ids_chunks_the_inlined_id_list(): # bound, same membership subquery. assert len(set(client.selects)) == 1 # And the published-head check runs once, not once per chunk. - heads = [sql for sql, _, _ in client.calls if "max(index_version)" in sql] + heads = [sql for sql, _, _ in client.calls if _HEAD_READ in sql] assert len(heads) == 1 @@ -838,7 +895,7 @@ def test_get_by_ids_matches_commit_membership_on_store_and_pack(): sql = client.selects[0] # Pack identity is (store_id, pack_id); matching pack_id alone would let # the same UUID committed by a second store slip inside a pinned snapshot. - assert "(store_id, pack_id) IN (SELECT store_id, pack_id FROM" in sql + assert "USING (store_id, pack_id)" in sql def test_search_matches_commit_membership_on_store_and_pack(): @@ -847,7 +904,7 @@ def test_search_matches_commit_membership_on_store_and_pack(): catalog.search(CaptureQuery(limit=10)) sql = client.selects[0] - assert "(store_id, pack_id) IN (SELECT store_id, pack_id FROM" in sql + assert "USING (store_id, pack_id)" in sql def test_get_by_ids_rejects_an_unpublished_watermark(): @@ -875,7 +932,7 @@ def _raw_row_catalog(row: tuple) -> ClickHouseCaptureCatalog: original = client.execute def execute(query, params=None, **kwargs): - if "max(index_version)" in query: + if _HEAD_READ in query: return original(query, params, **kwargs) client.calls.append((" ".join(query.split()), params, kwargs)) return [row] diff --git a/tests/test_clickhouse_snapshot_live.py b/tests/test_clickhouse_snapshot_live.py index ec4b9e141..9d078405d 100644 --- a/tests/test_clickhouse_snapshot_live.py +++ b/tests/test_clickhouse_snapshot_live.py @@ -934,6 +934,94 @@ def test_two_packs_describing_one_capture_both_survive_a_merge(second_descriptio assert at_fresh[0].locator == copied[0].locator +@pytest.mark.parametrize( + "second_description", (_copied_to_another_store, _retried_into_a_new_pack) +) +def test_a_replayed_pack_does_not_flip_a_pinned_read(second_description): + """Re-indexing a published pack must not re-promote it over a newer one. + + Found by model checking the publish protocol. Pass A publishes P1 and dies + before ``commit_packs`` -- the "redundant work next pass" the indexer + accepts. Pass B publishes P2, a second pack describing the same capture, + and a reader pins that watermark and resolves the capture to P2. A later + pass does not find P1 in the inventory, re-indexes it and writes P1's rows + at a fresh, higher version, then dies before publishing. P1 is still a + member of the pinned snapshot, and ranked on its rows' own version it + outranked P2 there -- permanently, since nothing ever rewrites those rows, + and a merge then leaves only the higher version. + + The same rows arrive by other routes (an outcome-unknown publish that + landed; a conflict whose ``commit_packs`` failed; a rebuild running beside + the live indexer), so the reader has to rank a pack by when its PUBLISH + reached the watermark, not by when its rows were written. + + And by its FIRST publish, not its newest: if the replay does publish, P1 + becomes a member again at a version above P2's, and ranked on that it + would win every head from then on -- the older pack superseding the newer + one on the ordinary crash-recovery path. A replay adds nothing new to the + catalog, so it must not move a pack's rank. + """ + original = synthetic_descriptors(3) + newer = second_description(original) + tenant = original[0].metadata.tenant_id + ids = [item.capture_id for item in original] + with _catalog() as (writer, reader, client, config): + # Pass A: published, never committed to the inventory. + version = writer.allocate_version() + writer.write_descriptors(original, index_version=version) + _publish(writer, version, refs=_refs(original)) + # Pass B: the newer pack, published and committed. + version = writer.allocate_version() + writer.write_descriptors(newer, index_version=version) + _publish(writer, version, refs=_refs(newer)) + _commit(writer, newer, version) + pinned = reader.current_watermark() + assert pinned == str(version) + + def resolved(watermark): + by_id = reader.get_by_ids(ids, tenant_id=tenant, watermark=watermark) + return {item.capture_id: item.locator for item in by_id} + + at_pin = resolved(pinned) + assert at_pin == {item.capture_id: item.locator for item in newer} + first = reader.search(CaptureQuery(limit=2, tenant_id=tenant)) + assert first.next_cursor is not None + + # The replay: P1's rows again, at a version above the pin, with no + # publish behind them. + replay = writer.allocate_version() + writer.write_descriptors(original, index_version=replay) + + assert resolved(pinned) == at_pin, "the replay flipped the pinned read" + page = reader.search(CaptureQuery(limit=10, tenant_id=tenant)) + assert page.watermark == pinned + assert page.items == newer + rest = reader.search( + CaptureQuery(limit=2, tenant_id=tenant, cursor=first.next_cursor) + ) + assert first.items + rest.items == newer + + # A merge leaves P1 only at the replay's version; still not a rank. + _merge(client, config) + assert resolved(pinned) == at_pin, "a merge flipped the pinned read" + + # Once the replay DOES publish P1, P1 is a member again at a version + # above P2's -- but it was first published below it, and that is its + # rank. P2 still wins at the new head, and the old pin is unmoved. + _publish(writer, replay, refs=_refs(original)) + head = reader.current_watermark() + assert head == str(replay) + assert resolved(head) == at_pin, "the replay's publish re-promoted P1" + page = reader.search(CaptureQuery(limit=10, tenant_id=tenant)) + assert page.watermark == head + assert page.items == newer + assert resolved(pinned) == at_pin + + _merge(client, config) + assert resolved(head) == at_pin, "a merge re-promoted P1" + assert resolved(pinned) == at_pin + + def test_a_pin_ignores_a_second_store_holding_the_same_pack_id(): """Pack identity is the PAIR, proven by behaviour rather than by SQL text. diff --git a/tests/test_native_capture_storage_live.py b/tests/test_native_capture_storage_live.py index 6d4b0ea4b..01218b5ba 100644 --- a/tests/test_native_capture_storage_live.py +++ b/tests/test_native_capture_storage_live.py @@ -1470,6 +1470,101 @@ def test_a_pass_whose_lease_changed_while_it_read_rereads_the_replay_guard( assert len(captures) == 4 +@pytest.mark.parametrize("after_conflict", [None, "transport", "lease_refused"]) +def test_a_conflicted_publish_reports_the_conflict_unless_the_lease_was_lost( + fake_s3, tmp_path, after_conflict): + """A publish that finds a second writer's row at its own version is + visible and must not be retried, so the indexer records its packs in the + inventory and raises kPublishConflict: a supervisor matching on it learns + that something else is writing the prefix. + + If that inventory INSERT then fails in transport, the conflict is still + what the pass reports, with the failure in its message (catalog.py's + `raise conflict from commit_failure`). Left to propagate, the transport + error replaced it. + + A lease refusal on that request is not a transport error, though. Under + the storage service every request runs behind the lease scope's hook, + which renews first once a renewal is due, and a rival claiming the lease + -- which is when a second writer turns up -- refuses it. index_bounded + tells a lost lease by its error kind and rethrows it, so that the pass it + cut short is owed (storage_service.cpp); relabelled a conflict, the loss + was handled as an ordinary failed batch. The refusal keeps its kind and + carries the conflict in its message. + + The conflict is real: a foreign watermark row lands at the pass's version + after the pass's own and before its owners read-back + (`conflict_at_publish`, in conformance_catalog's `index` op). + """ + from tests.test_native_catalog_lease_live import CatalogDriver, _open + + spool_root = tmp_path / "spool" + _stage(spool_root, range(4)) # two packs + store = _Driver(STORE_DRIVER) + try: + uploaded = store.call( + op="upload_pending", endpoint=fake_s3, bucket=BUCKET, + region=REGION, access=ACCESS, secret=SECRET, token=None, + insecure=True, connect_timeout=5, read_timeout=15, max_attempts=4, + store_id="s3", root=str(spool_root), spool_max_bytes=1 << 40, + limit=-1, max_workers=4, max_in_flight_bytes=1 << 30) + assert uploaded["ok"], uploaded + finally: + store.close() + refs = uploaded["refs"] + assert len(refs) == 2, refs + + with _catalog() as (client, catalog): + driver = CatalogDriver() + try: + _open(driver, catalog.table_prefix) + assert driver.call(op="ensure_schema")["ok"] + assert driver.call(op="acquire", holder="indexer")["ok"] + seam = {"conflict_at_publish": True} + if after_conflict is not None: + seam["after_conflict"] = after_conflict + result = driver.call( + op="index", refs=refs, endpoint=fake_s3, bucket=BUCKET, + region=REGION, access=ACCESS, secret=SECRET, insecure=True, + **seam) + finally: + driver.close() + + def table(name): + return f"`{DATABASE}`.`{catalog.table_prefix}_{name}`" + + # The conflict: two publishes at one version, the pass's own among + # them with its whole manifest, so its packs are visible. + conflicted = client.execute( + f"SELECT index_version FROM {table('index_watermark')} " + "GROUP BY index_version HAVING uniqExact(publish_id) = 2") + assert len(conflicted) == 1, conflicted + assert client.execute( + f"SELECT count() FROM {table('snapshot_manifest')} " + "WHERE index_version = %(version)s", + {"version": conflicted[0][0]}) == [(2,)] + recorded = client.execute( + f"SELECT count() FROM {table('pack_inventory_raw')}")[0][0] + + assert not result["ok"], result + message = result["message"] + if after_conflict is None: + assert result["error"] == "SnapshotPublishConflictError", result + assert recorded == 2, result + elif after_conflict == "transport": + assert result["error"] == "SnapshotPublishConflictError", result + assert "was published by this writer" in message, result + assert ("recording its packs in the inventory then failed too" + in message), result + assert "simulated connection reset" in message, result + assert recorded == 0, result + else: + assert result["error"] == "PublisherLeaseHeldError", result + assert "is contested" in message, result + assert "was published by this writer" in message, result + assert recorded == 0, result + + def _lease_head_read(request: bytes) -> bool: return b"SELECT term, toString(lease_id)" in request diff --git a/tests/test_native_reader_parity_live.py b/tests/test_native_reader_parity_live.py index 4c7f30883..85fe1d871 100644 --- a/tests/test_native_reader_parity_live.py +++ b/tests/test_native_reader_parity_live.py @@ -1075,10 +1075,10 @@ def test_get_by_ids_parity_and_watermark_validation(): def test_supersession_resolves_the_newest_pack(): """The same capture re-described by a later pack: newest wins, both sides. - The resolution order is (index_version, store_id, pack_id) — a later - version supersedes; within one version the highest (store_id, pack_id) - wins, a fixed choice so a selection resolved twice resolves to the - same bytes. + The resolution order is (member_version, store_id, pack_id, + index_version) — a pack first published at a later version supersedes; within + one version the highest (store_id, pack_id) wins, a fixed choice so a + selection resolved twice resolves to the same bytes. """ with _catalog() as (client, config, prefix): driver = CatalogDriver() @@ -1110,6 +1110,141 @@ def test_supersession_resolves_the_newest_pack(): driver.close() +def test_a_replayed_pack_does_not_flip_a_pinned_read_on_either_side(): + """Re-indexing a published pack must not re-promote it, native or Python. + + The old pack is published at 7 and never committed; a newer pack + describing the same captures is published at 8. A later pass re-indexes + the old pack and writes its rows at 9 without publishing (it crashed). + Both readers must still resolve snapshot 8 to the NEW pack: the old + pack's publish is 7, whatever version its rows were written at. And when + a replay does publish the old pack at 9, both must resolve head 9 to the + new pack too: a pack ranks by its FIRST publish, so replaying it cannot + move it above a pack published after it. See + test_clickhouse_snapshot_live.test_a_replayed_pack_does_not_flip_a_pinned_read. + """ + with _catalog() as (client, config, prefix): + driver = CatalogDriver() + try: + _open_helper(driver, prefix) + old_pack = str(uuid.uuid4()) + new_pack = str(uuid.uuid4()) + first = _descriptor_dicts(2, pack_id=old_pack) + second = _descriptor_dicts(2, pack_id=new_pack) + _publish_native(driver, prefix, first, 7) + _publish_native(driver, prefix, second, 8) + replayed = driver.call(op="write_descriptors", descriptors=first, + index_version=9) + assert replayed["ok"], replayed + assert driver.call(op="current_watermark")["watermark"] == "8" + + reader = _python_reader(client, config) + ids = [d["capture_id"] for d in first] + native_search = driver.call(op="search", limit=100) + native_ids = driver.call(op="get_by_ids", capture_ids=ids, + tenant_id="t", watermark="8") + page = _python_page_items(reader) + python_ids = reader.get_by_ids(ids, tenant_id="t", watermark="8") + assert _normalize(native_search["items"]) == _normalize(page.items) + assert _normalize(native_ids["items"]) == _normalize(python_ids) + for item in native_search["items"] + native_ids["items"]: + assert item[21] == new_pack, item # pack_id column + for item in page.items + python_ids: + assert item.locator.pack_id == new_pack + + # The replay publishes: the old pack is a member again, at 9. + published = driver.call( + op="publish_snapshot", index_version=9, + refs=[{"store_id": first[0]["store_id"], "pack_id": old_pack}], + published_at_ns=9, indexed_rows=len(first), indexed_packs=1) + assert published["ok"], published + assert driver.call(op="current_watermark")["watermark"] == "9" + for watermark in ("9", "8"): + native_ids = driver.call(op="get_by_ids", capture_ids=ids, + tenant_id="t", watermark=watermark) + python_ids = reader.get_by_ids(ids, tenant_id="t", + watermark=watermark) + assert _normalize(native_ids["items"]) == _normalize(python_ids) + assert [item[21] for item in native_ids["items"]] == ( + [new_pack] * len(ids)), (watermark, native_ids["items"]) + assert [item.locator.pack_id for item in python_ids] == ( + [new_pack] * len(ids)), watermark + native_search = driver.call(op="search", limit=100) + page = _python_page_items(reader) + assert page.watermark == "9" + assert _normalize(native_search["items"]) == _normalize(page.items) + for item in native_search["items"]: + assert item[21] == new_pack, item + finally: + driver.close() + + +@pytest.mark.parametrize("newer_part", ["written last", "written first"]) +def test_within_one_pack_the_row_a_merge_keeps_is_the_row_both_sides_read( + newer_part): + """A pack's rows at two versions resolve to the newer one, merged or not. + + A pack re-indexed after its publish has its rows at two index_versions, + and a merge keeps only the higher (ReplacingMergeTree(index_version)). + Both rows share member_version, store_id and pack_id, so the last + component of the resolution order, index_version, is all that makes a + pinned read pick the row the merge will keep. Without it the tie is + undefined, and a pin could read one row before a merge and the other + after it. Rows rendered from one pack are identical today, which hides + that; here the second rendering moves payload_offset so the pick shows. + Which row an undefined tie picks follows the order the parts are read + in, so the newer rows go into a part written after the older one and + into one written before it: with index_version dropped from + kResolutionOrder, the second case returns the older rows. The Python + order is also pinned by text, in test_clickhouse_capture_reader + (test_one_aggregate_on_a_total_order_resolves_both_query_sites). + """ + from dmi.storage.capture.clickhouse_reader import _RESOLVED + + with _catalog() as (client, config, prefix): + driver = CatalogDriver() + try: + _open_helper(driver, prefix) + pack = str(uuid.uuid4()) + first = _descriptor_dicts(2, pack_id=pack) + again = [dict(entry, payload_offset=entry["payload_offset"] + + 1_000_000) for entry in first] + if newer_part == "written last": + _publish_native(driver, prefix, first, 7) + written = driver.call(op="write_descriptors", descriptors=again, + index_version=9) + assert written["ok"], written + if newer_part == "written first": + _publish_native(driver, prefix, first, 7) + assert driver.call(op="current_watermark")["watermark"] == "7" + + reader = _python_reader(client, config) + ids = [entry["capture_id"] for entry in first] + offset_at = 5 + list(_RESOLVED).index("payload_offset") + expected = sorted(str(entry["payload_offset"]) for entry in again) + for phase in ("before a merge", "after a merge"): + if phase == "after a merge": + client.execute( + f"OPTIMIZE TABLE `{config.database}`." + f"`{prefix}_capture_raw` FINAL") + native_search = driver.call(op="search", limit=100) + native_ids = driver.call(op="get_by_ids", capture_ids=ids, + tenant_id="t", watermark="7") + assert native_search["ok"] and native_ids["ok"], phase + for rows in (native_search["items"], native_ids["items"]): + assert sorted(row[offset_at] for row in rows) == expected, ( + phase, rows) + page = _python_page_items(reader) + python_ids = reader.get_by_ids(ids, tenant_id="t", + watermark="7") + assert _normalize(native_search["items"]) == _normalize( + page.items), phase + assert _normalize(native_ids["items"]) == _normalize( + python_ids), phase + finally: + driver.close() + + # --- C2: hydration and core summary at parity -------------------------------- def _e2e_setup(fake_s3, prefix, record_count=3, **stage_kwargs):