From 7f0902973f24237526eaff658ffc5acb2a994631 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Tue, 22 Sep 2026 22:16:44 -0400 Subject: [PATCH 1/4] Resolve a search page's argMax for its own keys, not every row past the cursor search() ran one GROUP BY ... ORDER BY ... LIMIT over the raw capture table, so ClickHouse built the full 27-column argMax resolution tuple for every group 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 (EXPLAIN: 26/26 granules), so a 38-row page over a 198k-row catalog read all 198k rows and aggregated all of them: ~150 ms per page whatever the page size. Filtering alone was 22 ms and grouping the five key columns 37 ms; the argMax tuple was the other three quarters. The page's keys are now chosen first, by an inner query that groups only the key columns under the same filters and LIMIT, and the argMax is resolved for those keys alone. Both queries carry every filter, so the groups and each group's resolution are unchanged. Python and native readers take the same shape. Measured on a 204k-capture catalog, paging through one 196,608-capture run: limit 10000: 9.12s -> 6.51s (1.4x) limit 2000: 23.92s -> 10.56s (2.3x) limit 500: 87.29s -> 26.68s (3.3x) with identical ids in identical order at every page size. Not removed: each page still scans the rows past the cursor, because the keyset tuple remains invisible to the index; that floor (~68 ms/page here) needs index-usable cursor bounds, a separate change. The reader suites -- cpu, live reader, snapshot, end-to-end and native/Python reader parity -- pass 155/155 (154 before, plus the new guard), with the conformance drivers rebuilt so no parity case skipped. --- native/csrc/catalog/reader.cpp | 12 +++++++- src/dmi/storage/capture/clickhouse_reader.py | 18 ++++++++++-- tests/test_clickhouse_capture_reader.py | 29 ++++++++++++++++++++ 3 files changed, 55 insertions(+), 4 deletions(-) diff --git a/native/csrc/catalog/reader.cpp b/native/csrc/catalog/reader.cpp index b7f93e41a..c975a1c54 100644 --- a/native/csrc/catalog/reader.cpp +++ b/native/csrc/catalog/reader.cpp @@ -755,9 +755,19 @@ SearchPage NativeCaptureCatalog::search(const SearchFilters& filters) const { grouped += quoted(column); order += quoted(column); } + // Choose the page's keys first, then resolve the argMax tuple for those keys + // alone -- the same shape as the Python reader. One GROUP BY ... LIMIT built + // the full resolution tuple for EVERY group past the cursor before LIMIT kept + // limit + 1 of them, and the keyset tuple comparison is not usable by the + // primary-key index, so a page cost about the whole catalog whatever its + // size. Both queries carry every filter, so groups and resolution are + // unchanged. const std::vector rows = client_->execute( "SELECT " + projection() + " FROM " + qualified("capture_raw") + - " WHERE " + clauses + " GROUP BY " + grouped + " ORDER BY " + order + + " 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 + " LIMIT %(row_limit)s", params, bounded_read_settings()); diff --git a/src/dmi/storage/capture/clickhouse_reader.py b/src/dmi/storage/capture/clickhouse_reader.py index 68bb19d7b..2dc9cf743 100644 --- a/src/dmi/storage/capture/clickhouse_reader.py +++ b/src/dmi/storage/capture/clickhouse_reader.py @@ -319,11 +319,23 @@ 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 + # Choose the page's keys first, then resolve the argMax tuple for those + # keys alone. One GROUP BY ... LIMIT built the full resolution tuple for + # EVERY group past the cursor before LIMIT kept limit + 1 of them, and + # the keyset comparison above is not usable by the primary-key index, + # so each page cost about the whole catalog whatever its size. The inner + # query groups only the key columns under the same filters and LIMIT; + # the outer keeps every filter too, so groups and their resolution are + # unchanged. + key = ", ".join(quoted(name) for name in _SORT_KEY) + where = " AND ".join(clauses) sql = ( f"SELECT {self._projection()} FROM {self._qualified()} " - f"WHERE {' AND '.join(clauses)} " - f"GROUP BY {', '.join(quoted(name) for name in _SORT_KEY)} " - f"ORDER BY {', '.join(quoted(name) for name in _SORT_KEY)} " + f"WHERE {where} AND ({key}) IN (" + f"SELECT {key} FROM {self._qualified()} WHERE {where} " + f"GROUP BY {key} ORDER BY {key} LIMIT %(row_limit)s) " + f"GROUP BY {key} " + f"ORDER BY {key} " "LIMIT %(row_limit)s" ) rows = self._client.execute( diff --git a/tests/test_clickhouse_capture_reader.py b/tests/test_clickhouse_capture_reader.py index 5e55eaae8..3a07d9ad0 100644 --- a/tests/test_clickhouse_capture_reader.py +++ b/tests/test_clickhouse_capture_reader.py @@ -334,6 +334,35 @@ def test_search_groups_and_orders_by_the_sort_key(): assert f"ORDER BY {key}" in sql +def test_search_resolves_argmax_only_for_the_pages_own_keys(): + """The 27-column argMax runs over the page's keys, not every row past the cursor. + + A single GROUP BY ... ORDER BY ... LIMIT resolved the argMax tuple for EVERY + group past the cursor before LIMIT kept (limit + 1) of them. The keyset + tuple comparison is not usable by the primary-key index, so a 38-row page + over a 198k-row catalog read all 198k rows and built all their tuples: + ~150 ms per page whatever the page size, and a full scan paged at 10k read + the table 20 times. About three quarters of each page was that aggregate. + + The page's keys are chosen first by an inner query that groups only the five + key columns and carries the same filters and LIMIT; the argMax is then + computed for those keys alone. Groups and per-group resolution are + unchanged -- the outer query keeps every filter -- which the live parity + suites pin; this pins the shape that makes the cost track the page. + """ + catalog, client = _catalog(descriptors=synthetic_descriptors(1)) + + catalog.search(CaptureQuery(limit=10)) + + sql = client.selects[0] + key = "`tenant_id`, `experiment_id`, `run_id`, `captured_at_ns`, `capture_id`" + assert f"({key}) IN (SELECT {key} FROM" in sql, sql + inner = sql.split(f"({key}) IN (SELECT {key} FROM", 1)[1] + assert "argMax(" not in inner.split(")", 1)[0], "the key subquery must not aggregate the tuple" + assert "LIMIT %(row_limit)s)" in inner, "the key subquery must carry the page LIMIT" + assert sql.count("LIMIT %(row_limit)s") == 2, sql + + # --- pagination ------------------------------------------------------------- From e02e034a18e014ebe6dd4471417e3724ea1632d5 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Tue, 22 Sep 2026 22:28:29 -0400 Subject: [PATCH 2/4] Describe the two-phase search page in the capture storage design The design doc described search as one argMax over the non-grouped columns. That is still the resolution, but a page now computes it only for the keys an inner, key-only query selects under the same filters and LIMIT. This records why (the per-page cost measured before), why it is safe (the same immutability rule that already justifies the pre-aggregation filters), and what it leaves: each page still scans the rows past the cursor. --- docs/capture-storage-design.md | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/docs/capture-storage-design.md b/docs/capture-storage-design.md index 98d5125e1..293821a66 100644 --- a/docs/capture-storage-design.md +++ b/docs/capture-storage-design.md @@ -1566,6 +1566,19 @@ removed: because ClickHouse compares a tuple ordering argument through a generic `Field` once per row per aggregate. +A search page resolves that aggregate for **its own keys only**. An inner query +groups just the five sort-key columns under the page's filters and `LIMIT`, and +the outer query computes the `argMax` tuple for those keys. With a single +`GROUP BY ... LIMIT`, ClickHouse built the full 27-column tuple for every group +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, +so the groups and their resolution are unchanged: that is the same immutability +rule, below, that makes the pre-aggregation filters safe. Each page still scans +the rows past the cursor. Removing that needs index-usable cursor bounds. + Every descriptor field except the locator is immutable for a `(tenant_id, capture_id)`, which is what makes the pre-aggregation `WHERE` filters safe; the rule is written out in `clickhouse_reader`'s module From 02b70ffd3d30c3af880c51280e9a5778a772324e Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Wed, 23 Sep 2026 00:54:02 -0400 Subject: [PATCH 3/4] Keep the two-phase page native-only; guard it in the parity suite The Python storage capture path is being removed, so the Python half of this change is dropped: clickhouse_reader.py and its cpu shape test are back to main, and only the native reader keeps the two-phase page query. That also makes the guard stronger, not weaker. The new parity test walks the native reader in pages of 3 across a superseded version -- four captures re-described from a different pack at a later index_version, the case where argMax resolution decides the answer -- and compares every row with the Python reader, which still uses the single-phase shape. So the suite now checks the new query against the old one rather than against itself. It then reads the queries ClickHouse actually received from the native driver (HTTP, interface 2) out of system.query_log and requires the key-tuple subquery on each. Red first: with reader.cpp back to main the test fails on the missing `(keys) IN (SELECT keys FROM` subquery. An earlier draft asserted a bare ") IN (SELECT " and would have passed against the old query too, because the snapshot-membership filter already contains that text. Reader suites (cpu, live reader, snapshot, end-to-end, native parity) 155/155; cpu suite 2113 passed. --- docs/capture-storage-design.md | 8 ++- src/dmi/storage/capture/clickhouse_reader.py | 18 +---- tests/test_clickhouse_capture_reader.py | 29 -------- tests/test_native_reader_parity_live.py | 75 ++++++++++++++++++++ 4 files changed, 83 insertions(+), 47 deletions(-) diff --git a/docs/capture-storage-design.md b/docs/capture-storage-design.md index 293821a66..42e8c0dca 100644 --- a/docs/capture-storage-design.md +++ b/docs/capture-storage-design.md @@ -1566,7 +1566,7 @@ removed: because ClickHouse compares a tuple ordering argument through a generic `Field` once per row per aggregate. -A search page resolves that aggregate for **its own keys only**. An inner query +The native reader resolves that aggregate for **its own keys only**. An inner query groups just the five sort-key columns under the page's filters and `LIMIT`, and the outer query computes the `argMax` tuple for those keys. With a single `GROUP BY ... LIMIT`, ClickHouse built the full 27-column tuple for every group @@ -1576,8 +1576,10 @@ 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, so the groups and their resolution are unchanged: that is the same immutability -rule, below, that makes the pre-aggregation filters safe. Each page still scans -the rows past the cursor. Removing that needs index-usable cursor bounds. +rule, below, that makes the pre-aggregation filters safe. The Python reference +reader keeps the single-phase shape, so the parity suite compares the two +shapes directly. Each page still scans the rows past the cursor. Removing that +needs index-usable cursor bounds. Every descriptor field except the locator is immutable for a `(tenant_id, capture_id)`, which is what makes the pre-aggregation `WHERE` diff --git a/src/dmi/storage/capture/clickhouse_reader.py b/src/dmi/storage/capture/clickhouse_reader.py index 2dc9cf743..68bb19d7b 100644 --- a/src/dmi/storage/capture/clickhouse_reader.py +++ b/src/dmi/storage/capture/clickhouse_reader.py @@ -319,23 +319,11 @@ 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 - # Choose the page's keys first, then resolve the argMax tuple for those - # keys alone. One GROUP BY ... LIMIT built the full resolution tuple for - # EVERY group past the cursor before LIMIT kept limit + 1 of them, and - # the keyset comparison above is not usable by the primary-key index, - # so each page cost about the whole catalog whatever its size. The inner - # query groups only the key columns under the same filters and LIMIT; - # the outer keeps every filter too, so groups and their resolution are - # unchanged. - key = ", ".join(quoted(name) for name in _SORT_KEY) - where = " AND ".join(clauses) sql = ( f"SELECT {self._projection()} FROM {self._qualified()} " - f"WHERE {where} AND ({key}) IN (" - f"SELECT {key} FROM {self._qualified()} WHERE {where} " - f"GROUP BY {key} ORDER BY {key} LIMIT %(row_limit)s) " - f"GROUP BY {key} " - f"ORDER BY {key} " + f"WHERE {' AND '.join(clauses)} " + 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" ) rows = self._client.execute( diff --git a/tests/test_clickhouse_capture_reader.py b/tests/test_clickhouse_capture_reader.py index 3a07d9ad0..5e55eaae8 100644 --- a/tests/test_clickhouse_capture_reader.py +++ b/tests/test_clickhouse_capture_reader.py @@ -334,35 +334,6 @@ def test_search_groups_and_orders_by_the_sort_key(): assert f"ORDER BY {key}" in sql -def test_search_resolves_argmax_only_for_the_pages_own_keys(): - """The 27-column argMax runs over the page's keys, not every row past the cursor. - - A single GROUP BY ... ORDER BY ... LIMIT resolved the argMax tuple for EVERY - group past the cursor before LIMIT kept (limit + 1) of them. The keyset - tuple comparison is not usable by the primary-key index, so a 38-row page - over a 198k-row catalog read all 198k rows and built all their tuples: - ~150 ms per page whatever the page size, and a full scan paged at 10k read - the table 20 times. About three quarters of each page was that aggregate. - - The page's keys are chosen first by an inner query that groups only the five - key columns and carries the same filters and LIMIT; the argMax is then - computed for those keys alone. Groups and per-group resolution are - unchanged -- the outer query keeps every filter -- which the live parity - suites pin; this pins the shape that makes the cost track the page. - """ - catalog, client = _catalog(descriptors=synthetic_descriptors(1)) - - catalog.search(CaptureQuery(limit=10)) - - sql = client.selects[0] - key = "`tenant_id`, `experiment_id`, `run_id`, `captured_at_ns`, `capture_id`" - assert f"({key}) IN (SELECT {key} FROM" in sql, sql - inner = sql.split(f"({key}) IN (SELECT {key} FROM", 1)[1] - assert "argMax(" not in inner.split(")", 1)[0], "the key subquery must not aggregate the tuple" - assert "LIMIT %(row_limit)s)" in inner, "the key subquery must carry the page LIMIT" - assert sql.count("LIMIT %(row_limit)s") == 2, sql - - # --- pagination ------------------------------------------------------------- diff --git a/tests/test_native_reader_parity_live.py b/tests/test_native_reader_parity_live.py index 83d3e2aab..4109bdd19 100644 --- a/tests/test_native_reader_parity_live.py +++ b/tests/test_native_reader_parity_live.py @@ -2614,3 +2614,78 @@ def test_a_non_bmp_hook_name_indexes_intact(fake_s3): sink.close() store.close() driver.close() + + +def test_native_pages_resolve_argmax_for_their_own_keys_and_match_python(): + """The native page query resolves the argMax tuple for the page's keys only. + + A single GROUP BY ... LIMIT built the full 27-column resolution tuple for + every group past the cursor before LIMIT kept limit + 1 of them, and the + keyset tuple comparison is not usable by the primary-key index -- so each + page cost about the whole catalog whatever its size (~150 ms for a 38-row + page over 198k rows). The native reader now selects the page's keys in an + inner, key-only query under the same filters and LIMIT. + + Two things are pinned. The SHAPE, from the query ClickHouse actually + received: native reads arrive over HTTP (interface 2), and every argMax + query must carry the key subquery. And the RESULT across a superseded + version, walked in pages small enough to cross several cursors, against + the Python reader, which still uses the single-phase shape -- so the + parity suite checks the new shape against the old one. + """ + from dmi.storage.capture.clickhouse_reader import _RESOLVED + + with _catalog() as (client, config, prefix): + driver = CatalogDriver() + try: + _open_helper(driver, prefix) + first = _descriptor_dicts(10) + _publish_native(driver, prefix, first, 7) + # Re-describe four captures from a different pack at a later + # version: resolution must pick version 8's locator for those. + newer_pack = "018f0000-0000-7000-8000-00000000beef" + second = [dict(first[i], pack_id=newer_pack) for i in (0, 3, 6, 9)] + _publish_native(driver, prefix, second, 8) + + native, cursor = [], None + while True: + page = driver.call(op="search", limit=3, + **({"cursor": cursor} if cursor else {})) + assert page["ok"], page + native += page["items"] + cursor = page["next_cursor"] + if cursor is None: + break + reader = _python_reader(client, config) + python, cursor = [], None + while True: + page = _python_page_items(reader, limit=3, cursor=cursor) + python += page.items + cursor = page.next_cursor + if cursor is None: + break + + assert len(native) == 10 == len(python) + assert _normalize(native) == _normalize(python) + pack_at = 5 + list(_RESOLVED).index("pack_id") + moved = {row[4] for row in native if row[pack_at] == newer_pack} + assert moved == {"capture-0", "capture-3", "capture-6", "capture-9"} + + client.execute("SYSTEM FLUSH LOGS") + sent = [row[0] for row in client.execute( + "SELECT query FROM system.query_log WHERE type = 'QueryFinish' " + "AND interface = 2 AND query LIKE %(raw)s AND query LIKE '%%argMax%%' " + "AND event_time > now() - 600", + {"raw": f"%{prefix}_capture_raw%"})] + assert sent, "no native page query reached ClickHouse over HTTP" + # The KEY tuple feeds the subquery. A bare ") IN (SELECT " would also + # match the snapshot-membership filter `(store_id, pack_id) IN + # (SELECT ...)` that the old single-phase query already carried. + keys = "`tenant_id`,`experiment_id`,`run_id`,`captured_at_ns`,`capture_id`" + for query in sent: + assert f"({keys}) IN (SELECT {keys} FROM" in query, query + inner = query.split(f"({keys}) IN (SELECT {keys} FROM", 1)[1] + assert "argMax" not in inner.split(" GROUP BY ", 1)[0], query + assert query.count(" LIMIT ") == 2, query + finally: + driver.close() From 51060837165a60443d0fb1775f3f7835e8025d6d Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Sat, 26 Sep 2026 00:29:35 -0400 Subject: [PATCH 4/4] Pin that a page's two queries filter alike; state the doubled row read The two-phase page is correct only while its inner key query and its outer resolution query carry the same filters, and nothing tested that. With the inner query's snapshot membership dropped, every live test still passed, yet on a catalog holding an unpublished pack a walk ended after one page: the unpublished keys filled the inner LIMIT, the outer query discarded them, and a short page issues no cursor. The new parity test stages a published pack of ten captures with alternating hooks and an unpublished one -- written at version 8, never published, as a crashed indexer leaves it -- holding two keys between each pair of member keys, plus later re-descriptions of two members. Walked in pages of 2, with and without a hook filter, the native reader must return exactly the member captures in order, resolved to the published pack, and agree with the Python reader. Red first, each mutation applied to reader.cpp with conformance_catalog rebuilt: inner query without membership: unfiltered walk returned capture-0 alone, of 10 inner query without the hook filter: filtered walk returned 2 of 5 and in both the other 48 tests in the file still passed. Reverted: 49/49. Also: - reader.cpp no longer says the Python reader shares the shape; it has been single-phase since 02b70ff. - reader.cpp and the design doc state the read-guard cost. max_rows_to_read limits the whole statement and both queries read capture_raw, so with a selective filter a page can read up to twice the rows the single-phase query did: 1,102,131 -> 2,203,414 on the review's corpus, where a limit of 1,500,000 now refuses the page with Code 158; 34 -> 68 on the test corpus here. The query shape is unchanged. --- docs/capture-storage-design.md | 22 ++++++ native/csrc/catalog/reader.cpp | 28 ++++++-- tests/test_native_reader_parity_live.py | 91 +++++++++++++++++++++++++ 3 files changed, 135 insertions(+), 6 deletions(-) diff --git a/docs/capture-storage-design.md b/docs/capture-storage-design.md index 54949e4d6..4ce5b1a31 100644 --- a/docs/capture-storage-design.md +++ b/docs/capture-storage-design.md @@ -1592,6 +1592,28 @@ reader keeps the single-phase shape, so the parity suite compares the two shapes directly. Each page still scans the rows past the cursor. Removing that needs index-usable cursor bounds. +The two queries must filter alike. If the inner query drops a filter that the +outer one keeps, such as snapshot membership or a hook filter, its `LIMIT` fills +with keys the outer query then discards. The page comes back short, a short page +issues no cursor, and the walk ends early while reporting success. The parity +suite walks a corpus built for this: an unpublished pack whose keys interleave +with the member keys, crossed in pages of 2, with and without a hook filter. + +The second read counts against the read guard. `max_rows_to_read` limits the +whole statement, and both queries read `*_capture_raw`. The inner query reads +every row past the cursor that the filters do not prune. The outer query reads +every granule that holds one of the page's keys. A selective filter spreads +those keys across granules, so the outer read can approach a second full scan, +and a page can read up to **twice** the rows the single-phase query did. On one +corpus a page that read 1,102,131 rows now reads 2,203,414, and a +`max_rows_to_read` of 1,500,000 that the old shape passed refuses the new one +with Code 158 (`TOO_MANY_ROWS`). On the parity test's corpus, one granule per +part, 25.12 reports 68 rows read for a native page against 34 for the Python +reader's single-phase page. Size `ReaderConfig::max_rows_to_read` (`reader.h`) +for two passes over the rows past the cursor. At worst, the default of +50,000,000 then trips once about 25 million rows lie past the cursor, not 50 +million. + Every descriptor field except the locator is immutable for a `(tenant_id, capture_id)`, which is what makes the pre-aggregation `WHERE` filters safe; the rule is written out in `clickhouse_reader`'s module diff --git a/native/csrc/catalog/reader.cpp b/native/csrc/catalog/reader.cpp index ec87d158a..7f2d3342f 100644 --- a/native/csrc/catalog/reader.cpp +++ b/native/csrc/catalog/reader.cpp @@ -767,12 +767,28 @@ SearchPage NativeCaptureCatalog::search(const SearchFilters& filters) const { order += quoted(column); } // Choose the page's keys first, then resolve the argMax tuple for those keys - // alone -- the same shape as the Python reader. One GROUP BY ... LIMIT built - // the full resolution tuple for EVERY group past the cursor before LIMIT kept - // limit + 1 of them, and the keyset tuple comparison is not usable by the - // primary-key index, so a page cost about the whole catalog whatever its - // size. Both queries carry every filter, so groups and resolution are - // unchanged. + // alone. One GROUP BY ... LIMIT built the full resolution tuple for EVERY + // group past the cursor before LIMIT kept limit + 1 of them, and the keyset + // tuple comparison is not usable by the primary-key index, so a page cost + // 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. + // + // 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 + // every row past the cursor that the filters leave, the outer one each + // granule holding a page key. A selective filter spreads those keys across + // granules, so the outer read approaches a second full scan, and a page can + // read up to twice the rows the single-phase query did (1,102,131 -> + // 2,203,414 on one corpus, so a limit of 1,500,000 that passed the old + // 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 " + diff --git a/tests/test_native_reader_parity_live.py b/tests/test_native_reader_parity_live.py index 4109bdd19..4c7f30883 100644 --- a/tests/test_native_reader_parity_live.py +++ b/tests/test_native_reader_parity_live.py @@ -2689,3 +2689,94 @@ def test_native_pages_resolve_argmax_for_their_own_keys_and_match_python(): assert query.count(" LIMIT ") == 2, query finally: driver.close() + + +def test_a_page_walk_skips_unpublished_keys_without_ending_early(): + """Every member capture comes back once, however unpublished keys interleave. + + The two-phase page is correct only while its inner key query and its outer + resolution query filter alike. Let the inner one drop a filter the outer one + keeps -- snapshot membership, a hook filter -- and its LIMIT fills with keys + the outer query then discards: the page comes back short, a short page owes + no cursor, and the walk stops early while reporting success. The other tests + cannot see that, because every row they stage is a member that matches the + filter, so dropping the inner membership filter left them all passing. + + Here a published pack holds ten captures whose hooks alternate, and an + unpublished pack -- written at version 8 and never published, as a crashed + indexer leaves one -- holds two keys between each pair of member keys, all + under the filtered hook. It also re-describes two members, and only + membership keeps that later version from winning their argMax. Walked in + pages of 2 (an inner LIMIT of 3, so the first page alone crosses two + unpublished keys), with and without a hook filter, the native reader must + return exactly the member captures, in order and resolved to the published + pack, and agree with the Python reader. + """ + from dmi.storage.capture.clickhouse_reader import _RESOLVED + + with _catalog() as (client, config, prefix): + driver = CatalogDriver() + try: + _open_helper(driver, prefix) + published_pack = "018f0000-0000-7000-8000-00000000a001" + unpublished_pack = "018f0000-0000-7000-8000-00000000b002" + # captured_at_ns rises with the index, so the index is sort order. + members, unpublished = [], [] + for index, entry in enumerate(_descriptor_dicts(30)): + if index % 3 == 0: + hook = "resid_pre" if index % 6 == 0 else "other_hook" + members.append( + dict(entry, pack_id=published_pack, hook_name=hook)) + else: + unpublished.append( + dict(entry, pack_id=unpublished_pack, + hook_name="resid_pre")) + unpublished += [dict(members[i], pack_id=unpublished_pack) + for i in (1, 4)] + _publish_native(driver, prefix, members, 7) + written = driver.call(op="write_descriptors", + descriptors=unpublished, index_version=8) + assert written["ok"], written + assert driver.call(op="current_watermark")["watermark"] == "7" + + reader = _python_reader(client, config) + pages = len(members) + len(unpublished) + 2 # a walk that loops fails + + def walk_native(**filters): + rows, cursor = [], None + for _ in range(pages): + page = driver.call( + op="search", limit=2, **filters, + **({"cursor": cursor} if cursor else {})) + assert page["ok"], page + rows += page["items"] + cursor = page["next_cursor"] + if cursor is None: + return rows + raise AssertionError(f"native walk did not end: {filters}") + + def walk_python(**filters): + items, cursor = [], None + for _ in range(pages): + page = _python_page_items( + reader, limit=2, cursor=cursor, **filters) + items += page.items + cursor = page.next_cursor + if cursor is None: + return items + raise AssertionError(f"python walk did not end: {filters}") + + pack_at = 5 + list(_RESOLVED).index("pack_id") + for native_filters, python_filters, expected in ( + ({}, {}, [m["capture_id"] for m in members]), + ({"hook_names": ["resid_pre"]}, {"hook_names": ("resid_pre",)}, + [m["capture_id"] for m in members + if m["hook_name"] == "resid_pre"]), + ): + native = walk_native(**native_filters) + assert [row[4] for row in native] == expected, native_filters + assert {row[pack_at] for row in native} == {published_pack} + python = walk_python(**python_filters) + assert _normalize(native) == _normalize(python), native_filters + finally: + driver.close()