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 4ce5b1a31..9d50919db 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,7 +1605,8 @@ past the cursor before `LIMIT` kept `limit + 1` of them. The keyset comparison `(tenant_id, experiment_id, run_id, captured_at_ns, capture_id) > (...)` is not usable by the primary-key index, so each page cost about the whole catalog whatever its size. Measured on 25.12 over a 198k-row catalog, a 38-row page -took ~150 ms, three quarters of it that tuple. Both queries carry every filter, +took ~150 ms, three quarters of it that tuple. Both queries read the same +snapshot join (the one that supplies `member_version`) and carry every filter, so the groups and their resolution are unchanged: that is the same immutability rule, below, that makes the pre-aggregation filters safe. The Python reference reader keeps the single-phase shape, so the parity suite compares the two 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/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 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/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_reader_parity_live.py b/tests/test_native_reader_parity_live.py index 4c7f30883..b75bbf976 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,75 @@ 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() + + # --- C2: hydration and core summary at parity -------------------------------- def _e2e_setup(fake_s3, prefix, record_count=3, **stage_kwargs):