diff --git a/src/foundation/mem.c b/src/foundation/mem.c index 4c5b60ae7a..4f74229e15 100644 --- a/src/foundation/mem.c +++ b/src/foundation/mem.c @@ -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 diff --git a/src/foundation/mem.h b/src/foundation/mem.h index 22cf0b47aa..91a8ef57bb 100644 --- a/src/foundation/mem.h +++ b/src/foundation/mem.h @@ -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); diff --git a/src/pipeline/pass_parallel.c b/src/pipeline/pass_parallel.c index e85ad015f1..74e87599e1 100644 --- a/src/pipeline/pass_parallel.c +++ b/src/pipeline/pass_parallel.c @@ -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 @@ -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 @@ -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)))); diff --git a/tests/test_parallel.c b/tests/test_parallel.c index 85df7d55d1..678b5b26ec 100644 --- a/tests/test_parallel.c +++ b/tests/test_parallel.c @@ -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" @@ -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) { @@ -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);