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
11 changes: 6 additions & 5 deletions src/daemon/application.c
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down Expand Up @@ -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);
Expand Down
25 changes: 24 additions & 1 deletion src/mcp/mcp.c
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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);

Expand Down Expand Up @@ -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();
Expand Down
3 changes: 2 additions & 1 deletion src/pipeline/lsp_surface.c
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
37 changes: 31 additions & 6 deletions src/pipeline/pipeline.c
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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).
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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(). */
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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);
Expand Down
12 changes: 12 additions & 0 deletions src/pipeline/pipeline.h
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
28 changes: 16 additions & 12 deletions src/pipeline/pipeline_incremental.c
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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));
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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);

Expand Down Expand Up @@ -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 *));
Expand All @@ -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;
}
Expand Down Expand Up @@ -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;
}
Expand Down
6 changes: 3 additions & 3 deletions src/pipeline/pipeline_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
6 changes: 6 additions & 0 deletions tests/test_daemon_application.c
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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);
Expand Down
Loading
Loading