Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions src/foundation/mem.c
Original file line number Diff line number Diff line change
Expand Up @@ -594,7 +594,15 @@ size_t cbm_mem_allocator_committed(void) {
}

static _Atomic size_t g_peak_charged;
static _Atomic size_t g_charged_for_tests; /* 0 = live reading */
void cbm_mem_set_charged_for_tests(size_t bytes) {
atomic_store_explicit(&g_charged_for_tests, bytes, memory_order_relaxed);
}
size_t cbm_mem_charged(void) {
size_t pinned = atomic_load_explicit(&g_charged_for_tests, memory_order_relaxed);
if (pinned != 0) {
return pinned;
}
/* The OS number (phys_footprint on macOS, RSS elsewhere) is the charge for
* everything the process maps, and the memory core's own live bytes are its
* floor: macOS was measured under-reporting the OS number after
Expand Down
8 changes: 8 additions & 0 deletions src/foundation/mem.h
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,14 @@ size_t cbm_mem_budget(void);
* Never call from production code. */
void cbm_mem_set_budget_for_tests(size_t bytes);

/* TEST HOOK: pin cbm_mem_charged() (and so cbm_mem_over_budget()) to `bytes`;
* 0 restores the live reading. The charge is the process footprint, which a
* test cannot steer, so a gate keyed to it ("extraction ended just under the
* spill latch", #2184) is otherwise untestable. A pinned reading does not move
* cbm_mem_peak_charged(). Callers restore 0 before their assertions. Never call
* from production code. */
void cbm_mem_set_charged_for_tests(size_t bytes);

/* Returns true if current RSS exceeds the budget. */
bool cbm_mem_over_budget(void);

Expand Down
87 changes: 86 additions & 1 deletion src/pipeline/pass_parallel.c
Original file line number Diff line number Diff line change
Expand Up @@ -940,6 +940,88 @@ static int pp_spill_sweep(extract_ctx_t *ec, int worker_id) {
return parked;
}

/* ── Post-extraction projection (#2184) ───────────────────────────────
* Spill is entered DURING extraction, when the charge crosses the latch below.
* The phases after it -- registry build, cross-LSP prepare (all_defs, the
* per-language registries, surface rows, the module index) and resolve (edges
* into the graph buffer, cross-LSP appends into the cached results) -- cannot
* spill, and they grow the charge by a sizeable fraction of what extraction
* left. A run that ended extraction just under the latch therefore kept every
* result in memory and went over the budget later: openclaw, 8 workers,
* 4079 MB budget: 3678 MB charged at extraction end, 5829 MB in resolve.
*
* So extraction end projects that growth from the result counts and spills
* before handing over when charged + growth would cross the latch. A pure
* function of counts (O9): the same repo decides the same way on every run.
* Per-unit costs fitted 2026-09-25 (M5 Pro, release build, 8 workers) on the
* charge growth from the parallel_extract mark to the parallel_resolve mark:
*
* corpus files defs calls+usages growth model
* go 21882 735895 5,504,282 1307 MB 1909 MB
* django 4169 75419 591,350 231 MB 215 MB
* kotlin 5247 49669 421,788 129 MB 166 MB
* rust 818 23717 265,957 79 MB 76 MB
* php 2435 14691 137,969 47 MB 58 MB
* openclaw 48203 748829 >= 7,402,832 2151 MB >= 2419 MB
*
* It never under-reads by more than 7% (inside the budget/16 margin) and
* over-reads Go by 46% -- the safe side: an unneeded spill costs disk reads,
* never graph content (spilled and in-memory runs build the same graph). */
enum {
PP_POST_BYTES_PER_DEF = 1280, /* registry entry + LSP def + cross registries */
PP_POST_BYTES_PER_REF = 160, /* per call / usage: resolved edge + appends */
PP_POST_BYTES_PER_FILE = 8192, /* per-file tables: modules, surfaces, imports */
};

static size_t pp_post_extract_growth(const extract_ctx_t *ec, int64_t *defs_out,
int64_t *refs_out) {
int64_t defs = 0;
int64_t refs = 0;
for (int i = 0; i < ec->file_count; i++) {
const CBMFileResult *r = ec->result_cache[i];
if (r) {
defs += r->defs.count;
refs += (int64_t)r->calls.count + r->usages.count;
}
}
*defs_out = defs;
*refs_out = refs;
return (size_t)defs * PP_POST_BYTES_PER_DEF + (size_t)refs * PP_POST_BYTES_PER_REF +
(size_t)ec->file_count * PP_POST_BYTES_PER_FILE;
}

/* Enter spill mode at extraction end when the projected post-extraction growth
* would carry the charge over the latch; the final sweep then parks every
* cached result before registry build. No-op when spill is already on, not
* allowed for this owner, or no budget is set. */
static void pp_spill_if_projected_over(extract_ctx_t *ec) {
size_t budget = cbm_mem_budget();
if (budget == 0 || !ec->pctx || !ec->pctx->spill_allowed || !pp_spill_allowed(ec) ||
pp_spill_active(ec) || atomic_load_explicit(&ec->over_budget_abort, memory_order_relaxed)) {
return;
}
int64_t defs = 0;
int64_t refs = 0;
size_t growth = pp_post_extract_growth(ec, &defs, &refs);
size_t charged = cbm_mem_charged();
size_t line = budget - budget / PP_SPILL_EARLY_DIV;
bool spill = growth > line || charged > line - growth;
const size_t mb = (size_t)1024 * 1024;
char v[6][CBM_SZ_32];
snprintf(v[0], sizeof(v[0]), "%zu", charged / mb);
snprintf(v[1], sizeof(v[1]), "%zu", growth / mb);
snprintf(v[2], sizeof(v[2]), "%zu", line / mb);
snprintf(v[3], sizeof(v[3]), "%lld", (long long)defs);
snprintf(v[4], sizeof(v[4]), "%lld", (long long)refs);
snprintf(v[5], sizeof(v[5]), "%d", ec->file_count);
cbm_log_info("mem.post_extract.projection", "charged_mb", v[0], "growth_mb", v[1], "line_mb",
v[2], "defs", v[3], "refs", v[4], "files", v[5], "decision",
spill ? "spill" : "keep");
if (spill) {
pp_spill_enter(ec, "post_extract_projection");
}
}

/* Diagnostic (CBM_MEM_PHASES=1): where does the charge go between the
* near-budget latch and the first over-budget observation? One line per
* 256 MB step of the charge above its last logged value while spill mode is
Expand Down Expand Up @@ -1485,6 +1567,7 @@ int cbm_parallel_extract_ex(cbm_pipeline_ctx_t *ctx, const cbm_file_info_t *file
cbm_parallel_for_opts_t parallel_opts = {.max_workers = worker_count, .force_pthreads = false};
cbm_scale_begin(&ec.scale, "parallel_extract", (long)file_count);
cbm_parallel_for(worker_count, extract_worker, &ec, parallel_opts);
pp_spill_if_projected_over(&ec);
if (pp_spill_active(&ec) &&
!atomic_load_explicit(&ec.over_budget_abort, memory_order_relaxed)) {
/* Spill mode was entered, so results belong on disk: park every
Expand All @@ -1493,7 +1576,9 @@ int cbm_parallel_extract_ex(cbm_pipeline_ctx_t *ctx, const cbm_file_info_t *file
* run only on an over-budget observation; a run that latched early
* and then stayed under budget through extraction (kernel, 15 GB,
* 2026-09-14: 14,949 MB at this point, 44,797 results = 8 GB still
* cached) reached resolve with no headroom and aborted there. */
* cached) reached resolve with no headroom and aborted there.
* The projection above enters spill mode here too when the phases
* after extraction would carry the charge over the latch (#2184). */
int parked = pp_spill_sweep(&ec, 0);
cbm_log_info("mem.spill.final_sweep", "parked", itoa_log(parked), "charged_mb",
itoa_log((int)(cbm_mem_charged() / ((size_t)1024 * 1024))));
Expand Down
144 changes: 144 additions & 0 deletions tests/test_parallel.c
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
#include "discover/discover.h"
#include "foundation/platform.h"
#include "foundation/log.h"
#include "foundation/mem.h"
#include "cbm.h"
#include "result_spill.h"

Expand Down Expand Up @@ -616,6 +617,148 @@ TEST(parallel_spill_mode_builds_the_same_graph) {
PASS();
}

/* ── #2184: the post-extraction phases are charged before they run ─── */

/* Order-independent fingerprint of a whole graph: every node (label, QN, file,
* lines, properties) and every edge (type, source QN, target QN, properties)
* hashed and folded with a sum and an xor. Two graphs that differ in any field
* of any node or edge -- not just in a per-type count -- disagree here. */
typedef struct {
const cbm_gbuf_t *gb;
uint64_t sum;
uint64_t xr;
long count;
} graph_fp_t;

static uint64_t fp_mix(uint64_t h, const char *s) {
for (const unsigned char *p = (const unsigned char *)(s ? s : "\x01"); *p; p++) {
h = (h ^ *p) * 1099511628211ULL; /* FNV-1a */
}
return (h ^ 0xffU) * 1099511628211ULL; /* field separator */
}

static uint64_t fp_mix_int(uint64_t h, long long v) {
char buf[32];
snprintf(buf, sizeof(buf), "%lld", v);
return fp_mix(h, buf);
}

static void fp_fold(graph_fp_t *fp, uint64_t h) {
fp->sum += h;
fp->xr ^= h * 0x9E3779B97F4A7C15ULL;
fp->count++;
}

static void fp_visit_node(const cbm_gbuf_node_t *n, void *ud) {
uint64_t h = fp_mix(1469598103934665603ULL, "N");
h = fp_mix(h, n->label);
h = fp_mix(h, n->qualified_name);
h = fp_mix(h, n->file_path);
h = fp_mix_int(h, n->start_line);
h = fp_mix_int(h, n->end_line);
h = fp_mix(h, n->properties_json);
fp_fold(ud, h);
}

static void fp_visit_edge(const cbm_gbuf_edge_t *e, void *ud) {
graph_fp_t *fp = ud;
const cbm_gbuf_node_t *src = cbm_gbuf_find_by_id(fp->gb, e->source_id);
const cbm_gbuf_node_t *dst = cbm_gbuf_find_by_id(fp->gb, e->target_id);
uint64_t h = fp_mix(1469598103934665603ULL, "E");
h = fp_mix(h, e->type);
h = fp_mix(h, src ? src->qualified_name : NULL);
h = fp_mix(h, dst ? dst->qualified_name : NULL);
h = fp_mix(h, e->properties_json);
fp_fold(fp, h);
}

static graph_fp_t graph_fingerprint(const cbm_gbuf_t *gb) {
graph_fp_t fp = {.gb = gb};
cbm_gbuf_foreach_node(gb, fp_visit_node, &fp);
cbm_gbuf_foreach_edge(gb, fp_visit_edge, &fp);
return fp;
}

/* Mutator hook: runs between extraction and registry build, i.e. exactly where
* the post-extraction phases start. Records how many results are still held
* in memory there. */
static void count_cached_results(CBMFileResult **cache, int file_count, void *ud) {
int n = 0;
for (int i = 0; i < file_count; i++) {
n += cache[i] != NULL;
}
*(int *)ud = n;
}

/* #2184: spill was entered only DURING extraction, when the charge crossed
* budget - budget/16. A run that ended extraction just under that line kept
* every result in memory, and the phases that cannot spill (registry build,
* cross-LSP prepare, resolve) then grew the process past the budget -- openclaw
* at 8 workers, 4079 MB budget: 3678 MB at extraction end, 5829 MB in resolve.
* The extraction end now projects that growth from the result counts and
* spills first when budget - budget/16 would be crossed.
*
* The charge is pinned through the test seam (it is the process footprint
* otherwise): run A ends extraction ONE BYTE under the latch -- the extraction
* gate never fires, only the projection can spill; run B has the whole budget
* free -- the projection fits and nothing may spill. Both graphs must be
* identical, field for field, to each other and to the plain in-memory run. */
TEST(parallel_post_extract_projection_spills_before_resolve) {
if (ensure_parity_setup() != 0)
FAIL("setup failed");
cbm_discover_opts_t opts = {.mode = CBM_MODE_FULL};
cbm_file_info_t *files = NULL;
int file_count = 0;
ASSERT_EQ(cbm_discover(g_par_tmpdir, &opts, &files, &file_count), 0);
ASSERT_GT(file_count, 0);

const size_t saved_budget = cbm_mem_budget();
const size_t budget = (size_t)1024 * 1024 * 1024;
const size_t latch = budget - budget / 16; /* spill latches on charged > latch */
cbm_mem_set_budget_for_tests(budget);
g_harness_spill = true;

cbm_mem_set_charged_for_tests(latch);
int cached_near = -1;
cbm_gbuf_t *gb_near =
run_parallel_with_extract_opts_and_mutator("par-test", g_par_tmpdir, files, file_count, 2,
NULL, count_cached_results, &cached_near, false);
int64_t parked_near = g_harness_parked;

cbm_mem_set_charged_for_tests((size_t)1);
int cached_roomy = -1;
cbm_gbuf_t *roomy = run_parallel_with_extract_opts_and_mutator(
"par-test", g_par_tmpdir, files, file_count, 2, NULL, count_cached_results, &cached_roomy,
false);
int64_t parked_roomy = g_harness_parked;

cbm_mem_set_charged_for_tests(0);
g_harness_spill = false;
cbm_mem_set_budget_for_tests(saved_budget);
cbm_discover_free(files, file_count);
ASSERT(gb_near != NULL);
ASSERT(roomy != NULL);

/* A: every result went to disk before registry build. */
ASSERT_EQ((int)parked_near, file_count);
ASSERT_EQ(cached_near, 0);
/* B: the projection fits, so no store was opened and results stay cached. */
ASSERT_EQ((int)parked_roomy, -1);
ASSERT_GT(cached_roomy, 0);

graph_fp_t fp_near = graph_fingerprint(gb_near);
graph_fp_t fp_roomy = graph_fingerprint(roomy);
graph_fp_t fp_mem = graph_fingerprint(g_par_gbuf);
ASSERT_GT(fp_mem.count, 0);
ASSERT_EQ(fp_near.count, fp_mem.count);
ASSERT_EQ(fp_roomy.count, fp_mem.count);
ASSERT_TRUE(fp_near.sum == fp_mem.sum && fp_near.xr == fp_mem.xr);
ASSERT_TRUE(fp_roomy.sum == fp_mem.sum && fp_roomy.xr == fp_mem.xr);
cbm_gbuf_free(gb_near);
cbm_gbuf_free(roomy);
PASS();
}

/* ── Empty file list ──────────────────────────────────────────────── */

TEST(parallel_empty_files) {
Expand Down Expand Up @@ -4437,6 +4580,7 @@ SUITE(parallel) {
RUN_TEST(parallel_semantic_fixture_expected_counts);
RUN_TEST(parallel_total_edges);
RUN_TEST(parallel_spill_mode_builds_the_same_graph);
RUN_TEST(parallel_post_extract_projection_spills_before_resolve);
RUN_TEST(parallel_empty_files);
RUN_TEST(parallel_args_json_no_overflow);

Expand Down
Loading