{buckaroo_state.show_commands ? (
@@ -447,6 +464,7 @@ export function DFViewerInfiniteDS({
activeCol={activeCol}
setActiveCol={setActiveCol}
error_info={""}
+ stats_status={getStatsStatus(df_meta)}
/>
diff --git a/packages/buckaroo-js-core/src/components/DFViewerParts/DFViewerInfinite.tsx b/packages/buckaroo-js-core/src/components/DFViewerParts/DFViewerInfinite.tsx
index 7c8fa079e..4f4465ccb 100644
--- a/packages/buckaroo-js-core/src/components/DFViewerParts/DFViewerInfinite.tsx
+++ b/packages/buckaroo-js-core/src/components/DFViewerParts/DFViewerInfinite.tsx
@@ -7,7 +7,8 @@ import {
import * as _ from "lodash-es";
import { DFData, DFDataRow, DFViewerConfig, SDFT } from "./DFWhole";
-import { getCellRendererSelector, dfToAgrid, extractPinnedRows, extractSDFT } from "./gridUtils";
+import { getCellRendererSelector, dfToAgrid, extractPinnedRows, extractSDFT, getFieldVal } from "./gridUtils";
+import type { StatsStatus } from "../WidgetTypes";
import { AgGridReact } from "ag-grid-react"; // the AG Grid React Component
import {
@@ -21,6 +22,7 @@ import {
CellStyleModule,
ColumnAutoSizeModule,
PinnedRowModule,
+ RenderApiModule,
RowSelectionModule,
TooltipModule,
TextFilterModule,
@@ -45,6 +47,9 @@ ModuleRegistry.registerModules([
CellStyleModule,
ColumnAutoSizeModule,
PinnedRowModule,
+ // api.refreshCells lives here. Without it the call logs AG Grid error 200
+ // and does nothing.
+ RenderApiModule,
RowSelectionModule,
TooltipModule,
TextFilterModule,
@@ -150,6 +155,7 @@ export function DFViewerInfinite({
max_rows_in_configs,
view_name,
data_key,
+ stats_status,
}: {
data_wrapper: DatasourceOrRaw;
df_viewer_config: DFViewerConfig;
@@ -173,6 +179,10 @@ export function DFViewerInfinite({
// a rowId, even though their `index` values overlap (row 0 in main is a
// different record than row 0 in summary).
data_key?: string;
+ // df_meta.stats.status. While "pending" a pinned key with no value shows a
+ // placeholder row; when "not_computed" it is omitted. Undefined behaves
+ // as "complete".
+ stats_status?: StatsStatus;
}) {
/*
The idea is to do some pre-setup here for
@@ -233,6 +243,7 @@ export function DFViewerInfinite({
effectiveScheme={effectiveScheme}
view_name={view_name}
data_key={data_key}
+ stats_status={stats_status}
/>
)
@@ -250,6 +261,7 @@ export function DFViewerInfiniteInner({
effectiveScheme,
view_name,
data_key,
+ stats_status,
}: {
data_wrapper: DatasourceOrRaw;
df_viewer_config: DFViewerConfig;
@@ -266,6 +278,7 @@ export function DFViewerInfiniteInner({
effectiveScheme?: 'light' | 'dark';
view_name?: string;
data_key?: string;
+ stats_status?: StatsStatus;
}) {
/*
const lastProps = useRef(null);
@@ -344,8 +357,8 @@ export function DFViewerInfiniteInner({
// Always re-extract; upstream may mutate summary in-place without changing identity
// Memoize to ensure it updates when summary_stats_data changes
const topRowData = useMemo(
- () => extractPinnedRows(summary_stats_data, pinned_rows ? pinned_rows : []) as DFDataRow[],
- [summary_stats_data, pinned_rows]
+ () => extractPinnedRows(summary_stats_data, pinned_rows ? pinned_rows : [], stats_status) as DFDataRow[],
+ [summary_stats_data, pinned_rows, stats_status]
);
// Pinned rows are extracted and ready
@@ -442,6 +455,32 @@ export function DFViewerInfiniteInner({
// ignore until grid ready
}
}, [pinnedSig]);
+
+ // color_map reads histogram_bins from the grid context when a cell is
+ // painted, so cells that rendered before the bins arrived keep the
+ // neutral style. Repaint the color-mapped columns when their bins
+ // change after the first render. Bins that are already there at mount
+ // paint correctly, so the first run only records the signature.
+ const colorMapCols = useMemo(
+ () => df_viewer_config.column_config.flatMap((cc) =>
+ cc.color_map_config?.color_rule === "color_map"
+ ? [{ field: getFieldVal(cc), statsCol: cc.color_map_config.val_column }]
+ : []),
+ [df_viewer_config.column_config],
+ );
+ const colorMapSig = useMemo(() => {
+ if (colorMapCols.length === 0) return "";
+ const stats = extractSDFT(summary_stats_data);
+ return JSON.stringify(colorMapCols.map(
+ ({ field, statsCol }) => [field, statsCol === undefined ? undefined : stats[statsCol]?.histogram_bins]));
+ }, [colorMapCols, summary_stats_data]);
+ const colorMapSigRef = useRef(colorMapSig);
+ useEffect(() => {
+ if (colorMapSigRef.current === colorMapSig) return;
+ colorMapSigRef.current = colorMapSig;
+ if (colorMapCols.length === 0) return;
+ gridRef.current?.api?.refreshCells({ force: true, columns: colorMapCols.map((c) => c.field) });
+ }, [colorMapSig, colorMapCols]);
// Force update rowData when Raw data changes
const rawDataSig = useMemo(() => {
diff --git a/packages/buckaroo-js-core/src/components/DFViewerParts/SeriesSummaryTooltip.tsx b/packages/buckaroo-js-core/src/components/DFViewerParts/SeriesSummaryTooltip.tsx
index 907fef642..bfbfad3aa 100644
--- a/packages/buckaroo-js-core/src/components/DFViewerParts/SeriesSummaryTooltip.tsx
+++ b/packages/buckaroo-js-core/src/components/DFViewerParts/SeriesSummaryTooltip.tsx
@@ -24,11 +24,16 @@ export const getSimpleTooltip = (tooltipField:string) => {
// This should be possible with the tooltipValueGetter, but that
// wasn't working for some reason
- if (props.data.index === "histogram") {
+ if (props.data?.index === "histogram") {
return;
}
- const val = props.data[tooltipField].toString()
- return
{val}
;
+ // A pinned row whose stats have not arrived has no value for the
+ // column, and a row that has not loaded has no data: show nothing.
+ const raw = props.data?.[tooltipField];
+ if (raw === undefined || raw === null) {
+ return;
+ }
+ return
{raw.toString()}
;
};
return simpleTooltip;
}
diff --git a/packages/buckaroo-js-core/src/components/DFViewerParts/Styler.tsx b/packages/buckaroo-js-core/src/components/DFViewerParts/Styler.tsx
index ef83da6ae..fcdd0baff 100644
--- a/packages/buckaroo-js-core/src/components/DFViewerParts/Styler.tsx
+++ b/packages/buckaroo-js-core/src/components/DFViewerParts/Styler.tsx
@@ -64,13 +64,14 @@ export function colorMap(cmr: ColorMapRules) {
const summarys = params.context?.histogram_stats;
const statsCol = cmr.val_column; // || col_name;
+ // No bins is the normal state until the summary stats arrive, so
+ // fall back to the neutral style quietly; the grid repaints the
+ // column when they do.
if (statsCol === undefined || summarys === undefined){
- console.log("66 couldn't find stats_col")
return baseReturn;
}
const summary_stats_cell = summarys[statsCol];
if (summary_stats_cell === undefined || summary_stats_cell.histogram_bins === undefined ) {
- console.log("69 couldn't find summary_stats");
return baseReturn
}
const histogram_edges = summary_stats_cell.histogram_bins;
diff --git a/packages/buckaroo-js-core/src/components/DFViewerParts/gridUtils.test.ts b/packages/buckaroo-js-core/src/components/DFViewerParts/gridUtils.test.ts
index 01697ba2e..a32f4a7db 100644
--- a/packages/buckaroo-js-core/src/components/DFViewerParts/gridUtils.test.ts
+++ b/packages/buckaroo-js-core/src/components/DFViewerParts/gridUtils.test.ts
@@ -15,7 +15,8 @@ import {
import * as _ from "lodash-es";
import { DFData, DFViewerConfig, NormalColumnConfig, MultiIndexColumnConfig, PinnedRowConfig, ColumnConfig, FormatterArgs } from "./DFWhole";
import { getFormatter, getFloatFormatter, getCompactNumberFormatter, formatDuration, formatIsoDuration, getDurationFormatter } from './Displayer';
-import { ColDef, ICellRendererParams, ValueFormatterParams } from 'ag-grid-community';
+import { CellClassParams, ColDef, ICellRendererParams, ITooltipParams, ValueFormatterParams } from 'ag-grid-community';
+import { getSimpleTooltip } from './SeriesSummaryTooltip';
describe("testing utility functions in gridUtils ", () => {
// mostly sanity checks to help develop gridUtils
@@ -199,6 +200,18 @@ describe("testing utility functions in gridUtils ", () => {
]);
});
+ it("omits a required key with no value when the stats status is error, and keeps one that has a value", () => {
+ // An error is final for the state on screen, so no value is coming for
+ // the keys that are missing (rows-first c4).
+ const data: DFData = [{ index: "row1", value: 1 }];
+ const pinnedConfig: PinnedRowConfig[] = [
+ { primary_key_val: "row1", displayer_args: { displayer: "obj" } },
+ { primary_key_val: "missing", displayer_args: { displayer: "obj" } }
+ ];
+ expect(extractPinnedRows(data, pinnedConfig, "error")).toStrictEqual([{ index: "row1", value: 1 }]);
+ expect(extractPinnedRows([], pinnedConfig, "error")).toStrictEqual([]);
+ });
+
it("includes an optional `?`-prefixed pinned row when the unprefixed key exists in data", () => {
const data: DFData = [
{ index: "histogram_bins", value: 1 },
@@ -732,6 +745,70 @@ describe("testing multi index organiztion ", () => {
expect(children.length).toBe(2);
});
+});
+
+// Rows-first c0a: while summary stats are pending or not computed, pinned
+// cells have no value, and color_map has no histogram bins to read.
+describe("pinned cells without stats values (rows-first c0a)", () => {
+ it("the simple tooltip returns nothing for a valueless pinned cell instead of throwing", () => {
+ const tooltip = getSimpleTooltip("a");
+ // A pinned row with no value for the column: only the row label exists.
+ expect(() => tooltip({ data: { index: "dtype" } } as ITooltipParams)).not.toThrow();
+ expect(tooltip({ data: { index: "dtype" } } as ITooltipParams)).toBeUndefined();
+ // A null cell and a row that has not loaded (data undefined) are also valueless.
+ expect(() => tooltip({ data: { index: "mean", a: null } } as ITooltipParams)).not.toThrow();
+ expect(() => tooltip({ data: undefined } as unknown as ITooltipParams)).not.toThrow();
+ });
+
+ it("the simple tooltip still renders a cell that has a value", () => {
+ const tooltip = getSimpleTooltip("a");
+ const el = tooltip({ data: { index: 0, a: 5 } } as ITooltipParams);
+ expect(el).toBeDefined();
+ expect((el as any).props.children).toBe("5");
+ });
+
+ describe("color_map without histogram bins", () => {
+ const config: DFViewerConfig = {
+ pinned_rows: [],
+ left_col_configs: [],
+ column_config: [
+ {
+ col_name: "a",
+ header_name: "a",
+ displayer_args: { displayer: "obj" },
+ color_map_config: { color_rule: "color_map", map_name: "BLUE_TO_YELLOW", val_column: "a" },
+ },
+ ],
+ };
+ const cellStyleFor = (context: any, value: any = 3) => {
+ const colDef = dfToAgrid(config)[0] as ColDef;
+ const cellStyle = colDef.cellStyle as (p: CellClassParams) => Record;
+ return cellStyle({ context, data: { index: 0, a: value }, value, node: { rowPinned: undefined } } as unknown as CellClassParams);
+ };
+
+ let logSpy: jest.SpyInstance;
+ beforeEach(() => { logSpy = jest.spyOn(console, "log").mockImplementation(() => {}); });
+ afterEach(() => { logSpy.mockRestore(); });
-
+ it("returns the neutral style without logging when the stats column has no entry", () => {
+ expect(cellStyleFor({ histogram_stats: {} })).toEqual({ backgroundColor: "inherit" });
+ expect(logSpy).not.toHaveBeenCalled();
+ });
+
+ it("returns the neutral style without logging when the entry has no histogram_bins", () => {
+ expect(cellStyleFor({ histogram_stats: { a: { histogram_log_bins: [1, 2] } } })).toEqual({ backgroundColor: "inherit" });
+ expect(logSpy).not.toHaveBeenCalled();
+ });
+
+ it("returns the neutral style without logging when the context carries no histogram_stats", () => {
+ expect(cellStyleFor({})).toEqual({ backgroundColor: "inherit" });
+ expect(logSpy).not.toHaveBeenCalled();
+ });
+
+ it("colors the cell once bins exist", () => {
+ const style = cellStyleFor({ histogram_stats: { a: { histogram_bins: [1, 2, 3, 4, 5] } } });
+ expect(style.backgroundColor).not.toBe("inherit");
+ expect(style.backgroundColor).toBeDefined();
+ });
+ });
});
diff --git a/packages/buckaroo-js-core/src/components/DFViewerParts/gridUtils.ts b/packages/buckaroo-js-core/src/components/DFViewerParts/gridUtils.ts
index b963123e4..8d765f84e 100644
--- a/packages/buckaroo-js-core/src/components/DFViewerParts/gridUtils.ts
+++ b/packages/buckaroo-js-core/src/components/DFViewerParts/gridUtils.ts
@@ -38,6 +38,7 @@ import { getFormatterFromArgs, getCellRenderer, objFormatter, getFormatter } fro
import { CSSProperties, Dispatch, SetStateAction } from "react";
import { CommandConfigT } from "../CommandUtils";
import { KeyAwareSmartRowCache, PayloadArgs } from "./SmartRowCache";
+import type { StatsStatus } from "../WidgetTypes";
// for now colDef stuff with less than 3 implementantions should stay in this file
@@ -97,13 +98,33 @@ export function stripOptionalPinnedKey(key: string): string {
return isOptionalPinnedKey(key) ? key.slice(1) : key;
}
-export function extractPinnedRows(sdf: DFData, prc: PinnedRowConfig[]) {
+// Marks a pinned row that stands in for stats that have not arrived. The row
+// carries its key as `index`, so it keeps its label and gets a row id of its
+// own, and the cell renderer selector leaves its value cells empty.
+export const PENDING_STAT_ROW_KEY = "__stat_pending";
+
+const PendingStatCell = () => null;
+const pendingStatRenderer: CellRendererSelectorResult = { component: PendingStatCell };
+
+// `statsStatus` is df_meta.stats.status. A required key with no value is, by
+// status:
+// "pending" a placeholder row, which holds the pinned area's height
+// "not_computed", "error" omitted, since no value is coming
+// anything else undefined, as it was before df_meta.stats existed
+export function extractPinnedRows(sdf: DFData, prc: PinnedRowConfig[], statsStatus?: StatsStatus) {
const result: (DFData[number] | undefined)[] = [];
for (const cfg of prc) {
const raw = cfg.primary_key_val;
- const found = _.find(sdf, { index: stripOptionalPinnedKey(raw) });
- if (found === undefined && isOptionalPinnedKey(raw)) {
- continue;
+ const key = stripOptionalPinnedKey(raw);
+ const found = _.find(sdf, { index: key });
+ if (found === undefined) {
+ if (isOptionalPinnedKey(raw) || statsStatus === "not_computed" || statsStatus === "error") {
+ continue;
+ }
+ if (statsStatus === "pending") {
+ result.push({ index: key, [PENDING_STAT_ROW_KEY]: true });
+ continue;
+ }
}
result.push(found);
}
@@ -340,6 +361,9 @@ export function getCellRendererSelector(pinned_rows: PinnedRowConfig[], column_c
if (pk === undefined) {
return anyRenderer; // default renderer
}
+ if (_.get(params.node.data, PENDING_STAT_ROW_KEY) === true && params.column?.getColId() !== "index") {
+ return pendingStatRenderer; // a stat that has not arrived: leave the cell empty
+ }
const maybePrc: PinnedRowConfig | undefined = _.find(
pinned_rows,
(cfg) => stripOptionalPinnedKey(cfg.primary_key_val) === pk,
diff --git a/packages/buckaroo-js-core/src/components/StatusBar.stats.test.tsx b/packages/buckaroo-js-core/src/components/StatusBar.stats.test.tsx
new file mode 100644
index 000000000..7db44f0e8
--- /dev/null
+++ b/packages/buckaroo-js-core/src/components/StatusBar.stats.test.tsx
@@ -0,0 +1,152 @@
+/**
+ * StatusBar — summary stats status (rows-first c4).
+ *
+ * A session that reports df_meta.stats gets one extra, fixed-width column in
+ * the status bar showing where the stats stand: loading ("pending"), not
+ * computed (with a control that asks for them), error (with the reason) or
+ * ready. A session that does not report df_meta.stats gets the status bar it
+ * always had.
+ *
+ * AG Grid is stubbed to capture the props the status bar hands it; the cell
+ * renderer is rendered on its own.
+ */
+import "@testing-library/jest-dom";
+import { render, screen, fireEvent } from "@testing-library/react";
+
+// The props the status bar gave AG Grid on its last render.
+const mockGrid: { props: any } = { props: null };
+jest.mock("ag-grid-react", () => ({
+ AgGridReact: (props: any) => {
+ // React also calls this stub once with no props; keep the last real ones.
+ if (props) mockGrid.props = props;
+ return ;
+ },
+}));
+jest.mock("./useColorScheme", () => ({ useColorScheme: () => "light" }));
+
+import { StatusBar, StatsStatusCell } from "./StatusBar";
+import { BuckarooOptions, BuckarooState, DFMeta, DFMetaStats } from "./WidgetTypes";
+
+const baseMeta: DFMeta = { total_rows: 378, columns: 7, filtered_rows: 297, rows_shown: 297 };
+const options: BuckarooOptions = {
+ sampled: [],
+ cleaning_method: ["", "clean1"],
+ post_processing: ["", "post1"],
+ df_display: ["main", "summary"],
+ show_commands: ["0", "1"],
+};
+const bState: BuckarooState = {
+ sampled: false,
+ cleaning_method: false,
+ quick_command_args: {},
+ post_processing: false,
+ df_display: "main",
+ show_commands: false,
+};
+
+const renderBar = (dfMeta: DFMeta, onComputeStats?: () => void) =>
+ render(
+ {}}
+ buckarooOptions={options}
+ onComputeStats={onComputeStats}
+ />,
+ );
+
+const fields = (): string[] => mockGrid.props.columnDefs.map((c: any) => c.field);
+
+describe("StatusBar stats column", () => {
+ beforeEach(() => {
+ mockGrid.props = null;
+ });
+
+ it("adds no column and no row field when df_meta has no stats (every session today)", () => {
+ renderBar(baseMeta);
+ expect(fields()).not.toContain("stats");
+ expect(mockGrid.props.rowData[0]).not.toHaveProperty("stats");
+ });
+
+ it("adds a fixed-width stats column after the summary-view selector when df_meta.stats is present", () => {
+ const stats: DFMetaStats = { status: "pending", tier: "schema", gen: 1 };
+ renderBar({ ...baseMeta, stats });
+ const names = fields();
+ expect(names.indexOf("stats")).toBe(names.indexOf("df_display") + 1);
+
+ const column = mockGrid.props.columnDefs[names.indexOf("stats")];
+ expect(column.cellRenderer).toBe(StatsStatusCell);
+ // A fixed width, so the status changing never moves the other columns.
+ expect(typeof column.width).toBe("number");
+ expect(column.flex).toBeUndefined();
+ expect(mockGrid.props.rowData[0].stats).toBe(stats);
+ });
+
+ it("keeps the column for every status", () => {
+ for (const status of ["pending", "not_computed", "error", "complete"] as const) {
+ const { unmount } = renderBar({ ...baseMeta, stats: { status, gen: 1 } });
+ expect(fields()).toContain("stats");
+ unmount();
+ }
+ });
+
+ it("hands the compute callback to the cell renderer through the grid context", () => {
+ const onComputeStats = jest.fn();
+ renderBar({ ...baseMeta, stats: { status: "not_computed", gen: 1 } }, onComputeStats);
+ expect(mockGrid.props.context.onComputeStats).toBe(onComputeStats);
+ });
+});
+
+describe("StatsStatusCell", () => {
+ const cell = (value: DFMetaStats | undefined, onComputeStats?: () => void) =>
+ render();
+
+ it("pending: says the stats are being computed", () => {
+ cell({ status: "pending", gen: 1 });
+ const root = screen.getByTestId("stats-status");
+ expect(root).toHaveAttribute("data-stats-status", "pending");
+ expect(root).toHaveTextContent("Computing summary stats");
+ expect(root).toHaveAttribute("role", "status");
+ });
+
+ it("not_computed: offers a Compute summary stats button that calls the handler", () => {
+ const onComputeStats = jest.fn();
+ cell({ status: "not_computed", gen: 1 }, onComputeStats);
+ expect(screen.getByTestId("stats-status")).toHaveAttribute("data-stats-status", "not_computed");
+
+ fireEvent.click(screen.getByRole("button", { name: "Compute summary stats" }));
+ // Called with no arguments, not with the click event.
+ expect(onComputeStats).toHaveBeenCalledTimes(1);
+ expect(onComputeStats).toHaveBeenCalledWith();
+ });
+
+ it("not_computed: with no handler there is no button, only the label", () => {
+ cell({ status: "not_computed", gen: 1 });
+ expect(screen.queryByRole("button")).not.toBeInTheDocument();
+ expect(screen.getByTestId("stats-status")).toHaveTextContent("Summary stats not computed");
+ });
+
+ it("error: shows the reason", () => {
+ cell({ status: "error", gen: 1, reason: "stats_failed" });
+ const root = screen.getByTestId("stats-status");
+ expect(root).toHaveAttribute("data-stats-status", "error");
+ expect(root).toHaveTextContent("Stats error: stats_failed");
+ });
+
+ it("error: without a reason still says so", () => {
+ cell({ status: "error", gen: 1 });
+ expect(screen.getByTestId("stats-status")).toHaveTextContent("Stats error");
+ });
+
+ it("complete: says the stats are ready", () => {
+ cell({ status: "complete", tier: "full", gen: 1 });
+ const root = screen.getByTestId("stats-status");
+ expect(root).toHaveAttribute("data-stats-status", "complete");
+ expect(root).toHaveTextContent("Summary stats ready");
+ });
+
+ it("renders nothing without stats", () => {
+ const { container } = cell(undefined);
+ expect(container).toBeEmptyDOMElement();
+ });
+});
diff --git a/packages/buckaroo-js-core/src/components/StatusBar.tsx b/packages/buckaroo-js-core/src/components/StatusBar.tsx
index 1a08dedd8..e7c9be79e 100644
--- a/packages/buckaroo-js-core/src/components/StatusBar.tsx
+++ b/packages/buckaroo-js-core/src/components/StatusBar.tsx
@@ -4,7 +4,7 @@ import * as _ from "lodash-es";
import { AgGridReact } from "ag-grid-react"; // the AG Grid React Component
import { ColDef, GridApi, GridOptions } from "ag-grid-community";
import { basicIntFormatter } from "./DFViewerParts/Displayer";
-import { DFMeta } from "./WidgetTypes";
+import { DFMeta, DFMetaStats } from "./WidgetTypes";
import { BuckarooOptions } from "./WidgetTypes";
import { BuckarooState, BKeys } from "./WidgetTypes";
import { CustomCellEditorProps } from 'ag-grid-react';
@@ -307,6 +307,49 @@ export const SearchEditor = memo(({ value, onValueChange, stopEditing }: Custom
);
});
+/**
+ * Where the summary stats stand, as the server reports it in df_meta.stats:
+ * loading while they are pending, a control to ask for them while they are not
+ * computed, the reason when they failed. The cell always renders one line in a
+ * fixed-width column, so changing status moves nothing.
+ */
+export const StatsStatusCell = function (params: { value?: DFMetaStats; context?: { onComputeStats?: () => void } }) {
+ const stats = params.value;
+ if (stats === undefined) return null;
+ const onComputeStats = params.context?.onComputeStats;
+ const cell = (content: React.ReactNode, extra: React.HTMLAttributes = {}) => (
+
+ {content}
+
+ );
+ switch (stats.status) {
+ case "pending":
+ return cell(
+ <>
+
+ Computing summary stats…
+ >,
+ { role: "status", "aria-live": "polite" },
+ );
+ case "not_computed":
+ return onComputeStats ? (
+ cell(
+ ,
+ )
+ ) : (
+ cell("Summary stats not computed")
+ );
+ case "error": {
+ const text = stats.reason ? `Stats error: ${stats.reason}` : "Stats error";
+ return cell(text, { title: text });
+ }
+ default:
+ return cell("Summary stats ready");
+ }
+};
+
export function StatusBar({
dfMeta,
buckarooState,
@@ -316,6 +359,7 @@ export function StatusBar({
themeConfig,
inFlight,
componentConfig,
+ onComputeStats,
}: {
dfMeta: DFMeta;
buckarooState: BuckarooState;
@@ -334,6 +378,9 @@ export function StatusBar({
* Python's ComponentConfig TypedDict; cell renderers read them via
* params.context.componentConfig. */
componentConfig?: Record;
+ /** Sends a forced stats_request. The stats column shows it as a button while
+ * df_meta.stats.status is "not_computed"; without it there is no button. */
+ onComputeStats?: () => void;
}) {
if (false) {
console.log("heightOverride", heightOverride);
@@ -408,6 +455,18 @@ export function StatusBar({
width: 120,
cellRenderer: dfDisplayCell,
},
+ // Only for a session whose server reports df_meta.stats; a fixed width, so
+ // the status changing never moves the other columns.
+ ...(dfMeta.stats === undefined
+ ? []
+ : [{
+ field: "stats",
+ headerName: "stats",
+ headerTooltip: "Summary stats status",
+ width: 200,
+ cellDataType: false,
+ cellRenderer: StatsStatusCell,
+ }]),
/*
{
field: 'auto_clean',
@@ -471,7 +530,8 @@ export function StatusBar({
filtered_rows: basicIntFormatter.format(dfMeta.filtered_rows),
post_processing: buckarooState.post_processing,
show_commands: buckarooState.show_commands || "0",
- search: searchStr
+ search: searchStr,
+ ...(dfMeta.stats === undefined ? {} : { stats: dfMeta.stats }),
},
];
@@ -543,6 +603,7 @@ export function StatusBar({
setBuckarooState,
buckarooOptions,
componentConfig,
+ onComputeStats,
}}
>
diff --git a/packages/buckaroo-js-core/src/components/WidgetTypes.tsx b/packages/buckaroo-js-core/src/components/WidgetTypes.tsx
index e7f63015c..d15baac8e 100644
--- a/packages/buckaroo-js-core/src/components/WidgetTypes.tsx
+++ b/packages/buckaroo-js-core/src/components/WidgetTypes.tsx
@@ -1,11 +1,29 @@
+// Where the summary stats stand, as the server reports it in df_meta.stats.
+// "pending": a stats_update is expected. "not_computed": none will be sent
+// unless the user asks. A missing df_meta.stats means "complete", which is
+// what servers that predate the field send.
+export type StatsStatus = "complete" | "pending" | "not_computed" | "error";
+
+export interface DFMetaStats {
+ status: StatsStatus;
+ tier?: string;
+ reason?: string;
+ gen?: number;
+}
+
export interface DFMeta {
// static,
total_rows: number;
columns: number;
filtered_rows: number;
rows_shown: number;
+ // Absent when the server predates the two-message protocol.
+ stats?: DFMetaStats;
}
+export const getStatsStatus = (meta: DFMeta | undefined): StatsStatus =>
+ meta?.stats?.status ?? "complete";
+
export interface BuckarooOptions {
sampled: string[];
cleaning_method: string[];
diff --git a/packages/buckaroo-js-core/src/index.ts b/packages/buckaroo-js-core/src/index.ts
index 21e422a46..4a56e1810 100644
--- a/packages/buckaroo-js-core/src/index.ts
+++ b/packages/buckaroo-js-core/src/index.ts
@@ -18,6 +18,9 @@ import { BuckarooStaticTable } from './components/BuckarooStaticTable';
import { BuckarooServerView, buckarooWsUrl } from './server/BuckarooServerView';
import { BuckarooView } from './server/BuckarooView';
import { WebSocketModel } from './server/WebSocketModel';
+import { makeLatestDictDecoder } from './server/latestDictDecoder';
+import { withStatsCapability } from './server/StatsChannel';
+import { StateOrchestrator, requestStats } from './server/StateOrchestrator';
import { HistogramCell } from "./components/DFViewerParts/HistogramCell";
import { InfiniteEx } from "./components/DFViewerParts/TableInfinite";
@@ -60,6 +63,10 @@ export default {
BuckarooView,
buckarooWsUrl,
WebSocketModel,
+ makeLatestDictDecoder,
+ withStatsCapability,
+ StateOrchestrator,
+ requestStats,
};
// Named exports for direct imports
@@ -89,6 +96,10 @@ export {
BuckarooView,
buckarooWsUrl,
WebSocketModel,
+ makeLatestDictDecoder,
+ withStatsCapability,
+ StateOrchestrator,
+ requestStats,
};
export type { IModel } from './server/IModel';
diff --git a/packages/buckaroo-js-core/src/server/BuckarooServerView.caps.test.tsx b/packages/buckaroo-js-core/src/server/BuckarooServerView.caps.test.tsx
new file mode 100644
index 000000000..21e8b8421
--- /dev/null
+++ b/packages/buckaroo-js-core/src/server/BuckarooServerView.caps.test.tsx
@@ -0,0 +1,84 @@
+/**
+ * BuckarooServerView — capability advertisement (rows-first c2).
+ *
+ * The server records a client's capabilities from `?caps=` on the WebSocket
+ * URL, because it sends the first message before the client says anything. A
+ * client that merges `stats_update` must put it there, whatever URL the host
+ * passes in.
+ */
+import { render, cleanup, waitFor } from "@testing-library/react";
+import { BuckarooServerView } from "./BuckarooServerView";
+
+const capturedViewProps: any[] = [];
+
+jest.mock("./BuckarooView", () => ({
+ BuckarooView: (props: any) => {
+ capturedViewProps.push(props);
+ return ;
+ },
+ pickMode: (m: unknown) => (m === "buckaroo" ? "buckaroo" : "viewer"),
+}));
+
+jest.mock("./WebSocketModel", () => ({
+ WebSocketModel: class { constructor(_ws: any, _state: any) {} },
+}));
+
+class FakeWebSocket {
+ static instances: FakeWebSocket[] = [];
+ binaryType = "arraybuffer";
+ onopen: (() => void) | null = null;
+ onerror: ((e: any) => void) | null = null;
+ private listeners: Record void>> = {};
+ constructor(public url: string) {
+ FakeWebSocket.instances.push(this);
+ setTimeout(() => {
+ this.onopen?.();
+ setTimeout(() => {
+ const initial = {
+ type: "initial_state",
+ df_meta: { total_rows: 4, columns: 2, filtered_rows: 4, rows_shown: 4 },
+ df_data_dict: {},
+ df_display_args: {},
+ mode: "viewer",
+ };
+ this.listeners["message"]?.forEach((h) => h({ data: JSON.stringify(initial) } as any));
+ }, 0);
+ }, 0);
+ }
+ addEventListener(ev: string, h: (e: any) => void) {
+ (this.listeners[ev] ??= new Set()).add(h);
+ }
+ removeEventListener(ev: string, h: (e: any) => void) {
+ this.listeners[ev]?.delete(h);
+ }
+ close() {}
+}
+
+const origWebSocket = (globalThis as any).WebSocket;
+
+beforeAll(() => {
+ (globalThis as any).WebSocket = FakeWebSocket;
+});
+afterAll(() => {
+ (globalThis as any).WebSocket = origWebSocket;
+});
+afterEach(() => {
+ capturedViewProps.length = 0;
+ FakeWebSocket.instances.length = 0;
+ cleanup();
+});
+
+describe("BuckarooServerView advertises stats_update", () => {
+ it("opens the socket with ?caps=stats_update", async () => {
+ render();
+ await waitFor(() => expect(capturedViewProps.length).toBeGreaterThan(0));
+ expect(FakeWebSocket.instances).toHaveLength(1);
+ expect(FakeWebSocket.instances[0].url).toBe("ws://x/ws/s?caps=stats_update");
+ });
+
+ it("keeps the query string the host passed", async () => {
+ render();
+ await waitFor(() => expect(capturedViewProps.length).toBeGreaterThan(0));
+ expect(FakeWebSocket.instances[0].url).toBe("ws://x/ws/s?token=abc&caps=stats_update");
+ });
+});
diff --git a/packages/buckaroo-js-core/src/server/BuckarooServerView.tsx b/packages/buckaroo-js-core/src/server/BuckarooServerView.tsx
index 49fedbe3a..1687dddaa 100644
--- a/packages/buckaroo-js-core/src/server/BuckarooServerView.tsx
+++ b/packages/buckaroo-js-core/src/server/BuckarooServerView.tsx
@@ -2,6 +2,7 @@ import * as React from "react";
import { decodeDFDataDict } from "../components/DFViewerParts/resolveDFData";
import { WebSocketModel } from "./WebSocketModel";
+import { withStatsCapability } from "./StatsChannel";
import {
BuckarooView,
BuckarooServerMetadata,
@@ -106,7 +107,8 @@ export function BuckarooServerView({
(async () => {
try {
- ws = new WebSocket(wsUrl);
+ // The server reads capabilities from the URL at open.
+ ws = new WebSocket(withStatsCapability(wsUrl));
ws.binaryType = "arraybuffer";
await new Promise((resolve, reject) => {
diff --git a/packages/buckaroo-js-core/src/server/BuckarooView.stats.test.tsx b/packages/buckaroo-js-core/src/server/BuckarooView.stats.test.tsx
new file mode 100644
index 000000000..9bb495202
--- /dev/null
+++ b/packages/buckaroo-js-core/src/server/BuckarooView.stats.test.tsx
@@ -0,0 +1,89 @@
+/**
+ * BuckarooView — the Compute summary stats control (rows-first c4).
+ *
+ * While df_meta.stats says the stats are not computed, the status bar offers a
+ * control that asks the server for them. BuckarooView hands the widget the
+ * callback, and the callback sends `stats_request {force: true}` through
+ * whatever IModel the host gave it.
+ */
+import { render, cleanup, act } from "@testing-library/react";
+import { BuckarooView } from "./BuckarooView";
+import type { IModel } from "./IModel";
+
+const mockWidgetProps: any[] = [];
+jest.mock("../components/BuckarooWidgetInfinite", () => ({
+ BuckarooInfiniteWidget: (props: any) => {
+ mockWidgetProps.push(props);
+ return ;
+ },
+ DFViewerInfiniteDS: () => ,
+ getKeySmartRowCache: jest.fn(() => ({ __stub: "row-cache" })),
+}));
+
+function makeFakeModel(state: Record): { model: IModel; sent: any[] } {
+ const sent: any[] = [];
+ const model: IModel = {
+ send: (msg) => { sent.push(msg); },
+ get: (k) => state[k],
+ set: (k, v) => { state[k] = v; },
+ save_changes: () => {},
+ on: () => {},
+ off: () => {},
+ };
+ return { model, sent };
+}
+
+const displayArgs = {
+ main: { df_viewer_config: { pinned_rows: [], left_col_configs: [], column_config: [] }, summary_stats_key: "all_stats" },
+};
+const metaWith = (stats?: Record) => ({
+ total_rows: 3, columns: 1, filtered_rows: 3, rows_shown: 3,
+ ...(stats === undefined ? {} : { stats }),
+});
+
+const mountBuckaroo = async (state: Record) => {
+ const { model, sent } = makeFakeModel(state);
+ await act(async () => {
+ render();
+ });
+ return { model, sent, props: () => mockWidgetProps[mockWidgetProps.length - 1] };
+};
+
+afterEach(() => {
+ mockWidgetProps.length = 0;
+ cleanup();
+});
+
+describe("BuckarooView on_compute_stats (rows-first c4)", () => {
+ it("hands the widget a callback that sends a forced stats_request for the gen on screen", async () => {
+ const { sent, props } = await mountBuckaroo({
+ df_meta: metaWith({ status: "not_computed", tier: "schema", gen: 9 }),
+ df_data_dict: {},
+ df_display_args: displayArgs,
+ });
+ expect(typeof props().on_compute_stats).toBe("function");
+
+ props().on_compute_stats();
+ expect(sent).toEqual([{ type: "stats_request", stats_gen: 9, scope: "raw", force: true }]);
+ });
+
+ it("sends nothing when the model's df_meta carries no stats.gen", async () => {
+ const { sent, props } = await mountBuckaroo({
+ df_meta: metaWith(),
+ df_data_dict: {},
+ df_display_args: displayArgs,
+ });
+ expect(typeof props().on_compute_stats).toBe("function");
+ props().on_compute_stats();
+ expect(sent).toEqual([]);
+ });
+
+ it("sends no request on its own: asking is the scheduler's job, and a session that is not pending is left alone", async () => {
+ const { sent } = await mountBuckaroo({
+ df_meta: metaWith({ status: "not_computed", tier: "schema", gen: 9 }),
+ df_data_dict: {},
+ df_display_args: displayArgs,
+ });
+ expect(sent).toEqual([]);
+ });
+});
diff --git a/packages/buckaroo-js-core/src/server/BuckarooView.test.tsx b/packages/buckaroo-js-core/src/server/BuckarooView.test.tsx
index 8a38aa6cf..796e1db75 100644
--- a/packages/buckaroo-js-core/src/server/BuckarooView.test.tsx
+++ b/packages/buckaroo-js-core/src/server/BuckarooView.test.tsx
@@ -9,18 +9,34 @@
import { render, cleanup, act } from "@testing-library/react";
import { BuckarooView } from "./BuckarooView";
import type { IModel } from "./IModel";
+import { decodeDFDataDict } from "../components/DFViewerParts/resolveDFData";
// Stub the heavy widget surfaces — this test exercises the injection
// wiring, not AG-Grid. The widget components instantiate AgGridReact which
// is fragile under jsdom; the stub keeps the test focused on the model
// contract.
+//
+// The viewer stub records the props BuckarooView hands it, so the rows-first
+// tests below can see which df_meta / df_data_dict reached the widget.
+const mockViewerProps: any[] = [];
jest.mock("../components/BuckarooWidgetInfinite", () => ({
BuckarooInfiniteWidget: () => ,
- DFViewerInfiniteDS: () => ,
+ DFViewerInfiniteDS: (props: any) => {
+ mockViewerProps.push(props);
+ return ;
+ },
getKeySmartRowCache: jest.fn(() => ({ __stub: "row-cache" })),
}));
-function makeFakeModel(): { model: IModel; events: Map>; sent: any[] } {
+// Wrap decodeDFDataDict in a jest.fn that defaults to the real decoder, so the
+// existing tests run unchanged and the rows-first tests can count and delay
+// decodes.
+jest.mock("../components/DFViewerParts/resolveDFData", () => {
+ const actual = jest.requireActual("../components/DFViewerParts/resolveDFData");
+ return { ...actual, decodeDFDataDict: jest.fn(actual.decodeDFDataDict) };
+});
+
+function makeFakeModel(): { model: IModel; events: Map>; sent: any[]; state: Record } {
const events = new Map>();
const state: Record = {};
const sent: any[] = [];
@@ -35,10 +51,13 @@ function makeFakeModel(): { model: IModel; events: Map>; s
},
off: (e, h) => { events.get(e)?.delete(h); },
};
- return { model, events, sent };
+ return { model, events, sent, state };
}
-afterEach(() => cleanup());
+afterEach(() => {
+ mockViewerProps.length = 0;
+ cleanup();
+});
describe("BuckarooView (injectable IModel — #759)", () => {
it("renders the viewer widget when given a fake IModel + initialState — no WebSocket needed", async () => {
@@ -102,3 +121,110 @@ describe("BuckarooView (injectable IModel — #759)", () => {
expect(onMetadata).toHaveBeenCalledWith({ path: "/data/sales.parquet", rows: 42 }, "tell me about sales");
});
});
+
+// Rows-first c0a: a second message can follow the first at any moment, so the
+// view has to cope with changes that land while it is still wiring itself up.
+describe("BuckarooView two-message hardening (rows-first c0a)", () => {
+ const mockDecode = decodeDFDataDict as jest.Mock;
+ const realDecode = jest.requireActual("../components/DFViewerParts/resolveDFData").decodeDFDataDict;
+
+ const displayArgs = {
+ main: { df_viewer_config: { pinned_rows: [], left_col_configs: [], column_config: [] }, summary_stats_key: "all_stats" },
+ };
+ const metaWith = (total_rows: number) => ({ total_rows, columns: 1, filtered_rows: total_rows, rows_shown: total_rows });
+ const emit = (events: Map>, name: string, ...args: unknown[]) => {
+ for (const h of Array.from(events.get(name) ?? [])) h(...args);
+ };
+ const lastProps = () => mockViewerProps[mockViewerProps.length - 1];
+
+ beforeEach(() => {
+ mockDecode.mockReset();
+ mockDecode.mockImplementation(realDecode);
+ });
+
+ it("applies a change that reached the model before the effect subscribed", async () => {
+ // initialState is what the host held when it built the view. A second
+ // initial_state then landed on the model while React was committing,
+ // so its change:* events had no listener yet.
+ const { model, state } = makeFakeModel();
+ state.df_meta = metaWith(99);
+ const initialState = { df_meta: metaWith(1), df_data_dict: {}, df_display_args: displayArgs };
+
+ await act(async () => {
+ render();
+ });
+
+ expect(lastProps().df_meta.total_rows).toBe(99);
+ });
+
+ it("decodes a df_data_dict that reached the model before the effect subscribed", async () => {
+ const { model, state } = makeFakeModel();
+ const raw = { format: "mock", id: "N" };
+ state.df_data_dict = { main: raw };
+ mockDecode.mockImplementation(async (dict: any) => ({ main: [{ index: 0, from: dict.main.id }] }));
+ const initialState = { df_meta: metaWith(1), df_data_dict: {}, df_display_args: displayArgs };
+
+ await act(async () => {
+ render();
+ });
+
+ expect(lastProps().df_data_dict.main).toEqual([{ index: 0, from: "N" }]);
+ });
+
+ it("decodes one initial_state that carries metadata once, and renders its df_data_dict once", async () => {
+ const { model, events, state } = makeFakeModel();
+ const initialState = { df_meta: metaWith(1), df_data_dict: {}, df_display_args: displayArgs };
+ await act(async () => {
+ render();
+ });
+ const dictsBefore = new Set(mockViewerProps.map((p) => p.df_data_dict));
+ mockDecode.mockClear();
+ mockDecode.mockImplementation(async (dict: any) => ({ main: [{ index: 0, from: dict.main.id }] }));
+
+ // The order WebSocketModel emits for a full frame: one change per key,
+ // then "metadata". The model state already holds the new values.
+ const frame = {
+ df_meta: metaWith(7),
+ df_data_dict: { main: { format: "mock", id: "F" } },
+ df_display_args: displayArgs,
+ metadata: { path: "/data/f.parquet", rows: 7 },
+ };
+ Object.assign(state, frame);
+ await act(async () => {
+ emit(events, "change:df_meta", frame.df_meta);
+ emit(events, "change:df_data_dict", frame.df_data_dict);
+ emit(events, "change:df_display_args", frame.df_display_args);
+ emit(events, "metadata", frame.metadata, undefined);
+ });
+
+ expect(mockDecode).toHaveBeenCalledTimes(1);
+ const dictsAfter = new Set(mockViewerProps.map((p) => p.df_data_dict));
+ expect(dictsAfter.size - dictsBefore.size).toBe(1);
+ expect(lastProps().df_meta.total_rows).toBe(7);
+ });
+
+ it("applies the newer df_data_dict when decodes complete out of order", async () => {
+ const { model, events } = makeFakeModel();
+ const initialState = { df_meta: metaWith(1), df_data_dict: {}, df_display_args: displayArgs };
+ await act(async () => {
+ render();
+ });
+
+ const finish: Record void> = {};
+ mockDecode.mockImplementation(
+ (dict: any) =>
+ new Promise((resolve) => {
+ finish[dict.main.id] = () => resolve({ main: [{ index: 0, from: dict.main.id }] });
+ }),
+ );
+ await act(async () => {
+ emit(events, "change:df_data_dict", { main: { format: "mock", id: "older" } });
+ emit(events, "change:df_data_dict", { main: { format: "mock", id: "newer" } });
+ });
+ // The newer decode finishes first, then the stale one.
+ await act(async () => { finish["newer"](); });
+ await act(async () => { finish["older"](); });
+
+ expect(lastProps().df_data_dict.main).toEqual([{ index: 0, from: "newer" }]);
+ });
+});
diff --git a/packages/buckaroo-js-core/src/server/BuckarooView.tsx b/packages/buckaroo-js-core/src/server/BuckarooView.tsx
index 3ede8298e..d69eab528 100644
--- a/packages/buckaroo-js-core/src/server/BuckarooView.tsx
+++ b/packages/buckaroo-js-core/src/server/BuckarooView.tsx
@@ -1,15 +1,16 @@
import * as React from "react";
import { BuckarooInfiniteWidget, DFViewerInfiniteDS, getKeySmartRowCache } from "../components/BuckarooWidgetInfinite";
-import { decodeDFDataDict } from "../components/DFViewerParts/resolveDFData";
import { DFMeta, BuckarooState, BuckarooOptions } from "../components/WidgetTypes";
import { CommandConfigT } from "../components/CommandUtils";
import { Operation } from "../components/OperationUtils";
import { OperationResult, baseOperationResults } from "../components/DependentTabs";
-import { DFData, DFDataOrPayload } from "../components/DFViewerParts/DFWhole";
+import { DFData } from "../components/DFViewerParts/DFWhole";
import { IDisplayArgs } from "../components/DFViewerParts/gridUtils";
import { stampLayoutType, isFitContentLayout } from "../components/DFViewerParts/displayArgsUtils";
import { IModel } from "./IModel";
+import { makeLatestDictDecoder, RawDFDataDict } from "./latestDictDecoder";
+import { requestStats } from "./StateOrchestrator";
export type BuckarooServerMode = "viewer" | "buckaroo";
@@ -161,6 +162,22 @@ export function BuckarooView({
const onMetadataRef = React.useRef(onMetadata);
React.useEffect(() => { onMetadataRef.current = onMetadata; }, [onMetadata]);
+ // Every df_data_dict that reaches the view (the seed, change:df_data_dict,
+ // the metadata handler, the catch-up on subscribe) goes through one
+ // decoder, so a frame decodes once and the newest decode wins. The seed is
+ // already applied when it needs no resolution.
+ const dictDecoderRef = React.useRef<((raw: RawDFDataDict) => void) | null>(null);
+ if (dictDecoderRef.current === null) {
+ dictDecoderRef.current = makeLatestDictDecoder(
+ (d) => {
+ setDfDataDict(d as Record);
+ setDataReady(true);
+ },
+ initialNeedsResolution ? undefined : (initialState.df_data_dict as RawDFDataDict),
+ );
+ }
+ const loadDfDataDict = dictDecoderRef.current;
+
// Resolve any parquet-encoded payloads in df_data_dict. Pre-resolved
// dicts (e.g. when BuckarooServerView already ran decodeDFDataDict)
// pass through unchanged, so this is cheap in the common case. Skip
@@ -168,19 +185,13 @@ export function BuckarooView({
// for the BuckarooServerView path.
React.useEffect(() => {
if (!initialNeedsResolution) return;
- let cancelled = false;
- const dict = initialState.df_data_dict as Record | undefined;
+ const dict = initialState.df_data_dict as RawDFDataDict;
if (!dict) {
setDataReady(true);
return;
}
- decodeDFDataDict(dict).then((d) => {
- if (cancelled) return;
- setDfDataDict(d as Record);
- setDataReady(true);
- });
- return () => { cancelled = true; };
- }, [initialState, initialNeedsResolution]);
+ loadDfDataDict(dict);
+ }, [initialState, initialNeedsResolution, loadDfDataDict]);
// Fire onMetadata for the initial payload, matching BuckarooServerView's
// pre-split behavior.
@@ -202,8 +213,7 @@ export function BuckarooView({
const onMeta = (metadata: BuckarooServerMetadata, prompt?: string) => {
onMetadataRef.current?.(metadata, prompt);
setDfMeta((model.get("df_meta") as DFMeta | undefined) ?? { ...DEFAULT_DF_META, total_rows: metadata?.rows ?? 0 });
- decodeDFDataDict((model.get("df_data_dict") as Record | undefined) ?? {})
- .then((d) => setDfDataDict(d as Record));
+ loadDfDataDict(model.get("df_data_dict") as RawDFDataDict);
setDfDisplayArgs((model.get("df_display_args") as Record | undefined) ?? {});
setBuckarooStateLocal((model.get("buckaroo_state") as BuckarooState | undefined) ?? DEFAULT_BUCKAROO_STATE);
setBuckarooOptions((model.get("buckaroo_options") as BuckarooOptions | undefined) ?? DEFAULT_BUCKAROO_OPTIONS);
@@ -212,9 +222,7 @@ export function BuckarooView({
setOperations((model.get("operations") as Operation[] | undefined) ?? []);
};
const onDfMeta = (v: DFMeta) => setDfMeta(v);
- const onDfDataDict = (v: Record) => {
- decodeDFDataDict(v).then((d) => setDfDataDict(d as Record));
- };
+ const onDfDataDict = (v: RawDFDataDict) => loadDfDataDict(v);
const onDfDisplayArgs = (v: Record) => setDfDisplayArgs(v);
const onBState = (v: BuckarooState) => setBuckarooStateLocal(v);
const onBOpts = (v: BuckarooOptions) => setBuckarooOptions(v);
@@ -232,6 +240,25 @@ export function BuckarooView({
model.on("change:operation_results", onOpRes);
model.on("change:operations", onOps);
+ // A change:* emitted after the model was built and before this effect
+ // ran had no listener. The model holds the latest value of every key,
+ // so read each one now that the handlers are in place. A value that is
+ // the one already held is a no-op (same reference).
+ const catchUp: Array<[string, (v: any) => void]> = [
+ ["df_meta", onDfMeta],
+ ["df_data_dict", onDfDataDict],
+ ["df_display_args", onDfDisplayArgs],
+ ["buckaroo_state", onBState],
+ ["buckaroo_options", onBOpts],
+ ["command_config", onCmdCfg],
+ ["operation_results", onOpRes],
+ ["operations", onOps],
+ ];
+ for (const [key, apply] of catchUp) {
+ const v = model.get(key);
+ if (v !== undefined) apply(v);
+ }
+
return () => {
model.off("metadata", onMeta);
model.off("change:df_meta", onDfMeta);
@@ -243,7 +270,7 @@ export function BuckarooView({
model.off("change:operation_results", onOpRes);
model.off("change:operations", onOps);
};
- }, [model]);
+ }, [model, loadDfDataDict]);
const onBuckarooState = React.useCallback>>((newState) => {
const resolved = typeof newState === "function"
@@ -258,6 +285,12 @@ export function BuckarooView({
model.save_changes();
}, [model]);
+ // The status bar's Compute summary stats button, shown while the server
+ // reports the stats as not computed.
+ const onComputeStats = React.useCallback(() => {
+ requestStats(model, { force: true });
+ }, [model]);
+
// gridUtils honors component_config.layoutType. Stamp it per entry so the
// prop wins when provided. autoHeight=undefined → server value left intact.
const effectiveDisplayArgs = React.useMemo(
@@ -296,6 +329,7 @@ export function BuckarooView({
on_buckaroo_state={onBuckarooState}
buckaroo_options={buckarooOptions}
src={src}
+ on_compute_stats={onComputeStats}
/>
) : (
>();
+
+ constructor(public state: Record) {}
+
+ get(key: string) {
+ return this.state[key];
+ }
+ set(key: string, value: any) {
+ this.state[key] = value;
+ this.emit(`change:${key}`, value);
+ }
+ send(msg: any) {
+ this.sent.push(msg);
+ }
+ on(event: string, handler: (...args: any[]) => void) {
+ if (!this.handlers.has(event)) this.handlers.set(event, new Set());
+ this.handlers.get(event)!.add(handler);
+ }
+ off(event: string, handler: (...args: any[]) => void) {
+ this.handlers.get(event)?.delete(handler);
+ }
+ emit(event: string, ...args: any[]) {
+ for (const h of Array.from(this.handlers.get(event) ?? [])) h(...args);
}
- last(): Record {
- return JSON.parse(this.sent[this.sent.length - 1]);
+ listenerCount() {
+ return Array.from(this.handlers.values()).reduce((n, set) => n + set.size, 0);
}
- byType(type: string): Record[] {
- return this.sent
- .map((s) => JSON.parse(s))
- .filter((m) => m.type === type);
+ /** A full frame: each key lands and fires its change event in turn, as
+ * WebSocketModel's initial_state branch does. */
+ frame(msg: Record) {
+ for (const [key, value] of Object.entries(msg)) this.set(key, value);
}
}
-describe("StateOrchestrator", () => {
- let ws: FakeWs;
- let orch: StateOrchestrator;
+const meta = (stats?: Record) => ({
+ total_rows: 3, columns: 2, filtered_rows: 3, rows_shown: 3,
+ ...(stats === undefined ? {} : { stats }),
+});
+const pending = (gen: number) => ({ status: "pending", tier: "schema", gen });
+const complete = (gen: number) => ({ status: "complete", tier: "full", gen });
+const dict = (rows: any[] = []) => ({ all_stats: rows });
+const statRow = (stat: string) => ({ index: stat, level_0: stat, a: 1 });
+// Typed loosely: the server's buckaroo_state also carries keys (search_string)
+// that BuckarooState does not declare.
+const bState = (over: Record = {}): any => ({
+ sampled: false, cleaning_method: false, quick_command_args: {}, post_processing: false,
+ df_display: "main", show_commands: false, ...over,
+});
+
+const request = (gen: number, extra: Record = {}) => ({
+ type: "stats_request", stats_gen: gen, scope: "raw", ...extra,
+});
+
+const makeModel = (stats?: Record) =>
+ new FakeModel({ df_meta: meta(stats), df_data_dict: dict(), buckaroo_state: bState() });
+
+/** The first rows reached the client: WebSocketModel emits msg:custom once it
+ * has paired an infinite_resp with its parquet frame. */
+const rowsArrived = (model: FakeModel) =>
+ model.emit("msg:custom", { type: "infinite_resp", key: { start: 0, end: 3 }, length: 3 }, []);
+
+const start = (model: FakeModel, opts: Record = {}) => {
+ const orchestrator = new StateOrchestrator({ model, ...opts });
+ orchestrator.start();
+ return orchestrator;
+};
- beforeEach(() => {
- jest.useFakeTimers();
- ws = new FakeWs();
- orch = new StateOrchestrator({ ws, minDebounceMs: 10, maxDebounceMs: 5000 });
+// With the defaults, a state change waits 2 x 250 ms before it asks again.
+const DEBOUNCE = 500;
+const FIRST_PAINT_TIMEOUT = 1500;
+
+// Runs due timers and every promise continuation they leave behind.
+const tick = (ms = 0) => jest.advanceTimersByTimeAsync(ms);
+
+beforeEach(() => {
+ jest.useFakeTimers();
+});
+
+afterEach(() => {
+ jest.useRealTimers();
+});
+
+describe("when nothing is pending", () => {
+ it("requests nothing for a session whose df_meta has no stats (every session today)", async () => {
+ const model = makeModel();
+ start(model);
+ rowsArrived(model);
+ await tick(10_000);
+ expect(model.sent).toEqual([]);
+ });
+
+ it.each(["complete", "not_computed", "error"])("requests nothing when the status is %s", async (status) => {
+ const model = makeModel({ status, tier: "schema", gen: 3 });
+ start(model);
+ rowsArrived(model);
+ await tick(10_000);
+ expect(model.sent).toEqual([]);
});
+});
- afterEach(() => {
- orch.dispose();
- jest.useRealTimers();
+describe("the first request", () => {
+ it("goes out after the first infinite_resp, not before", async () => {
+ const model = makeModel(pending(3));
+ start(model);
+ await tick(100);
+ expect(model.sent).toEqual([]);
+
+ rowsArrived(model);
+ await tick();
+ expect(model.sent).toEqual([request(3)]);
});
- it("starts with token 0", () => {
- expect(orch.currentToken).toBe(0);
+ it("goes out anyway when no rows come (an empty frame, the summary view)", async () => {
+ const model = makeModel(pending(3));
+ start(model);
+ await tick(FIRST_PAINT_TIMEOUT - 1);
+ expect(model.sent).toEqual([]);
+
+ await tick(2);
+ expect(model.sent).toEqual([request(3)]);
});
- it("onStateChange bumps token and ships state_change", () => {
- orch.onStateChange({ search_string: "x" });
- expect(orch.currentToken).toBe(1);
- const msg = ws.last();
- expect(msg.type).toBe("state_change");
- expect(msg.state_token).toBe(1);
- expect((msg.new_state as Record).search_string).toBe("x");
+ it("carries no force flag", async () => {
+ const model = makeModel(pending(3));
+ start(model);
+ rowsArrived(model);
+ await tick();
+ expect(model.sent[0]).not.toHaveProperty("force");
});
- it("schedules compute_stat_group after the debounce", () => {
- orch.onStateChange({ search_string: "PIZZA" });
- // Only the immediate state_change has been sent so far.
- expect(ws.byType("compute_stat_group")).toHaveLength(0);
+ it("is sent once, however many row responses follow", async () => {
+ const model = makeModel(pending(3));
+ start(model);
+ rowsArrived(model);
+ rowsArrived(model);
+ await tick();
+ rowsArrived(model);
+ await tick(10_000);
+ expect(model.sent).toEqual([request(3)]);
+ });
+});
+
+describe("one request per reply", () => {
+ const afterFirstRequest = async () => {
+ const model = makeModel(pending(3));
+ const orchestrator = start(model);
+ rowsArrived(model);
+ await tick();
+ return { model, orchestrator };
+ };
- // Default baseline is 500ms × 2× = 1000ms; clamped to maxDebounceMs=5000.
- jest.advanceTimersByTime(999);
- expect(ws.byType("compute_stat_group")).toHaveLength(0);
+ it("asks again for each reply that leaves the stats pending, and stops at the final one", async () => {
+ const { model } = await afterFirstRequest();
+ expect(model.sent).toEqual([request(3)]);
- jest.advanceTimersByTime(1);
- const reqs = ws.byType("compute_stat_group");
- expect(reqs).toHaveLength(1);
- expect(reqs[0].scope).toBe("filt");
- expect(reqs[0].group).toBe("aggregate");
- expect(reqs[0].state_token).toBe(1);
+ // No reply yet: nothing more goes out, however long the server takes.
+ await tick(10_000);
+ expect(model.sent).toEqual([request(3)]);
+
+ // A partial reply is a new df_data_dict under the same df_meta, which is
+ // what StatsChannel does for a stats_update that is not final.
+ model.set("df_data_dict", dict([statRow("mean")]));
+ await tick();
+ expect(model.sent).toEqual([request(3), request(3)]);
+ await tick(10_000);
+ expect(model.sent).toHaveLength(2);
+
+ model.set("df_data_dict", dict([statRow("mean"), statRow("std")]));
+ await tick();
+ expect(model.sent).toHaveLength(3);
+
+ // The final reply also sets df_meta.stats to complete.
+ model.set("df_data_dict", dict([statRow("mean"), statRow("std"), statRow("max")]));
+ model.set("df_meta", meta(complete(3)));
+ await tick(10_000);
+ expect(model.sent).toHaveLength(3);
});
- it("back-to-back state_changes cancel the previous debounce timer", () => {
- orch.onStateChange({ search_string: "P" });
- jest.advanceTimersByTime(500);
- orch.onStateChange({ search_string: "PI" });
- // The first timer would have fired at t=1000ms; the second
- // resets it so at t=999ms from the SECOND call (= 1499 overall)
- // no compute_stat_group has fired yet.
- jest.advanceTimersByTime(998);
- expect(ws.byType("compute_stat_group")).toHaveLength(0);
+ it("reads a reply the same way whichever of the two events comes first", async () => {
+ const { model } = await afterFirstRequest();
+ model.set("df_meta", meta(complete(3)));
+ model.set("df_data_dict", dict([statRow("mean")]));
+ await tick(10_000);
+ expect(model.sent).toEqual([request(3)]);
+ });
- // The second timer fires.
- jest.advanceTimersByTime(2);
- const reqs = ws.byType("compute_stat_group");
- expect(reqs).toHaveLength(1);
- expect(reqs[0].state_token).toBe(2); // second state_change's token
+ it.each([
+ ["df_meta then df_data_dict", ["df_meta", "df_data_dict"]],
+ ["df_data_dict then df_meta", ["df_data_dict", "df_meta"]],
+ ])("does not take a full frame for the same state as a reply (%s)", async (_label, order) => {
+ // A search term that changes only the highlight comes back as a full
+ // initial_state for the same stats_gen, with a new df_meta and a new dict.
+ const { model } = await afterFirstRequest();
+ const full: Record = { df_meta: meta(pending(3)), df_data_dict: dict([statRow("dtype")]) };
+ model.frame(Object.fromEntries(order.map((k) => [k, full[k]])));
+ await tick(10_000);
+ expect(model.sent).toEqual([request(3)]);
});
- it("rapid typing produces just one aggregate request, with the latest token", () => {
- // Simulate 5 keystrokes at 100ms intervals — typical typing cadence.
- for (let i = 0; i < 5; i++) {
- orch.onStateChange({ search_string: "P".repeat(i + 1) });
- jest.advanceTimersByTime(100);
- }
- // Default debounce = 1000ms after the last keystroke. No aggregate yet.
- expect(ws.byType("compute_stat_group")).toHaveLength(0);
-
- jest.advanceTimersByTime(1000);
- const reqs = ws.byType("compute_stat_group");
- // Exactly one aggregate request, with the 5th (final) token.
- expect(reqs).toHaveLength(1);
- expect(reqs[0].state_token).toBe(5);
- });
-
- it("onStatGroupResult with matching token updates the baseline", () => {
- orch.onStateChange({ search_string: "x" });
- const applied = orch.onStatGroupResult({
- type: "stat_group_result",
- state_token: 1,
- scope: "filt",
- group: "aggregate",
- elapsed_ms: 7500,
- });
- expect(applied).toBe(true);
-
- // Next debounce is 2× 7500 = 15000ms, clamped to maxDebounceMs=5000.
- expect(orch.computeDebounce("filt")).toBe(5000);
- });
-
- it("onStatGroupResult with stale token is silently dropped", () => {
- orch.onStateChange({ search_string: "x" });
- orch.onStateChange({ search_string: "xy" });
- // Token is now 2.
- const applied = orch.onStatGroupResult({
- type: "stat_group_result",
- state_token: 1, // stale
- scope: "filt",
- group: "aggregate",
- elapsed_ms: 9999,
- });
- expect(applied).toBe(false);
- // Baseline unchanged — debounce stays at the default.
- expect(orch.computeDebounce("filt")).toBe(1000);
+ it("stops when the server reports an error", async () => {
+ const { model } = await afterFirstRequest();
+ model.set("df_meta", meta({ status: "error", tier: "schema", gen: 3, reason: "stats_failed" }));
+ model.set("df_data_dict", dict([statRow("mean")]));
+ await tick(10_000);
+ expect(model.sent).toEqual([request(3)]);
});
- it("computeDebounce respects minDebounceMs floor", () => {
- const o = new StateOrchestrator({
- ws: new FakeWs(),
- minDebounceMs: 500,
- maxDebounceMs: 3000,
- multiplier: 2,
- });
- o.onStatGroupResult({
- type: "stat_group_result",
- state_token: 0,
- scope: "filt",
- group: "aggregate",
- elapsed_ms: 10, // 2× 10 = 20, well under floor
- });
- expect(o.computeDebounce("filt")).toBe(500);
+ it("stops when the server says the stats are not computed", async () => {
+ const { model } = await afterFirstRequest();
+ model.set("df_meta", meta({ status: "not_computed", tier: "schema", gen: 3 }));
+ await tick(10_000);
+ expect(model.sent).toEqual([request(3)]);
});
+});
- it("computeDebounce respects maxDebounceMs ceiling", () => {
- const o = new StateOrchestrator({
- ws: new FakeWs(),
- minDebounceMs: 200,
- maxDebounceMs: 3000,
- multiplier: 2,
- });
- o.onStatGroupResult({
- type: "stat_group_result",
- state_token: 0,
- scope: "filt",
- group: "aggregate",
- elapsed_ms: 6000, // 2× = 12000, hits ceiling
- });
- expect(o.computeDebounce("filt")).toBe(3000);
+describe("a state change", () => {
+ const pendingAt = async (gen: number) => {
+ const model = makeModel(pending(gen));
+ const orchestrator = start(model);
+ rowsArrived(model);
+ await tick();
+ return { model, orchestrator };
+ };
+
+ // The server answers a dataflow change with a frame for the next stats_gen,
+ // then the grid's refetch brings rows. They are separate messages, so the
+ // scheduler has read the frame by the time the rows arrive.
+ const nextFrame = async (model: FakeModel, gen: number) => {
+ model.frame({ df_meta: meta(pending(gen)), df_data_dict: dict() });
+ await tick();
+ rowsArrived(model);
+ };
+
+ it("stops the requests for the old state, then asks once for the new one", async () => {
+ const { model } = await pendingAt(3);
+ expect(model.sent).toEqual([request(3)]);
+
+ // The user changes a dataflow field while the request is out.
+ model.set("buckaroo_state", bState({ post_processing: "log_scale" }));
+ await tick(10);
+ // The reply to the old request lands. It is not followed by another.
+ model.set("df_data_dict", dict([statRow("mean")]));
+ await tick(10);
+ await nextFrame(model, 4);
+
+ await tick(DEBOUNCE - 1);
+ expect(model.sent).toEqual([request(3)]);
+ await tick(1);
+ expect(model.sent).toEqual([request(3), request(4)]);
+ await tick(10_000);
+ expect(model.sent).toHaveLength(2);
+ });
+
+ it("reads a frame whose dict comes before its df_meta as the next state, not as a reply", async () => {
+ const { model } = await pendingAt(3);
+ // The frame for the next gen carries its dict before its df_meta.
+ model.frame({ df_data_dict: dict(), df_meta: meta(pending(4)) });
+ await tick();
+ rowsArrived(model);
+ await tick(10_000);
+ expect(model.sent).toEqual([request(3), request(4)]);
+ });
+
+ it.each([
+ ["post_processing", { post_processing: "log_scale" }],
+ ["cleaning_method", { cleaning_method: "aggressive" }],
+ ["quick_command_args", { quick_command_args: { search: ["x"] } }],
+ ])("a %s change cancels a request that is waiting out its delay", async (_field, change) => {
+ const { model } = await pendingAt(3);
+ await nextFrame(model, 4);
+ await tick(300);
+
+ model.set("buckaroo_state", bState(change));
+ await tick(300);
+ // 600 ms in: the request that was due at 500 ms never went out.
+ expect(model.sent).toEqual([request(3)]);
+
+ await nextFrame(model, 5);
+ await tick(DEBOUNCE - 1);
+ expect(model.sent).toEqual([request(3)]);
+ await tick(1);
+ expect(model.sent).toEqual([request(3), request(5)]);
+ });
+
+ it.each([
+ ["search_string (the #1015 path)", { search_string: "x" }],
+ ["df_display", { df_display: "summary" }],
+ ["show_commands", { show_commands: "1" }],
+ ["sampled", { sampled: "sample" }],
+ ])("a %s-only change is skipped", async (_label, change) => {
+ const { model } = await pendingAt(3);
+ await nextFrame(model, 4);
+ await tick(300);
+
+ model.set("buckaroo_state", bState(change));
+ await tick(200);
+ // The request is still due at 500 ms, as if nothing had changed.
+ expect(model.sent).toEqual([request(3), request(4)]);
+ });
+
+ it("a search_string-only change does not interrupt a chain of replies", async () => {
+ const { model } = await pendingAt(3);
+ model.set("buckaroo_state", bState({ search_string: "x" }));
+ model.set("df_data_dict", dict([statRow("mean")]));
+ await tick();
+ expect(model.sent).toEqual([request(3), request(3)]);
});
- it("initialAggregateMs seeds the baseline before any observed compute", () => {
- const o = new StateOrchestrator({
- ws: new FakeWs(),
- initialAggregateMs: { filt: 250 },
- minDebounceMs: 10,
- maxDebounceMs: 5000,
- multiplier: 2,
+ it("a frame that repeats the same dataflow state is not a state change", async () => {
+ // Every full frame carries buckaroo_state back as the client sent it.
+ const { model } = await pendingAt(3);
+ model.set("buckaroo_state", bState());
+ model.set("df_data_dict", dict([statRow("mean")]));
+ await tick();
+ expect(model.sent).toEqual([request(3), request(3)]);
+ });
+
+ it("waits out the delay for a state change made after the earlier stats completed", async () => {
+ const { model } = await pendingAt(3);
+ await tick(250); // a request that took as long as the default assumes, so the delay is DEBOUNCE
+ model.set("df_data_dict", dict([statRow("mean")]));
+ model.set("df_meta", meta(complete(3))); // the final reply
+ await tick();
+
+ // The first state asked at once; this one is a change to it.
+ model.set("buckaroo_state", bState({ post_processing: "log_scale" }));
+ await nextFrame(model, 4);
+ await tick(DEBOUNCE - 1);
+ expect(model.sent).toEqual([request(3)]);
+ await tick(1);
+ expect(model.sent).toEqual([request(3), request(4)]);
+ });
+
+ it("waits out the delay for the first state change of a model that started with its stats complete", async () => {
+ const model = makeModel(complete(3));
+ start(model);
+ rowsArrived(model);
+ await tick(10_000);
+ expect(model.sent).toEqual([]);
+
+ model.set("buckaroo_state", bState({ quick_command_args: { search: ["a"] } }));
+ await nextFrame(model, 4);
+ await tick(DEBOUNCE - 1);
+ expect(model.sent).toEqual([]);
+ await tick(1);
+ expect(model.sent).toEqual([request(4)]);
+ });
+
+ it("asks again for the same state when the server never answers the change with a frame", async () => {
+ const { model } = await pendingAt(3);
+ model.set("buckaroo_state", bState({ post_processing: "log_scale" }));
+ await tick(FIRST_PAINT_TIMEOUT + DEBOUNCE);
+ expect(model.sent).toEqual([request(3), request(3)]);
+ });
+
+ it("waits 2 x the last request's time, within the limits, before the next state's request", async () => {
+ const model = makeModel(pending(3));
+ const orchestrator = start(model, { minDebounceMs: 100, maxDebounceMs: 2000 });
+ expect(orchestrator.computeDebounce()).toBe(500);
+
+ rowsArrived(model);
+ await tick(400);
+ model.set("df_data_dict", dict([statRow("mean")])); // the reply, 400 ms after the request
+ await tick();
+ expect(orchestrator.computeDebounce()).toBe(800);
+
+ model.set("buckaroo_state", bState({ post_processing: "log_scale" }));
+ await nextFrame(model, 4);
+ await tick(799);
+ expect(model.sent).toHaveLength(2);
+ await tick(1);
+ expect(model.sent).toHaveLength(3);
+ expect(model.sent[2]).toEqual(request(4));
+ });
+
+ it("measures a request once: later events do not stretch the delay", async () => {
+ const model = makeModel(pending(3));
+ const orchestrator = start(model, { minDebounceMs: 100, maxDebounceMs: 20_000 });
+ rowsArrived(model);
+ await tick(400);
+ model.set("df_data_dict", dict([statRow("mean")]));
+ model.set("df_meta", meta(complete(3))); // the final reply
+ await tick();
+ expect(orchestrator.computeDebounce()).toBe(800);
+
+ // A full frame for the finished state, long afterwards.
+ await tick(10_000);
+ model.frame({ df_meta: meta(complete(3)), df_data_dict: dict() });
+ await tick();
+ expect(orchestrator.computeDebounce()).toBe(800);
+ });
+
+ it.each([
+ [10, 100],
+ [400, 800],
+ [5000, 2000],
+ ])("computeDebounce after a %i ms request is %i ms (floor 100, ceiling 2000)", async (elapsed, expected) => {
+ const model = makeModel(pending(3));
+ const orchestrator = start(model, { minDebounceMs: 100, maxDebounceMs: 2000 });
+ rowsArrived(model);
+ await tick(elapsed);
+ model.set("df_data_dict", dict([statRow("mean")]));
+ await tick();
+ expect(orchestrator.computeDebounce()).toBe(expected);
+ });
+});
+
+describe("requestStats", () => {
+ it("sends a stats_request for the gen the model shows", () => {
+ const model = makeModel(pending(7));
+ expect(requestStats(model)).toBe(true);
+ expect(model.sent).toEqual([request(7)]);
+ });
+
+ it("force sends force: true, which the Compute summary stats control uses", () => {
+ const model = makeModel({ status: "not_computed", tier: "schema", gen: 7 });
+ expect(requestStats(model, { force: true })).toBe(true);
+ expect(model.sent).toEqual([request(7, { force: true })]);
+ });
+
+ it("sends nothing when df_meta carries no stats.gen", () => {
+ const model = makeModel();
+ expect(requestStats(model)).toBe(false);
+ expect(requestStats(model, { force: true })).toBe(false);
+ expect(model.sent).toEqual([]);
+ });
+});
+
+describe("touchesDataflow", () => {
+ it("is true for each field the server reruns the dataflow for", () => {
+ expect(touchesDataflow(bState(), bState({ post_processing: "x" }))).toBe(true);
+ expect(touchesDataflow(bState(), bState({ cleaning_method: "x" }))).toBe(true);
+ expect(touchesDataflow(bState(), bState({ quick_command_args: { search: ["x"] } }))).toBe(true);
+ });
+
+ it("is false for the others, and for an equal quick_command_args that is a new object", () => {
+ expect(touchesDataflow(bState(), bState({ search_string: "x" }))).toBe(false);
+ expect(touchesDataflow(bState(), bState({ df_display: "summary" }))).toBe(false);
+ expect(touchesDataflow(bState({ quick_command_args: { search: ["x"] } }), bState({ quick_command_args: { search: ["x"] } }))).toBe(false);
+ });
+
+ it("is false when there is no earlier state to compare with", () => {
+ expect(touchesDataflow(undefined, bState({ post_processing: "x" }))).toBe(false);
+ });
+});
+
+describe("start and stop", () => {
+ it("stop removes every listener and cancels the pending request", async () => {
+ const model = makeModel(pending(3));
+ const orchestrator = start(model);
+ expect(model.listenerCount()).toBeGreaterThan(0);
+
+ rowsArrived(model);
+ orchestrator.stop();
+ expect(model.listenerCount()).toBe(0);
+ await tick(10_000);
+ expect(model.sent).toEqual([]);
+ });
+
+ it("start is idempotent: a second call adds no listeners", () => {
+ const model = makeModel(pending(3));
+ const orchestrator = start(model);
+ const listeners = model.listenerCount();
+ orchestrator.start();
+ expect(model.listenerCount()).toBe(listeners);
+ });
+
+ it("can start again after a stop", async () => {
+ const model = makeModel(pending(3));
+ const orchestrator = start(model);
+ orchestrator.stop();
+ orchestrator.start();
+ rowsArrived(model);
+ await tick();
+ expect(model.sent).toEqual([request(3)]);
+ });
+});
+
+// The scheduler is started by WebSocketModel, so a session reached through
+// BuckarooServerView or the standalone page gets it with no wiring of its own.
+describe("wired into WebSocketModel", () => {
+ class FakeSocket {
+ readyState = 1; // WebSocket.OPEN
+ onmessage: ((e: MessageEvent) => void) | null = null;
+ sent: any[] = [];
+ send(data: string) {
+ this.sent.push(JSON.parse(data));
+ }
+ deliver(msg: object) {
+ this.onmessage?.({ data: JSON.stringify(msg) } as MessageEvent);
+ }
+ deliverBinary() {
+ this.onmessage?.({ data: new ArrayBuffer(8) } as MessageEvent);
+ }
+ }
+
+ const update = (gen: number, stat: string, final: boolean) => ({
+ type: "stats_update",
+ stats_gen: gen,
+ scope: "raw",
+ tier: "full",
+ final,
+ payload: { format: "json", layout: "wide", data: [{ index: stat, level_0: stat, a: 1 }] },
+ elapsed_ms: 5,
+ });
+
+ const makeSocketModel = (stats?: Record) => {
+ const ws = new FakeSocket();
+ const model = new WebSocketModel(ws as unknown as WebSocket, {
+ df_meta: meta(stats),
+ df_data_dict: { all_stats: [{ index: "dtype", level_0: "dtype", a: "int64" }] },
+ buckaroo_state: bState(),
});
- expect(o.computeDebounce("filt")).toBe(500); // 2× 250
- // Unrelated scope still uses fallback.
- expect(o.computeDebounce("clean")).toBe(1000); // 2× 500 (default)
+ return { ws, model };
+ };
+ const rowsFromServer = (ws: FakeSocket) => {
+ ws.deliver({ type: "infinite_resp", key: { start: 0, end: 3 }, length: 3 });
+ ws.deliverBinary();
+ };
+
+ it("requests the stats after the first rows, once per reply, and stops at the final one", async () => {
+ const { ws, model } = makeSocketModel(pending(3));
+ await tick(100);
+ expect(ws.sent).toEqual([]);
+
+ rowsFromServer(ws);
+ await tick();
+ expect(ws.sent).toEqual([request(3)]);
+
+ ws.deliver(update(3, "mean", false));
+ await tick();
+ expect(ws.sent).toEqual([request(3), request(3)]);
+
+ ws.deliver(update(3, "std", true));
+ await tick(10_000);
+ expect(ws.sent).toHaveLength(2);
+ expect(model.get("df_meta").stats.status).toBe("complete");
+ expect(model.get("df_data_dict").all_stats.map((r: any) => r.index)).toEqual(["dtype", "mean", "std"]);
});
- it("dispose cancels all pending aggregate timers", () => {
- orch.onStateChange({ search_string: "x" });
- orch.dispose();
- jest.advanceTimersByTime(10_000);
- expect(ws.byType("compute_stat_group")).toHaveLength(0);
+ it("moves on to the next gen when a state change frame arrives", async () => {
+ const { ws } = makeSocketModel(pending(3));
+ rowsFromServer(ws);
+ await tick();
+ expect(ws.sent).toEqual([request(3)]);
+
+ ws.deliver({ type: "initial_state", df_meta: meta(pending(4)), df_data_dict: dict() });
+ await tick();
+ rowsFromServer(ws);
+ await tick(DEBOUNCE);
+ expect(ws.sent).toEqual([request(3), request(4)]);
});
- it("scopesForAggregate parameter overrides default ['filt']", () => {
- orch.onStateChange({ search_string: "x" }, { scopesForAggregate: ["filt", "clean"] });
- jest.advanceTimersByTime(10_000);
- const reqs = ws.byType("compute_stat_group");
- const scopes = reqs.map((r) => r.scope).sort();
- expect(scopes).toEqual(["clean", "filt"]);
+ it("sends nothing for a model whose df_meta has no stats", async () => {
+ const { ws } = makeSocketModel();
+ rowsFromServer(ws);
+ await tick(10_000);
+ expect(ws.sent).toEqual([]);
});
});
diff --git a/packages/buckaroo-js-core/src/server/StateOrchestrator.ts b/packages/buckaroo-js-core/src/server/StateOrchestrator.ts
index e096ac606..e03650470 100644
--- a/packages/buckaroo-js-core/src/server/StateOrchestrator.ts
+++ b/packages/buckaroo-js-core/src/server/StateOrchestrator.ts
@@ -1,153 +1,263 @@
/**
- * Client-side scheduler for the JS-driven progressive-stats protocol.
+ * Client-side scheduler for the stats wire (rows-first c4).
*
- * On a state_change, ships the cheap (scalar) request immediately so
- * df_meta + scalar pinned rows update fast. Schedules the expensive
- * (aggregate) compute per scope behind an adaptive debounce — by
- * default 2× the last observed aggregate compute time for that scope,
- * clamped to ``[minDebounceMs, maxDebounceMs]``.
+ * A server that defers a session's summary stats sends a first frame whose
+ * `df_meta.stats.status` is "pending", then waits to be asked. This scheduler
+ * asks: it sends `stats_request {stats_gen, scope: "raw"}` once the first rows
+ * have arrived, sends another for each reply that leaves the stats pending (a
+ * `stats_update` that is not final), and stops when the status leaves
+ * "pending". Nothing is requested unless `df_meta.stats` says pending, so a
+ * session whose server reports no stats is never touched.
*
- * Token-based cancellation: every state_change bumps an internal
- * token; results carrying a stale token are dropped on arrival.
+ * It watches the model and sends through it, so any IModel works:
*
- * See plans/js-driven-stat-debounce.md for the full protocol design.
+ * change:df_meta the status and stats_gen of the state on screen
+ * change:df_data_dict a reply merged into all_stats (StatsChannel's job)
+ * change:buckaroo_state a state change made on this client
+ * msg:custom an infinite_resp, the first rows
*
- * This module is transport-agnostic — pass any object with a
- * ``send(s: string)`` method. Existing buckaroo WS connections fit.
+ * The merge is not here. A reply reaches this class only as a change on the
+ * model: a new df_data_dict under the same df_meta is a partial update, and a
+ * status other than "pending" is the end.
+ *
+ * A state change is a change to a dataflow field of buckaroo_state (the fields
+ * the server reruns the dataflow for) and the frame that answers it carries the
+ * next stats_gen. Either one starts the wait again with a delay, so a burst of
+ * changes (typing in the search box) asks once, for the state it ends on. A
+ * change to anything else, a search_string for one, leaves the schedule alone.
+ *
+ * `stats_gen` is the server's token for the state a reply describes. A reply for
+ * another gen is dropped by StatsChannel, so no client-side token is kept here.
*/
+import { BuckarooState, DFMeta } from "../components/WidgetTypes";
+import { IModel } from "./IModel";
+
+/** What the scheduler needs of a model. */
+export type StatsModel = Pick;
-export type ScopeName = "raw" | "clean" | "filt";
-export type CostGroup = "scalar" | "aggregate";
+/** The fields of buckaroo_state the server reruns the dataflow for, and so
+ * bumps stats_gen on. Mirrors _DATAFLOW_FIELDS in
+ * buckaroo/server/websocket_handler.py. */
+export const DATAFLOW_STATE_FIELDS = ["post_processing", "cleaning_method", "quick_command_args"] as const;
+
+/** Whether `next` differs from `prev` in a dataflow field. False when there is
+ * no earlier state to compare with. */
+export function touchesDataflow(prev: BuckarooState | undefined, next: BuckarooState | undefined): boolean {
+ if (prev === undefined || next === undefined) return false;
+ return DATAFLOW_STATE_FIELDS.some((field) => JSON.stringify(prev[field]) !== JSON.stringify(next[field]));
+}
-export interface WsLike {
- send(message: string): void;
+export interface StatsRequestOptions {
+ /** Ask for stats the server did not plan to compute. The "Compute summary
+ * stats" control sends this. */
+ force?: boolean;
+}
+
+/**
+ * Send a `stats_request` for the stats_gen of the state the model shows.
+ * Returns false, and sends nothing, when its df_meta carries no stats.gen.
+ */
+export function requestStats(model: Pick, opts: StatsRequestOptions = {}): boolean {
+ const gen = (model.get("df_meta") as DFMeta | undefined)?.stats?.gen;
+ if (typeof gen !== "number") return false;
+ model.send({ type: "stats_request", stats_gen: gen, scope: "raw", ...(opts.force ? { force: true } : {}) });
+ return true;
}
export interface OrchestratorOptions {
- ws: WsLike;
- /** Lower bound on the per-scope debounce. Default 200 ms. */
+ model: StatsModel;
+ /** Lower bound on the delay before asking for a new state's stats. Default 200 ms. */
minDebounceMs?: number;
- /** Upper bound on the per-scope debounce. Default 3000 ms. */
+ /** Upper bound on that delay. Default 3000 ms. */
maxDebounceMs?: number;
- /** Multiplier on the last observed aggregate compute time. Default 2. */
+ /** Multiplier on the last observed request time. Default 2. */
multiplier?: number;
- /** Initial aggregate baseline per scope (used until we observe a real one). */
- initialAggregateMs?: Partial>;
+ /** Request time assumed until one has been observed. Default 250 ms. */
+ initialRequestMs?: number;
+ /** How long to wait for the first rows before asking anyway (an empty
+ * frame, the summary view, and a grid that never fetches send none).
+ * Default 1500 ms. */
+ firstPaintTimeoutMs?: number;
}
-export interface StatGroupResult {
- type: "stat_group_result";
- state_token: number;
- scope: ScopeName;
- group: CostGroup;
- elapsed_ms: number;
- stats?: unknown;
-}
+type Timer = ReturnType;
export class StateOrchestrator {
- private token = 0;
- private aggregateTimers = new Map>();
- private lastAggregateMs = new Map();
- private readonly ws: WsLike;
+ private readonly model: StatsModel;
private readonly minDebounceMs: number;
private readonly maxDebounceMs: number;
private readonly multiplier: number;
- /** Fallback baseline before any real aggregate compute has been observed. */
- private readonly defaultBaselineMs = 500;
+ private readonly initialRequestMs: number;
+ private readonly firstPaintTimeoutMs: number;
+
+ private started = false;
+ // The stats_gen being driven; undefined while nothing is pending.
+ private gen: number | undefined;
+ // A request is out and its reply has not been seen.
+ private inFlight = false;
+ // Delay before the next request: 0 for the first state and for each reply
+ // in a chain, the debounce after a state change.
+ private delayMs = 0;
+ // The state the model held at start has been read, so a pending state after
+ // it is a state change, whatever the first one was.
+ private began = false;
+ private sentAt = 0;
+ private lastRequestMs: number | undefined;
+ private requestTimer: Timer | undefined;
+ private paintTimer: Timer | undefined;
+ private syncQueued = false;
+ private seenMeta: unknown;
+ private seenDict: unknown;
+ private seenState: BuckarooState | undefined;
constructor(opts: OrchestratorOptions) {
- this.ws = opts.ws;
+ this.model = opts.model;
this.minDebounceMs = opts.minDebounceMs ?? 200;
this.maxDebounceMs = opts.maxDebounceMs ?? 3000;
this.multiplier = opts.multiplier ?? 2;
- if (opts.initialAggregateMs) {
- for (const [scope, ms] of Object.entries(opts.initialAggregateMs)) {
- if (ms != null) this.lastAggregateMs.set(scope as ScopeName, ms);
- }
- }
+ this.initialRequestMs = opts.initialRequestMs ?? 250;
+ this.firstPaintTimeoutMs = opts.firstPaintTimeoutMs ?? 1500;
}
- /**
- * Current state-change token. Tests inspect this; production
- * code doesn't usually need it.
- */
- get currentToken(): number {
- return this.token;
+ /** Start watching the model, and adopt the state it already holds. */
+ start(): void {
+ if (this.started) return;
+ this.started = true;
+ this.seenMeta = this.model.get("df_meta");
+ this.seenDict = this.model.get("df_data_dict");
+ this.seenState = this.model.get("buckaroo_state");
+ this.model.on("change:df_meta", this.onModelChange);
+ this.model.on("change:df_data_dict", this.onModelChange);
+ this.model.on("change:buckaroo_state", this.onState);
+ this.model.on("msg:custom", this.onMessage);
+ this.sync();
+ this.began = true;
+ }
+
+ /** Stop watching and cancel anything scheduled. Call on unmount. */
+ stop(): void {
+ if (!this.started) return;
+ this.started = false;
+ this.model.off("change:df_meta", this.onModelChange);
+ this.model.off("change:df_data_dict", this.onModelChange);
+ this.model.off("change:buckaroo_state", this.onState);
+ this.model.off("msg:custom", this.onMessage);
+ this.standDown();
+ this.began = false;
}
/**
- * Compute the debounce delay (ms) for a given scope based on
- * the last observed aggregate compute time, clamped to
- * ``[minDebounceMs, maxDebounceMs]``.
+ * The delay (ms) before asking for a new state's stats: the last observed
+ * request time times the multiplier, clamped to
+ * `[minDebounceMs, maxDebounceMs]`. A request that took longer means the
+ * server was busy, so the next one waits longer.
*/
- computeDebounce(scope: ScopeName): number {
- const last = this.lastAggregateMs.get(scope) ?? this.defaultBaselineMs;
- const raw = last * this.multiplier;
+ computeDebounce(): number {
+ const raw = (this.lastRequestMs ?? this.initialRequestMs) * this.multiplier;
return Math.max(this.minDebounceMs, Math.min(this.maxDebounceMs, raw));
}
- /**
- * Drive a single user-initiated state change.
- *
- * 1. Bump the state token.
- * 2. Cancel any pending aggregate timers from the previous change.
- * 3. Ship the ``state_change`` message (server will reply with
- * the scalar stats fast).
- * 4. For each scope expected to have aggregate work, schedule a
- * debounced ``compute_stat_group`` request. Timer fires the
- * request only if no further state_change has bumped the
- * token meanwhile.
- */
- onStateChange(
- newState: Record,
- opts?: { scopesForAggregate?: ScopeName[] },
- ): void {
- const token = ++this.token;
- for (const t of this.aggregateTimers.values()) clearTimeout(t);
- this.aggregateTimers.clear();
-
- this.ws.send(JSON.stringify({
- type: "state_change",
- state_token: token,
- new_state: newState,
- }));
-
- const scopes = opts?.scopesForAggregate ?? ["filt"];
- for (const scope of scopes) {
- const delay = this.computeDebounce(scope);
- const tid = setTimeout(() => {
- // If a newer state_change bumped the token while we
- // were waiting, drop the request — sending it would
- // produce a stat_group_aborted from the server anyway.
- if (token !== this.token) return;
- this.ws.send(JSON.stringify({
- type: "compute_stat_group",
- state_token: token,
- scope,
- group: "aggregate",
- }));
- }, delay);
- this.aggregateTimers.set(scope, tid);
+ // A frame fires one change event per key, in the order the server wrote
+ // them (df_data_dict before df_meta), so the model is read once they have
+ // all landed, in a microtask, not at the first event.
+ private readonly onModelChange = (): void => {
+ if (this.syncQueued) return;
+ this.syncQueued = true;
+ void Promise.resolve().then(() => {
+ this.syncQueued = false;
+ if (this.started) this.sync();
+ });
+ };
+
+ private readonly onState = (next?: BuckarooState): void => {
+ const prev = this.seenState;
+ this.seenState = next ?? this.model.get("buckaroo_state");
+ if (this.gen === undefined || !touchesDataflow(prev, this.seenState)) return;
+ this.begin(this.gen);
+ };
+
+ private readonly onMessage = (msg?: { type?: string }): void => {
+ if (msg?.type === "infinite_resp") this.markPainted();
+ };
+
+ private sync(): void {
+ const meta = this.model.get("df_meta") as DFMeta | undefined;
+ const dict = this.model.get("df_data_dict");
+ const metaChanged = meta !== this.seenMeta;
+ const dictChanged = dict !== this.seenDict;
+ this.seenMeta = meta;
+ this.seenDict = dict;
+
+ const stats = meta?.stats;
+ if (stats?.status !== "pending" || typeof stats.gen !== "number") {
+ // Complete, not computed, an error, or a server that reports no stats.
+ this.standDown();
+ } else if (stats.gen !== this.gen) {
+ this.begin(stats.gen);
+ } else if (this.inFlight && dictChanged && !metaChanged) {
+ // A new df_data_dict under the same df_meta is a stats_update that
+ // is not final (a full frame for this state replaces both). One
+ // request per reply: ask again.
+ this.inFlight = false;
+ this.noteRequestTime();
+ this.arm();
}
}
- /**
- * Process a ``stat_group_result`` from the server. Stale results
- * (mismatched token) are silently ignored. Successful results
- * update the per-scope aggregate baseline used by the next
- * debounce.
- *
- * Returns true if the result was applied, false if it was stale.
- */
- onStatGroupResult(msg: StatGroupResult): boolean {
- if (msg.state_token !== this.token) return false;
- this.lastAggregateMs.set(msg.scope, msg.elapsed_ms);
- return true;
+ // Wait for rows, then ask, for `gen`. The state the model starts with asks as
+ // soon as rows are up; every later one is a state change and waits out the
+ // debounce.
+ private begin(gen: number): void {
+ this.clearTimers();
+ this.gen = gen;
+ this.inFlight = false;
+ this.delayMs = this.began ? this.computeDebounce() : 0;
+ this.paintTimer = setTimeout(() => {
+ this.paintTimer = undefined;
+ this.markPainted();
+ }, this.firstPaintTimeoutMs);
+ }
+
+ private standDown(): void {
+ if (this.inFlight) this.noteRequestTime();
+ this.clearTimers();
+ this.gen = undefined;
+ this.inFlight = false;
+ }
+
+ // The first rows are up (or the wait for them timed out): ask.
+ private markPainted(): void {
+ if (this.paintTimer !== undefined) {
+ clearTimeout(this.paintTimer);
+ this.paintTimer = undefined;
+ }
+ this.arm();
+ }
+
+ private arm(): void {
+ if (this.inFlight || this.gen === undefined || this.requestTimer !== undefined) return;
+ this.requestTimer = setTimeout(() => {
+ this.requestTimer = undefined;
+ this.fire();
+ }, this.delayMs);
+ }
+
+ private fire(): void {
+ if (requestStats(this.model)) {
+ this.inFlight = true;
+ this.sentAt = Date.now();
+ this.delayMs = 0;
+ }
+ }
+
+ private noteRequestTime(): void {
+ this.lastRequestMs = Date.now() - this.sentAt;
}
- /** Cancel all pending aggregate timers. Call on widget unmount. */
- dispose(): void {
- for (const t of this.aggregateTimers.values()) clearTimeout(t);
- this.aggregateTimers.clear();
+ private clearTimers(): void {
+ if (this.requestTimer !== undefined) clearTimeout(this.requestTimer);
+ if (this.paintTimer !== undefined) clearTimeout(this.paintTimer);
+ this.requestTimer = undefined;
+ this.paintTimer = undefined;
}
}
diff --git a/packages/buckaroo-js-core/src/server/StatsChannel.test.ts b/packages/buckaroo-js-core/src/server/StatsChannel.test.ts
new file mode 100644
index 000000000..9d2caa839
--- /dev/null
+++ b/packages/buckaroo-js-core/src/server/StatsChannel.test.ts
@@ -0,0 +1,401 @@
+/**
+ * StatsChannel — the client half of the stats wire (rows-first c2).
+ *
+ * A capable client receives a stats-free `initial_state` (df_meta.stats.status
+ * "pending"), asks for the stats and merges the `stats_update` that answers,
+ * keyed by `stats_gen`. These tests drive a WebSocketModel with a fake socket,
+ * the way the server's frames reach it.
+ */
+import { WebSocketModel } from "./WebSocketModel";
+import { withStatsCapability } from "./StatsChannel";
+import { decodeDFData } from "../components/DFViewerParts/resolveDFData";
+
+// A wide summary-stats envelope (parquet_b64, layout "wide") as the server
+// sends one; the decoder tests use the same fixture.
+// eslint-disable-next-line @typescript-eslint/no-var-requires
+const wideFixture = require("../components/DFViewerParts/test-fixtures/summary_stats_parquet_b64.json");
+
+// The real decoder, except that an envelope carrying `hold` waits on `gate`,
+// so a test can deliver a frame while a decode is in flight.
+let gate: Promise = Promise.resolve();
+jest.mock("../components/DFViewerParts/resolveDFData", () => {
+ const actual = jest.requireActual("../components/DFViewerParts/resolveDFData");
+ return {
+ ...actual,
+ decodeDFData: jest.fn(async (env: any, buffers?: DataView[]) => {
+ if (env && env.hold) await gate;
+ return actual.decodeDFData(env, buffers);
+ }),
+ };
+});
+
+class FakeSocket {
+ readyState = 1; // WebSocket.OPEN
+ onmessage: ((e: MessageEvent) => void) | null = null;
+ sent: any[] = [];
+ send(data: string) {
+ this.sent.push(JSON.parse(data));
+ }
+ deliver(msg: object) {
+ this.onmessage?.({ data: JSON.stringify(msg) } as MessageEvent);
+ }
+}
+
+// Lets every promise continuation (the payload decode) run.
+const settle = () => new Promise((resolve) => setTimeout(resolve, 0));
+
+const row = (stat: string, cells: Record) => ({ index: stat, level_0: stat, ...cells });
+
+// What the schema tier ships: identity rows only.
+const schemaStats = () => [
+ row("dtype", { a: "int64", b: "float64", c: "object" }),
+ row("length", { a: 3, b: 3, c: 3 }),
+];
+
+const metaFor = (stats?: Record) => ({
+ total_rows: 3, columns: 3, filtered_rows: 3, rows_shown: 3,
+ ...(stats === undefined ? {} : { stats }),
+});
+
+const pending = (gen: number) => ({ status: "pending", tier: "schema", gen });
+
+const frame = (gen: number | undefined, allStats: any = schemaStats()) => ({
+ type: "initial_state",
+ df_meta: metaFor(gen === undefined ? undefined : pending(gen)),
+ df_data_dict: { all_stats: allStats },
+});
+
+// The server's payload is a wide DFEnvelope; the json format decodes to the
+// same row shape without a parquet fixture.
+const update = (gen: number, rows: any[], extra: object = {}) => ({
+ type: "stats_update",
+ stats_gen: gen,
+ scope: "raw",
+ tier: "full",
+ final: true,
+ payload: { format: "json", layout: "wide", data: rows },
+ elapsed_ms: 12.5,
+ ...extra,
+});
+
+// `gen` null builds a model whose df_meta carries no stats, as an old server's does.
+function makeModel(gen: number | null = 3, allStats: any = schemaStats()) {
+ const ws = new FakeSocket();
+ const model = new WebSocketModel(ws as unknown as WebSocket, {
+ df_meta: metaFor(gen === null ? undefined : pending(gen)),
+ df_data_dict: { all_stats: allStats },
+ });
+ const events: { event: string; value: any }[] = [];
+ for (const key of ["df_data_dict", "df_meta"]) {
+ model.on(`change:${key}`, (value: any) => events.push({ event: `change:${key}`, value }));
+ }
+ return { ws, model, events };
+}
+
+describe("withStatsCapability", () => {
+ it("adds ?caps=stats_update to a bare URL", () => {
+ expect(withStatsCapability("ws://localhost:8700/ws/sales")).toBe("ws://localhost:8700/ws/sales?caps=stats_update");
+ });
+
+ it("appends to an existing query string", () => {
+ expect(withStatsCapability("ws://h/ws/s?token=abc")).toBe("ws://h/ws/s?token=abc&caps=stats_update");
+ });
+
+ it("extends an existing caps value", () => {
+ expect(withStatsCapability("ws://h/ws/s?caps=other")).toBe("ws://h/ws/s?caps=other,stats_update");
+ });
+
+ it("leaves a URL that already advertises the capability alone", () => {
+ expect(withStatsCapability("ws://h/ws/s?caps=stats_update")).toBe("ws://h/ws/s?caps=stats_update");
+ expect(withStatsCapability("ws://h/ws/s?caps=a,stats_update")).toBe("ws://h/ws/s?caps=a,stats_update");
+ });
+
+ it("keeps the fragment last", () => {
+ expect(withStatsCapability("ws://h/ws/s#frag")).toBe("ws://h/ws/s?caps=stats_update#frag");
+ });
+});
+
+describe("stats_update merge", () => {
+ it("key-merges the payload's columns into all_stats and keeps the other columns", async () => {
+ const { ws, model } = makeModel(3);
+ ws.deliver(update(3, [
+ row("length", { a: 3, b: 5 }),
+ row("mean", { a: 2, b: 4.5 }),
+ ]));
+ await settle();
+ expect(model.get("df_data_dict").all_stats).toEqual([
+ row("dtype", { a: "int64", b: "float64", c: "object" }),
+ row("length", { a: 3, b: 5, c: 3 }),
+ row("mean", { a: 2, b: 4.5 }),
+ ]);
+ });
+
+ it("does not let a null in the payload erase a value already merged", async () => {
+ const { ws, model } = makeModel(3);
+ // The wide pivot pads a stat a column did not carry with null.
+ ws.deliver(update(3, [row("dtype", { a: null, b: "float32" })]));
+ await settle();
+ const dtype = model.get("df_data_dict").all_stats.find((r: any) => r.index === "dtype");
+ expect(dtype).toEqual(row("dtype", { a: "int64", b: "float32", c: "object" }));
+ });
+
+ it("assigns a new df_data_dict and leaves the previous objects untouched", async () => {
+ const { ws, model, events } = makeModel(3);
+ const before = model.get("df_data_dict");
+ const beforeStats = before.all_stats;
+ const snapshot = JSON.parse(JSON.stringify(beforeStats));
+ ws.deliver(update(3, [row("length", { a: 99 }), row("mean", { a: 2 })]));
+ await settle();
+ const after = model.get("df_data_dict");
+ expect(after).not.toBe(before);
+ expect(after.all_stats).not.toBe(beforeStats);
+ expect(beforeStats).toEqual(snapshot);
+ const dictEvents = events.filter((e) => e.event === "change:df_data_dict");
+ expect(dictEvents).toHaveLength(1);
+ expect(dictEvents[0].value).toBe(after);
+ });
+
+ it("merges onto an all_stats that arrived as an undecoded envelope", async () => {
+ const { ws, model } = makeModel(3);
+ // A later initial_state hands the model the dict as the server sent it.
+ ws.deliver(frame(4, { format: "json", layout: "wide", data: schemaStats() }));
+ ws.deliver(update(4, [row("mean", { a: 2, b: 4.5, c: null })]));
+ await settle();
+ const stats = model.get("df_data_dict").all_stats;
+ expect(Array.isArray(stats)).toBe(true);
+ expect(stats.map((r: any) => r.index)).toEqual(["dtype", "length", "mean"]);
+ });
+
+ it("keeps the other df_data_dict keys", async () => {
+ const ws = new FakeSocket();
+ const model = new WebSocketModel(ws as unknown as WebSocket, {
+ df_meta: metaFor(pending(3)),
+ df_data_dict: { all_stats: schemaStats(), empty: [], main: [{ a: 1 }] },
+ });
+ ws.deliver(update(3, [row("mean", { a: 2 })]));
+ await settle();
+ const dict = model.get("df_data_dict");
+ expect(dict.empty).toEqual([]);
+ expect(dict.main).toEqual([{ a: 1 }]);
+ expect(dict.all_stats).toHaveLength(3);
+ });
+
+ it("merges a wide parquet_b64 payload as the server sends it", async () => {
+ const { ws, model } = makeModel(3, [row("orig_col_name", { a: "first" })]);
+ const decoded: any[] = await decodeDFData(wideFixture);
+ expect(decoded.length).toBeGreaterThan(1);
+ ws.deliver(update(3, [], { payload: wideFixture }));
+ await settle();
+ const stats = model.get("df_data_dict").all_stats;
+ expect(stats[0]).toEqual(row("orig_col_name", { a: "first" }));
+ for (const decodedRow of decoded) {
+ expect(stats).toContainEqual(decodedRow);
+ }
+ });
+
+ it("builds all_stats when the model holds no dict yet", async () => {
+ const ws = new FakeSocket();
+ const model = new WebSocketModel(ws as unknown as WebSocket, { df_meta: metaFor(pending(3)) });
+ ws.deliver(update(3, [row("mean", { a: 2 })]));
+ await settle();
+ expect(model.get("df_data_dict")).toEqual({ all_stats: [row("mean", { a: 2 })] });
+ });
+
+ it("adds all_stats to a dict that has none", async () => {
+ const ws = new FakeSocket();
+ const model = new WebSocketModel(ws as unknown as WebSocket, {
+ df_meta: metaFor(pending(3)),
+ df_data_dict: { main: [{ a: 1 }] },
+ });
+ ws.deliver(update(3, [row("mean", { a: 2 })]));
+ await settle();
+ expect(model.get("df_data_dict")).toEqual({ main: [{ a: 1 }], all_stats: [row("mean", { a: 2 })] });
+ });
+
+ it("applies updates that arrive back to back, in order", async () => {
+ const { ws, model } = makeModel(3);
+ ws.deliver(update(3, [row("mean", { a: 2 })], { final: false }));
+ ws.deliver(update(3, [row("mean", { b: 4.5 }), row("max", { a: 9 })]));
+ await settle();
+ const stats = model.get("df_data_dict").all_stats;
+ expect(stats.find((r: any) => r.index === "mean")).toEqual(row("mean", { a: 2, b: 4.5 }));
+ expect(stats.find((r: any) => r.index === "max")).toEqual(row("max", { a: 9 }));
+ });
+});
+
+describe("stats_update order", () => {
+ it("merges updates in arrival order when an earlier payload decodes slower", async () => {
+ let release: () => void = () => {};
+ gate = new Promise((resolve) => { release = resolve; });
+ try {
+ const { ws, model, events } = makeModel(3);
+ const held = { format: "json", layout: "wide", data: [row("mean", { a: 1 })], hold: true };
+ ws.deliver(update(3, [], { final: false, payload: held }));
+ ws.deliver(update(3, [row("mean", { a: 2 })]));
+ await settle(); // the second payload has decoded; the first is held
+ release();
+ await settle();
+ // The final update goes last: its value stands and the status
+ // completes only after both merges.
+ const mean = model.get("df_data_dict").all_stats.find((r: any) => r.index === "mean");
+ expect(mean.a).toBe(2);
+ expect(events.map((e) => e.event)).toEqual(["change:df_data_dict", "change:df_data_dict", "change:df_meta"]);
+ } finally {
+ release();
+ gate = Promise.resolve();
+ }
+ });
+});
+
+describe("stats_update and df_meta.stats", () => {
+ it("a final update marks the stats complete at the update's tier", async () => {
+ const { ws, model, events } = makeModel(3);
+ const metaBefore = model.get("df_meta");
+ ws.deliver(update(3, [row("mean", { a: 2 })]));
+ await settle();
+ const meta = model.get("df_meta");
+ expect(meta).not.toBe(metaBefore);
+ expect(meta.stats).toEqual({ status: "complete", tier: "full", gen: 3 });
+ expect(meta.total_rows).toBe(3);
+ expect(events.filter((e) => e.event === "change:df_meta")).toHaveLength(1);
+ });
+
+ it("a non-final update merges but leaves the status pending", async () => {
+ const { ws, model, events } = makeModel(3);
+ ws.deliver(update(3, [row("mean", { a: 2 })], { final: false }));
+ await settle();
+ expect(model.get("df_data_dict").all_stats).toHaveLength(3);
+ expect(model.get("df_meta").stats.status).toBe("pending");
+ expect(events.filter((e) => e.event === "change:df_meta")).toHaveLength(0);
+ });
+});
+
+describe("stats_gen", () => {
+ it("drops a stats_update whose stats_gen is not the expected one", async () => {
+ const { ws, model, events } = makeModel(3);
+ ws.deliver(update(2, [row("mean", { a: 2 })]));
+ ws.deliver(update(3, [row("max", { a: 9 })]));
+ await settle();
+ expect(model.get("df_data_dict").all_stats.map((r: any) => r.index)).toEqual(["dtype", "length", "max"]);
+ expect(events.filter((e) => e.event === "change:df_data_dict")).toHaveLength(1);
+ });
+
+ it("drops every stats_update when the server reported no stats", async () => {
+ const { ws, model, events } = makeModel(null);
+ ws.deliver(update(0, [row("mean", { a: 2 })]));
+ await settle();
+ expect(model.get("df_data_dict").all_stats).toHaveLength(2);
+ expect(events).toHaveLength(0);
+ });
+
+ it("clears the expectation when an initial_state carries no df_meta.stats", async () => {
+ const { ws, model } = makeModel(3);
+ ws.deliver(frame(undefined));
+ ws.deliver(update(3, [row("mean", { a: 2 })]));
+ await settle();
+ expect(model.get("df_data_dict").all_stats).toHaveLength(2);
+ });
+
+ it("advances the expected gen on a broadcast initial_state with no reply_seq", async () => {
+ const { ws, model } = makeModel(3);
+ ws.deliver(frame(4));
+ ws.deliver(update(3, [row("mean", { a: 2 })]));
+ await settle();
+ expect(model.get("df_data_dict").all_stats).toHaveLength(2);
+ ws.deliver(update(4, [row("mean", { a: 2 })]));
+ await settle();
+ expect(model.get("df_data_dict").all_stats).toHaveLength(3);
+ });
+
+ it("discards a merge when an initial_state for a newer gen arrives while it decodes", async () => {
+ const { ws, model } = makeModel(3);
+ ws.deliver(update(3, [row("mean", { a: 2 })]));
+ // Still decoding: the model has not merged anything yet.
+ ws.deliver(frame(4, [row("dtype", { a: "int32" })]));
+ await settle();
+ expect(model.get("df_data_dict").all_stats).toEqual([row("dtype", { a: "int32" })]);
+ expect(model.get("df_meta").stats).toEqual(pending(4));
+ ws.deliver(update(4, [row("mean", { a: 2 })]));
+ await settle();
+ expect(model.get("df_data_dict").all_stats).toEqual([row("dtype", { a: "int32" }), row("mean", { a: 2 })]);
+ });
+
+ it("merges onto the new dict when a same-gen initial_state replaces it while the old one decodes", async () => {
+ let release: () => void = () => {};
+ gate = new Promise((resolve) => { release = resolve; });
+ try {
+ const held = { format: "json", layout: "wide", data: schemaStats(), hold: true };
+ const { ws, model } = makeModel(3, held);
+ ws.deliver(update(3, [row("mean", { a: 2 })]));
+ await settle(); // the update is now waiting on the held decode
+ expect((decodeDFData as jest.Mock).mock.calls.some(([env]) => env === held)).toBe(true);
+ ws.deliver(frame(3, [row("dtype", { a: "int32" })]));
+ release();
+ await settle();
+ expect(model.get("df_data_dict").all_stats).toEqual([
+ row("dtype", { a: "int32" }),
+ row("mean", { a: 2 }),
+ ]);
+ expect(model.get("df_meta").stats).toEqual({ status: "complete", tier: "full", gen: 3 });
+ } finally {
+ release();
+ gate = Promise.resolve();
+ }
+ });
+});
+
+describe("stats_aborted", () => {
+ it("marks the stats failed when the run for the expected gen failed", async () => {
+ const { ws, model } = makeModel(3);
+ ws.deliver({ type: "stats_aborted", stats_gen: 3, current_gen: 3, scope: "raw", reason: "error" });
+ await settle();
+ expect(model.get("df_meta").stats).toEqual({ status: "error", tier: "schema", gen: 3, reason: "stats_failed" });
+ });
+
+ it("marks the stats not computed when the server says they cannot be requested", async () => {
+ const { ws, model } = makeModel(3);
+ ws.deliver({ type: "stats_aborted", stats_gen: 3, current_gen: 3, scope: "raw", reason: "not_requestable" });
+ await settle();
+ expect(model.get("df_meta").stats.status).toBe("not_computed");
+ });
+
+ it("changes nothing for a stale reply, an unsupported scope or a session with no data", async () => {
+ const { ws, events } = makeModel(3);
+ for (const reason of ["stale", "unsupported_scope", "no_data"]) {
+ ws.deliver({ type: "stats_aborted", stats_gen: 3, current_gen: 4, scope: "raw", reason });
+ }
+ await settle();
+ expect(events).toHaveLength(0);
+ });
+
+ it("ignores a reply for a gen the client has left", async () => {
+ const { ws, model } = makeModel(3);
+ ws.deliver({ type: "stats_aborted", stats_gen: 2, current_gen: 3, scope: "raw", reason: "error" });
+ await settle();
+ expect(model.get("df_meta").stats.status).toBe("pending");
+ });
+});
+
+describe("other messages", () => {
+ it("ignores a type it does not know and keeps applying frames", async () => {
+ const { ws, model, events } = makeModel(3);
+ expect(() => ws.deliver({ type: "something_new", stats_gen: 3 })).not.toThrow();
+ ws.deliver({ type: "error", message: "boom" });
+ await settle();
+ expect(events).toHaveLength(0);
+ ws.deliver(frame(4, [row("dtype", { a: "int32" })]));
+ expect(model.get("df_data_dict").all_stats).toEqual([row("dtype", { a: "int32" })]);
+ expect(model.get("df_meta").stats).toEqual(pending(4));
+ });
+
+ it("still pairs an infinite_resp with the binary frame that follows it", () => {
+ const { ws, model } = makeModel(3);
+ const seen: [any, DataView[]][] = [];
+ model.on("msg:custom", (msg: any, buffers: DataView[]) => seen.push([msg, buffers]));
+ ws.deliver({ type: "infinite_resp", key: { start: 0, end: 5 }, length: 5 });
+ ws.onmessage?.({ data: new ArrayBuffer(8) } as MessageEvent);
+ expect(seen).toHaveLength(1);
+ expect(seen[0][0].type).toBe("infinite_resp");
+ expect(seen[0][1][0].byteLength).toBe(8);
+ });
+});
diff --git a/packages/buckaroo-js-core/src/server/StatsChannel.ts b/packages/buckaroo-js-core/src/server/StatsChannel.ts
new file mode 100644
index 000000000..0ead005ed
--- /dev/null
+++ b/packages/buckaroo-js-core/src/server/StatsChannel.ts
@@ -0,0 +1,210 @@
+/**
+ * StatsChannel — the client half of the stats wire (rows-first c2).
+ *
+ * A client that advertises `?caps=stats_update` gets a first `initial_state`
+ * whose `df_meta.stats` says the stats are pending, asks for them with
+ * `stats_request {stats_gen, scope}`, and receives either
+ *
+ * stats_update {stats_gen, scope, tier, final, payload, elapsed_ms}
+ * stats_aborted {stats_gen, current_gen?, scope, reason}
+ *
+ * `payload` is an inline wide DFEnvelope holding `all_stats`. `stats_gen` is the
+ * server's counter for the state the stats describe; it rides on every
+ * `initial_state` as `df_meta.stats.gen`, and a reply for any other gen is for
+ * a state the client has left. Sending the request is the scheduler's job, not
+ * this module's.
+ *
+ * Merge semantics are WebSocket-only: `WebSocketModel` hands every frame to
+ * `handle()` first, and Jupyter's widget sets `df_data_dict` whole.
+ */
+import { decodeDFData } from "../components/DFViewerParts/resolveDFData";
+import { DFData, DFDataOrPayload } from "../components/DFViewerParts/DFWhole";
+import { DFMeta, DFMetaStats } from "../components/WidgetTypes";
+import { IModel } from "./IModel";
+
+/** The capability this client advertises, as one value of `?caps=` on the
+ * WebSocket URL: it merges `stats_update` messages. The server records it per
+ * connection when the socket opens, since it sends the first message before
+ * the client can say anything. */
+export const STATS_UPDATE_CAP = "stats_update";
+
+const decodeQueryValue = (value: string): string => {
+ try {
+ return decodeURIComponent(value);
+ } catch {
+ return value;
+ }
+};
+
+/** `wsUrl` with `caps=stats_update` added: a new query on a bare URL, a new
+ * parameter after an existing query, or a comma-joined value when the host
+ * already passes `caps`. The fragment stays last and other parameters are
+ * left as the host wrote them. */
+export function withStatsCapability(wsUrl: string): string {
+ const hashAt = wsUrl.indexOf("#");
+ const fragment = hashAt === -1 ? "" : wsUrl.slice(hashAt);
+ const beforeFragment = hashAt === -1 ? wsUrl : wsUrl.slice(0, hashAt);
+ const queryAt = beforeFragment.indexOf("?");
+ const path = queryAt === -1 ? beforeFragment : beforeFragment.slice(0, queryAt);
+ const params = queryAt === -1 ? [] : beforeFragment.slice(queryAt + 1).split("&").filter((p) => p !== "");
+
+ const capsAt = params.findIndex((p) => p === "caps" || p.startsWith("caps="));
+ if (capsAt === -1) {
+ params.push(`caps=${STATS_UPDATE_CAP}`);
+ } else {
+ const caps = decodeQueryValue(params[capsAt].slice("caps=".length))
+ .split(",")
+ .map((cap) => cap.trim())
+ .filter((cap) => cap !== "");
+ if (!caps.includes(STATS_UPDATE_CAP)) caps.push(STATS_UPDATE_CAP);
+ params[capsAt] = `caps=${caps.join(",")}`;
+ }
+ return `${path}?${params.join("&")}${fragment}`;
+}
+
+export interface StatsUpdateMessage {
+ type: "stats_update";
+ stats_gen: number;
+ scope?: string;
+ tier?: string;
+ final?: boolean;
+ payload?: DFDataOrPayload;
+ elapsed_ms?: number;
+}
+
+export interface StatsAbortedMessage {
+ type: "stats_aborted";
+ /** The request's gen. */
+ stats_gen?: number;
+ /** The session's gen, omitted when the session has no data. */
+ current_gen?: number;
+ scope?: string;
+ reason?: "stale" | "unsupported_scope" | "not_requestable" | "error" | "no_data";
+}
+
+/**
+ * Key-merge a stats payload into `all_stats`. Both are row-per-stat tables
+ * (`{index: , : , ...}`). For each stat row the payload
+ * names, every column it carries replaces that cell, columns it does not carry
+ * keep theirs, and a stat the table lacks is appended.
+ *
+ * A `null` in the payload never replaces a value already merged: the wide
+ * pivot fills a stat a column did not carry in this message with `null`.
+ *
+ * Returns new row objects and a new array. `base` may be the decoder's cached
+ * array, so nothing in it is touched.
+ */
+export function mergeStatRows(base: DFData, update: DFData): DFData {
+ const merged = base.slice();
+ const at = new Map();
+ merged.forEach((row, i) => at.set(row.index, i));
+
+ for (const updateRow of update) {
+ const i = at.get(updateRow.index);
+ if (i === undefined) {
+ at.set(updateRow.index, merged.length);
+ merged.push({ ...updateRow });
+ continue;
+ }
+ const row = { ...merged[i] };
+ for (const [column, value] of Object.entries(updateRow)) {
+ if (column === "index" || column === "level_0") continue;
+ if (value === null && row[column] != null) continue;
+ row[column] = value;
+ }
+ merged[i] = row;
+ }
+ return merged;
+}
+
+const genOf = (meta: DFMeta | undefined): number | undefined => {
+ const gen = meta?.stats?.gen;
+ return typeof gen === "number" ? gen : undefined;
+};
+
+/** What the channel needs of a model. */
+type StatsModel = Pick;
+
+export class StatsChannel {
+ // Updates are applied one at a time: each reads the dict the previous one
+ // wrote, so a later update cannot overwrite an earlier one's merge.
+ private applying: Promise = Promise.resolve();
+
+ constructor(private model: StatsModel) {}
+
+ /**
+ * The stats_gen of the state the client is showing, or `undefined` when the
+ * server reports none; a `stats_request` carries it. It is read off the
+ * model's `df_meta` each time, so it starts from the frame the model was
+ * built from and every applied `initial_state` (a broadcast frame with no
+ * `reply_seq` included) moves it. A `df_meta` with no `stats.gen` leaves
+ * nothing expected, so a server that stops reporting stats cannot have a
+ * late `stats_update` merged. Reading the model, not the message, keeps this
+ * correct under any scheme that drops stale `initial_state` frames.
+ */
+ get expectedGen(): number | undefined {
+ return genOf(this.model.get("df_meta"));
+ }
+
+ /** Consume a stats message. Returns false for every other message type,
+ * which the model handles (or ignores) as before. */
+ handle(msg: { type?: string }): boolean {
+ if (msg.type === "stats_update") {
+ this.receiveUpdate(msg as StatsUpdateMessage);
+ return true;
+ }
+ if (msg.type === "stats_aborted") {
+ this.receiveAborted(msg as StatsAbortedMessage);
+ return true;
+ }
+ return false;
+ }
+
+ private receiveUpdate(msg: StatsUpdateMessage): void {
+ if (msg.stats_gen !== this.expectedGen) return;
+ this.applying = this.applying
+ .then(() => this.applyUpdate(msg))
+ .catch((e) => console.error("[StatsChannel] stats_update failed:", e));
+ }
+
+ private async applyUpdate(msg: StatsUpdateMessage): Promise {
+ const update = await decodeDFData(msg.payload);
+ for (;;) {
+ // A frame may have moved the gen on while something decoded.
+ if (msg.stats_gen !== this.expectedGen) return;
+ const dict: Record | null | undefined = this.model.get("df_data_dict");
+ // The dict is decoded when the seed built it and raw when a later
+ // initial_state did.
+ const base = await decodeDFData(dict?.all_stats);
+ // A frame replaced the dict (or moved the gen) while it decoded:
+ // start over from the new one.
+ if (dict !== this.model.get("df_data_dict") || msg.stats_gen !== this.expectedGen) continue;
+ this.model.set("df_data_dict", { ...dict, all_stats: mergeStatRows(base, update) });
+ if (msg.final) {
+ this.replaceStats((stats) => ({ status: "complete", tier: msg.tier ?? stats.tier, gen: stats.gen }));
+ }
+ return;
+ }
+ }
+
+ private receiveAborted(msg: StatsAbortedMessage): void {
+ // Only a reply to a request for the state on screen says anything about
+ // it. A `stale` reply is answered by the initial_state that carries the
+ // new gen; taking `current_gen` without that frame would merge stats
+ // for a state the client has not seen into the one it shows.
+ if (msg.stats_gen !== this.expectedGen) return;
+ if (msg.reason === "error") {
+ this.replaceStats((stats) => ({ ...stats, status: "error", reason: "stats_failed" }));
+ } else if (msg.reason === "not_requestable") {
+ this.replaceStats((stats) => ({ ...stats, status: "not_computed" }));
+ }
+ }
+
+ /** Replace `df_meta.stats` in a new `df_meta`, the reference that c0a's
+ * `inFlight` rule and the pinned rows' placeholders key on. */
+ private replaceStats(change: (stats: DFMetaStats) => DFMetaStats): void {
+ const meta: DFMeta | undefined = this.model.get("df_meta");
+ if (!meta?.stats) return;
+ this.model.set("df_meta", { ...meta, stats: change(meta.stats) });
+ }
+}
diff --git a/packages/buckaroo-js-core/src/server/WebSocketModel.ts b/packages/buckaroo-js-core/src/server/WebSocketModel.ts
index 1e1fd4a2d..0d249081d 100644
--- a/packages/buckaroo-js-core/src/server/WebSocketModel.ts
+++ b/packages/buckaroo-js-core/src/server/WebSocketModel.ts
@@ -13,21 +13,32 @@
* Binary protocol (matching anywidget's msg + buffers pattern):
* Server sends a JSON text frame (infinite_resp), then a binary frame (Parquet).
* This class pairs them and emits "msg:custom" with (msg, [DataView]).
+ *
+ * stats_update and stats_aborted frames go to `stats` (see StatsChannel), and
+ * `scheduler` asks for the stats a session defers (see StateOrchestrator).
*/
+import { StatsChannel } from "./StatsChannel";
+import { StateOrchestrator } from "./StateOrchestrator";
+
export class WebSocketModel {
private ws: WebSocket;
private pendingMsg: any = null;
private handlers: Map> = new Map();
private state: Record;
private pendingChanges: Set = new Set();
+ readonly stats: StatsChannel;
+ readonly scheduler: StateOrchestrator;
constructor(ws: WebSocket, initialState: Record) {
this.state = { ...initialState };
this.ws = ws;
+ this.stats = new StatsChannel(this);
+ this.scheduler = new StateOrchestrator({ model: this });
this.ws.onmessage = (event: MessageEvent) => {
if (typeof event.data === "string") {
const msg = JSON.parse(event.data);
+ if (this.stats.handle(msg)) return;
if (msg.type === "infinite_resp") {
// Expect a following binary frame — stash this JSON
@@ -60,6 +71,9 @@ export class WebSocketModel {
}
}
};
+
+ // Idle unless df_meta.stats says pending, which no default session does.
+ this.scheduler.start();
}
send(msg: any): void {
diff --git a/packages/buckaroo-js-core/src/server/latestDictDecoder.ts b/packages/buckaroo-js-core/src/server/latestDictDecoder.ts
new file mode 100644
index 000000000..16c4c3ef9
--- /dev/null
+++ b/packages/buckaroo-js-core/src/server/latestDictDecoder.ts
@@ -0,0 +1,34 @@
+import { decodeDFDataDict } from "../components/DFViewerParts/resolveDFData";
+import { DFData, DFDataOrPayload } from "../components/DFViewerParts/DFWhole";
+
+export type RawDFDataDict = Record | undefined | null;
+
+/**
+ * Decode df_data_dict values arriving from a model and hand `apply` only the
+ * newest result.
+ *
+ * A full `initial_state` reaches a view twice, as `change:df_data_dict` and
+ * then as `metadata`, and both handlers read the same dict off the model. The
+ * returned function ignores a dict it has already seen (same reference), so
+ * that frame decodes once. Decoding is async, so a slow decode of an earlier
+ * dict can finish after a later one; each call takes a token, and a result is
+ * applied only if no later call has started.
+ *
+ * `seen` is a dict the caller has already applied (the view's seed), so the
+ * first read of that same dict off the model does not decode it again.
+ */
+export function makeLatestDictDecoder(
+ apply: (decoded: Record) => void,
+ seen?: RawDFDataDict,
+): (raw: RawDFDataDict) => void {
+ let latest = 0;
+ let lastRaw: RawDFDataDict = seen;
+ return (raw) => {
+ if (raw === lastRaw) return;
+ lastRaw = raw;
+ const mine = ++latest;
+ decodeDFDataDict(raw).then((decoded) => {
+ if (mine === latest) apply(decoded);
+ });
+ };
+}
diff --git a/packages/buckaroo-js-core/src/stories/StatsPendingPinnedRows.stories.tsx b/packages/buckaroo-js-core/src/stories/StatsPendingPinnedRows.stories.tsx
new file mode 100644
index 000000000..22e3e0c44
--- /dev/null
+++ b/packages/buckaroo-js-core/src/stories/StatsPendingPinnedRows.stories.tsx
@@ -0,0 +1,128 @@
+/**
+ * Story for the two-message protocol's client states (rows-first c0a).
+ *
+ * The server sends the first message once the schema and row count are known,
+ * and the summary stats later or never. `df_meta.stats.status` says which:
+ *
+ * - "pending" pinned keys with no value show a placeholder row
+ * - "not_computed" pinned keys with no value are omitted
+ * - "complete" stats are present; also the meaning of a missing
+ * `df_meta.stats`, as servers without the field send
+ *
+ * Column `a` is color-mapped, so the story also shows that its cells restyle
+ * when the histogram bins arrive. Used by stats-pending-pinned-rows.spec.ts.
+ */
+import type { Meta, StoryObj } from "@storybook/react";
+import React, { useMemo, useState } from "react";
+import { DFViewerInfiniteDS } from "../components/BuckarooWidgetInfinite";
+import { DFData, DFViewerConfig } from "../components/DFViewerParts/DFWhole";
+import { IDisplayArgs } from "../components/DFViewerParts/gridUtils";
+import { KeyAwareSmartRowCache, PayloadResponse } from "../components/DFViewerParts/SmartRowCache";
+import { DFMeta } from "../components/WidgetTypes";
+
+type Status = "pending" | "not_computed" | "complete";
+const STATUSES: Status[] = ["pending", "not_computed", "complete"];
+
+const mainData: DFData = [
+ { index: 0, a: 1, b: "x" },
+ { index: 1, a: 2, b: "y" },
+ { index: 2, a: 3, b: "z" },
+ { index: 3, a: 4, b: "w" },
+ { index: 4, a: 5, b: "v" },
+];
+
+const completeStats: DFData = [
+ { index: "dtype", a: "int64", b: "object" },
+ { index: "mean", a: 3, b: "N/A" },
+ { index: "histogram_bins", a: [0, 1, 2, 3, 4, 5], b: [] },
+];
+
+const viewerConfig: DFViewerConfig = {
+ column_config: [
+ {
+ col_name: "a",
+ header_name: "a",
+ displayer_args: { displayer: "obj" },
+ color_map_config: { color_rule: "color_map", map_name: "BLUE_TO_YELLOW", val_column: "a" },
+ },
+ { col_name: "b", header_name: "b", displayer_args: { displayer: "obj" } },
+ ],
+ left_col_configs: [{ col_name: "index", header_name: "index", displayer_args: { displayer: "obj" } }],
+ pinned_rows: [
+ { primary_key_val: "dtype", displayer_args: { displayer: "obj" } },
+ { primary_key_val: "mean", displayer_args: { displayer: "obj" } },
+ ],
+};
+
+const displayArgs: Record = {
+ main: { data_key: "main", df_viewer_config: viewerConfig, summary_stats_key: "all_stats" },
+};
+
+const StatsPendingPinnedRowsInner: React.FC = () => {
+ const [status, setStatus] = useState("pending");
+
+ const src = useMemo(() => {
+ const cache = new KeyAwareSmartRowCache((pa) => {
+ const resp: PayloadResponse = {
+ key: pa,
+ data: mainData.slice(pa.start, Math.min(pa.end, mainData.length)),
+ length: mainData.length,
+ };
+ setTimeout(() => cache.addPayloadResponse(resp), 10);
+ });
+ return cache;
+ }, []);
+
+ const df_meta = useMemo(
+ () =>
+ ({
+ total_rows: mainData.length,
+ columns: 2,
+ filtered_rows: mainData.length,
+ rows_shown: mainData.length,
+ stats: { status },
+ }) as DFMeta,
+ [status],
+ );
+ const df_data_dict = useMemo(
+ () => ({
+ main: [] as DFData,
+ all_stats: status === "complete" ? completeStats : ([] as DFData),
+ empty: [] as DFData,
+ }),
+ [status],
+ );
+
+ return (
+