Skip to content

feat(stats): resumable stat units and the StatRun store (rows-first s4) - #1026

Draft
paddymul wants to merge 16 commits into
adr-002-rows-first-stats-deliveryfrom
feat/rowsfirst-s4-stat-units
Draft

paddymul wants to merge 16 commits into
adr-002-rows-first-stats-deliveryfrom
feat/rowsfirst-s4-stat-units

Conversation

@paddymul

@paddymul paddymul commented Oct 4, 2026 •

Copy link
Copy Markdown
Collaborator

Stacked on #1024 (feat/rowsfirst-s3-stats-wire-server), which is stacked on #1022 and #1021. This PR is based on main so the repo's Checks workflow runs on it (its pull_request trigger uses branches: "*", 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 are dde83c16 (a characterization test that passes on the base), 667ae34d (failing tests) and 799583f7 (implementation), then 746ae8c5 (failing tests) and fa05adba (fix) from a review round on column names; read only those five.

Problem

The stats classes compute everything in one call: DfStatsV2 and PlDfStatsV2 run StatPipeline.process_df in their constructor, and XorqDfStatsV2 runs XorqStatPipeline.process_table. #1024's stats_request is 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 of measurements-phase0.md.

Approach

  • Units. A new module, pluggable_analysis_framework/stat_units.py, defines StatUnit (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"), and rewrite_sd, which re-keys an original-name dict to the rewritten a, b, c names the way XorqDataflow._get_summary_sd did inline. StatPipeline and XorqStatPipeline each gain plan(state), new_accumulator(state) and run(unit, acc), and so do the three stats classes (DfStatsV2, PlDfStatsV2, XorqDfStatsV2), which delegate to their pipeline. plan and new_accumulator compute nothing. process_df and process_table are these three called in order, so the inline path and a caller running units one at a time cannot differ.
  • pandas and polars. One unit per column that is not skipped, in column order, run through the existing process_column.
  • xorq. The scalar batch is the first unit (length, null_count, min, max, mean, std, approx_median, distinct_count, plus the stats that are pure functions of them, so histogram_bins comes with it and color_map needs no histogram query). Then one histogram GROUP BY per column, with the columns the caller names first. A stat that needs an XorqExpr or XorqExecute, 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_impl are the two kinds of unit. The histogram queries are built and executed as before, so their snapshot-cache keys are unchanged.
  • Column-chunk split of the batch. 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 sets stat_chunk_cells on XorqServerDataflow, 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.
  • What decides that a source may be split. The maintainer removed an op-tree scan classifier on purpose: 17f031e9 (fix(stats): materialize the stat source only for chained sources (#915) #916) gated source materialization on the op tree, and fde5213b (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 with stat_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 a Read of 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.py holds 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 in prefer, 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 in acc.errors and the unit finishes. StatCursor is 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. DataStreamHandler creates one per connection (stats_cursor), so the work is done once and each client reads the whole list at its own pace.
  • On the session. SessionState.stat_runs maps (stats_gen, scope) to its run. begin_stats_generation clears it, so a run lives until the generation changes. start_stat_run(session, scope="raw") in stats_wire.py returns the run for the current generation, creating it from the dataflow if absent (planning only, no query). It analyzes the frame assign_full_stats does. Nothing in the server calls it yet.
  • Dataflow. CustomizableDataflow.build_stats(processed_df, run=True) builds the stats class the way _get_summary_sd always did, and _get_summary_sd now calls it; run=False builds it without running, for plan/run. XorqDataflow overrides it to pass cache_storage, chunk_cells and rows. The three stats classes take run=True, and build_stats passes run only when it is False, so a DFStatsClass written before the keyword existed still builds.

Column names in a request

A column group, StatState.priority and StatRun's prefer hint are written in one of two kinds of name: the original column names, or the rewritten a, b, c names 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 columns c, b, a rewrites them to a, b, c, so c is the original name of the first column and the rewritten name of the last. A group of ('c',) planned the units of both, and prefer=('a',) ran the unit of the first column. A review found this; no server path calls these yet, but the columns hint of the next phase's stats_request will.

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 says rewritten.

StatState(namespace=...) sets it for the group and the priority names, and StatRun.next_unit and run_next take it for prefer. An unknown namespace raises ValueError. 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 takes the state instead of the priority tuple because of this. skip_columns matching 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_df over them), xorq_stat_pipeline.py (the same, XorqAccumulator, chunk_cells, chunk_refusal; process_table over them), df_stats_v2.py (run=False, state, plan/run).
  • buckaroo/dataflow/dataflow.py (build_stats, the DfStats protocol), buckaroo/xorq_buckaroo.py (build_stats, stat_chunk_cells, rewrite_sd), buckaroo/server/xorq_loading.py (the stat_chunk_cells keyword).
  • buckaroo/server/stat_run.py (new), session.py (stat_runs, cleared in begin_stats_generation), stats_wire.py (start_stat_run), websocket_handler.py (stats_cursor).
  • Review round on column names: stat_units.py (resolve_names, StatState.namespace, prioritized(state, pairs)), stat_pipeline.py and xorq_stat_pipeline.py (the prioritized call), stat_run.py (namespace on next_unit and run_next).

Tests

Five commits: dde83c16 (characterization tests that pass on the base), 667ae34d (failing tests) and 799583f7 (implementation), then 746ae8c5 (failing tests for the column-name namespaces) and fa05adba (the fix).

  • Characterization, in test_paf_v2.py and test_xorq_stats_v2.py: process_df equals the process_column loop 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 shape histogram documents (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 equals process_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, plan sends no query, the batch is one query and each histogram unit one more, the batch fragment carries histogram_bins and no histogram, the union equals process_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) and test_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 by assemble_merged_sd with init_sd, a cleaning op, a search filter, column_config_overrides and a post-processing sd override, equal the dataflow's merged_sd; a column group returns only those columns; skip_stat_columns columns get no unit; stat_chunk_cells splits a parquet-backed dataflow's batch and a joined one stays whole. The helpers shared by those two files sit beside _scope_inputs in scoped_summary_stats_test.py.

  • StatRun and 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_sd equals the full-stats summary_sd, cursors read one list each at their own pace, a cursor starts again on another run, no thread, timer or IOLoop callback is created, a raising unit fails the run) and TestStatsWire in test_load_expr.py (a run started from a deferred session is keyed by the generation and sends no query, running it ends at the summary_sd an inline session computes, a dataflow-field change and /reload_expr drop 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.py uses a frame with columns c, 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, rewritten and original read 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.py checks the same on the xorq plan (the batch columns and one histogram unit per column picked). test_data_loading_polars.py checks next_unit and run_next on pandas and polars. Of the 32 new tests, 27 fail on 746ae8c5 alone (the wrong columns picked, or the missing namespace argument) and 5 pin names that do not collide and pass on both commits. All nine Python / Test jobs failed on 746ae8c5.

On 667ae34d (tests only) all nine Python / Test jobs 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_default passes 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) and deploy was skipped. That includes all nine Python / Test jobs (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) and deploy was skipped, including all nine Python / Test jobs.

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. basedpyright on the scoped files reports 0 errors and 0 warnings. After fa05adba the 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_generation not dropping runs, a cursor that does not advance, a skipped column's length not set, fragments that alias the accumulator's dicts, priority ignored, and process_df bypassing the units. A ninth, StatRun.next_unit ignoring a unit's prerequisites, survives: no test asks for a unit whose prerequisite has not run (the prefer tests 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/run as 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.

Fixture Backend Units Longest unit (3 runs) Whole run, median
small, 52,814 rows (1.4M cells) pandas 27 0.010-0.011 s (column:float3) 0.145 s
polars 27 0.004-0.005 s 0.074 s
xorq, one batch 28 0.061-0.069 s (the batch) 0.325 s
xorq, chunk_cells=1M (18 columns per chunk) 29 0.048-0.050 s 0.323 s
medium, 2,000,000 rows (54M cells) pandas 27 0.011 s 0.175 s
polars 27 0.005 s 0.078 s
xorq, one batch 28 0.858-0.899 s (the batch) 1.624 s
xorq, chunk_cells=60M 28 0.848-0.866 s (still one batch: 27 columns x 2M rows is under 60M cells) 1.618 s
xorq, chunk_cells=12M (6 columns per chunk) 32 0.243-0.245 s (batch:0) 1.610 s

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

  • Everything new is opt-in or not called. stat_chunk_cells and chunk_cells default to off; run=True is the default of the stats classes; start_stat_run, StatRun and stats_cursor are not used by any request path yet, so stats_request and complete_stats behave as in feat(server): stats_request, stats_update and df_meta.stats on deferred sessions (rows-first s3) #1024, through the same units.
  • The default path computes what it did. I compared the summary dicts of DfStatsV2, PlDfStatsV2 and XorqDfStatsV2 (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, between 9bfebfb9 and this branch, with a script outside the repo. All were identical, modulo the order of tied histogram entries and polars' mode over tied counts, which differ between two runs of the same code (two runs of the base disagree there, so the script sorts histograms and drops mode). The characterization tests above pin the same on the base.
  • process_df keeps its early return for an empty frame and its perf summary; process_table keeps the stat.xorq.total and stat.xorq.batch_aggregate spans, the cache counters and the firstpull.summary_stats span on the dataflow, which the existing telemetry tests cover.
  • No existing test changed. One import line in test_data_loading_polars.py grew (_scope_inputs plus the new shared helpers); everything else in the test files is added.

Deviations from the plan

  • Fragments are keyed by original column name on all three backends, as the plan writes them, but process_df still returns its summary keyed by rewritten name and process_table by original name, as they always did. rewrite_sd converts a run's results to the rewritten keys the dataflow carries, and XorqDataflow._get_summary_sd now uses it instead of its inline loop.
  • plan(state) takes a StatState, 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 and prefer are written in a namespace (any, original or rewritten, default any; see "Column names in a request"), because a client only knows the rewritten names; skip_columns keeps the matching each backend always had (original or rewritten on pandas and polars, original on xorq).
  • A skipped column gets no unit, but it keeps its entry in acc.sd(), with the structural keys the pipeline always gave it (orig_col_name and rewritten_col_name, and on xorq dtype, the batch's length and the None placeholders), because process_df and process_table always returned one. It is in no fragment.
  • Visible-columns-first is StatState.priority for the plan and StatRun.next_unit(prefer) for a run, since the plan leaves the source of "visible" to the request.
  • The xorq batch also computes the stats that are pure functions of its results (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. With chunk_cells every chunk precedes every histogram unit, as the plan's unit order says, and each histogram unit names its chunk as its prerequisite.
  • The flag is stat_chunk_cells on XorqServerDataflow (and chunk_cells on the pipeline and the stats class), a cell count and not a boolean, because chunks are sized by cells. It is not a /load_expr body field yet.
  • The guard is an allow-list of one source kind, where the plan's wording ("scan-backed sources") invites a classifier. See "What decides that a source may be split".
  • StatRun has 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_stats and the DfStats protocol gain run and plan/new_accumulator/run (the plan says "each stats class gains two methods"; new_accumulator is the third, because run(unit, acc) needs somewhere to start).

Not in this PR

  • The request branch that runs units under a time budget, the remaining field of stats_update and final per client: stats_request still runs the whole run, through the same units.
  • Final assignment and snapshot refresh generalized to a run that finished in several requests, and the config upgrade.
  • Per-unit telemetry spans (stats.unit) and the run's cache and timing counters.
  • The client, the count memoization and the scalar-tier unit sets of plan 3 phase 4.
  • A /load_expr body field for stat_chunk_cells, and any default for it.
  • A test for StatRun.next_unit's prerequisite filter (see Tests).

🤖 Generated with Claude Code

paddymul and others added 13 commits October 3, 2026 22:51
…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>
@github-actions

github-actions Bot commented Oct 4, 2026 •

Copy link
Copy Markdown
Contributor

📦 TestPyPI package published

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.dev37186186027

or 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.dev37186186027

MCP server for Claude Code

claude 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>
@paddymul

paddymul commented Oct 4, 2026

Copy link
Copy Markdown
Collaborator Author

Review round on this PR: one finding, fixed in fa05adba with its tests in 746ae8c5.

Finding (medium): a column name was read as an original name and as a rewritten name at once. A column group, StatState.priority and StatRun's prefer hint matched each 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. On a frame with columns c, b, a (rewritten a, b, c), StatState(df, columns=('c',)) planned the units of c and a, priority=('a',) planned c before a, and StatRun.next_unit(prefer=('a',)) returned the unit of c. A frame b, a, c with columns=('a',) returned both b and a. The values in the fragments were right; the selection and the order were wrong. I reproduced it with the reviewer's prefer script and its edge script before changing anything. The PR body described accepting both kinds of name but not this collision, and it now has a section on it. No request path calls these yet, but the columns hint of the next phase's stats_request will.

How it was addressed.

  • 746ae8c5 adds the tests alone: 32 new tests in test_paf_v2.py, test_xorq_stats_v2.py and test_data_loading_polars.py. Locally 27 of them failed on the tests-only commit (the wrong columns picked, or the new namespace argument missing) and 5 pass because they pin names that do not collide. All nine Python / Test jobs failed on that commit before the fix was pushed.
  • fa05adba adds resolve_names, which reads each name once, in a namespace the caller states: rewritten (every name is a rewritten name, which is what a client's hint is), original, or any. In 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, so names that never collided pick what they picked before and no existing test changed. StatState(namespace=...) sets it for the group and priority names, and StatRun.next_unit and run_next take namespace for prefer. An unknown namespace raises ValueError. Priority and prefer names are resolved against every column of the frame rather than the group's columns, so a name picks the same column whatever group is asked for. prioritized now takes the state because of that.

I chose an explicit namespace over a fixed rule because neither order of precedence is right for both callers: a client's c means the rewritten name, a Python caller's c means the column called c. The request branch of the next phase should pass namespace="rewritten" for the client's hint.

Not changed. skip_columns still matches original or rewritten names on pandas and polars and original names on xorq. That is the matching those backends have on main, so it is outside this PR.

Verification. After the fix, the reviewer's two scripts pick the single column asked for (prefer picks a for rewritten a, a group of ('c',) plans only c, a group of ('a',) on b, a, c plans only a). The full unit suite gives 1393 passed and 5 skipped locally, the previous 1361 plus the 32 new tests, with ruff and paddy_format clean. On fa05adba all 28 checks completed: 27 succeeded (one is the Read the Docs status) and deploy was skipped, including all nine Python / Test jobs. Nothing was declined.

paddymul added a commit that referenced this pull request Oct 6, 2026
…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>
paddymul added a commit that referenced this pull request Oct 6, 2026
…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>
@paddymul
paddymul force-pushed the adr-002-rows-first-stats-delivery branch from c500780 to f61df3d Compare October 6, 2026 18:27

This branch was successfully deployed

1 active deployment
testpypi — fa05adba Deployed Oct 4, 2026 by paddymul via Publish to TestPyPI #1690
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant