diff --git a/src/daemon/application.c b/src/daemon/application.c index 22c62029b..d1f78841b 100644 --- a/src/daemon/application.c +++ b/src/daemon/application.c @@ -2199,10 +2199,10 @@ static char *application_auto_index_args(cbm_daemon_application_t *application, return NULL; } yyjson_mut_doc_set_root(document, root); - char *args = yyjson_mut_obj_add_strcpy(document, root, "repo_path", root_path) && - application_index_args_add_policy(application, document, root) - ? yyjson_mut_write(document, 0, NULL) - : NULL; + bool encoded = yyjson_mut_obj_add_strcpy(document, root, "repo_path", root_path) && + application_index_args_add_policy(application, document, root) && + yyjson_mut_obj_add_bool(document, root, "_cbm_background", true); + char *args = encoded ? yyjson_mut_write(document, 0, NULL) : NULL; yyjson_mut_doc_free(document); return args; } @@ -3894,7 +3894,8 @@ static int application_background_index(cbm_daemon_application_t *application, } yyjson_mut_doc_set_root(document, root); bool encoded = yyjson_mut_obj_add_strcpy(document, root, "repo_path", canonical_root) && - application_index_args_add_policy(application, document, root); + application_index_args_add_policy(application, document, root) && + yyjson_mut_obj_add_bool(document, root, "_cbm_background", true); char *default_project = cbm_project_name_from_path(canonical_root); bool custom_project = project_name[0] && (!default_project || strcmp(default_project, project_name) != 0); diff --git a/src/mcp/mcp.c b/src/mcp/mcp.c index c877d5110..31efa9420 100644 --- a/src/mcp/mcp.c +++ b/src/mcp/mcp.c @@ -11109,7 +11109,8 @@ static char *index_run_supervised_path(cbm_mcp_server_t *srv, const char *root_p yyjson_mut_val *root = yyjson_mut_obj(doc); yyjson_mut_doc_set_root(doc, root); if (!yyjson_mut_obj_add_strcpy(doc, root, "repo_path", root_path) || - !cbm_mcp_index_policy_add_to_args(doc, root, &policy)) { + !cbm_mcp_index_policy_add_to_args(doc, root, &policy) || + !yyjson_mut_obj_add_bool(doc, root, "_cbm_background", true)) { yyjson_mut_doc_free(doc); return NULL; } @@ -11369,6 +11370,22 @@ static bool index_root_owner_resolve(const char *repo_path, char **owner_out, ch return false; } +/* Background CPU headroom is for a refresh of a graph that already serves. + * A project's first index keeps every core, even when it was started + * automatically: the caller is waiting on a graph that does not exist yet. + * An empty or unreadable file is not a committed index. Both sites that + * honour automatic indexing — handle_index_repository's _cbm_background + * flag and the in-process autoindex_thread — go through this so they + * cannot drift. */ +static bool project_has_committed_index(const char *project) { + if (!project || !project[0]) { + return false; + } + char db_path[CBM_SZ_1K]; + project_db_path(project, db_path, sizeof(db_path)); + return db_path[0] && project_db_is_servable(project, db_path); +} + /* The three heap strings handle_index_repository owns from * cbm_mcp_get_string_arg / resolved_repo_path_from_project_arg. One release * point keeps the dozen early-return paths in step; free(NULL) is a no-op, so @@ -11498,6 +11515,7 @@ static char *handle_index_repository(cbm_mcp_server_t *srv, const char *args) { char *repo_path = cbm_mcp_get_string_arg(args, "repo_path"); char *mode_str = cbm_mcp_get_string_arg(args, "mode"); char *name_override = cbm_mcp_get_string_arg(args, "name"); + bool background = cbm_mcp_get_bool_arg(args, "_cbm_background"); cbm_normalize_path_sep(repo_path); if (!repo_path) { @@ -11669,6 +11687,10 @@ static char *handle_index_repository(cbm_mcp_server_t *srv, const char *args) { return cbm_mcp_text_result("invalid project name", true); } free(name_override); + /* Headroom follows the final project name, decided before artifact + * bootstrap can create the first local copy of this index. */ + cbm_pipeline_set_background(p, background && + project_has_committed_index(cbm_pipeline_project_name(p))); cbm_pipeline_set_persistence(p, persistence); cbm_pipeline_set_resource_policy(p, &resource_policy); @@ -18106,6 +18128,7 @@ static void *autoindex_thread(void *arg) { cbm_log_warn("autoindex.err", "msg", "pipeline_create_failed"); return NULL; } + cbm_pipeline_set_background(p, project_has_committed_index(cbm_pipeline_project_name(p))); /* Block until any concurrent pipeline finishes */ cbm_pipeline_lock(); diff --git a/src/pipeline/lsp_surface.c b/src/pipeline/lsp_surface.c index 4c3cde336..95ed7f63b 100644 --- a/src/pipeline/lsp_surface.c +++ b/src/pipeline/lsp_surface.c @@ -231,8 +231,9 @@ int cbm_lsp_surface_build_rows(const cbm_pipeline_ctx_t *ctx, const char *projec .rows = rows, }; atomic_init(&job.failed, false); + int workers = cbm_pipeline_worker_count(ctx ? ctx->pipeline : NULL); cbm_parallel_for(file_count, surface_row_one, &job, - (cbm_parallel_for_opts_t){.max_workers = 0, .force_pthreads = false}); + (cbm_parallel_for_opts_t){.max_workers = workers, .force_pthreads = false}); if (atomic_load_explicit(&job.failed, memory_order_relaxed)) { cbm_store_free_lsp_surfaces(rows, file_count); /* untouched rows are all NULL */ return -1; diff --git a/src/pipeline/pipeline.c b/src/pipeline/pipeline.c index ba41b877c..be881942e 100644 --- a/src/pipeline/pipeline.c +++ b/src/pipeline/pipeline.c @@ -71,6 +71,18 @@ static inline void *intptr_to_ptr(intptr_t v) { * Atomic spinlock: 0 = free, 1 = locked. */ static atomic_int g_pipeline_busy = 0; +#if defined(CBM_INCREMENTAL_TEST_API) && CBM_INCREMENTAL_TEST_API +static atomic_int g_last_worker_count = 0; + +void cbm_pipeline_worker_count_test_reset(void) { + atomic_store(&g_last_worker_count, 0); +} + +int cbm_pipeline_worker_count_test_last(void) { + return atomic_load(&g_last_worker_count); +} +#endif + #if defined(CBM_INCREMENTAL_TEST_API) && CBM_INCREMENTAL_TEST_API static atomic_bool g_persist_test_fail_after_stage_dump = false; static atomic_bool g_persist_test_cancel_after_predump = false; @@ -192,6 +204,7 @@ struct cbm_pipeline { bool persistence; /* write .codebase-memory/graph.db.zst after indexing */ cbm_index_resource_policy_t resource_policy; cbm_index_resource_violation_t resource_violation; + bool background; /* reserve CPU headroom when no user is waiting */ /* Snapshot of the artifact export failure of THIS run (set only by * export_after_publish failure, zeroed at run start, cleared on success). @@ -400,6 +413,12 @@ const char *cbm_pipeline_export_error(const cbm_pipeline_t *p) { return p ? p->export_error : ""; } +void cbm_pipeline_set_background(cbm_pipeline_t *p, bool background) { + if (p) { + p->background = background; + } +} + bool cbm_pipeline_set_project_name(cbm_pipeline_t *p, const char *name) { if (!p || !name || !name[0]) { return false; @@ -613,12 +632,18 @@ void cbm_pipeline_set_committed_counts(cbm_pipeline_t *p, int nodes, int edges) * crasher; a parallel re-run would race the marker. Honour that override * everywhere the worker count drives the parallel/sequential decision, so the * whole extraction phase collapses to the deterministic sequential path. */ -static int effective_worker_count(bool initial) { +int cbm_pipeline_worker_count(const cbm_pipeline_t *p) { const char *st = getenv("CBM_INDEX_SINGLE_THREAD"); + int workers; if (st && st[0] == '1') { - return 1; + workers = 1; + } else { + workers = cbm_default_worker_count(!p || !p->background); } - return cbm_default_worker_count(initial); +#if defined(CBM_INCREMENTAL_TEST_API) && CBM_INCREMENTAL_TEST_API + atomic_store(&g_last_worker_count, workers); +#endif + return workers; } /* Resolve the DB path for this pipeline. Caller must free(). */ @@ -2590,7 +2615,7 @@ static int run_githistory(cbm_pipeline_t *p, cbm_pipeline_ctx_t *ctx) { gh_compute_arg_t gh_arg = {.repo_path = ctx->repo_path, .result = &gh_result}; if (p->mode != CBM_MODE_FAST) { - if (effective_worker_count(true) > SKIP_ONE) { + if (cbm_pipeline_worker_count(p) > SKIP_ONE) { if (cbm_thread_create(&gh_thread, 0, gh_compute_thread_fn, &gh_arg) == 0) { gh_threaded = true; } @@ -2690,7 +2715,7 @@ static int run_extraction_phase(cbm_pipeline_t *p, cbm_pipeline_ctx_t *ctx, return CBM_NOT_FOUND; } - int worker_count = effective_worker_count(true); + int worker_count = cbm_pipeline_worker_count(p); CBM_PROF_START(t_extract_total); int rc = (worker_count > SKIP_ONE && file_count > MIN_FILES_FOR_PARALLEL) ? run_parallel_pipeline(p, ctx, files, file_count, worker_count, &t) @@ -2785,7 +2810,7 @@ static int cbm_pipeline_run_staged(cbm_pipeline_t *p) { rc = mode_promoted ? cbm_pipeline_build_fresh_semantic_manifest(p, p->project_name, &baseline_manifest, &baseline_count) - : cbm_pipeline_build_semantic_manifest(p->project_name, p->repo_path, files, + : cbm_pipeline_build_semantic_manifest(p, p->project_name, p->repo_path, files, file_count, p->excluded_dirs, p->excluded_count, &p->git_ctx, p->userconfig, &baseline_manifest, &baseline_count); diff --git a/src/pipeline/pipeline.h b/src/pipeline/pipeline.h index c1d571ecd..e52513673 100644 --- a/src/pipeline/pipeline.h +++ b/src/pipeline/pipeline.h @@ -59,6 +59,18 @@ void cbm_pipeline_set_resource_policy(cbm_pipeline_t *p, const cbm_index_resourc /* Copy the exact discovery violation from the most recent run. */ void cbm_pipeline_get_resource_violation(const cbm_pipeline_t *p, cbm_index_resource_violation_t *violation); +/* Mark work that runs without a waiting user. Background pipelines reserve + * CPU headroom; foreground pipelines retain the initial-index all-core policy. */ +void cbm_pipeline_set_background(cbm_pipeline_t *p, bool background); + +/* Resolve the worker policy for this pipeline, including environment and + * crash-recovery overrides. */ +int cbm_pipeline_worker_count(const cbm_pipeline_t *p); + +#if defined(CBM_INCREMENTAL_TEST_API) && CBM_INCREMENTAL_TEST_API +void cbm_pipeline_worker_count_test_reset(void); +int cbm_pipeline_worker_count_test_last(void); +#endif /* Snapshot of the artifact export failure of the last cbm_pipeline_run, or "" * when the run succeeded / did not reach post-publish export. Used to diff --git a/src/pipeline/pipeline_incremental.c b/src/pipeline/pipeline_incremental.c index 154518658..8c348c3f7 100644 --- a/src/pipeline/pipeline_incremental.c +++ b/src/pipeline/pipeline_incremental.c @@ -413,9 +413,9 @@ static void *manifest_hash_worker(void *arg) { enum { MANIFEST_PARALLEL_MIN_FILES = 64 }; -int cbm_pipeline_build_semantic_manifest(const char *project, const char *repo_path, - const cbm_file_info_t *files, int file_count, - char **excluded_dirs, int excluded_count, +int cbm_pipeline_build_semantic_manifest(const cbm_pipeline_t *p, const char *project, + const char *repo_path, const cbm_file_info_t *files, + int file_count, char **excluded_dirs, int excluded_count, const cbm_git_context_t *git_ctx, const cbm_userconfig_t *userconfig, cbm_file_hash_t **out, int *out_count) { @@ -456,7 +456,7 @@ int cbm_pipeline_build_semantic_manifest(const char *project, const char *repo_p } struct timespec t_hash; cbm_clock_gettime(CLOCK_MONOTONIC, &t_hash); - int hash_workers = cbm_default_worker_count(true); + int hash_workers = cbm_pipeline_worker_count(p); if (rc == 0 && file_count >= MANIFEST_PARALLEL_MIN_FILES && hash_workers > SKIP_ONE) { char (*shas)[CBM_SHA256_HEX_LEN + 1] = malloc((size_t)file_count * sizeof(*shas)); int64_t *mtimes = calloc((size_t)file_count, sizeof(int64_t)); @@ -634,9 +634,9 @@ int cbm_pipeline_build_fresh_semantic_manifest(cbm_pipeline_t *p, const char *pr &fresh_ignored_total); } if (rc == 0) { - rc = cbm_pipeline_build_semantic_manifest(project, repo_path, fresh_files, fresh_file_count, - fresh_excluded, fresh_excluded_count, - &fresh_git_ctx, fresh_userconfig, out, out_count); + rc = cbm_pipeline_build_semantic_manifest( + p, project, repo_path, fresh_files, fresh_file_count, fresh_excluded, + fresh_excluded_count, &fresh_git_ctx, fresh_userconfig, out, out_count); } cbm_set_user_lang_config(previous_userconfig); cbm_git_context_free(&fresh_git_ctx); @@ -1212,7 +1212,7 @@ static int run_extract_resolve(cbm_pipeline_ctx_t *ctx, cbm_file_info_t *changed * a full build takes, which is what makes its output converge. */ #define MIN_FILES_FOR_PARALLEL_INCR 50 - int worker_count = cbm_default_worker_count(true); + int worker_count = cbm_pipeline_worker_count(ctx->pipeline); bool use_parallel = closure != NULL || (worker_count > SKIP_ONE && ci > MIN_FILES_FOR_PARALLEL_INCR); @@ -1590,7 +1590,7 @@ static int closure_probe_surfaces(cbm_pipeline_t *p, const char *project, _Atomic int64_t probe_ids; atomic_init(&probe_ids, cbm_gbuf_next_id(probe_gbuf)); rc = cbm_parallel_extract(&probe_ctx, probe_files, probe_count, cache, &probe_ids, - cbm_default_worker_count(true)); + cbm_pipeline_worker_count(p)); } if (rc == 0) { char **def_modules = (char **)calloc((size_t)probe_count, sizeof(char *)); @@ -1602,8 +1602,12 @@ static int closure_probe_surfaces(cbm_pipeline_t *p, const char *project, if (def_modules && def_starts) { defs = cbm_pxc_collect_all_defs(NULL, &probe_arena, cache, probe_files, probe_count, project, def_modules, &def_count, def_starts); - rc = cbm_lsp_surface_build_rows(NULL, project, cache, probe_files, probe_count, defs, - def_starts, out_rows, out_count); + /* Extract left pipeline NULL so file errors are not recorded twice. + * Surface rows only borrow it for the background worker policy. */ + probe_ctx.pipeline = p; + rc = cbm_lsp_surface_build_rows(&probe_ctx, project, cache, probe_files, probe_count, + defs, def_starts, out_rows, out_count); + probe_ctx.pipeline = NULL; } else { rc = -1; } @@ -2061,7 +2065,7 @@ static int run_closure_delta(cbm_pipeline_t *p, const char *db_path, const char elig[elig_count++] = row; } } - int workers = cbm_default_worker_count(true); + int workers = cbm_pipeline_worker_count(p); if (workers < 1) { workers = 1; } diff --git a/src/pipeline/pipeline_internal.h b/src/pipeline/pipeline_internal.h index 7f266eae9..c2a87e0bf 100644 --- a/src/pipeline/pipeline_internal.h +++ b/src/pipeline/pipeline_internal.h @@ -808,9 +808,9 @@ int cbm_pipeline_run_incremental(cbm_pipeline_t *p, const char *db_path, cbm_fil #define CBM_SEMANTIC_INPUT_GLOBAL_CONFIG CBM_SEMANTIC_INPUT_PREFIX "global-extension-config-v1" #define CBM_SEMANTIC_INPUT_PROJECT_CONFIG CBM_SEMANTIC_INPUT_PREFIX "project-extension-config-v1" -int cbm_pipeline_build_semantic_manifest(const char *project, const char *repo_path, - const cbm_file_info_t *files, int file_count, - char **excluded_dirs, int excluded_count, +int cbm_pipeline_build_semantic_manifest(const cbm_pipeline_t *p, const char *project, + const char *repo_path, const cbm_file_info_t *files, + int file_count, char **excluded_dirs, int excluded_count, const cbm_git_context_t *git_ctx, const cbm_userconfig_t *userconfig, cbm_file_hash_t **out, int *out_count); diff --git a/tests/test_daemon_application.c b/tests/test_daemon_application.c index 7f4c9cc28..0bec333a8 100644 --- a/tests/test_daemon_application.c +++ b/tests/test_daemon_application.c @@ -1375,6 +1375,7 @@ typedef struct { * itself instead, which is the edge a reader on another thread needs. */ atomic_bool args_captured[APP_FAKE_MAX_ATTEMPTS]; size_t memory_budgets[APP_FAKE_MAX_ATTEMPTS]; + bool background_requests[APP_FAKE_MAX_ATTEMPTS]; } app_fake_worker_context_t; typedef struct { @@ -1442,6 +1443,8 @@ static int app_fake_worker_start(void *opaque, const char *args_json, size_t mem args_json ? args_json : ""); atomic_store(&context->args_captured[worker->attempt], true); context->memory_budgets[worker->attempt] = memory_budget_bytes; + context->background_requests[worker->attempt] = + cbm_mcp_get_bool_arg(args_json, "_cbm_background"); if (marker_file) { (void)snprintf(context->marker_paths[worker->attempt], APP_TEST_PATH_CAP, "%s", marker_file); @@ -2202,6 +2205,7 @@ TEST(daemon_application_initialize_coalesces_auto_index_for_full_sessions) { ASSERT_TRUE(first_initialized); ASSERT_TRUE(first_owned); ASSERT_TRUE(auto_policy_propagated); + ASSERT_TRUE(fake.background_requests[0]); ASSERT_TRUE(second_initialized); ASSERT_TRUE(coalesced); ASSERT_TRUE(restricted_disconnect_kept_job); @@ -3634,6 +3638,7 @@ TEST(daemon_application_request_cancel_preserves_persistent_watch_and_session) { ASSERT_TRUE(thread_started); ASSERT_TRUE(subscribed); ASSERT_TRUE(worker_started); + ASSERT_FALSE(fixture.fake.background_requests[0]); ASSERT_TRUE(request_returned); ASSERT_TRUE(thread_joined); ASSERT_EQ(request.status, CBM_DAEMON_RUNTIME_APPLICATION_CANCELLED); @@ -3931,6 +3936,7 @@ TEST(daemon_application_final_cancel_drains_admitted_watcher_job) { ASSERT_TRUE(fixture_ready); ASSERT_TRUE(thread_started); ASSERT_TRUE(worker_started); + ASSERT_TRUE(fixture.fake.background_requests[0]); ASSERT_TRUE(job_active); ASSERT_EQ(watches_after_cancel, 0); ASSERT_TRUE(thread_joined); diff --git a/tests/test_mcp.c b/tests/test_mcp.c index 22af972af..f42f0dc6a 100644 --- a/tests/test_mcp.c +++ b/tests/test_mcp.c @@ -20006,6 +20006,159 @@ TEST(autoindex_limit_guards_git_root_issue713) { PASS(); } +/* A servable project row is a committed index. user_version stays 0, so a + * later run cannot take the incremental no-op and skip recording workers. */ +static bool mcp_publish_committed_index(const char *root) { + char *project = root ? cbm_project_name_from_path(root) : NULL; + cbm_store_t *store = project ? cbm_store_open(project) : NULL; + bool published = store && cbm_store_upsert_project(store, project, root) == CBM_STORE_OK; + cbm_store_close(store); + free(project); + return published; +} + +#ifdef CBM_ENABLE_TEST_SEAMS +typedef struct { + const char *root; + bool published; +} mcp_committed_seed_t; + +/* maybe_auto_index returns before the thread when .db already + * exists. Publish the committed index from the count hook, which runs after + * that check and before autoindex_thread chooses a width. */ +static void mcp_autoindex_publish_committed(void *context) { + mcp_committed_seed_t *seed = context; + if (seed) { + seed->published = mcp_publish_committed_index(seed->root); + } +} +#endif + +/* First index, including one started automatically, uses every core. Headroom + * applies only once a servable project database already exists. + * automatic: in-process auto-index (initialize → autoindex_thread). + * Otherwise index_repository with _cbm_background, the daemon/watcher path. */ +static int mcp_assert_inprocess_worker_policy(bool committed, bool automatic) { + char cache[256]; + char repo[512]; + (void)snprintf(cache, sizeof(cache), "%s/cbm-worker-policy-%d%d-XXXXXX", cbm_tmpdir(), + committed ? 1 : 0, automatic ? 1 : 0); + bool cache_ready = cbm_mkdtemp(cache) != NULL; + (void)snprintf(repo, sizeof(repo), "%s/repo", cache); + char source[640]; + (void)snprintf(source, sizeof(source), "%s/main.py", repo); + bool repo_ready = cache_ready && th_mkdir_p(repo) == 0 && + th_write_file(source, "def background_index():\n return True\n") == 0; + + mcp_test_env_backup_t environment[] = { + {.name = "CBM_CACHE_DIR"}, + {.name = "CBM_WORKERS"}, + {.name = "CBM_INDEX_SINGLE_THREAD"}, + }; + bool environment_saved = true; + for (size_t i = 0; i < sizeof(environment) / sizeof(environment[0]); i++) { + const char *value = getenv(environment[i].name); + environment[i].present = value != NULL; + environment[i].value = value ? strdup(value) : NULL; + environment_saved = environment_saved && (!value || environment[i].value); + } + bool environment_ready = environment_saved && cbm_setenv("CBM_CACHE_DIR", cache, 1) == 0 && + cbm_unsetenv("CBM_WORKERS") == 0 && + cbm_unsetenv("CBM_INDEX_SINGLE_THREAD") == 0; + + bool seeded = !committed; + char old_cwd[CBM_SZ_4K] = {0}; + bool cwd_ready = true; + cbm_config_t *config = NULL; + bool config_ready = true; + cbm_mcp_server_t *server = NULL; + char *response = NULL; + if (automatic) { + cwd_ready = repo_ready && environment_ready && cbm_getcwd(old_cwd, sizeof(old_cwd)) && + cbm_chdir(repo) == 0; + config = cwd_ready ? cbm_config_open(cache) : NULL; + config_ready = config && cbm_config_set(config, CBM_CONFIG_AUTO_INDEX, "true") == 0 && + cbm_config_set(config, CBM_CONFIG_AUTO_WATCH, "false") == 0; + cbm_pipeline_worker_count_test_reset(); + server = config_ready ? cbm_mcp_server_new(NULL) : NULL; + if (server) { + cbm_mcp_server_set_config(server, config); +#ifdef CBM_ENABLE_TEST_SEAMS + mcp_committed_seed_t seed = {.root = repo, .published = !committed}; + if (committed) { + cbm_mcp_server_set_auto_index_count_test_hook( + server, mcp_autoindex_publish_committed, &seed); + } +#endif + response = cbm_mcp_server_handle( + server, "{\"jsonrpc\":\"2.0\",\"id\":1,\"method\":\"initialize\",\"params\":{}}"); +#ifdef CBM_ENABLE_TEST_SEAMS + seeded = seed.published; +#endif + cbm_mcp_server_free(server); /* joins the automatic index thread */ + } + } else { + if (committed && repo_ready && environment_ready) { + seeded = mcp_publish_committed_index(repo); + } + cbm_pipeline_worker_count_test_reset(); + server = repo_ready && environment_ready && seeded ? cbm_mcp_server_new(NULL) : NULL; + if (server) { + char repo_json[512]; + char args[CBM_SZ_1K]; + (void)snprintf(repo_json, sizeof(repo_json), "%s", repo); + cbm_normalize_path_sep(repo_json); + (void)snprintf(args, sizeof(args), "{\"_cbm_background\":true,\"repo_path\":\"%s\"}", + repo_json); + response = cbm_mcp_handle_tool(server, "index_repository", args); + cbm_mcp_server_free(server); + } + } + int selected_workers = cbm_pipeline_worker_count_test_last(); + bool server_ready = server != NULL; + bool response_ready = response != NULL; + + free(response); + cbm_config_close(config); + if (automatic && cwd_ready) { + (void)cbm_chdir(old_cwd); + } + /* Capture while CBM_WORKERS is still unset. Restoring first makes the + * expected count follow the lane override instead of this policy. */ + int expected = cbm_default_worker_count(!committed); + mcp_test_restore_env(environment, sizeof(environment) / sizeof(environment[0])); + bool cleaned = !cache_ready || th_rmtree(cache) == 0; + + ASSERT_TRUE(cache_ready); + ASSERT_TRUE(repo_ready); + ASSERT_TRUE(environment_saved); + ASSERT_TRUE(environment_ready); + ASSERT_TRUE(cwd_ready); + ASSERT_TRUE(config_ready); + ASSERT_TRUE(seeded); + ASSERT_TRUE(server_ready); + ASSERT_TRUE(response_ready); + ASSERT_EQ(selected_workers, expected); + ASSERT_TRUE(cleaned); + PASS(); +} + +TEST(mcp_auto_index_in_process_fresh_project_uses_full_width) { + return mcp_assert_inprocess_worker_policy(false, true); +} + +TEST(mcp_auto_index_in_process_existing_index_keeps_headroom) { + return mcp_assert_inprocess_worker_policy(true, true); +} + +TEST(mcp_background_flag_fresh_project_uses_full_width) { + return mcp_assert_inprocess_worker_policy(false, false); +} + +TEST(mcp_background_flag_existing_index_keeps_headroom) { + return mcp_assert_inprocess_worker_policy(true, false); +} + /* ══════════════════════════════════════════════════════════════════ * #853 — auto_watch=false must ALSO gate the SUPERVISED fresh-index * watcher registration (keystone × #849 merge interaction) @@ -21361,6 +21514,10 @@ SUITE(mcp) { RUN_TEST(autoindex_limit_guards_non_git_root_issue713); RUN_TEST(autoindex_limit_admits_non_git_root_under_limit_issue713); RUN_TEST(autoindex_limit_guards_git_root_issue713); + RUN_TEST(mcp_auto_index_in_process_fresh_project_uses_full_width); + RUN_TEST(mcp_auto_index_in_process_existing_index_keeps_headroom); + RUN_TEST(mcp_background_flag_fresh_project_uses_full_width); + RUN_TEST(mcp_background_flag_existing_index_keeps_headroom); } /* Kept separate so daemon-coordination regressions can be iterated without diff --git a/tests/test_pipeline.c b/tests/test_pipeline.c index 56be6b88e..c541db5e2 100644 --- a/tests/test_pipeline.c +++ b/tests/test_pipeline.c @@ -117,6 +117,120 @@ TEST(pipeline_create_free) { PASS(); } +TEST(pipeline_background_worker_policy_preserves_overrides) { + const char *old_workers = getenv("CBM_WORKERS"); + char *saved_workers = old_workers ? strdup(old_workers) : NULL; + const char *old_single = getenv("CBM_INDEX_SINGLE_THREAD"); + char *saved_single = old_single ? strdup(old_single) : NULL; + cbm_unsetenv("CBM_WORKERS"); + cbm_unsetenv("CBM_INDEX_SINGLE_THREAD"); + + cbm_pipeline_t *foreground = cbm_pipeline_new("/some/path", NULL, CBM_MODE_FULL); + cbm_pipeline_t *background = cbm_pipeline_new("/some/path", NULL, CBM_MODE_FULL); + ASSERT_NOT_NULL(foreground); + ASSERT_NOT_NULL(background); + cbm_pipeline_set_background(background, true); + + int foreground_default = cbm_pipeline_worker_count(foreground); + int background_default = cbm_pipeline_worker_count(background); + bool defaults_match = foreground_default == cbm_default_worker_count(true) && + background_default == cbm_default_worker_count(false) && + foreground_default >= 1 && background_default >= 1; + + cbm_setenv("CBM_WORKERS", "3", 1); + bool env_override = cbm_pipeline_worker_count(foreground) == 3 && + cbm_pipeline_worker_count(background) == 3; + + cbm_setenv("CBM_INDEX_SINGLE_THREAD", "1", 1); + bool recovery_override = cbm_pipeline_worker_count(foreground) == 1 && + cbm_pipeline_worker_count(background) == 1; + + cbm_pipeline_free(foreground); + cbm_pipeline_free(background); + saved_workers ? cbm_setenv("CBM_WORKERS", saved_workers, 1) : cbm_unsetenv("CBM_WORKERS"); + saved_single ? cbm_setenv("CBM_INDEX_SINGLE_THREAD", saved_single, 1) + : cbm_unsetenv("CBM_INDEX_SINGLE_THREAD"); + free(saved_workers); + free(saved_single); + + ASSERT_TRUE(defaults_match); + ASSERT_TRUE(env_override); + ASSERT_TRUE(recovery_override); + PASS(); +} + +TEST(pipeline_background_policy_preserves_index_results) { + ASSERT_EQ(setup_test_repo(), 0); + /* 64 fillers plus setup_test_repo's three files exceed MIN_FILES_FOR_PARALLEL + * (50). Below that threshold both policies take the sequential path, so the + * comparison stays green even with the background policy removed. */ + for (int i = 0; i < 64; i++) { + char name[32]; + char body[96]; + (void)snprintf(name, sizeof(name), "pad%02d.go", i); + (void)snprintf(body, sizeof(body), "package main\n\nfunc Pad%02d() int {\n\treturn %d\n}\n", + i, i); + if (th_write_file(TH_PATH(g_tmpdir, name), body) != 0) { + teardown_test_repo(); + FAIL("failed to pad fixture past MIN_FILES_FOR_PARALLEL"); + } + } + + /* An inherited CBM_WORKERS pin makes both runs use the same count. Drop it + * so this compares an all-core index against a headroom index. */ + const char *old_workers = getenv("CBM_WORKERS"); + char *saved_workers = old_workers ? strdup(old_workers) : NULL; + const char *old_single = getenv("CBM_INDEX_SINGLE_THREAD"); + char *saved_single = old_single ? strdup(old_single) : NULL; + if ((old_workers && !saved_workers) || (old_single && !saved_single)) { + free(saved_workers); + free(saved_single); + teardown_test_repo(); + FAIL("failed to save worker env"); + } + cbm_unsetenv("CBM_WORKERS"); + cbm_unsetenv("CBM_INDEX_SINGLE_THREAD"); + + char foreground_db[512]; + char background_db[512]; + (void)snprintf(foreground_db, sizeof(foreground_db), "%s/foreground.db", g_tmpdir); + (void)snprintf(background_db, sizeof(background_db), "%s/background.db", g_tmpdir); + + cbm_pipeline_t *foreground = cbm_pipeline_new(g_tmpdir, foreground_db, CBM_MODE_FULL); + cbm_pipeline_t *background = cbm_pipeline_new(g_tmpdir, background_db, CBM_MODE_FULL); + int foreground_rc = -1; + int background_rc = -1; + int foreground_nodes = -1; + int foreground_edges = -1; + int background_nodes = -1; + int background_edges = -1; + bool created = foreground && background; + if (created) { + cbm_pipeline_set_background(background, true); + foreground_rc = cbm_pipeline_run(foreground); + background_rc = cbm_pipeline_run(background); + cbm_pipeline_get_committed_counts(foreground, &foreground_nodes, &foreground_edges); + cbm_pipeline_get_committed_counts(background, &background_nodes, &background_edges); + } + + cbm_pipeline_free(foreground); + cbm_pipeline_free(background); + teardown_test_repo(); + saved_workers ? cbm_setenv("CBM_WORKERS", saved_workers, 1) : cbm_unsetenv("CBM_WORKERS"); + saved_single ? cbm_setenv("CBM_INDEX_SINGLE_THREAD", saved_single, 1) + : cbm_unsetenv("CBM_INDEX_SINGLE_THREAD"); + free(saved_workers); + free(saved_single); + + ASSERT_TRUE(created); + ASSERT_EQ(foreground_rc, 0); + ASSERT_EQ(background_rc, 0); + ASSERT_GT(foreground_nodes, 0); + ASSERT_EQ(background_nodes, foreground_nodes); + ASSERT_EQ(background_edges, foreground_edges); + PASS(); +} + TEST(pipeline_null_repo) { cbm_pipeline_t *p = cbm_pipeline_new(NULL, NULL, CBM_MODE_FULL); ASSERT_NULL(p); @@ -5380,8 +5494,8 @@ TEST(pipeline_semantic_manifest_rejects_non_directory_root) { cbm_file_hash_t *manifest = NULL; int manifest_count = -1; - int rc = cbm_pipeline_build_semantic_manifest("manifest-fail-closed", root_path, NULL, 0, NULL, - 0, NULL, NULL, &manifest, &manifest_count); + int rc = cbm_pipeline_build_semantic_manifest(NULL, "manifest-fail-closed", root_path, NULL, 0, + NULL, 0, NULL, NULL, &manifest, &manifest_count); cbm_pipeline_free_semantic_manifest(manifest, manifest_count > 0 ? manifest_count : 0); th_rmtree(tmp); @@ -16284,6 +16398,8 @@ SUITE(pipeline) { RUN_TEST(pipeline_lock_release_allows_contender); /* Lifecycle */ RUN_TEST(pipeline_create_free); + RUN_TEST(pipeline_background_worker_policy_preserves_overrides); + RUN_TEST(pipeline_background_policy_preserves_index_results); RUN_TEST(pipeline_null_repo); RUN_TEST(pipeline_free_null); RUN_TEST(pipeline_cancel);