From 439f465a8ff5049c95464f24a3f42601d3e10e38 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 20:27:06 +0000 Subject: [PATCH 1/2] Rank a capture's packs by publish version, not descriptor version A pass that re-indexes an already-published pack (crash before commit_packs, an outcome-unknown publish that landed, a conflict whose commit_packs failed, a rebuild beside the live indexer) writes that pack's descriptor rows at a fresh, higher version before publishing anything. The reader ranked a capture's candidate rows by the row's own index_version, and membership bounds packs rather than rows, so those rows outranked a newer pack describing the same capture inside snapshots that were already pinned. If the replay never published, that stayed true. Found by TLA+ model checking (VersionPublish Replay/ReplayCrash/ ReplayCrashPinned). The Python and C++ readers now join the snapshot's packs with the newest version their publish reached the watermark at (clickhouse_sql. member_versions, the same paired-manifest rows as membership_predicate) and resolve on (member_version, store_id, pack_id, index_version). The public view's DDL is unchanged. The C++ indexer also gets Python's guard on the conflict path: a commit_packs failure no longer replaces kPublishConflict. --- docs/capture-storage-design.md | 26 +++- .../catalog-differential-review-2026-09-01.md | 2 + native/csrc/catalog/indexer.cpp | 27 +++- native/csrc/catalog/reader.cpp | 54 ++++--- native/csrc/catalog/reader.h | 2 +- src/dmi/storage/capture/catalog.py | 29 ++-- src/dmi/storage/capture/clickhouse_reader.py | 134 ++++++++++++------ src/dmi/storage/capture/clickhouse_sql.py | 37 ++++- src/dmi/storage/capture/pack.py | 6 +- tests/test_capture_review_findings.py | 4 +- tests/test_clickhouse_capture_reader.py | 90 +++++++++--- tests/test_clickhouse_snapshot_live.py | 74 ++++++++++ tests/test_native_reader_parity_live.py | 50 ++++++- 13 files changed, 418 insertions(+), 117 deletions(-) diff --git a/docs/capture-storage-design.md b/docs/capture-storage-design.md index cfd0f4cfc..420107dc1 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 publish 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,20 @@ 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 + publish reached the watermark, at or below the pin: 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. `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 diff --git a/docs/catalog-differential-review-2026-09-01.md b/docs/catalog-differential-review-2026-09-01.md index 03425346a..6ae3eb225 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 publish reached the watermark at (`clickhouse_reader._snapshot`), and `indexer.cpp` now has this guard too. + **Recommendation:** ```python except SnapshotPublishConflictError as conflict: diff --git a/native/csrc/catalog/indexer.cpp b/native/csrc/catalog/indexer.cpp index 9281ea00c..f7e7c08c4 100644 --- a/native/csrc/catalog/indexer.cpp +++ b/native/csrc/catalog/indexer.cpp @@ -287,7 +287,8 @@ IndexResultData NativeIndexer::index(const std::vector& refs) { // the inventory, so a pack recorded there but never made visible is // skipped forever AND invisible. Only a lost VERSION race is retried // here, repaired by allocating higher and rewriting the descriptors at - // the winning version (supersession ranks by index_version). + // the winning version (supersession ranks by the pack's membership + // version, not by the rows' index_version; see catalog.py _publish). uint64_t attempts = 0; for (; attempts < static_cast(config_.max_publish_attempts); ++attempts) { @@ -298,10 +299,30 @@ IndexResultData NativeIndexer::index(const std::vector& refs) { all_rows.size(), indexed.size()); } catch (const CatalogError& e) { if (e.kind() != CatalogError::Kind::kPublishRace) { - if (e.kind() == CatalogError::Kind::kPublishConflict) { + if (e.kind() == CatalogError::Kind::kPublishConflict && + !indexed.empty()) { // Visible, so skippable: record the packs before propagating. - if (!indexed.empty()) { + try { writer_->commit_packs(RenderPackRows(indexed), version); + } catch (const std::exception& commit_failure) { + // The conflict is the finding (catalog.py's `raise conflict + // from commit_failure`). Left to propagate, the commit's + // transport error replaced kPublishConflict, so a supervisor + // matching on it never saw the second writer, and the visible + // packs stayed out of the inventory for the next pass to + // re-publish. + throw CatalogError( + CatalogError::Kind::kPublishConflict, + std::string(e.what()) + + " (recording its packs in the inventory then failed " + "too: " + + commit_failure.what() + ")"); + } catch (...) { + throw CatalogError( + CatalogError::Kind::kPublishConflict, + std::string(e.what()) + + " (recording its packs in the inventory then failed " + "too)"); } } throw; diff --git a/native/csrc/catalog/reader.cpp b/native/csrc/catalog/reader.cpp index a83ff32ae..c181839f7 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 +// 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. 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,21 @@ 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 became a member + // at, 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, max(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 +709,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 +723,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 +734,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 +743,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 +766,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 @@ -767,8 +782,9 @@ SearchPage NativeCaptureCatalog::search(const SearchFilters& filters) const { order += quoted(column); } const std::vector rows = client_->execute( - "SELECT " + projection() + " FROM " + qualified("capture_raw") + - " WHERE " + clauses + " GROUP BY " + grouped + " ORDER BY " + order + + "SELECT " + projection() + " FROM " + snapshot() + + (clauses.empty() ? "" : " WHERE " + clauses) + " GROUP BY " + + grouped + " ORDER BY " + order + " LIMIT %(row_limit)s", params, bounded_read_settings()); @@ -899,7 +915,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. @@ -915,9 +931,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/src/dmi/storage/capture/catalog.py b/src/dmi/storage/capture/catalog.py index d137b3539..b33d8f917 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 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..5d8b26fdb 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,30 @@ 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 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 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 +457,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 +499,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 +510,30 @@ 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 published later supersedes an earlier one. It is the + version the pack's 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. 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 +545,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 +574,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..c7f6ccfb8 100644 --- a/src/dmi/storage/capture/clickhouse_sql.py +++ b/src/dmi/storage/capture/clickhouse_sql.py @@ -128,12 +128,41 @@ 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 newest 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. + """ + return ( + f"SELECT store_id, pack_id, max(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..981d0d558 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 publish 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..e6cad51c6 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 publish 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..2437e29bc 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 takes a `max(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,67 @@ 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. 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, max(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 +494,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 +513,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 +696,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 +892,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 +901,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 +929,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..441a664aa 100644 --- a/tests/test_clickhouse_snapshot_live.py +++ b/tests/test_clickhouse_snapshot_live.py @@ -934,6 +934,80 @@ 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. + """ + 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, that is the newest publish and + # wins at the new head -- while the old pin still resolves to P2. + _publish(writer, replay, refs=_refs(original)) + assert resolved(str(replay)) == { + item.capture_id: item.locator for item in original + } + 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_reader_parity_live.py b/tests/test_native_reader_parity_live.py index 83d3e2aab..eb3fa6b4a 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 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,48 @@ 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. 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 + finally: + driver.close() + + # --- C2: hydration and core summary at parity -------------------------------- def _e2e_setup(fake_s3, prefix, record_count=3, **stage_kwargs): From 59c28ce0dd6596d92c3bd6abf1430b55efa0b15d Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 19:46:23 -0400 Subject: [PATCH 2/2] Rank a pack by its first publish, not its newest The reader ranked a capture's packs on the newest paired publish at or below the pin, max(index_version) over the manifest. A replay that does publish moves that rank: pass A publishes P1 at v1 and dies before commit_packs, pass B publishes P2 at v2 describing the same capture, and a later pass that re-indexes P1 publishes it again at v3. From v3 on every head resolved the capture back to P1, the older pack, on the normal crash-recovery path. Pins stayed correct; heads did not. Both readers now take min(index_version): a pack's rank is fixed once it is first published, so a replay, published or not, cannot move it. Pins are stable either way, since a later publish lands above the pin. A genuine re-capture is a new pack and a mirror is another store, so both still get a fresh first publish. The Python and C++ member_versions subqueries change together and stay identical. The live replay test and the native/Python parity test now publish the replay and require P2 at the new head, before and after a merge; both failed with max. EXPLAIN indexes=1 plans for the probe corpus (2M rows, 20k packs) are identical under min and max, as are rows read, and timings are within run-to-run noise. --- docs/benchmarks.md | 3 +- docs/capture-storage-design.md | 18 ++++++++--- docs/catalog-descriptor-key.md | 27 ++++++++++++---- .../catalog-differential-review-2026-09-01.md | 2 +- native/csrc/catalog/reader.cpp | 15 ++++----- src/dmi/storage/capture/catalog.py | 2 +- src/dmi/storage/capture/clickhouse_reader.py | 23 ++++++++------ src/dmi/storage/capture/clickhouse_sql.py | 15 +++++++-- src/dmi/storage/capture/pack.py | 2 +- tests/test_capture_review_findings.py | 2 +- tests/test_clickhouse_capture_reader.py | 9 ++++-- tests/test_clickhouse_snapshot_live.py | 24 +++++++++++--- tests/test_native_reader_parity_live.py | 31 +++++++++++++++++-- 13 files changed, 129 insertions(+), 44 deletions(-) diff --git a/docs/benchmarks.md b/docs/benchmarks.md index 024979c2f..0301eab7b 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 diff --git a/docs/capture-storage-design.md b/docs/capture-storage-design.md index 420107dc1..b53b51779 100644 --- a/docs/capture-storage-design.md +++ b/docs/capture-storage-design.md @@ -383,7 +383,7 @@ 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 `(member_version, store_id, -pack_id, index_version)`, newest publish first -- become the pack that capture +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 @@ -1568,8 +1568,9 @@ removed: 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 - publish reached the watermark, at or below the pin: the manifest rows paired - with the watermark log, joined in on `(store_id, pack_id)`. It is NOT a + 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 @@ -1577,8 +1578,15 @@ removed: 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. `index_version` is kept as the last component, - so within one pack the row a merge keeps is also the row a read resolves. + 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 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 6ae3eb225..dec9e8565 100644 --- a/docs/catalog-differential-review-2026-09-01.md +++ b/docs/catalog-differential-review-2026-09-01.md @@ -144,7 +144,7 @@ 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 publish reached the watermark at (`clickhouse_reader._snapshot`), and `indexer.cpp` now has this guard too. +*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 diff --git a/native/csrc/catalog/reader.cpp b/native/csrc/catalog/reader.cpp index c181839f7..7d9618719 100644 --- a/native/csrc/catalog/reader.cpp +++ b/native/csrc/catalog/reader.cpp @@ -33,10 +33,10 @@ constexpr const char* kProjection[] = { "payload_offset", "stored_length", "decoded_length", "codec", "payload_checksum"}; // clickhouse_reader._RESOLUTION_ORDER: a pack ranks by the version its -// 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. index_version last picks, within one pack, the row a -// merge keeps. +// 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)"; @@ -481,13 +481,14 @@ NativeCaptureCatalog::bounded_read_settings() const { 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 became a member - // at, which is what kResolutionOrder ranks on. + // 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 ( qualified("capture_raw") + - " INNER JOIN (SELECT store_id, pack_id, max(index_version) AS " + " 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 + diff --git a/src/dmi/storage/capture/catalog.py b/src/dmi/storage/capture/catalog.py index b33d8f917..a9bd0c30f 100644 --- a/src/dmi/storage/capture/catalog.py +++ b/src/dmi/storage/capture/catalog.py @@ -529,7 +529,7 @@ def _publish( 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 publish reached the watermark at + 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 diff --git a/src/dmi/storage/capture/clickhouse_reader.py b/src/dmi/storage/capture/clickhouse_reader.py index 5d8b26fdb..028cf9ce7 100644 --- a/src/dmi/storage/capture/clickhouse_reader.py +++ b/src/dmi/storage/capture/clickhouse_reader.py @@ -422,7 +422,7 @@ def _qualified(self, table: str | None = None) -> str: def _snapshot(self) -> str: """The descriptor rows of the packs inside the snapshot, each carrying - the version its pack became a member at. + 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 @@ -437,7 +437,10 @@ def _snapshot(self) -> str: 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. + 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 @@ -524,13 +527,15 @@ def _projection() -> str: 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 published later supersedes an earlier one. It is the - version the pack's 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. Within a version the winner is the highest + 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. diff --git a/src/dmi/storage/capture/clickhouse_sql.py b/src/dmi/storage/capture/clickhouse_sql.py index c7f6ccfb8..2ad008839 100644 --- a/src/dmi/storage/capture/clickhouse_sql.py +++ b/src/dmi/storage/capture/clickhouse_sql.py @@ -153,16 +153,27 @@ def membership_predicate(manifest: str, watermark: str, *, bounded: bool) -> str def member_versions(manifest: str, watermark: str) -> str: """The packs inside the snapshot at ``%(watermark)s``, one row each, with - the newest version at which a publish that reached the watermark made the + 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, max(index_version) AS {MEMBER_VERSION} " + 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 981d0d558..3161c2801 100644 --- a/src/dmi/storage/capture/pack.py +++ b/src/dmi/storage/capture/pack.py @@ -628,7 +628,7 @@ def reject_a_foreign_tenant(ref: PackRef, tenants: Iterable[str]) -> None: 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 ``(member_version, store_id, - pack_id, index_version)``, newest publish first -- become the pack that + 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 e6cad51c6..97e20b2fe 100644 --- a/tests/test_capture_review_findings.py +++ b/tests/test_capture_review_findings.py @@ -402,7 +402,7 @@ 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 - `(member_version, store_id, pack_id, index_version)`, newest publish first, + `(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( diff --git a/tests/test_clickhouse_capture_reader.py b/tests/test_clickhouse_capture_reader.py index 2437e29bc..d51b035fb 100644 --- a/tests/test_clickhouse_capture_reader.py +++ b/tests/test_clickhouse_capture_reader.py @@ -40,7 +40,7 @@ # The published-head read, told apart from the descriptor reads -- whose -# membership join also takes a `max(index_version)`, per pack -- by its shape. +# membership join also aggregates `index_version`, per pack -- by its shape. _HEAD_READ = "SELECT max(index_version) FROM" @@ -301,7 +301,10 @@ def test_a_replayed_pack_ranks_by_its_publish_not_by_its_rewritten_rows(): 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. The live suite runs the scenario itself + 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) @@ -315,7 +318,7 @@ def test_a_replayed_pack_ranks_by_its_publish_not_by_its_rewritten_rows(): manifest = "`default`.`dmi_snapshot_manifest`" watermark = "`default`.`dmi_index_watermark`" members = ( - "INNER JOIN (SELECT store_id, pack_id, max(index_version) AS member_version " + "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) " diff --git a/tests/test_clickhouse_snapshot_live.py b/tests/test_clickhouse_snapshot_live.py index 441a664aa..9d078405d 100644 --- a/tests/test_clickhouse_snapshot_live.py +++ b/tests/test_clickhouse_snapshot_live.py @@ -954,6 +954,12 @@ def test_a_replayed_pack_does_not_flip_a_pinned_read(second_description): 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) @@ -999,12 +1005,20 @@ def resolved(watermark): _merge(client, config) assert resolved(pinned) == at_pin, "a merge flipped the pinned read" - # Once the replay DOES publish P1, that is the newest publish and - # wins at the new head -- while the old pin still resolves to P2. + # 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)) - assert resolved(str(replay)) == { - item.capture_id: item.locator for item in 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 diff --git a/tests/test_native_reader_parity_live.py b/tests/test_native_reader_parity_live.py index eb3fa6b4a..f78295891 100644 --- a/tests/test_native_reader_parity_live.py +++ b/tests/test_native_reader_parity_live.py @@ -1076,7 +1076,7 @@ def test_supersession_resolves_the_newest_pack(): """The same capture re-described by a later pack: newest wins, both sides. The resolution order is (member_version, store_id, pack_id, - index_version) — a pack published at a later version supersedes; within + 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. """ @@ -1117,7 +1117,10 @@ def test_a_replayed_pack_does_not_flip_a_pinned_read_on_either_side(): 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. See + 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): @@ -1148,6 +1151,30 @@ def test_a_replayed_pack_does_not_flip_a_pinned_read_on_either_side(): 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()