Repository navigation
Conversation
…the _handle_widget_change split (rows-first s1) A default-tier XorqBuckarooWidget and BuckarooWidget publish df_data_dict, then df_display_args, then the rest of the widget_args_tuple observers on a search change, and merged_sd carries the full stat set. These pass on main and pin the behaviour the split must keep. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
…-handler fields (rows-first s1) A schema-tier XorqServerDataflow should match full stats on pinned_rows, data_key, summary_stats_key and (except the stats-derived minWidth) column_config, issue no data query besides the cached count, and keep init_sd hints and sorted windows working. A pending state must write nothing under a full-tier cache key, and a later full assignment must reach merged_sd for the raw, clean and filt scopes. assemble_merged_sd must equal merged_sd, and _handle_widget_change must be built from separately callable all_stats and display-args builders. /load_expr and /reload_expr accept stats_tier and stats_delivery, replay them on reload, and keep them out of the warm short-circuit's has_config tuple. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
… s1)
Add a dataflow-level stats_tier ("full" default, "schema"). The xorq schema
tier builds identity and typing for every column from the expression's
schema, with no data query beyond the cached row count, so a dataflow
constructs in milliseconds rather than the stats' hundreds.
The tier is part of _scope_cache_key, so a schema entry is never read as a
full one, and _populate_sd_cache stores summary_sd under the filt key only
if it was computed for the current frame, klass list and tier. add_analysis
no longer builds DFStatsClass outside the hook when the tier is not full.
The merged_sd observer body is extracted as the pure assemble_merged_sd,
and _handle_widget_change is split into _build_df_data_dict and
_build_df_display_args.
/load_expr and /reload_expr accept stats_tier and stats_delivery, stored on
the session beside dataflow_kwargs and replayed on reload. They stay out of
the has_config tuple; the warm short-circuit compares the stored pair.
stats_delivery="deferred" builds the schema-tier dataflow and publishes it.
Defaults (full, inline) leave behaviour unchanged.
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
…ation test (rows-first s1) The Max Versions jobs resolve pandas 3, which reports a string column's dtype as 'str' where pandas 2 says 'object'. The characterization test asserted 'object'. Verified in a Max Versions environment (pandas 3.0.6, polars 1.44.2, xorq 0.4.5): the unit suite passes. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
…g on a skipped column (rows-first s1) A column in skip_stat_columns gets only name, dtype and length from the full-tier pipeline, so its _type comes from init_sd. The schema tier layers the schema-derived _type and is_* keys over it, so an int64 column that init_sd types as float merges as integer and renders with zero fraction digits instead of the float displayer init_sd asked for. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
…ows-first s1) _get_schema_sd never read skip_stat_columns, so a skipped column's schema-derived _type and is_* keys overrode init_sd's _type once merged. The full tier gives a skipped column only name, dtype and length. The schema tier now does the same, so init_sd's _type decides the displayer at both tiers. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
…er (rows-first s2) ServerDataflow and PolarsServerDataflow with stats_tier="schema" should publish the display state full stats give (column_config without stats-derived keys, pinned_rows including a host-supplied one, data_key, summary_stats_key) with no stat computed on the data, still apply init_sd, serve sorted windows, take a later full assignment into merged_sd for every scope, and assemble to the same sd as merged_sd. Both backends run through the same parametrized class, plus a pandas test that pins how an object column is typed from its dtype. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
ServerDataflow and PolarsServerDataflow now implement the _get_schema_sd hook,
so stats_tier="schema" builds a dataflow that types every column from its dtype
and runs no stat on the data. schema_sd (stat_pipeline.py) builds the sd the way
process_df shapes it: an empty frame gives {}, a skipped column keeps only its
names. pandas applies the existing typing_stats to a zero-row slice and derives
_type through the _type stat; polars factors pl_dtype_typing out of
pl_typing_stats and feeds it the dtype.
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
…tats_request (rows-first s3)
On a deferred /load_expr session a stats_request {stats_gen, scope} should
return a stats_update with the matching stats_gen whose inline wide payload
equals the all_stats an inline session sends, a stale stats_gen should get
stats_aborted and run no query, and /load_expr and /reload_expr should bump the
generation. df_meta.stats should be injected on every frame and survive a
dataflow-field change, which returns the session to the schema tier. With a
caps client and a legacy client on one session, the legacy client should keep
getting complete messages through the websocket broadcast, the /load_expr,
/reload_expr, /load and /load_compare pushes and the highlight overlay, while
the caps client gets a stats-free frame and then pulls a stats_update. A spy
telemetry sink should see a stats.request span.
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
…eration edge cases (rows-first s3) Four more cases for the stats wire format, kept in their own commit so each is seen failing on CI before the implementation lands. Returning to a state whose stats were completed once is answered from summary_stats_cache with no query. Completing the stats keeps the session's component_config on the refreshed display config. A warm /load_expr, which rebuilds nothing, leaves stats_gen alone. A stats_request on a session with no data is answered with stats_aborted. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
…ed sessions (rows-first s3)
A client that advertises ?caps=stats_update gets a stats-free initial_state on a
deferred /load_expr session (df_meta.stats.status "pending") and pulls the stats
with stats_request {stats_gen, scope}. The reply is a stats_update carrying the
dataflow's all_stats as an inline wide envelope, or stats_aborted when the
generation is stale. The request is the whole run: one synchronous call that
computes the full stats, writes the full-tier summary_stats_cache entry, assigns
summary_sd and refreshes the session snapshot through one helper.
stats_gen is a server-owned counter bumped by every load handler and by a state
change that touches a dataflow field, which also returns a deferred session to
the schema tier. df_meta.stats is injected by build_state_message from the
session, since the dataflow rebuilds df_meta wholesale. Every send site goes
through build_state_message_for, so a client without the capability still gets
complete messages (its missing stats run synchronously first); broadcast_state
replaces the five copies of the send loop and sends to capable clients first.
The stats_request branch binds the session's telemetry sink and emits a
stats.request span.
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
…s (rows-first s4) Both pass on the current code. process_df is written out as the process_column loop it is, and each xorq histogram query's snapshot-cache key is compared with the key of a query built in the test, so the refactor into resumable units that follows can be checked against something that does not go through the units. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
…tore (rows-first s4) Each stats class gets plan(state) and run(unit, acc): the fragments of the planned units, assembled through assemble_merged_sd, must equal the full-stats merged_sd on pandas, polars and xorq with init_sd, cleaning and overrides; a column group returns only its columns; skip_stat_columns columns get no unit; the xorq batch can be split by column chunk behind a default-off flag and is refused for anything but a plain parquet scan; the histogram cache keys stay put. A StatRun held on the session, keyed by (stats_gen, scope), keeps an append-only fragment list and accumulator, runs a unit only when asked, is dropped when stats_gen changes, and each connection keeps its own cursor. test_the_batch_is_one_query_by_default passes on the current code: it pins the default the flag must not change. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
📦 TestPyPI package publishedpip install --index-strategy unsafe-best-match --index-url https://test.pypi.org/simple/ --extra-index-url https://pypi.org/simple/ buckaroo==0.15.9.dev37186186027or with uv: uv pip install --index-strategy unsafe-best-match --index-url https://test.pypi.org/simple/ --extra-index-url https://pypi.org/simple/ buckaroo==0.15.9.dev37186186027MCP server for Claude Codeclaude mcp add buckaroo-table -- uvx --from "buckaroo[mcp]==0.15.9.dev37186186027" --index-strategy unsafe-best-match --index-url https://test.pypi.org/simple/ --extra-index-url https://pypi.org/simple/ buckaroo-table📖 Docs preview🎨 Storybook preview |
StatPipeline and XorqStatPipeline get plan(state), new_accumulator(state) and run(unit, acc); process_df and process_table are those three in order, so the inline path and a caller running units one at a time compute the same stats. pandas and polars plan one unit per non-skipped column. xorq plans the scalar batch (which also yields histogram_bins), then one histogram query per column, visible columns first. The batch can be cut into column chunks sized by cells when a host sets stat_chunk_cells, and only for a plain parquet scan: any other source falls back to the single batch and logs why. A StatRun, kept on the session by (stats_gen, scope), holds the planned units, an append-only fragment list and the accumulator. It has no thread, timer or callback, begin_stats_generation drops it, and each WebSocket connection keeps its own StatCursor into the list. CustomizableDataflow.build_stats is the one place a stats class is built, with run=False for plan/run callers. Two of the new tests are corrected here: the default-path test now ignores the PERVERSE_DF self-check units, and the xorq comparisons round floats to nine digits because a parallel aggregate adds its partial sums in varying order. The new-module imports in the tests move to module level. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
… once (rows-first s4)
A column group, the priority columns and StatRun's prefer hint accept
original or rewritten (a, b, c) names and read each name in both. When an
original name equals another column's rewritten name, a group returns the
columns of both, priority ranks both first, and prefer picks the unit of
whichever comes first in the frame.
On a frame whose columns are c, b, a (rewritten a, b, c):
- StatState(df, columns=('c',)) plans the units of c and a.
- priority=('a',) plans c before a.
- StatRun.next_unit(prefer=('a',)) picks the unit of c.
The new tests also pin the namespace argument the fix adds, so a client that
holds the rewritten names can ask for exactly those: StatState(namespace=)
and StatRun.next_unit/run_next(namespace=), with "any" (the default),
"original" and "rewritten".
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
…rst s4) A column group, StatState.priority and StatRun's prefer hint matched a name against both the original and the rewritten (a, b, c) column names. When an original name equals another column's rewritten name, one name picked two columns: a group returned columns that were not asked for, and prefer ran the unit of whichever column came first in the frame. resolve_names reads each name once. In "any", the default and what the code did for names that do not collide, a name that is an original column name picks that column and only a name that is not one is read as a rewritten name. StatState(namespace=) and StatRun.next_unit/run_next(namespace=) also take "original" and "rewritten", so a caller that holds the client's rewritten names says so and gets exactly those columns. Priority and prefer names are resolved against every column of the frame, not only the columns in the group, so a name picks the same column whatever group is asked for. prioritized now takes the state, since it needs the frame's columns and the namespace to read the priority names. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
|
Review round on this PR: one finding, fixed in Finding (medium): a column name was read as an original name and as a rewritten name at once. A column group, How it was addressed.
I chose an explicit namespace over a fixed rule because neither order of precedence is right for both callers: a client's Not changed. Verification. After the fix, the reviewer's two scripts pick the single column asked for ( |
…d the compare reset under a stats policy (rows-first p33) Cases found untested after the first tests commit: a scalar tier named with inline delivery builds a schema dataflow, and /load_compare clears the policy a session held. Rebased off the closed unit PRs (#1026, #1028): the assertions on unit tiers and on df_display_args in the final reply are dropped with the code they tested. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
…irst p33) /load_expr and /reload_expr accept stats_tier auto, full, scalar or schema, stored with the pair and kept out of has_config, so the warm short-circuit holds. When the dataflow is built at the schema tier the handler resolves the policy (resolve_stats_policy) from the count it already has, stores it on the session and starts the stats generation from it: a target below full is not_computed with the policy's reason. df_meta.stats reports the policy as tier_target and estimate, plus auto_request, requestable, omitted_keys, approx_keys and demand_columns where they differ from their documented defaults. It is applied at WebSocket open only for a client that sends ?caps=stats_update,stats_ondemand. A client with stats_update only is told the session is pending and pulls the stats, and a client with no caps gets them at connect, as for any deferred session. A tier the host named (scalar, schema) reaches every client as before. A session on an explicit stats_tier full within the ceiling sends the message it always has. /load keeps resolving to full: it does not read the field and clears a policy left on the session. Rebased onto #1024 without the unit PRs (#1026, #1028): stats_request is the whole run of #1024, so handle_stats_request now takes the client to decide stats_to_pull, and the incremental and display-config parts are gone. Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
c500780 to
f61df3d
Compare
Stacked on #1024 (
feat/rowsfirst-s3-stats-wire-server), which is stacked on #1022 and #1021. This PR is based onmainso the repo's Checks workflow runs on it (itspull_requesttrigger usesbranches: "*", which does not match a base branch containing a slash). The diff therefore includes the commits of #1021, #1022 and #1024 until they merge. The commits of this phase aredde83c16(a characterization test that passes on the base),667ae34d(failing tests) and799583f7(implementation), then746ae8c5(failing tests) andfa05adba(fix) from a review round on column names; read only those five.Problem
The stats classes compute everything in one call:
DfStatsV2andPlDfStatsV2runStatPipeline.process_dfin their constructor, andXorqDfStatsV2runsXorqStatPipeline.process_table. #1024'sstats_requestis therefore a whole run, one synchronous call that blocks the loop for as long as the stats take. On xorq that is one batch aggregate over every column plus one histogram query per column, and the batch alone took 162.5 s on the 78M-row telemetry entry. To bound the length of a request, a run has to be cut into units that can be run one at a time. A session also has to keep the pieces that have finished, so that two clients, or a client that reconnects, do not recompute them or lose them.Phase and plan references
Rows-first s4, from
buckaroo2-reports/plans/: plan 1 (01-rows-first-stats-separate-plumbing.md) section 9 "Phase 3" and plan 2 (02-rows-first-xorq-and-lazy-polars.md) section 7 "Phase 2", with plan 2 sections 4.1 (units, run object) and 4.2 (the xorq batch split), and gate 2 ofmeasurements-phase0.md.Approach
pluggable_analysis_framework/stat_units.py, definesStatUnit(id, columns, phase, prerequisite unit ids, cost class),StatState(the frame or expression, the columns that get no unit, an optional column group, columns to plan first, an optional row count, and the namespace the group and priority names are written in), the fragment type{orig_col: {stat: value}},StatAccumulator, and three helpers:merge_fragments,resolve_names, which turns requested column names into original names (see "Column names in a request"), andrewrite_sd, which re-keys an original-name dict to the rewrittena, b, cnames the wayXorqDataflow._get_summary_sddid inline.StatPipelineandXorqStatPipelineeach gainplan(state),new_accumulator(state)andrun(unit, acc), and so do the three stats classes (DfStatsV2,PlDfStatsV2,XorqDfStatsV2), which delegate to their pipeline.planandnew_accumulatorcompute nothing.process_dfandprocess_tableare these three called in order, so the inline path and a caller running units one at a time cannot differ.process_column.length,null_count,min,max,mean,std,approx_median,distinct_count, plus the stats that are pure functions of them, sohistogram_binscomes with it andcolor_mapneeds no histogram query). Then one histogram GROUP BY per column, with the columns the caller names first. A stat that needs anXorqExprorXorqExecute, or reads one that does, belongs to its column's histogram unit, so a user-defined stat that queries the table lands there without a special case. The two phases of the old_process_table_implare the two kinds of unit. The histogram queries are built and executed as before, so their snapshot-cache keys are unchanged.XorqStatPipeline(..., chunk_cells=N)cuts the batch into one aggregate per chunk of columns, sized by cells (rows x columns, with at least one column per chunk). It is off unless a host setsstat_chunk_cellsonXorqServerDataflow, and it needs the row count (StatState.rows; the dataflow passes its cached_expr_count, and without a count the batch stays whole). Phase 0 gate 2 found the split by column cheap on a parquet scan (the chunks sum to 0.90-1.16x the single batch at 10-12M rows, peak RSS halves, the longest unit falls from 1.1-1.4 s to 0.20-0.26 s for 6-column chunks), not worth it by stat class (1.3-1.4x, every piece rescans) and costly on a CSV scan (2.1-2.2x, every chunk reparses the file). So only the column split exists.17f031e9(fix(stats): materialize the stat source only for chained sources (#915) #916) gated source materialization on the op tree, andfde5213b(refactor(stats): remove buckaroo-managed source materialization #929) deleted the materialization with it, because the worthiness heuristic and the machinery around it "made the system hard to reason about". That classifier chose behaviour for every source. This one only vetoes a host's request. The host opts in withstat_chunk_cells, which is its statement that it wants chunks for this source.XorqStatPipeline.chunk_refusal(table)then refuses the split unless the expression is aReadof a parquet file, or a projection of plain columns of one. Anything else (a join, an aggregate, a diff, a filter, a derived column, a cache node, a CSV read, a memtable, and a table registered on a backend, whose format is not recorded) keeps the single batch and logs why at info level. The guard cannot turn chunking on, it has one allow-list entry, and an expression it does not recognise is refused. If the host's declaration and the guard disagree, the guard wins. Nothing classifies an expression to decide what a default does.StatRun.buckaroo/server/stat_run.pyholds the run a session keeps for one(stats_gen, scope): the planned units, an append-only list of the fragments in the order they finished, the accumulator the units read, and a status (pending,complete,error).run_next(prefer=(), namespace="any")runs one unit, the first whose prerequisites have run, or the first of those that covers a column named inprefer, and returns its fragment. It has no thread, timer or callback; nothing runs unless a caller asks. A unit that raises fails the run (status == "error", the exception propagates once, and the run is not retried); a stat that fails inside a unit is not that, it is an entry inacc.errorsand the unit finishes.StatCursoris one client's position in the list:take(run)returns the fragments it has not seen and advances, and a cursor handed a different run starts again at the beginning.DataStreamHandlercreates one per connection (stats_cursor), so the work is done once and each client reads the whole list at its own pace.SessionState.stat_runsmaps(stats_gen, scope)to its run.begin_stats_generationclears it, so a run lives until the generation changes.start_stat_run(session, scope="raw")instats_wire.pyreturns the run for the current generation, creating it from the dataflow if absent (planning only, no query). It analyzes the frameassign_full_statsdoes. Nothing in the server calls it yet.CustomizableDataflow.build_stats(processed_df, run=True)builds the stats class the way_get_summary_sdalways did, and_get_summary_sdnow calls it;run=Falsebuilds it without running, forplan/run.XorqDataflowoverrides it to passcache_storage,chunk_cellsandrows. The three stats classes takerun=True, andbuild_statspassesrunonly when it isFalse, so aDFStatsClasswritten before the keyword existed still builds.Column names in a request
A column group,
StatState.priorityandStatRun'spreferhint are written in one of two kinds of name: the original column names, or the rewrittena, b, cnames a client holds. The first version of this PR read each name as both. That breaks when an original name equals another column's rewritten name. A frame with columnsc, b, arewrites them toa, b, c, socis the original name of the first column and the rewritten name of the last. A group of('c',)planned the units of both, andprefer=('a',)ran the unit of the first column. A review found this; no server path calls these yet, but thecolumnshint of the next phase'sstats_requestwill.resolve_names(pairs, names, namespace)now reads each name once, in the namespace the caller states:rewritten: every name is a rewritten name. This is what the request branch passes for a client's hint.original: every name is an original name.any(the default): a name that is an original column name picks that column, and only a name that is not one is read as a rewritten name. Names that do not collide pick what they picked before, so the existing tests are unchanged. A name that does collide is read as the original, which is wrong for a client; that is why a client saysrewritten.StatState(namespace=...)sets it for the group and the priority names, andStatRun.next_unitandrun_nexttake it forprefer. An unknown namespace raisesValueError. Priority and prefer names are resolved against every column of the frame, not only the columns in the group, so a name picks the same column whatever group is asked for.prioritizedtakes the state instead of the priority tuple because of this.skip_columnsmatching is not part of this change; it keeps the matching each backend had on main.What changes
buckaroo/pluggable_analysis_framework/stat_units.py(new),stat_pipeline.py(plan,new_accumulator,run;process_dfover them),xorq_stat_pipeline.py(the same,XorqAccumulator,chunk_cells,chunk_refusal;process_tableover them),df_stats_v2.py(run=False,state,plan/run).buckaroo/dataflow/dataflow.py(build_stats, theDfStatsprotocol),buckaroo/xorq_buckaroo.py(build_stats,stat_chunk_cells,rewrite_sd),buckaroo/server/xorq_loading.py(thestat_chunk_cellskeyword).buckaroo/server/stat_run.py(new),session.py(stat_runs, cleared inbegin_stats_generation),stats_wire.py(start_stat_run),websocket_handler.py(stats_cursor).stat_units.py(resolve_names,StatState.namespace,prioritized(state, pairs)),stat_pipeline.pyandxorq_stat_pipeline.py(theprioritizedcall),stat_run.py(namespaceonnext_unitandrun_next).Tests
Five commits:
dde83c16(characterization tests that pass on the base),667ae34d(failing tests) and799583f7(implementation), then746ae8c5(failing tests for the column-name namespaces) andfa05adba(the fix).Characterization, in
test_paf_v2.pyandtest_xorq_stats_v2.py:process_dfequals theprocess_columnloop it replaced (pandas and polars, with and without a skipped column, plus the per-column error list), and each xorq histogram query keeps its snapshot-cache key, checked against keys of queries built in the test in the shapehistogramdocuments (the batch plus one per column, nothing else).Units, in the files that already test each backend:
test_paf_v2.py(pandas and polars pipelines: one unit per column, a fragment keyed by original name, the union of fragments equalsprocess_df, a column group by original or rewritten name, columns planned first, an empty frame, an error recorded on the accumulator),test_xorq_stats_v2.py(the batch first then one histogram unit per column, visible columns first,plansends no query, the batch is one query and each histogram unit one more, the batch fragment carrieshistogram_binsand nohistogram, the union equalsprocess_table, a failing stat is reported once and the run goes on, a column group),skip_columns_test.py(a skipped column gets no unit on the three backends and still has its entry).The column-chunk split, in
test_xorq_stats_v2.py(TestColumnChunkSplit): off by default, chunks sized by cells with at least one column, an unknown row count keeps one batch, batches precede histograms, visible columns fill the first chunk, the split equals the single batch (and queries are one per chunk), the histogram cache keys keep their keys and order, and nine kinds of source are refused (a join, an aggregate, a diff, a CSV read, a filter, a derived column, a cache node, a memtable and a backend table) while a plain parquet scan and a column projection of it are allowed. A refusal is logged and gives the single batch and the same stats.Through a dataflow, in
test_data_loading_polars.py(pandas and polars) andtest_load_expr.py(xorq):build_stats(run=False)computes nothing; the default run goes through the units (every stat query is issued inside a unit); the fragments of each scope's units, assembled byassemble_merged_sdwithinit_sd, a cleaning op, a search filter,column_config_overridesand a post-processing sd override, equal the dataflow'smerged_sd; a column group returns only those columns;skip_stat_columnscolumns get no unit;stat_chunk_cellssplits a parquet-backed dataflow's batch and a joined one stays whole. The helpers shared by those two files sit beside_scope_inputsinscoped_summary_stats_test.py.StatRunand the session:test_data_loading_polars.py(key, unit plan, fragments appended in completion order and never changed,prefer, completion, the accumulator equals the merged fragments,raw_sdequals the full-statssummary_sd, cursors read one list each at their own pace, a cursor starts again on another run, no thread, timer orIOLoopcallback is created, a raising unit fails the run) andTestStatsWireintest_load_expr.py(a run started from a deferred session is keyed by the generation and sends no query, running it ends at thesummary_sdan inline session computes, a dataflow-field change and/reload_exprdrop it while a typed search term does not, each connection has its own cursor).Column-name namespaces (
746ae8c5), in the files that already test each layer.test_paf_v2.pyuses a frame with columnsc, b, a: a group or a priority name that is an original name and also another column's rewritten name picks the original column only,rewrittenandoriginalread every name in their namespace, a rewritten name that is no original name still picks its column, and an unknown namespace raises.test_xorq_stats_v2.pychecks the same on the xorq plan (the batch columns and one histogram unit per column picked).test_data_loading_polars.pychecksnext_unitandrun_nexton pandas and polars. Of the 32 new tests, 27 fail on746ae8c5alone (the wrong columns picked, or the missingnamespaceargument) and 5 pin names that do not collide and pass on both commits. All ninePython / Testjobs failed on746ae8c5.On
667ae34d(tests only) all ninePython / Testjobs failed, including 3.14 and Windows, where the xorq test file is skipped but the pandas and polars tests are not; 86 of the new tests fail locally, each on a missing module or attribute (stat_units,stat_run,plan,run,build_stats,chunk_refusal,chunk_cells,start_stat_run,stat_chunk_cells). CI logs were not read (the REST log API is rate-limited for this account), so the reason rests on the local run of the same commit.test_the_batch_is_one_query_by_defaultpasses on that commit by design: it pins the default the flag must not change.CI on
799583f7: all 28 checks completed, 27 succeeded (one of them the Read the Docs status) anddeploywas skipped. That includes all ninePython / Testjobs (3.11 to 3.14, Max Versions 3.11 to 3.14, Windows),Python / Lint,Python / Typecheck, the JS job, the wheel build and the Playwright jobs.CI on
fa05adba: all 28 checks completed, 27 succeeded (one of them the Read the Docs status) anddeploywas skipped, including all ninePython / Testjobs.Locally the full unit suite (
pytest ./tests/unit -m "not slow") gives 1361 passed and 5 skipped, which is the 1268 of #1024 plus 6 characterization tests and 87 new ones. In a Max Versions environment (pandas 3.0.6, polars 1.44.2, xorq 0.4.5, numpy 2.5.3, resolved as CI does) it gives the same 1361 passed.basedpyrighton the scoped files reports 0 errors and 0 warnings. Afterfa05adbathe unit suite gives 1393 passed and 5 skipped locally, the 1361 plus the 32 namespace tests.Eight deliberate regressions of the implementation each make at least one test fail: the batch no longer first, a chunk guard that always allows,
begin_stats_generationnot dropping runs, a cursor that does not advance, a skipped column's length not set, fragments that alias the accumulator's dicts,priorityignored, andprocess_dfbypassing the units. A ninth,StatRun.next_unitignoring a unit's prerequisites, survives: no test asks for a unit whose prerequisite has not run (theprefertests run pandas units, which have none). I did not add one after the implementation because it would not have been seen failing on CI first; the request branch of the next phase, which drives a xorq run with a column hint, is where it belongs.Measurements
The longest single unit, each unit timed alone through
plan/runas a request would run it, three repeats, Apple M4 Pro with other agents running (load average about 4). The fixture is synthetic: 27 columns (4 ints, 6 floats with nulls, 2 low-cardinality ints, 8 strings of 5 to 5,000 distinct values, 3 booleans, 4 timestamps). xorq reads a parquet file. pandas and polars analyze the 50,000-row sample their stats classes take above 1M cells, so their units do not grow with the file.column:float3)chunk_cells=1M(18 columns per chunk)chunk_cells=60Mchunk_cells=12M(6 columns per chunk)batch:0)On xorq the longest histogram unit is 0.038-0.039 s on the medium fixture (one outlier run at 0.093 s), so with the batch cut into 6-column chunks the longest unit is a batch chunk at about 0.24 s, and the chunks add up to what the single batch costs (1.61 s against 1.62 s for the whole run). That matches gate 2: 0.20-0.26 s for 6-column chunks at 10-12M rows. Against the 250 ms budget the plans propose, a 6-column chunk fits at 2M rows, just under it; the plans' 60M cells per unit is a size where this fixture needs no split. These are the figures for parquet; a CSV-backed source is refused by the guard and keeps its single batch, whose cost is the 54-215 s the telemetry shows, which is plan 3's tier policy and not this phase's.
Why default behaviour is unchanged
stat_chunk_cellsandchunk_cellsdefault to off;run=Trueis the default of the stats classes;start_stat_run,StatRunandstats_cursorare not used by any request path yet, sostats_requestandcomplete_statsbehave as in feat(server): stats_request, stats_update and df_meta.stats on deferred sessions (rows-first s3) #1024, through the same units.DfStatsV2,PlDfStatsV2andXorqDfStatsV2(on a parquet scan and on a memtable, each with and without a skipped column, on a mixed-dtype fixture with nulls, strings, booleans, timestamps and a categorical) and the xorq snapshot-cache keys of a run, between9bfebfb9and this branch, with a script outside the repo. All were identical, modulo the order of tied histogram entries and polars'modeover tied counts, which differ between two runs of the same code (two runs of the base disagree there, so the script sorts histograms and dropsmode). The characterization tests above pin the same on the base.process_dfkeeps its early return for an empty frame and its perf summary;process_tablekeeps thestat.xorq.totalandstat.xorq.batch_aggregatespans, the cache counters and thefirstpull.summary_statsspan on the dataflow, which the existing telemetry tests cover.test_data_loading_polars.pygrew (_scope_inputsplus the new shared helpers); everything else in the test files is added.Deviations from the plan
process_dfstill returns its summary keyed by rewritten name andprocess_tableby original name, as they always did.rewrite_sdconverts a run's results to the rewritten keys the dataflow carries, andXorqDataflow._get_summary_sdnow uses it instead of its inline loop.plan(state)takes aStatState, which the plan does not define. It carries what the plan lists as the unit's inputs, plus the skipped columns, an optional column group, columns to plan first and the row count. A group, the priority columns andpreferare written in a namespace (any,originalorrewritten, defaultany; see "Column names in a request"), because a client only knows the rewritten names;skip_columnskeeps the matching each backend always had (original or rewritten on pandas and polars, original on xorq).acc.sd(), with the structural keys the pipeline always gave it (orig_col_nameandrewritten_col_name, and on xorqdtype, the batch'slengthand theNoneplaceholders), becauseprocess_dfandprocess_tablealways returned one. It is in no fragment.StatState.priorityfor the plan andStatRun.next_unit(prefer)for a run, since the plan leaves the source of "visible" to the request.histogram_bins,non_null_count,nan_per,distinct_per), so the batch fragment is complete for everything that needs no histogram query, as plan 2 section 4.2 describes. Withchunk_cellsevery chunk precedes every histogram unit, as the plan's unit order says, and each histogram unit names its chunk as its prerequisite.stat_chunk_cellsonXorqServerDataflow(andchunk_cellson the pipeline and the stats class), a cell count and not a boolean, because chunks are sized by cells. It is not a/load_exprbody field yet.StatRunhas the plan's unit list, fragment list, accumulator, creation time and status, but not the cache and timing counters of plan 2 section 4.5, which belong with the request branch's spans.build_statsand theDfStatsprotocol gainrunandplan/new_accumulator/run(the plan says "each stats class gains two methods";new_accumulatoris the third, becauserun(unit, acc)needs somewhere to start).Not in this PR
remainingfield ofstats_updateandfinalper client:stats_requeststill runs the whole run, through the same units.stats.unit) and the run's cache and timing counters./load_exprbody field forstat_chunk_cells, and any default for it.StatRun.next_unit's prerequisite filter (see Tests).🤖 Generated with Claude Code