diff --git a/plans/ADR-007-tallyman-owned-materialization.md b/plans/ADR-007-tallyman-owned-materialization.md new file mode 100644 index 0000000..404cd5d --- /dev/null +++ b/plans/ADR-007-tallyman-owned-materialization.md @@ -0,0 +1,498 @@ +# ADR: Tallyman owns result materialization (no xorq cache nodes in builds) + +- **Status:** Proposed (2026-09-18, revised 2026-09-20 in the grilling + session, which added the governing rule, decision D10, and the resolution + recorded under D5). Supersedes two decisions of + `plans/ADR-006-read-path-loads-builds.md`: its D4 (chaining inlines the + parent's cache node) and its D8 (the manifest records the snapshot key and + reads assert it). Five other ADR-006 decisions keep their intent: D5 (the + canonical sort), D6 (a missing build is a hard error), D7 (verification runs + in production and is loud), D10 (an unfaithful heal wipes the entry's + Buckaroo state) and D12 (unfaithful entries are pinned and badged). The last + three attach to `ensure_materialized`, which is this ADR's D5. +- **Reading decision labels:** a bare label such as "D5" in this document + always means this ADR's own decision. Another ADR's decision is always + written with its ADR number and a few words saying what it decides. +- **Context:** the 2026-09-18 cache audit (tallyman @ `a748ea6`, buckaroo + 0.15.4, xorq 0.3.26). One finding is filed upstream of tallyman: + buckaroo-data/buckaroo#972 (`/load_expr` has no `cache_dir`). The direction + was set by Paddy the same day: "I want to depend on xorq as little as + possible for caching." +- **Affected code:** `src/tallyman_xorq/source_cache.py` (`rewrite_for_build`), + `src/tallyman_xorq/result_cache.py` (`_resolve_result_plan`, + `cached_result_expr`, `entry_graph_expr`, `baked_snapshot_path`, + `_cached_node_path`, `_assert_recorded_snapshot_key`, `_verify_self_heal`), + `src/tallyman_xorq/build.py` (the execute-once step, `build.py:496-557`), + `src/tallyman_xorq/io.py` (`tracked_expr_from_alias`, + `pinned_expr_from_alias`), `src/tallyman_xorq/portable.py` + (`rewrite_cache_dirs`), `src/tallyman_core/manifest.py` (`snapshot_key`), + `src/tallyman_companion/buckaroo_lifecycle.py` (`load_session`), + `docs/system-contract.md`. +- **Related ADRs:** `plans/ADR-002-source-identity-content-hash.md` (content in + the path; D3 reuses the device), `plans/ADR-003-result-cache-cost-rubric.md` + (its motivating misclassification; see Consequences), + `plans/ADR-008-row-order-of-reads.md` and + `plans/ADR-009-digest-stability.md` (the other two hash- or digest-changing + decisions that share this ADR's corpus rebuild). +- **Evidence:** `scripts/spike_bare_read_chaining.py` (results under D3), and + the audit measurements quoted in the Problem section. + +## Terms + +- **Entry:** one catalog computation, stored under its content hash. +- **Build:** the `xorq_build/` directory inside an entry, which is the entry's + computation graph written to disk by xorq. +- **Replay a build:** load that directory back into a live expression and + execute it. +- **Materialize:** run an entry's computation once and write the result to a + parquet file. That file is the entry's **snapshot**, and it lives under + `compute_cache/result_cache/`. +- **Worthy entry:** an entry tallyman materializes, because its graph does work + that is expensive or that cannot inherit a row order (an aggregate, join, + sort, window function, UDF, union). **Cheap entry:** one it does not, whose + small plan re-runs on every read. +- **Cache node:** xorq's `CachedNode`, a marker inside a graph meaning "look for + a file with this key, and if it is missing run the graph below me and write + it". +- **Chaining:** a recipe building on another entry through + `tracked_expr_from_alias` or `pinned_expr_from_alias`. +- **Heal:** re-create a snapshot that is missing from disk by re-running the + entry's build. +- **Session:** one grid's state inside the Buckaroo server process. +- **View build:** a build whose whole graph is one step, "read this parquet + file". +- **Live diff:** the compare grid shown when two versions are diffed. A + **promoted diff** is one that has been saved as a catalog entry. +- **Checkpoint:** tallyman's step that zips new entries and makes one git + commit in the catalog repository. + +## Problem + +Tallyman's result cache is xorq's cache. `rewrite_for_build` +(`source_cache.py:241-243`) wraps every worthy expression in +`.cache(cache=ParquetSnapshotCache(...))`, and that one call decides everything +else: the file's name is xorq's snapshot key, its directory is a base path that +xorq does not serialize, its writer is xorq's `ParquetStorage`, it is written as +a side effect of `loaded.count().execute()` (`build.py:540`, +`result_cache.py:733`), a missing ancestor is regenerated by a nested cache node +noticing its own file is gone, and Buckaroo is handed a build that contains all +of this. + +The audit found the following, each a consequence of that arrangement. + +1. **Snapshots escape the project.** xorq serializes a cache node's + `relative_path` and never its base path, and `load_expr(cache_dir=...)` + redirects only the root cache node. + - Buckaroo's `/load_expr` calls `load_expr(build_dir)` with no `cache_dir`, + so the viewer never reads tallyman's snapshot. The first view re-executes + the whole graph inside Buckaroo and writes a second copy under + `~/.cache/xorq/result_cache/`. On the dev machine that directory held 68 + files and 14 GB, with the same keys and sizes as the project's + `compute_cache` (buckaroo#972). + - `build.py:496` loads the new build with `cache_dir` but without + `rewrite_cache_dirs`, so nested ancestor nodes resolve to `~/.cache/xorq` + while a child builds. Every child build re-executes its expensive parent + and leaves a duplicate there. Reproduced in scratch homes where Buckaroo + never ran. +2. **The writer is unsafe under concurrency.** `ParquetStorage` writes through + a fixed `.parquet.tmp`. Four concurrent cold writers of one key gave + three `FileNotFoundError`s and a corrupt final file, and because a hit is + decided by file existence the corrupt file is then served forever. + `_heal_lock` covers threads in one process only, builds take no lock, and + FastMCP runs sync tools in a threadpool, so parallel tool calls do build + concurrently. +3. **The file shape is poor.** One row group per DataFusion batch (at most + 8,192 rows), Snappy. A real 3.68 GB snapshot has 9,525 row groups and a + 44.9 MB footer. `cached_result_expr` constructs a fresh read on every call + (`result_cache.py:742`), so that footer is opened again for every page + request, and every call registers another table in the shared backend + (40 reads, 40 tables). +4. **The bytes depend on how the node was executed.** xorq stamps provenance + metadata only when the cache node is the root of the executed expression: + 1,234 bytes against 858 for the same rows. Build and heal agree today only + because both happen to call `count()`. +5. **Inlined chaining misclassifies children.** Under ADR-006 decision D4 + (chaining inlines the parent's cache node) a child's + graph contains its parent's Aggregate and canonical Sort, so a filter over a + cached aggregate is itself worthy and writes a full copy. This is ADR-003's + motivating bug. Graphs also grow with chain depth, which is the cost the + pre-#74 per-entry parquet boundary existed to avoid (#82). +6. **Tallyman carries code whose only job is to aim xorq's cache.** + `rewrite_cache_dirs`, the contract's "every loader must supply `cache_dir`", + `manifest.snapshot_key` with `_assert_recorded_snapshot_key` (ADR-006 + decision D8, the snapshot-key tripwire), + and `_cached_node_path`. + +xorq is frozen upstream for this project, so every item above needs a +tallyman-side workaround for as long as xorq's cache is in the path. + +## Governing rule + +Stated by Paddy in the grilling session (2026-09-19 and 2026-09-20), and the +reason several decisions below are consequences and not separate choices: + +> Buckaroo is a very good displayer. It runs queries only for summary stats, +> sorting and paging. Anything more is done by tallyman first. If something +> tallyman asked Buckaroo to display does not exist, that is on tallyman. + +> When the MCP asks tallyman to create an entry, tallyman should run the query +> and materialize the parquet if necessary immediately. Tallyman shouldn't call +> Buckaroo to display an entry until the original query has finished. + +> I want a cohesive system that works reliably, then we can worry about speed +> problems as they come up. We don't have a cohesive system now. + +That last statement sets the priority for all three ADRs of this set (this one, +`plans/ADR-008-row-order-of-reads.md` and `plans/ADR-009-digest-stability.md`): +where a uniform rule and a faster special case compete, the uniform rule is the +decision and the faster path is noted for later. + +So tallyman runs an entry's computation, or a diff, to completion before it +asks Buckaroo to show anything. No other process runs an entry's expensive +computation, writes result files, or repairs tallyman's cache. Besides being +simpler, this puts every failure of a computation in tallyman's process, where +it can be logged and reported, and none inside a grid query in Buckaroo. + +## Decisions + +### D1. Builds carry no cache nodes + +`rewrite_for_build` keeps the in-memory-read rejection and the canonical sort +(ADR-006 decision D5, so the sort stays inside the build) and stops calling +`.cache()`: + +```python +if _is_worthy_expr(expr): + expr = _canonical_sorted(expr) +``` + +The source-read injection (`source_cache.py:215-234`) is deleted with it. It has +no producer today: `deferred_read_csv` is a build error and no JSON reader +exists. A recipe that arrives already containing a `CachedNode` becomes a build +error next to the in-memory check, because a node with default storage would +write under `~/.cache/xorq`. + +`source_cache.py:229` and `:243` are the only places in `src/` that create a +cache node, so this one change removes xorq's cache from tallyman's data path. +xorq remains the expression, build, load, hashing and execution layer. + +Removing the node alone would turn every worthy entry into a recompute entry +(`_resolve_result_plan` treats "no `CachedNode` on top" as recompute, +`result_cache.py:642-649`), so D2 to D6 land in the same change. + +*Rejected:* patch ADR-006 decision D4 (inlined chaining) in place. Deep-rewrite +the build's load (a one-line +change, verified to stop the build leak with no hash change), wait for +buckaroo#972, put a cross-process lock around xorq's writer, and pre-seed +xorq's files with a better writer (a hit is a bare existence check, so tallyman +can write `.parquet` itself). This keeps every hash and needs no rebuild. +It also keeps the file's name as xorq's tokenization of the graph (so the +snapshot-key tripwire of ADR-006 decision D8 stays), keeps item 5, and leaves each fix as a workaround for +behaviour in a dependency that cannot be changed from here. + +### D2. A snapshot's location is a function of the content hash + +`snapshot_path(project, content_hash)` returns +`compute_cache/result_cache/.parquet`. Whether an entry has a +snapshot comes from the manifest's `cache_worthy`, as it does now. + +`cached_result_expr` keeps its name and signature (ADR-006 decision D2, "keep +the facade"). For a worthy +entry it returns one bare read of `snapshot_path`, memoized per +`(project, content_hash)` for the life of the process, so repeated reads stop +registering tables and stop re-opening the footer. A worthy entry whose file +exists is served without loading its build at all. For a cheap entry it returns +the loaded graph, as now. + +Deleted: `manifest.snapshot_key`, `_assert_recorded_snapshot_key`, +`_cached_node_path`, `rewrite_cache_dirs`, and the `cache_dir` argument. With +the path computed from the hash there is no second derivation for a tripwire to +compare against. + +*Rejected:* `-.parquet`, which would make a child name +its parent's exact bytes. An unfaithful heal (a heal whose result does not +match the recorded digest, ADR-006 decision D12) would then write a +file no existing child can find, and a flagged condition would become a hard +failure for the whole subtree. The digest stays in the manifest, where verify +already looks. + +### D3. Chaining through a worthy parent is a bare read of its snapshot + +`tracked_expr_from_alias` and `pinned_expr_from_alias` return +`deferred_read_parquet(snapshot_path(parent))` for a worthy parent, after +making sure the file exists (D5). For a cheap parent they return the parent's +loaded graph, which by D1 contains no cache nodes. `entry_graph_expr` and +`cached_result_expr` become the same function. + +The child's hash covers the literal path of the parent's snapshot, and that +path contains the parent's content hash. A child's identity is therefore a +function of its parent's identity, which is the device ADR-002 already uses for +sources: if you want content identity, put it in the path. + +Measured with `scripts/spike_bare_read_chaining.py` (raw xorq, no cache nodes): + +| Question | Result | +| --- | --- | +| Child hash, snapshot rewritten with other bytes at the same path | unchanged | +| Child hash, same bytes at another path | changes | +| `build_expr` of a child while the parent's snapshot is absent | raises `FileNotFoundError: local path does not exist` | +| `load_expr` of the child's build while the snapshot is absent | succeeds | +| Executing it in that state | raises `ValueError: At least one path is required` | +| Parent re-materialized from its own build, then the child executed | digest unchanged; child rows match the reference (19,777) | +| `classify_build` of a filter + computed column over the snapshot | cheap | +| `classify_build` of the same child with the parent graph inlined | worthy (`ops:Aggregate,Sort,SortKey`) | +| Child `expr.yaml` | carries the literal path; 3,434 bytes against 6,903 inlined | +| Files written under `XORQ_CACHE_DIR` | none | + +So the graph is cut at every worthy entry, a child of an aggregate is cheap, +and the literal path in `expr.yaml` means `make_portable_inplace` handles it +with the existing `${TALLYMAN_PROJECT_ROOT}` placeholder. + +*Rejected:* keep inlining the parent's graph without its cache node. Every read +of every descendant would re-run the parent's expensive subgraph. + +### D4. One writer, used by the build and by every heal + +`materialize(project, content_hash)` loads the entry's build, executes it as a +record-batch stream, and writes the snapshot itself: + +- a unique temp name in the destination directory, then `os.replace`; +- under a per-hash cross-process `flock`, with the existence check repeated + inside the lock, so a second process waits and then finds the file; +- it numbers the rows as it writes them, in a last column named `__row_order` + (decision D2 of `plans/ADR-008-row-order-of-reads.md`); +- it returns the digest of what it wrote. + +The build's execute-once step and every heal call this function, so the +contract's I2 ("result bytes are manufactured exactly once; a self-heal is +reproduce-and-verify") holds because there is one routine that manufactures +bytes. The file format and the digest definition are ADR-009's. + +*Rejected:* keep `ParquetStorage` behind a tallyman lock. That fixes the race +and keeps items 3 and 4. + +### D5. One entry point makes files exist: `ensure_materialized` + +`ensure_materialized(project, content_hash)` guarantees that every snapshot an +entry's plan reads is on disk before anything executes: + +1. If the entry is worthy and its snapshot exists, return. No build is loaded. +2. Otherwise load the entry's build and collect the snapshot paths its `Read` + nodes point at (any read under `compute_cache/result_cache/`). The list is + kept with the loaded plan in the existing LRU. +3. For each of those that is missing, recurse on the hash in its file name. +4. If the entry is worthy, `materialize` it. + +Every file it writes is verified against the manifest's `result_digest` before +it is served. A mismatch takes the existing path: a durable `unfaithful_heal` +record and the SSE event (ADR-006 decision D7, loud verification), the +stat-cache wipe and session eviction (ADR-006 decision D10), and the pin and +badge (ADR-006 decision D12). + +Callers: the canonical read (`cached_result_expr`, on every call, where step 1 +or a handful of `stat` calls is the whole cost, and which covers diff +composition), chaining at mint time (D3), `load_session` (D6), and the verify +sweep. + +ADR-006 decision D4 rejected bare-read chaining because "builds stay +non-self-contained and the pre-heal choreography stays load-bearing forever". +What that decision bought was a build that repairs its own ancestors when +something other than tallyman executes it cold. Buckaroo is the only such +thing, and the comment at `buckaroo_lifecycle.py:570-575` says the design +relied on it: "Buckaroo's replay of the build regenerates any evicted snapshot +through ordinary cache mechanics on first query." Under the governing rule +Buckaroo never does that, so the property has no user and nothing is given up. +This was question 1 of the grilling session, resolved 2026-09-20. + +The old pre-heal was also weaker than this function. It was a +`cached_result_expr` call with a discarded result at one call site, and +ancestors were healed only because reads then re-executed recipes. Here the +requirement is a precondition of the one canonical read that every in-process +consumer already uses (contract invariant I3, "one read semantics"), and it is +computed from the build's own reads, so it cannot drift from what execution +opens. + +*Rejected:* derive the required snapshots from `manifest.parents`. It needs only +JSON reads, but it depends on the recorded edges being complete and on each +parent's recorded worthiness matching what chaining did when the child was +minted. Those edges are recalc policy and are not complete today: a promoted +diff's recipe calls `build_diff_expr`, which reads both sides through +`cached_result_expr` and records no parent edge at all. The build's reads are +exactly what the engine will open. + +### D6. Buckaroo is handed something that already exists + +This is the governing rule applied to entry grids. + +For a worthy entry `load_session` calls `ensure_materialized`, then posts a +view build of the snapshot, written once to a stable per-entry directory +(Buckaroo's stat-cache keys include the build directory's path, which is why +the expanded build already lives at a stable path). Buckaroo never executes an +aggregate, join or sort on tallyman's behalf and never writes a snapshot, and +tallyman stops depending on buckaroo#972, whose premise (teach Buckaroo where +tallyman's cache is) the rule contradicts. The grid and `/api/data` read the +same file (contract invariant I5, "one question, one path"). + +For a cheap entry it calls `ensure_materialized` and posts the entry's own +expanded build. A cheap build is a view in the database sense: a stored +definition (filter these rows, keep these columns) over files that exist. +Buckaroo's stats, sorts and pages run through it. Tallyman has already executed +that plan once, in full, at build time, so an error in it has already surfaced +in tallyman. + +Paddy confirmed this reading on 2026-09-20 (question 4 of the grilling +session). The words that decide it are "materialize the parquet if necessary": +a file is written when the entry is created and only for a worthy entry, and +nothing is written when an entry is viewed. The alternative, writing a file the +first time a cheap entry is opened so that Buckaroo only ever reads one file, +was set aside. It costs a wait on first open and a full copy per viewed +revision; the audit measured 19 GB of cache against 779 MB of data when every +CSV revision wrote a copy. + +The timing half of the rule already holds for creation. The build executes the +entry once before it writes the manifest, a worthy entry's snapshot is written +in that step (D4), and an entry with no manifest is treated as absent, so +Buckaroo cannot be asked to display an entry whose query is still running. The +same ordering now covers a snapshot that was deleted later: `load_session` +waits for `ensure_materialized` before it posts anything. + +Deleting a snapshot (the Cache page, a reset prune, a future budget eviction) +ends every live Buckaroo session whose plan reads it first, using the +session-eviction hook that ADR-006 decision D10 introduced. The next +`/api/session` re-materializes and opens a new session. Without this, a page +request against a deleted file would fail inside Buckaroo with the +`At least one path is required` error above, which is exactly the kind of +failure the rule says belongs to tallyman. + +This resembles what #104 removed: #102's viewer build over +`/result.parquet`. #104's objection was two materialized copies per +entry and two read paths. Here there is one copy and the view build only points +at it. + +*Rejected:* post the entry's own build for worthy entries too. With no cache +node in it, Buckaroo would re-run the full computation for every page and stat +query. + +### D7. The cold state is an empty `compute_cache` + +The contract's cold seam today is the `cache_dir` argument, and the standing +tests exercise it by deleting `compute_cache/` and reading again +(`tests/test_lineage_faithful_reads.py`). That remains the test: with +`compute_cache/` removed, the canonical read must reproduce every snapshot the +entry needs, each with its recorded digest. The `cache_dir` parameter leaves +the contract, since nothing in a build resolves through it any more. + +### D8. A sentinel test keeps xorq's cache out of the path + +With `XORQ_CACHE_DIR` pointing at an empty sentinel directory, a build, a +chained child build, a view, an eviction and a heal must leave the sentinel +empty. The test fails on `main` today (finding 1), so it belongs in the +failing-tests commit, and it outlives this change as the check that no xorq +cache node has crept back in. + +### D9. One change, one rebuild + +Removing the cache node changes the hash of every worthy entry, and bare-read +chaining changes every child's. This lands together with ADR-008's change to +`tallyman_read_csv` and ADR-009's digest definition, behind a single corpus +rebuild, with the failing tests committed and seen red first. After the +rebuild, `~/.cache/xorq/result_cache` (14 GB) and the older leaks under +`~/.cache/xorq/parquet/` (2.6 GB) can be deleted by hand. + +### D10. Every diff is built as an entry before it is displayed + +Live and promoted diffs already build the same expression +(`build_compare_expr`). The live path posts it to Buckaroo unmaterialized +(`_build_compare_expr`, `app.py:418-441`), so the outer join runs again for +every page, sort and stat query in the diff grid, and a failure of the join +surfaces inside Buckaroo. The promote path writes a recipe that calls +`build_diff_expr(a_hash, b_hash, keys)` and runs the normal build. + +Every diff now takes the promote path. When a diff view is opened, tallyman +builds the diff entry and waits for the build to finish. A diff contains a +join, so it is worthy and is materialized, and Buckaroo is then handed a view +build of the finished file (D6) with the diff's display configuration. The page +shows a "building diff" state while it waits. Promoting a diff is reduced to +putting an alias on an entry that already exists. + +- **An unnamed diff entry is ephemeral.** The checkpoint zips and commits every + complete directory under `entries/`, aliased or not (`zip_pending_entries`, + `catalog.py:160-180`), so an unnamed diff stored there would put every diff + ever viewed into the catalog's git history. Ephemeral entries live under + `compute_cache/ephemeral_entries//`, where everything is + already defined as deletable at any time, and they can be rebuilt from the + two hashes and the keys. Promote moves the directory into `entries/`, sets + the alias and checkpoints. The content hash is computed from the graph, so + the move does not change it, and the snapshot path stays the same. +- **Parents.** A worthy side is read from its snapshot. A cheap side is not + copied first: because the diff is itself materialized, each side is read + exactly once, while the diff is written. +- **Sessions.** The same entry can be opened as a diff, with the diff display + classes, or as a plain entry, so the Buckaroo session key includes the view + kind and is no longer the content hash alone. +- **Retired with this:** `_build_compare_expr` and its temp build directory, + the separate diff-session bookkeeping (`diff_session_is_loaded`, + `mark_diff_session_loaded`), and `diff_stat_cache/`, since a diff's Buckaroo + stats become its entry's own. The audit finding that every recalc wipes all + of `diff_stat_cache/` goes with it. + +*Rejected:* keep posting the join and materialize nothing. It shows a first +page sooner, and it breaks the governing rule in both directions: Buckaroo runs +tallyman's join, repeatedly, and tallyman never learns whether it succeeded. + +## Consequences + +- **Retired:** the `.cache()` call and the source-read injection in + `rewrite_for_build`; `rewrite_cache_dirs`; `_cached_node_path`; + `manifest.snapshot_key` and `_assert_recorded_snapshot_key`; the `baked` / + `recompute` plan split keyed on a `CachedNode`; the thread-only `_heal_lock` + and the `(FileNotFoundError, ValueError)` retry around xorq's shared temp + file; `entry_graph_expr` as a separate function. +- **ADR-006:** its D4 (inlined chaining) and D8 (snapshot-key tripwire) are + superseded. Its D2's "snapshot path derived from the loaded expression" + becomes `snapshot_path`. Its D3 (rebind composition onto the default backend) + is still needed for cheap parents and for diff composition. Its D5 (canonical + sort) and D6 (a missing build is a hard error) are unchanged. Its D7, D10 and + D12 (loud verification, the Buckaroo-state wipe, the pin and badge) attach to + `ensure_materialized`. +- **`docs/system-contract.md`** needs rewriting in Part 1 ยง4 (xorq caching + becomes background, not mechanism), "Content hash" (no cache injection in the + hashed expression; parents appear as snapshot paths), "Manifest" (drop + `snapshot_key`), "Worthiness", write path steps 2 and 5, the read path, + "Chaining", and the `[^preheal]` footnote, which currently describes + bare-read chaining as a retired workaround. +- **`plans/remove-ondemand-result-parquet.md`:** its premise that xorq's + snapshot is the single materialized copy is replaced. Tallyman's snapshot is + the single copy. +- **ADR-003:** chained descendants of an expensive entry are cheap, so its + motivating lineage stops filling the cache when the iterations chain off the + join. A revision that restates the join in its own recipe is still worthy + and still writes a copy, so the budget and eviction half of ADR-003 remains + and should be rewritten against this design. +- **CSV lineages:** a child of a `tallyman_read_csv` entry reads the parent's + snapshot and no longer inherits its Sort, so revisions stop baking one full + sorted copy each. The root entry's own copy is ADR-008's subject. +- **Buckaroo's stat keys** for a worthy entry become a function of one read of + one content-addressed path. +- **Source identity `salt` mode:** `rewrite_for_build` returns early under + `salt` because xorq's path-only snapshot keys would collide. A snapshot named + by the entry's content hash has no such collision, so the early return should + become unnecessary. Not tested. +- **Cost accepted:** a worthy parent's snapshot must exist before a child can + be built. A build also no longer repairs its own ancestors when something + outside tallyman executes it, and under the governing rule nothing does. + +## Open questions + +1. **Deep cheap chains.** Nothing cuts the graph between cheap entries. Run + `tests/test_perf_chain_depth.py` against this design and decide whether a + node-count or `compile_seconds` threshold should make an otherwise cheap + entry worthy. +2. **An unfaithful parent's descendants.** Descendants of an entry whose heal + failed verification were computed from bytes that no longer exist. Nothing + flags them today, and nothing here does either. +3. **Eviction policy.** D6 says what eviction must do to live sessions. Which + snapshots to evict, and when, stays with the ADR-003 rewrite. +4. **Garbage collection of ephemeral entries.** D10 says where they live and + that they are deletable. When to delete them belongs with the eviction + policy of open question 3. diff --git a/scripts/spike_bare_read_chaining.py b/scripts/spike_bare_read_chaining.py new file mode 100644 index 0000000..cf5044e --- /dev/null +++ b/scripts/spike_bare_read_chaining.py @@ -0,0 +1,117 @@ +"""ADR-007 evidence: chaining through a bare read of the parent's snapshot, with no xorq cache node anywhere. + +Raw xorq plus tallyman's ``classify_build``; no tallyman build machinery. A worthy parent (an aggregate, canonically +sorted) is built with no cache node and materialized by a plain parquet write to ``/.parquet``. +A child reads that path with ``deferred_read_parquet`` and adds a filter and a computed column. + +Questions, in the order printed: + +1. What does the child's build hash depend on: the snapshot's path string, or its bytes/mtime? +2. Can a child be built while the parent's snapshot is absent? +3. Does ``load_expr`` of the child's build need the snapshot? What does executing without it raise? +4. After the parent is re-materialized from its own build, does the child execute to the right answer? +5. Is the child classified cheap? (Under today's inlined chaining the same child is worthy.) +6. Does the child's ``expr.yaml`` carry the literal path, so the portable-path placeholder applies to it? +7. Did anything get written under xorq's global cache directory? + + uv run python scripts/spike_bare_read_chaining.py +""" + +from __future__ import annotations + +import hashlib +import os +import tempfile +from pathlib import Path + +HOME = Path(tempfile.mkdtemp(prefix="spike_bare_read_")) +os.environ["XORQ_CACHE_DIR"] = str(HOME / "_global_xorq") # must be set before xorq is imported + +import numpy as np # noqa: E402 +import pyarrow as pa # noqa: E402 +import pyarrow.parquet as pq # noqa: E402 +import xorq.api as xo # noqa: E402 +from xorq.ibis_yaml.compiler import build_expr, load_expr # noqa: E402 + +from tallyman_xorq.result_cache import classify_build # noqa: E402 + +N = 400_000 +BUILDS = HOME / "builds" +SNAPSHOTS = HOME / "compute_cache" / "result_cache" + + +def materialize(expr, dest: Path) -> str: + tmp = dest.with_name(f".{dest.name}.{os.getpid()}.tmp") + pq.write_table(expr.to_pyarrow(), tmp, compression="zstd") + os.replace(tmp, dest) + return hashlib.sha256(dest.read_bytes()).hexdigest()[:16] + + +def child_of(snapshot: Path, schema=None): + parent = xo.deferred_read_parquet(str(snapshot), schema=schema) + return parent.filter(parent.n > 10).mutate(r=parent.s / parent.n) + + +def attempt(fn) -> str: + try: + return str(fn()) + except Exception as exc: # noqa: BLE001 - the spike reports whatever is raised + return f"raises {type(exc).__name__}: {str(exc)[:90]}" + + +def main() -> None: + SNAPSHOTS.mkdir(parents=True) + source = HOME / "cas" / "deadbeef.parquet" + source.parent.mkdir() + rng = np.random.default_rng(7) + pq.write_table( + pa.table({"id": np.arange(N), "g": rng.integers(0, 20_000, N), "v": rng.integers(0, 1000, N)}), source + ) + + t = xo.deferred_read_parquet(str(source)) + parent = t.group_by("g").agg(n=t.count(), s=t.v.sum()).order_by(["g", "n", "s"]) + parent_build = Path(build_expr(parent, builds_dir=BUILDS)) + snapshot = SNAPSHOTS / f"{parent_build.name}.parquet" + first_digest = materialize(load_expr(parent_build), snapshot) + print(f"parent {parent_build.name}: {classify_build(parent_build)}, snapshot digest {first_digest}\n") + + same_path = Path(build_expr(child_of(snapshot), builds_dir=BUILDS)).name + original = snapshot.read_bytes() + pq.write_table(pa.table({"g": [1, 2], "n": [99, 98], "s": [5, 6]}), snapshot) # other rows, other size and mtime + other_bytes = Path(build_expr(child_of(snapshot), builds_dir=BUILDS)).name + elsewhere = SNAPSHOTS / "0123456789ab.parquet" + elsewhere.write_bytes(original) + other_path = Path(build_expr(child_of(elsewhere), builds_dir=BUILDS)).name + print(f"1. child hash: {same_path}") + print(f" same path with other bytes: {other_bytes}; same bytes at another path: {other_path}") + + snapshot.unlink() + print(f"2. compose with snapshot absent, no schema: {attempt(lambda: child_of(snapshot).schema().names)}") + with_schema = attempt(lambda: build_expr(child_of(snapshot, parent.schema()), builds_dir=BUILDS)) + print(f" build with snapshot absent, schema given: {with_schema}") + + child_build = BUILDS / same_path + print(f"3. load_expr(child) with snapshot absent: {attempt(lambda: type(load_expr(child_build)).__name__)}") + print(f" execute with snapshot absent: {attempt(lambda: load_expr(child_build).count().execute())}") + + again = materialize(load_expr(parent_build), snapshot) + rows = int(load_expr(child_build).count().execute()) + want = int(parent.filter(parent.n > 10).count().execute()) + print(f"4. parent re-materialized from its build: digest unchanged={again == first_digest}") + print(f" child rows {rows}, expected {want}") + + inlined = Path(build_expr(parent.filter(parent.n > 10).mutate(r=parent.s / parent.n), builds_dir=BUILDS)) + print(f"5. classify_build(bare-read child) = {classify_build(child_build)}") + print(f" classify_build(same child, parent graph inlined) = {classify_build(inlined)}") + + yaml_text = (child_build / "expr.yaml").read_text() + inlined_size = len((inlined / "expr.yaml").read_text()) + print(f"6. literal snapshot path in child expr.yaml: {str(snapshot) in yaml_text}") + print(f" child expr.yaml is {len(yaml_text):,} bytes; with the parent graph inlined it is {inlined_size:,}") + + leaked = sorted(str(f) for f in (HOME / "_global_xorq").rglob("*.parquet")) + print(f"7. parquet files under XORQ_CACHE_DIR: {leaked or 'none'}") + + +if __name__ == "__main__": + main()