Skip to content

Let each scatter shard decide the file cache against its own node (T-630) - #454

Merged
chasers merged 2 commits into
mainfrom
t-630-shard-file-cache-decision
Oct 4, 2026
Merged

chasers merged 2 commits into
mainfrom
t-630-shard-file-cache-decision

Conversation

@chasers

@chasers chasers commented Oct 4, 2026 •

Copy link
Copy Markdown
Owner

TL;DR: Each scatter shard now decides the file cache against its own node's cache. Before, the coordinator decided for every shard, so a warmed table bypassed the cache on every run.

Tracker: T-630. Follow-up to T-629 (PR 453).

The bug

  1. The coordinator plans the job. It checks which sealed files its own cache index has read.
  2. Scatter sends the files to shards on other nodes. Each shard writes its blocks to its own node's cache.
  3. The coordinator's index never sees those blocks. The table still looks cold to it.
  4. The coordinator sends its verdict (:bypassed) to every shard. Shards with a warm cache skip it too.
  5. A fileCache: true warming run fills the peers' caches, but the next auto run bypasses again.

What changed

  • The coordinator sends the job's mode (auto, on, off), not its verdict.
  • Shards on one node share one cache directory. So the request carries node_files: the sealed files (url and bytes) of every shard on that node.
  • Each shard checks node_files against its own node's index (FileCache.shard_decision/3).
  • So every shard on a node makes the same decision.
  • A node whose files add up to no more than the threshold always uses the cache, as a single-engine job does.
  • Each shard returns its decision. The job reports the combined decision (FileCache.combined/1): used when any shard read through its cache.

Metrics

Metric Change
smolquery_query_file_cache_shards_total{decision} ✅ New. One count per shard that answered.
smolquery_query_file_cache_jobs_total{decision} Counts after the job runs, done or failed. A scattered job counts by its shards' combined decision.

Rolling upgrade

  • The request still carries the coordinator's verdict as file_cache.
  • An old shard ignores the new keys and follows that verdict, as before. Its reply has no decision, so the coordinator counts it as that verdict.
  • A new shard with no mode (from an old coordinator) follows the verdict too.

How to review

  1. lib/smolquery/query_service/file_cache.ex: shard_decision/3 and combined/1.
  2. lib/smolquery/query_service/partial_worker.ex: cache_decision/2.
  3. lib/smolquery/query_service/scatter.ex: node_files/1, request keys, gathered decisions.
  4. lib/smolquery/query_service/runner.ex: count_cache_decision/3 and scatter_cache/2.

Tests

  • ✅ Unit: shard_decision/3 against a swept index (warm, cold, empty, forced, no cache). combined/1.
  • ✅ Scatter integration: cold cache → every shard bypasses. Warm cache → every shard uses it. file_cache: false → every shard is off.
  • ✅ Scatter integration: a cold cache under the threshold → every shard uses it.
  • ✅ Client integration: a job that fails after planning still counts its decision.
  • ✅ Scatter integration: a shard gets a stale :bypassed verdict and decides :used from its own index. This test fails without the fix.

Checks

  • ✅ mix precommit, mix ci, mix dialyzer pass locally.
  • ✅ mix test --include integration for the scatter and client suites.

Review

/code-review high found 10 candidates. Fixed in "Review of T-630":

  • ✅ A per-shard multiplier gave mixed decisions on one node. Now node_files.
  • ✅ Old shards reported null. Now they count as the coordinator's verdict.
  • ✅ Failed jobs were not counted. Now they are.
  • ✅ The shard metric counted before the shard ran. Now it counts after.
  • ✅ job.scatter no longer repeats file_cache.
  • ✅ One test helper for the cache block name (Smolquery.Test.FileCacheFixture).
  • Not changed: the planner's segment_files call. The coordinator still needs its verdict for its own engine and the fallback.

Watch out

  • ⚠️ The client integration suite can fail with commit_conflict (the sqlite catalog is locked while engines attach). It happened in 2 of 6 local runs, including in a test this PR does not touch. CI passed.

  • The test cluster is one node, so the tests use local workers. The cross-node case is proven by the stale-verdict test, not by a real peer.

  • Shards balance by file index, not by file. A change in membership moves files between nodes, so some cache misses still happen after a roll.

🤖 Generated with Claude Code

…630)

T-629 made the coordinator decide for every shard of a scattered query,
from its own node's cache index. The shards read on other nodes, and
their blocks land in those nodes' caches, which that index never sees.
So a table warmed through its shards still looked cold to the
coordinator, and every auto run bypassed the cache again, on every
shard. A `fileCache: true` warming run could not change that.

Each shard now decides for itself (FileCache.shard_decision/4):

- The coordinator sends the job's mode (auto, true, false) as
  file_cache_mode, and node_shards, how many of the job's shards run on
  that node.
- A sealed shard unit carries its "bytes". On auto, the shard sums the
  bytes of its sealed files that its own node's index has never read,
  times node_shards, since shards on one node fill one directory. Over
  bypass_bytes it is :bypassed; otherwise :used.
- The reply carries the shard's decision. The job reports the combined
  decision (FileCache.combined/1): :used when any shard read through
  its cache. job.scatter carries it too.
- smolquery_query_file_cache_shards_total{decision} counts each shard's
  decision. smolquery_query_file_cache_jobs_total now counts once the
  job has run, so a scattered job counts by what its shards did.

Rolling upgrade: the request still carries the coordinator's verdict as
file_cache, so a shard on an older node behaves as before, and a new
shard that gets no mode from an older coordinator follows that verdict.

Tests: shard_decision against a swept index (warm, cold, hot, forced,
no cache, node_shards); combined; scatter integration with a cold cache
(every shard bypasses), a warm cache (every shard uses it, and false
turns it off), and a shard called with a stale :bypassed verdict that
decides :used from its own index.
@chasers
chasers added this pull request to stack #456 October 4, 2026 01:45
Code review of PR 454 (/code-review high). Fixed:

- A shard multiplied its own cold bytes by node_shards. Round-robin
  shards differ in size, so shards on one node could split between used
  and bypassed, and a table under bypass_bytes in total could bypass,
  which a single-engine run never does. The request now carries
  node_files, the sealed files of every shard on that node, and
  shard_decision/3 sums their cold bytes, so every shard on a node
  decides alike and a small table always uses the cache. node_files is
  built once per job, not counted per shard.
- A shard on a pre-T-630 node replies without file_cache. It followed
  the coordinator's verdict, so it now counts as that verdict, not nil.
- The jobs metric fired only at the end of a successful execute, so a
  job that failed after planning went uncounted. It now fires for done
  and failed jobs alike. A scattered job's event no longer carries the
  coordinator's uncached_bytes, which came from another node's index.
- The shard metric fired before the shard ran, and not at all on the
  old-coordinator path. It now fires once the partial answers, on both.
- job.scatter no longer repeats file_cache; the runner pops it off.
- The cache block name lives in one test helper,
  Smolquery.Test.FileCacheFixture, instead of three copies.

Not changed: the planner's segment_files call for a scattered job. The
coordinator's verdict still sets its own engine and the single-engine
fallback, so the call is not wasted.

Tests: shard_decision/3 over node files; a cold cache under the
threshold is used by every shard; a job that fails after planning
counts its decision.
@chasers
chasers merged commit 03b0fa4 into main Oct 4, 2026
7 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant