From e173c31065345d4d022ba377e3f18a733a4017f5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=88=98=E5=86=B2?= Date: Wed, 19 Aug 2026 14:39:38 +0800 Subject: [PATCH 1/3] feat(index): add worker resource watchdogs MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Discovery limits bound what indexing accepts, not what it then costs. A repository well inside those bounds can still exhaust the host through parser memory, or simply never finish, and a supervised worker that hangs leaves the parent waiting with nothing to report. Add index_max_rss_mb and index_max_duration_seconds, enforced by the parent against the worker process tree rather than the worker process alone, so a runaway child cannot hide behind a small parent. Resident memory is sampled through the platform interface on macOS, Linux and Windows. Crossing a limit terminates the tree and yields one trusted, structured terminal result that attributes the failure to the resource that caused it. Both limits default to off. A measurement that cannot be taken fails the attempt instead of passing it: a watchdog that quietly stops watching is worse than no watchdog at all. The shell fixture that stands in for the supervisor names the two new keys. The worker accepts only a policy that spells out every key it knows, which is what keeps a stale supervisor from starting a worker it cannot bound. Signed-off-by: 刘冲 --- README.md | 4 +- docs/CONFIGURATION.md | 10 +- docs/INDEX_RESOURCE_LIMITS.md | 48 +++++- scripts/test-runtime.sh | 3 +- src/cli/cli.c | 10 +- src/daemon/application.c | 19 ++- src/foundation/index_policy.c | 126 +++++++++----- src/foundation/index_policy.h | 11 ++ src/foundation/subprocess.c | 286 ++++++++++++++++++++++++++++++++ src/foundation/subprocess.h | 18 ++ src/mcp/index_supervisor.c | 256 ++++++++++++++++++++++++++-- src/mcp/index_supervisor.h | 28 +++- src/mcp/mcp.c | 114 ++++++++++--- src/mcp/mcp.h | 1 + src/mcp/mcp_internal.h | 4 + src/pipeline/pipeline.c | 5 +- tests/test_daemon_application.c | 133 +++++++++++++-- tests/test_index_policy.c | 72 +++++++- tests/test_index_supervisor.c | 269 ++++++++++++++++++++++++++++++ tests/test_mcp.c | 108 ++++++++++++ tests/test_subprocess.c | 47 ++++++ 21 files changed, 1462 insertions(+), 110 deletions(-) diff --git a/README.md b/README.md index 66b50fe08..a70aa7805 100644 --- a/README.md +++ b/README.md @@ -738,10 +738,12 @@ codebase-memory-mcp config set auto_watch false # don't register backgr codebase-memory-mcp config set watcher_enabled false # stop the watcher thread entirely (default: true) codebase-memory-mcp config set index_max_files 250000 # optional per-index source-file limit codebase-memory-mcp config set index_max_source_mb 16384 # optional per-index source-size limit +codebase-memory-mcp config set index_max_rss_mb 8192 # optional worker-tree current RSS limit +codebase-memory-mcp config set index_max_duration_seconds 3600 # optional total worker duration codebase-memory-mcp config reset auto_index # reset to default ``` -The two `index_max_*` settings default to `off`. Exceeding one fails the complete +The four `index_max_*` settings default to `off`. Exceeding one fails the complete index attempt rather than publishing a partial graph; an existing serving index is preserved. See [Index resource limits](docs/INDEX_RESOURCE_LIMITS.md). diff --git a/docs/CONFIGURATION.md b/docs/CONFIGURATION.md index 038256976..79b74e7e1 100644 --- a/docs/CONFIGURATION.md +++ b/docs/CONFIGURATION.md @@ -91,6 +91,8 @@ Current keys: | `watcher_enabled` | `true` | Master switch for the background watcher subsystem. Set `false` to stop the watcher from starting at all — no poll thread and no project registration. Reindex manually with `index_repository` when disabled. | | `index_max_files` | `off` | Optional maximum number of accepted source files in one discovery run. | | `index_max_source_mb` | `off` | Optional maximum accepted source size in MiB in one discovery run. | +| `index_max_rss_mb` | `off` | Optional maximum current RSS in MiB for the complete contained index-worker process tree (`64..1048576`). | +| `index_max_duration_seconds` | `off` | Optional maximum total worker duration in seconds (`1..86400`). | > **`watcher_enabled` vs `auto_watch`.** `watcher_enabled` controls whether the > watcher *subsystem* starts at all (the background poll thread). `auto_watch` is @@ -118,11 +120,13 @@ Current keys: > `auto_index` still runs, and `index_repository` stays available for manual > reindexing. -The two `index_max_*` settings are independent and disabled by default. They +The four `index_max_*` settings are independent and disabled by default. They apply to explicit indexing, automatic indexing, and watcher re-indexing, but not to `cross-repo-intelligence`, which does not scan repository source files. -Equality is allowed; exceeding either setting fails the complete index request -and preserves any previously serving database. See +Equality is allowed; exceeding any setting fails the complete index request and +preserves any previously serving database. Worker RSS covers descendants and is +not the same as the internal `CBM_MEM_BUDGET_MB` allocation budget. Total +duration is independent of the existing 15-minute no-log-progress timeout. See [Index resource limits](INDEX_RESOURCE_LIMITS.md) for counting, validation, and error-response details. diff --git a/docs/INDEX_RESOURCE_LIMITS.md b/docs/INDEX_RESOURCE_LIMITS.md index a3d57e015..66be4ecf4 100644 --- a/docs/INDEX_RESOURCE_LIMITS.md +++ b/docs/INDEX_RESOURCE_LIMITS.md @@ -1,8 +1,8 @@ # Index resource limits Index resource limits are optional operator controls for repositories whose -discovery breadth is not known in advance. They are disabled by default so -existing large-repository workloads retain their current behavior. +discovery breadth or worker runtime is not known in advance. They are disabled +by default so existing large-repository workloads retain their current behavior. ## Discovery settings @@ -23,6 +23,39 @@ Values use base-10 integers. MiB means 1,048,576 bytes. Empty values, zero, negative values, suffixes, trailing characters, and values outside the stated ranges are rejected without changing the stored value. +## Worker settings + +| Key | Default | Accepted value | Protects | +|---|---:|---:|---| +| `index_max_rss_mb` | `off` | `off` or `64..1048576` | Current RSS of the complete worker process tree | +| `index_max_duration_seconds` | `off` | `off` or `1..86400` | Total worker wall-clock duration | + +```bash +codebase-memory-mcp config set index_max_rss_mb 8192 +codebase-memory-mcp config set index_max_duration_seconds 3600 +``` + +RSS is the current resident memory of the contained worker and every descendant, +not the worker's allocation budget and not peak memory. This hard watchdog is +separate from the internal `CBM_MEM_BUDGET_MB` soft budget. The supervisor +samples RSS at most once every 250 milliseconds so the watchdog does not turn +full process-table enumeration into a busy loop. + +Duration uses a monotonic clock from successful spawn. It is independent of the +existing 15-minute quiet timeout: continuous log progress does not reset total +duration, while the quiet timeout continues to identify a worker that stops +making progress. + +Equality is allowed. The first RSS or elapsed-duration observation above its +limit starts the existing graceful-to-force process-tree shutdown. CBM reports +terminal only after the tree is quiescent or a bounded containment failure is +explicitly surfaced. Resource termination is not retried and does not +quarantine a source file. + +If RSS is enabled and three consecutive probes cannot obtain any trustworthy +tree measurement while the root worker is still running, CBM fails closed with +`code=resource_probe_failed`. + ## Counting and failure semantics `index_max_files` counts a file only after it passes directory pruning, ignore @@ -55,6 +88,12 @@ The previous database remains available because publication occurs only after a complete discovery and successful staged build. If no previous database exists, `serving_index_preserved` is false. +Worker limit failures use the same shape with `stage=worker`, +`resource=rss_bytes` and `unit=bytes`, or `resource=duration_ms` and +`unit=milliseconds`. RSS measurement failures use `code=resource_probe_failed` +and omit `observed`, `limit`, and `unit` because no trustworthy observation was +available. + ## Trust and compatibility Limits are read from the CLI-managed `_config.db`; they are not MCP request @@ -64,5 +103,6 @@ parent policy. These settings do not replace or increase `auto_index_limit`, change the 512 MiB single-file cap, alter workspace-root authorization, or affect -`cross-repo-intelligence`. With both settings `off`, discovery follows the -existing unbounded path. +`cross-repo-intelligence`. With all settings `off`, discovery follows the +existing path and the supervisor performs no periodic RSS probe or total-duration +termination. diff --git a/scripts/test-runtime.sh b/scripts/test-runtime.sh index 5d2c0c3ae..e6e60b15a 100644 --- a/scripts/test-runtime.sh +++ b/scripts/test-runtime.sh @@ -112,7 +112,8 @@ _cbm_test_runtime_daemon() { # cbm_mcp_index_policy_add_to_args: every key in cbm_index_policy_key_at, and # nothing else. cbm_test_index_worker_policy_json() { - printf '%s' '"_cbm_index_policy":{"index_max_files":"off","index_max_source_mb":"off"}' + printf '%s' '"_cbm_index_policy":{"index_max_files":"off","index_max_source_mb":"off"' + printf '%s' ',"index_max_rss_mb":"off","index_max_duration_seconds":"off"}' } cbm_test_runtime_cleanup() { diff --git a/src/cli/cli.c b/src/cli/cli.c index b59a894a6..66817c902 100644 --- a/src/cli/cli.c +++ b/src/cli/cli.c @@ -7449,6 +7449,8 @@ static const config_key_def_t CONFIG_KEYS[] = { {CBM_CONFIG_UI_PORT, "9749", "Port for the graph UI listener when enabled"}, {CBM_INDEX_CONFIG_MAX_FILES, "off", "Max accepted source files per index, or off"}, {CBM_INDEX_CONFIG_MAX_SOURCE_MB, "off", "Max accepted source MiB per index, or off"}, + {CBM_INDEX_CONFIG_MAX_RSS_MB, "off", "Max worker process-tree RSS MiB, or off"}, + {CBM_INDEX_CONFIG_MAX_DURATION_SECONDS, "off", "Max worker duration in seconds, or off"}, }; /* #1558: ui_enabled and ui_port were reachable ONLY by hand-editing @@ -7475,8 +7477,12 @@ static bool config_key_is_ui(const char *key) { } static bool config_key_is_index_policy(const char *key) { - return key && (strcmp(key, CBM_INDEX_CONFIG_MAX_FILES) == 0 || - strcmp(key, CBM_INDEX_CONFIG_MAX_SOURCE_MB) == 0); + for (size_t index = 0; key && index < cbm_index_policy_key_count(); index++) { + if (strcmp(key, cbm_index_policy_key_at(index)) == 0) { + return true; + } + } + return false; } static int config_index_policy_write(cbm_config_t *config, const char *key, const char *value) { diff --git a/src/daemon/application.c b/src/daemon/application.c index 748b8bb75..f8922e36f 100644 --- a/src/daemon/application.c +++ b/src/daemon/application.c @@ -305,9 +305,18 @@ static int application_worker_start_default(void *context, const char *args_json const char *quarantine_file, cbm_daemon_application_worker_t *worker_out) { (void)context; + cbm_index_resource_policy_t resource_policy; + char error[CBM_SZ_256] = {0}; + if (!cbm_mcp_index_policy_from_internal_args(args_json, &resource_policy, error, + sizeof(error))) { + cbm_log_error("daemon.index.policy", "error", error); + *worker_out = NULL; + return -1; + } cbm_index_worker_handle_t *worker = NULL; - int result = cbm_index_worker_start(args_json, memory_budget_bytes, false, marker_file, - quarantine_file, &worker); + int result = + cbm_index_worker_start_with_policy(args_json, memory_budget_bytes, &resource_policy, false, + marker_file, quarantine_file, &worker); *worker_out = worker; return result; } @@ -1284,6 +1293,12 @@ static application_attempt_decision_t application_consume_attempt( application_attempt_free(attempt); return APPLICATION_ATTEMPT_DECISION_SUCCESS; } + if (disposition == CBM_MCP_SUPERVISED_RESULT_RESOURCE_FAILURE) { + execution->response = + cbm_mcp_index_worker_resource_response(job->args_json, &attempt->result); + application_attempt_free(attempt); + return APPLICATION_ATTEMPT_DECISION_STOP; + } if (disposition == CBM_MCP_SUPERVISED_RESULT_UNSAFE_TERMINAL) { execution->unsafe_terminal = true; execution->supervision_failed = diff --git a/src/foundation/index_policy.c b/src/foundation/index_policy.c index 0d2b986bf..80aa2f8a6 100644 --- a/src/foundation/index_policy.c +++ b/src/foundation/index_policy.c @@ -1,21 +1,38 @@ #include "foundation/index_policy.h" +#include #include #include -static const char *const INDEX_POLICY_KEYS[] = { - CBM_INDEX_CONFIG_MAX_FILES, - CBM_INDEX_CONFIG_MAX_SOURCE_MB, +typedef struct { + const char *key; + size_t field_offset; + uint64_t minimum; + uint64_t maximum; + uint64_t multiplier; +} index_policy_metadata_t; + +static const index_policy_metadata_t INDEX_POLICY_METADATA[] = { + {CBM_INDEX_CONFIG_MAX_FILES, offsetof(cbm_index_resource_policy_t, max_files), 1, + CBM_INDEX_MAX_FILES_VALUE, 1}, + {CBM_INDEX_CONFIG_MAX_SOURCE_MB, offsetof(cbm_index_resource_policy_t, max_source_bytes), 1, + CBM_INDEX_MAX_SOURCE_MB_VALUE, CBM_INDEX_MIB_BYTES}, + {CBM_INDEX_CONFIG_MAX_RSS_MB, offsetof(cbm_index_resource_policy_t, max_rss_bytes), + CBM_INDEX_MIN_RSS_MB_VALUE, CBM_INDEX_MAX_RSS_MB_VALUE, CBM_INDEX_MIB_BYTES}, + {CBM_INDEX_CONFIG_MAX_DURATION_SECONDS, offsetof(cbm_index_resource_policy_t, max_duration_ms), + 1, CBM_INDEX_MAX_DURATION_SECONDS_VALUE, UINT64_C(1000)}, }; -static void set_error(char *error, size_t error_size, const char *key, uint64_t maximum) { +static void set_error(char *error, size_t error_size, const char *key, uint64_t minimum, + uint64_t maximum) { if (error && error_size > 0) { - (void)snprintf(error, error_size, "%s must be off or an integer from 1 to %llu", key, - (unsigned long long)maximum); + (void)snprintf(error, error_size, "%s must be off or an integer from %llu to %llu", key, + (unsigned long long)minimum, (unsigned long long)maximum); } } -static bool parse_bounded_uint64(const char *value, uint64_t maximum, uint64_t *parsed) { +static bool parse_bounded_uint64(const char *value, uint64_t minimum, uint64_t maximum, + uint64_t *parsed) { if (!value || !value[0] || !parsed) { return false; } @@ -30,7 +47,7 @@ static bool parse_bounded_uint64(const char *value, uint64_t maximum, uint64_t * } result = result * 10U + digit; } - if (result == 0) { + if (result < minimum) { return false; } *parsed = result; @@ -44,20 +61,39 @@ void cbm_index_policy_init(cbm_index_resource_policy_t *policy) { } bool cbm_index_policy_enabled(const cbm_index_resource_policy_t *policy) { + if (!policy) { + return false; + } + for (size_t index = 0; index < cbm_index_policy_key_count(); index++) { + const cbm_index_limit_u64_t *limit = + (const cbm_index_limit_u64_t *)((const unsigned char *)policy + + INDEX_POLICY_METADATA[index].field_offset); + if (limit->enabled) { + return true; + } + } + return false; +} + +bool cbm_index_policy_discovery_enabled(const cbm_index_resource_policy_t *policy) { return policy && (policy->max_files.enabled || policy->max_source_bytes.enabled); } +bool cbm_index_policy_worker_enabled(const cbm_index_resource_policy_t *policy) { + return policy && (policy->max_rss_bytes.enabled || policy->max_duration_ms.enabled); +} + size_t cbm_index_policy_key_count(void) { - return sizeof(INDEX_POLICY_KEYS) / sizeof(INDEX_POLICY_KEYS[0]); + return sizeof(INDEX_POLICY_METADATA) / sizeof(INDEX_POLICY_METADATA[0]); } const char *cbm_index_policy_key_at(size_t index) { - return index < cbm_index_policy_key_count() ? INDEX_POLICY_KEYS[index] : NULL; + return index < cbm_index_policy_key_count() ? INDEX_POLICY_METADATA[index].key : NULL; } const char *cbm_index_policy_default_value(const char *key) { for (size_t index = 0; index < cbm_index_policy_key_count(); index++) { - if (key && strcmp(key, INDEX_POLICY_KEYS[index]) == 0) { + if (key && strcmp(key, INDEX_POLICY_METADATA[index].key) == 0) { return "off"; } } @@ -70,34 +106,33 @@ bool cbm_index_policy_set(cbm_index_resource_policy_t *policy, const char *key, error[0] = '\0'; } if (!policy || !key || !value) { - set_error(error, error_size, key ? key : "index resource limit", 0); + set_error(error, error_size, key ? key : "index resource limit", 0, 0); return false; } - cbm_index_limit_u64_t *target = NULL; - uint64_t maximum = 0; - uint64_t multiplier = 1; - if (strcmp(key, CBM_INDEX_CONFIG_MAX_FILES) == 0) { - target = &policy->max_files; - maximum = CBM_INDEX_MAX_FILES_VALUE; - } else if (strcmp(key, CBM_INDEX_CONFIG_MAX_SOURCE_MB) == 0) { - target = &policy->max_source_bytes; - maximum = CBM_INDEX_MAX_SOURCE_MB_VALUE; - multiplier = CBM_INDEX_MIB_BYTES; - } else { - set_error(error, error_size, key, 0); + const index_policy_metadata_t *metadata = NULL; + for (size_t index = 0; index < cbm_index_policy_key_count(); index++) { + if (strcmp(key, INDEX_POLICY_METADATA[index].key) == 0) { + metadata = &INDEX_POLICY_METADATA[index]; + break; + } + } + if (!metadata) { + set_error(error, error_size, key, 0, 0); return false; } + cbm_index_limit_u64_t *target = + (cbm_index_limit_u64_t *)((unsigned char *)policy + metadata->field_offset); cbm_index_limit_u64_t candidate = {0}; if (strcmp(value, "off") != 0) { uint64_t parsed = 0; - if (!parse_bounded_uint64(value, maximum, &parsed)) { - set_error(error, error_size, key, maximum); + if (!parse_bounded_uint64(value, metadata->minimum, metadata->maximum, &parsed)) { + set_error(error, error_size, key, metadata->minimum, metadata->maximum); return false; } candidate.enabled = true; - candidate.value = parsed * multiplier; + candidate.value = parsed * metadata->multiplier; } *target = candidate; return true; @@ -108,18 +143,21 @@ bool cbm_index_policy_format(const cbm_index_resource_policy_t *policy, const ch if (!policy || !key || !out || out_size == 0) { return false; } - const cbm_index_limit_u64_t *limit = NULL; - uint64_t divisor = 1; - if (strcmp(key, CBM_INDEX_CONFIG_MAX_FILES) == 0) { - limit = &policy->max_files; - } else if (strcmp(key, CBM_INDEX_CONFIG_MAX_SOURCE_MB) == 0) { - limit = &policy->max_source_bytes; - divisor = CBM_INDEX_MIB_BYTES; - } else { + const index_policy_metadata_t *metadata = NULL; + for (size_t index = 0; index < cbm_index_policy_key_count(); index++) { + if (strcmp(key, INDEX_POLICY_METADATA[index].key) == 0) { + metadata = &INDEX_POLICY_METADATA[index]; + break; + } + } + if (!metadata) { return false; } + const cbm_index_limit_u64_t *limit = + (const cbm_index_limit_u64_t *)((const unsigned char *)policy + metadata->field_offset); int length = limit->enabled - ? snprintf(out, out_size, "%llu", (unsigned long long)(limit->value / divisor)) + ? snprintf(out, out_size, "%llu", + (unsigned long long)(limit->value / metadata->multiplier)) : snprintf(out, out_size, "off"); return length >= 0 && (size_t)length < out_size; } @@ -130,6 +168,10 @@ const char *cbm_index_resource_name(cbm_index_resource_t resource) { return "files"; case CBM_INDEX_RESOURCE_SOURCE_BYTES: return "source_bytes"; + case CBM_INDEX_RESOURCE_RSS_BYTES: + return "rss_bytes"; + case CBM_INDEX_RESOURCE_DURATION_MS: + return "duration_ms"; case CBM_INDEX_RESOURCE_NONE: default: return "unknown"; @@ -137,7 +179,13 @@ const char *cbm_index_resource_name(cbm_index_resource_t resource) { } const char *cbm_index_resource_unit(cbm_index_resource_t resource) { - return resource == CBM_INDEX_RESOURCE_FILES ? "files" : "bytes"; + if (resource == CBM_INDEX_RESOURCE_FILES) { + return "files"; + } + if (resource == CBM_INDEX_RESOURCE_DURATION_MS) { + return "milliseconds"; + } + return "bytes"; } const char *cbm_index_resource_config_key(cbm_index_resource_t resource) { @@ -146,6 +194,10 @@ const char *cbm_index_resource_config_key(cbm_index_resource_t resource) { return CBM_INDEX_CONFIG_MAX_FILES; case CBM_INDEX_RESOURCE_SOURCE_BYTES: return CBM_INDEX_CONFIG_MAX_SOURCE_MB; + case CBM_INDEX_RESOURCE_RSS_BYTES: + return CBM_INDEX_CONFIG_MAX_RSS_MB; + case CBM_INDEX_RESOURCE_DURATION_MS: + return CBM_INDEX_CONFIG_MAX_DURATION_SECONDS; case CBM_INDEX_RESOURCE_NONE: default: return "index_resource_limit"; diff --git a/src/foundation/index_policy.h b/src/foundation/index_policy.h index 6a6c499fc..617a6ab7d 100644 --- a/src/foundation/index_policy.h +++ b/src/foundation/index_policy.h @@ -7,9 +7,14 @@ #define CBM_INDEX_CONFIG_MAX_FILES "index_max_files" #define CBM_INDEX_CONFIG_MAX_SOURCE_MB "index_max_source_mb" +#define CBM_INDEX_CONFIG_MAX_RSS_MB "index_max_rss_mb" +#define CBM_INDEX_CONFIG_MAX_DURATION_SECONDS "index_max_duration_seconds" #define CBM_INDEX_MAX_FILES_VALUE UINT64_C(10000000) #define CBM_INDEX_MAX_SOURCE_MB_VALUE UINT64_C(1048576) +#define CBM_INDEX_MIN_RSS_MB_VALUE UINT64_C(64) +#define CBM_INDEX_MAX_RSS_MB_VALUE UINT64_C(1048576) +#define CBM_INDEX_MAX_DURATION_SECONDS_VALUE UINT64_C(86400) #define CBM_INDEX_MIB_BYTES UINT64_C(1048576) typedef struct { @@ -20,12 +25,16 @@ typedef struct { typedef struct { cbm_index_limit_u64_t max_files; cbm_index_limit_u64_t max_source_bytes; + cbm_index_limit_u64_t max_rss_bytes; + cbm_index_limit_u64_t max_duration_ms; } cbm_index_resource_policy_t; typedef enum { CBM_INDEX_RESOURCE_NONE = 0, CBM_INDEX_RESOURCE_FILES, CBM_INDEX_RESOURCE_SOURCE_BYTES, + CBM_INDEX_RESOURCE_RSS_BYTES, + CBM_INDEX_RESOURCE_DURATION_MS, } cbm_index_resource_t; typedef struct { @@ -36,6 +45,8 @@ typedef struct { void cbm_index_policy_init(cbm_index_resource_policy_t *policy); bool cbm_index_policy_enabled(const cbm_index_resource_policy_t *policy); +bool cbm_index_policy_discovery_enabled(const cbm_index_resource_policy_t *policy); +bool cbm_index_policy_worker_enabled(const cbm_index_resource_policy_t *policy); size_t cbm_index_policy_key_count(void); const char *cbm_index_policy_key_at(size_t index); diff --git a/src/foundation/subprocess.c b/src/foundation/subprocess.c index 83f788937..b14a74fbc 100644 --- a/src/foundation/subprocess.c +++ b/src/foundation/subprocess.c @@ -8,9 +8,11 @@ #include "compat.h" /* cbm_nanosleep */ #include "compat_fs.h" #include "log.h" +#include "mem_core.h" #include "platform.h" /* cbm_now_ms */ #include "sanitized.h" /* CBM_SANITIZED — spawn-retry budget */ +#include #include #include #include @@ -18,14 +20,17 @@ #ifdef _WIN32 #include +#include #include "win_utf8.h" /* cbm_utf8_to_wide — spawn the worker with a wide command line so a * non-ASCII repo path survives CreateProcess (#423/#20) */ #include /* free */ #else +#include #include #include #include #ifdef __APPLE__ +#include #include extern char **environ; #endif @@ -437,6 +442,287 @@ struct cbm_subprocess { #endif }; +static void cbm_rss_add_saturated(uint64_t *total, uint64_t value) { + if (value > UINT64_MAX - *total) { + *total = UINT64_MAX; + } else { + *total += value; + } +} + +#ifdef CBM_ENABLE_TEST_SEAMS +uint64_t cbm_subprocess_rss_sum_for_testing(const uint64_t *values, size_t count) { + uint64_t total = 0; + for (size_t index = 0; values && index < count; index++) { + cbm_rss_add_saturated(&total, values[index]); + } + return total; +} +#endif + +#ifdef _WIN32 +static cbm_proc_tree_rss_status_t cbm_subprocess_tree_rss_platform(cbm_subprocess_t *process, + uint64_t *rss_bytes) { + DWORD capacity = 32; + JOBOBJECT_BASIC_PROCESS_ID_LIST *processes = NULL; + for (;;) { + size_t bytes = sizeof(*processes) + (size_t)(capacity - 1) * sizeof(ULONG_PTR); + if (bytes > UINT32_MAX) { + cbm_free(CBM_MEM_CLASS_OTHER, processes); + return CBM_PROC_TREE_RSS_ERROR; + } + JOBOBJECT_BASIC_PROCESS_ID_LIST *grown = cbm_realloc(CBM_MEM_CLASS_OTHER, processes, bytes); + if (!grown) { + cbm_free(CBM_MEM_CLASS_OTHER, processes); + return CBM_PROC_TREE_RSS_ERROR; + } + processes = grown; + ZeroMemory(processes, bytes); + processes->NumberOfAssignedProcesses = capacity; + if (QueryInformationJobObject(process->job, JobObjectBasicProcessIdList, processes, + (DWORD)bytes, NULL)) { + break; + } + if (GetLastError() != ERROR_MORE_DATA || capacity > (UINT32_MAX / 2U)) { + cbm_free(CBM_MEM_CLASS_OTHER, processes); + return CBM_PROC_TREE_RSS_ERROR; + } + capacity *= 2U; + } + + uint64_t total = 0; + DWORD measured = 0; + bool root_failed = false; + for (DWORD index = 0; index < processes->NumberOfProcessIdsInList; index++) { + DWORD pid = (DWORD)processes->ProcessIdList[index]; + HANDLE member = OpenProcess(PROCESS_QUERY_INFORMATION | PROCESS_VM_READ, FALSE, pid); + if (!member) { + root_failed = root_failed || pid == process->process_id; + continue; + } + PROCESS_MEMORY_COUNTERS memory; + ZeroMemory(&memory, sizeof(memory)); + if (GetProcessMemoryInfo(member, &memory, sizeof(memory))) { + cbm_rss_add_saturated(&total, (uint64_t)memory.WorkingSetSize); + measured++; + } else if (pid == process->process_id) { + root_failed = true; + } + CloseHandle(member); + } + DWORD listed = processes->NumberOfProcessIdsInList; + cbm_free(CBM_MEM_CLASS_OTHER, processes); + if (root_failed || (measured == 0 && listed > 0)) { + return CBM_PROC_TREE_RSS_ERROR; + } + if (measured == 0) { + return CBM_PROC_TREE_RSS_EMPTY; + } + *rss_bytes = total; + return CBM_PROC_TREE_RSS_OK; +} +#elif defined(__APPLE__) +static cbm_proc_tree_rss_status_t cbm_subprocess_tree_rss_platform(cbm_subprocess_t *process, + uint64_t *rss_bytes) { + int capacity = proc_listallpids(NULL, 0); + if (capacity <= 0 || capacity > INT_MAX / (int)sizeof(pid_t) - 64) { + return CBM_PROC_TREE_RSS_ERROR; + } + capacity += 64; + pid_t *pids = NULL; + int count = 0; + for (;;) { + pid_t *grown = cbm_realloc(CBM_MEM_CLASS_OTHER, pids, (size_t)capacity * sizeof(*pids)); + if (!grown) { + cbm_free(CBM_MEM_CLASS_OTHER, pids); + return CBM_PROC_TREE_RSS_ERROR; + } + pids = grown; + count = proc_listallpids(pids, capacity * (int)sizeof(*pids)); + if (count <= 0) { + cbm_free(CBM_MEM_CLASS_OTHER, pids); + return CBM_PROC_TREE_RSS_ERROR; + } + if (count < capacity) { + break; + } + if (capacity > INT_MAX / (int)sizeof(*pids) / 2) { + cbm_free(CBM_MEM_CLASS_OTHER, pids); + return CBM_PROC_TREE_RSS_ERROR; + } + capacity *= 2; + } + + uint64_t total = 0; + int measured = 0; + bool root_failed = false; + for (int index = 0; index < count; index++) { + struct proc_bsdinfo info; + if (pids[index] <= 0) { + continue; + } + if (proc_pidinfo(pids[index], PROC_PIDTBSDINFO, 0, &info, sizeof(info)) != + (int)sizeof(info)) { + root_failed = root_failed || pids[index] == process->pid; + continue; + } + if ((pid_t)info.pbi_pgid != process->pgid) { + continue; + } + struct rusage_info_v2 usage; + if (proc_pid_rusage(pids[index], RUSAGE_INFO_V2, (rusage_info_t *)&usage) != 0) { + root_failed = root_failed || pids[index] == process->pid; + continue; + } + cbm_rss_add_saturated(&total, usage.ri_resident_size); + measured++; + } + cbm_free(CBM_MEM_CLASS_OTHER, pids); + if (root_failed) { + return CBM_PROC_TREE_RSS_ERROR; + } + if (measured == 0) { + return CBM_PROC_TREE_RSS_EMPTY; + } + *rss_bytes = total; + return CBM_PROC_TREE_RSS_OK; +} +#else +static bool cbm_proc_pid_name(const char *name) { + if (!name || !name[0]) { + return false; + } + for (const unsigned char *cursor = (const unsigned char *)name; *cursor; cursor++) { + if (*cursor < '0' || *cursor > '9') { + return false; + } + } + return true; +} + +static bool cbm_linux_proc_stat(const char *pid_name, pid_t *group, int64_t *rss_pages) { + char path[64]; + int length = snprintf(path, sizeof(path), "/proc/%s/stat", pid_name); + if (length <= 0 || length >= (int)sizeof(path)) { + return false; + } + FILE *stat_file = cbm_fopen(path, "rb"); + if (!stat_file) { + return false; + } + char stat_line[4096]; + bool read = fgets(stat_line, sizeof(stat_line), stat_file) != NULL; + (void)fclose(stat_file); + char *command_end = read ? strrchr(stat_line, ')') : NULL; + if (!command_end || command_end[1] != ' ') { + return false; + } + + char *save = NULL; + char *token = strtok_r(command_end + 2, " ", &save); + int field = 3; + bool have_group = false; + bool have_rss = false; + while (token) { + if (field == 5 || field == 24) { + char *end = NULL; + long long value = strtoll(token, &end, 10); + if (!end || *end != '\0') { + return false; + } + if (field == 5) { + *group = (pid_t)value; + have_group = true; + } else { + *rss_pages = (int64_t)value; + have_rss = true; + break; + } + } + token = strtok_r(NULL, " ", &save); + field++; + } + return have_group && have_rss; +} + +static cbm_proc_tree_rss_status_t cbm_subprocess_tree_rss_platform(cbm_subprocess_t *process, + uint64_t *rss_bytes) { + DIR *proc = opendir("/proc"); + if (!proc) { + return CBM_PROC_TREE_RSS_ERROR; + } + long page_size = sysconf(_SC_PAGESIZE); + if (page_size <= 0) { + (void)closedir(proc); + return CBM_PROC_TREE_RSS_ERROR; + } + + uint64_t total = 0; + int measured = 0; + bool root_failed = false; + struct dirent *entry; + while ((entry = readdir(proc)) != NULL) { + if (!cbm_proc_pid_name(entry->d_name)) { + continue; + } + pid_t group = 0; + int64_t pages = 0; + bool root_entry = strtol(entry->d_name, NULL, 10) == (long)process->pid; + if (!cbm_linux_proc_stat(entry->d_name, &group, &pages)) { + root_failed = root_failed || root_entry; + continue; + } + if (group != process->pgid || pages < 0) { + continue; + } + uint64_t page_count = (uint64_t)pages; + uint64_t bytes = page_count > UINT64_MAX / (uint64_t)page_size + ? UINT64_MAX + : page_count * (uint64_t)page_size; + cbm_rss_add_saturated(&total, bytes); + measured++; + } + (void)closedir(proc); + if (root_failed) { + return CBM_PROC_TREE_RSS_ERROR; + } + if (measured == 0) { + return CBM_PROC_TREE_RSS_EMPTY; + } + *rss_bytes = total; + return CBM_PROC_TREE_RSS_OK; +} +#endif + +cbm_proc_tree_rss_status_t cbm_subprocess_tree_rss_bytes(cbm_subprocess_t *process, + uint64_t *rss_bytes) { + if (!process || !rss_bytes) { + return CBM_PROC_TREE_RSS_ERROR; + } + *rss_bytes = 0; + return cbm_subprocess_tree_rss_platform(process, rss_bytes); +} + +bool cbm_subprocess_root_running(const cbm_subprocess_t *process) { + if (!process || process->root_reaped) { + return false; + } + int lifecycle = atomic_load_explicit(&process->lifecycle, memory_order_acquire); + return lifecycle == CBM_SUBPROCESS_ACTIVE || lifecycle == CBM_SUBPROCESS_CANCEL_REQUESTED; +} + +bool cbm_subprocess_termination_pending(const cbm_subprocess_t *process) { + return process && process->termination_started; +} + +bool cbm_subprocess_supervision_active(const cbm_subprocess_t *process) { + if (!process) { + return false; + } + int lifecycle = atomic_load_explicit(&process->lifecycle, memory_order_acquire); + return lifecycle == CBM_SUBPROCESS_ACTIVE || lifecycle == CBM_SUBPROCESS_CANCEL_REQUESTED; +} + static void cbm_subprocess_result_init(cbm_proc_result_t *result) { result->outcome = CBM_PROC_SPAWN_FAILED; result->exit_code = -1; diff --git a/src/foundation/subprocess.h b/src/foundation/subprocess.h index 54940de11..1179d97c4 100644 --- a/src/foundation/subprocess.h +++ b/src/foundation/subprocess.h @@ -22,6 +22,7 @@ #include #include /* size_t (cbm_build_win_cmdline) */ +#include /* How a supervised child ended. */ typedef enum { @@ -92,6 +93,12 @@ typedef enum { CBM_PROC_POLL_TERMINAL = 1 } cbm_proc_poll_t; +typedef enum { + CBM_PROC_TREE_RSS_ERROR = -1, + CBM_PROC_TREE_RSS_EMPTY = 0, + CBM_PROC_TREE_RSS_OK = 1, +} cbm_proc_tree_rss_status_t; + /* Spawn opts->bin and return immediately with a supervisor handle. On success, * *out owns the process until a terminal poll followed by destroy. Spawn copies * the option strings/argv it needs after return; log_ud remains caller-owned until @@ -122,6 +129,16 @@ int cbm_subprocess_spawn(const cbm_proc_opts_t *opts, cbm_subprocess_t **out); * tree when cancel_grace_ms elapses. Callers must keep polling to make progress. */ cbm_proc_poll_t cbm_subprocess_poll(cbm_subprocess_t *process, cbm_proc_result_t *out); +/* Read the current resident-set size of the complete contained process tree. + * OK returns an overflow-safe byte total, EMPTY means the owned tree currently + * has no observable members, and ERROR means no trustworthy measurement could + * be obtained. Individual processes that exit during enumeration are ignored. */ +cbm_proc_tree_rss_status_t cbm_subprocess_tree_rss_bytes(cbm_subprocess_t *process, + uint64_t *rss_bytes); +bool cbm_subprocess_root_running(const cbm_subprocess_t *process); +bool cbm_subprocess_termination_pending(const cbm_subprocess_t *process); +bool cbm_subprocess_supervision_active(const cbm_subprocess_t *process); + /* Record an explicit cancellation request without waiting. Safe to repeat and * safe to call from a cancellation thread while one owner thread polls. The * owner must stop cancellation producers before destroying the handle. true @@ -188,6 +205,7 @@ bool cbm_build_win_cmd_payload(char *buf, size_t cap, const char *cmd_executable * of hoping a loaded machine reproduces it. Test builds only. */ void cbm_subprocess_force_spawn_eagain_for_testing(int attempts); int cbm_subprocess_pending_spawn_eagain_for_testing(void); +uint64_t cbm_subprocess_rss_sum_for_testing(const uint64_t *values, size_t count); #endif #endif /* CBM_SUBPROCESS_H */ diff --git a/src/mcp/index_supervisor.c b/src/mcp/index_supervisor.c index 4cd69f642..93d3bef46 100644 --- a/src/mcp/index_supervisor.c +++ b/src/mcp/index_supervisor.c @@ -410,8 +410,56 @@ enum { INDEX_WORKER_SYNC_POLL_NS = 10000000, INDEX_WORKER_RELAY_LINES_PER_POLL = 64, INDEX_WORKER_RELAY_BYTES_PER_POLL = 64 * 1024, + INDEX_WORKER_RSS_PROBE_FAILURE_LIMIT = 3, + INDEX_WORKER_RSS_PROBE_INTERVAL_MS = 250, }; +typedef enum { + INDEX_WORKER_TERMINATION_NONE = 0, + INDEX_WORKER_TERMINATION_CANCEL_PENDING, + INDEX_WORKER_TERMINATION_CANCEL, + INDEX_WORKER_TERMINATION_RESOURCE, +} index_worker_termination_reason_t; + +#ifdef CBM_ENABLE_TEST_SEAMS +static cbm_index_supervisor_clock_fn g_resource_clock; +static cbm_index_supervisor_rss_fn g_resource_rss; +static void *g_resource_hook_context; + +void cbm_index_supervisor_set_resource_hooks_for_testing(cbm_index_supervisor_clock_fn clock_fn, + cbm_index_supervisor_rss_fn rss_fn, + void *context) { + g_resource_clock = clock_fn; + g_resource_rss = rss_fn; + g_resource_hook_context = context; +} + +void cbm_index_supervisor_reset_resource_hooks_for_testing(void) { + g_resource_clock = NULL; + g_resource_rss = NULL; + g_resource_hook_context = NULL; +} +#endif + +static uint64_t worker_resource_now_ms(void) { +#ifdef CBM_ENABLE_TEST_SEAMS + if (g_resource_clock) { + return g_resource_clock(g_resource_hook_context); + } +#endif + return cbm_now_ms(); +} + +static cbm_proc_tree_rss_status_t worker_resource_rss(cbm_subprocess_t *process, + uint64_t *rss_bytes) { +#ifdef CBM_ENABLE_TEST_SEAMS + if (g_resource_rss) { + return g_resource_rss(process, rss_bytes, g_resource_hook_context); + } +#endif + return cbm_subprocess_tree_rss_bytes(process, rss_bytes); +} + struct cbm_index_worker_handle { cbm_subprocess_t *process; char response_path[INDEX_WORKER_PATH_CAP]; @@ -421,6 +469,12 @@ struct cbm_index_worker_handle { long relay_tail_pos; bool process_terminal; cbm_proc_result_t process_result; + cbm_index_resource_policy_t resource_policy; + uint64_t started_ms; + uint64_t last_rss_probe_ms; + unsigned int rss_probe_failures; + bool rss_probe_started; + atomic_int termination_reason; atomic_bool terminal; cbm_index_worker_result_t result; }; @@ -599,7 +653,8 @@ static bool worker_unique_file(char *out, size_t out_size, const char *kind) { static bool worker_result_succeeded(const cbm_index_worker_result_t *result) { return result && result->outcome == CBM_PROC_CLEAN && !result->cancellation_requested && - result->tree_quiesced && !result->supervision_failed; + result->resource_violation.resource == CBM_INDEX_RESOURCE_NONE && + !result->resource_probe_failed && result->tree_quiesced && !result->supervision_failed; } static void worker_terminal_log(cbm_index_worker_handle_t *handle) { @@ -630,6 +685,13 @@ static void worker_terminal_log(cbm_index_worker_handle_t *handle) { } else if (handle->result.supervision_failed || !handle->result.tree_quiesced) { cbm_log_error("index.supervisor.containment_failed", "outcome", cbm_proc_outcome_str(handle->result.outcome), "log", handle->log_path); + } else if (handle->result.resource_probe_failed) { + cbm_log_error("index.supervisor.resource_probe_failed", "resource", "rss_bytes", "log", + handle->log_path); + } else if (handle->result.resource_violation.resource != CBM_INDEX_RESOURCE_NONE) { + cbm_log_warn("index.supervisor.resource_limit", "resource", + cbm_index_resource_name(handle->result.resource_violation.resource), "log", + handle->log_path); } else if (handle->result.cancellation_requested) { cbm_log_warn("index.supervisor.worker_cancelled", "outcome", cbm_proc_outcome_str(handle->result.outcome), "log", handle->log_path); @@ -664,10 +726,87 @@ size_t cbm_index_worker_job_memory_limit(size_t memory_budget_bytes) { return memory_budget_bytes > SIZE_MAX - headroom ? SIZE_MAX : memory_budget_bytes + headroom; } -int cbm_index_worker_start_with_log(const char *args_json, size_t memory_budget_bytes, - bool single_thread, const char *marker_file, - const char *quarantine_file, cbm_proc_log_cb log_callback, - void *log_context, cbm_index_worker_handle_t **handle_out) { +static void worker_request_resource_termination(cbm_index_worker_handle_t *handle, + cbm_index_resource_t resource, uint64_t observed, + uint64_t limit, bool probe_failed) { + int expected = INDEX_WORKER_TERMINATION_NONE; + if (!atomic_compare_exchange_strong_explicit(&handle->termination_reason, &expected, + INDEX_WORKER_TERMINATION_RESOURCE, + memory_order_acq_rel, memory_order_acquire)) { + return; + } + handle->result.resource_violation = (cbm_index_resource_violation_t){ + .resource = resource, .observed = observed, .limit = limit}; + handle->result.resource_probe_failed = probe_failed; + if (!cbm_subprocess_request_cancel(handle->process)) { + handle->result.resource_violation = (cbm_index_resource_violation_t){0}; + handle->result.resource_probe_failed = false; + expected = INDEX_WORKER_TERMINATION_RESOURCE; + (void)atomic_compare_exchange_strong_explicit(&handle->termination_reason, &expected, + INDEX_WORKER_TERMINATION_NONE, + memory_order_acq_rel, memory_order_acquire); + } +} + +static bool worker_check_resource_limits(cbm_index_worker_handle_t *handle) { + if (atomic_load_explicit(&handle->termination_reason, memory_order_acquire) != + INDEX_WORKER_TERMINATION_NONE || + cbm_subprocess_termination_pending(handle->process) || + !cbm_subprocess_supervision_active(handle->process)) { + return false; + } + uint64_t now = 0; + bool have_now = false; + bool rss_probe_failed = false; + if (handle->resource_policy.max_rss_bytes.enabled) { + now = worker_resource_now_ms(); + have_now = true; + bool probe_due = !handle->rss_probe_started || now < handle->last_rss_probe_ms || + now - handle->last_rss_probe_ms >= INDEX_WORKER_RSS_PROBE_INTERVAL_MS; + if (!probe_due) { + goto duration_check; + } + handle->rss_probe_started = true; + handle->last_rss_probe_ms = now; + uint64_t rss_bytes = 0; + cbm_proc_tree_rss_status_t status = worker_resource_rss(handle->process, &rss_bytes); + if (status == CBM_PROC_TREE_RSS_OK) { + handle->rss_probe_failures = 0; + if (rss_bytes > handle->resource_policy.max_rss_bytes.value) { + worker_request_resource_termination(handle, CBM_INDEX_RESOURCE_RSS_BYTES, rss_bytes, + handle->resource_policy.max_rss_bytes.value, + false); + return false; + } + } else if (status == CBM_PROC_TREE_RSS_ERROR) { + handle->rss_probe_failures++; + if (handle->rss_probe_failures >= INDEX_WORKER_RSS_PROBE_FAILURE_LIMIT) { + rss_probe_failed = true; + } + } else { + handle->rss_probe_failures = 0; + } + } +duration_check: + if (handle->resource_policy.max_duration_ms.enabled) { + if (!have_now) { + now = worker_resource_now_ms(); + } + uint64_t elapsed = now >= handle->started_ms ? now - handle->started_ms : 0; + if (elapsed > handle->resource_policy.max_duration_ms.value) { + worker_request_resource_termination(handle, CBM_INDEX_RESOURCE_DURATION_MS, elapsed, + handle->resource_policy.max_duration_ms.value, + false); + } + } + return rss_probe_failed; +} + +static int worker_start_internal(const char *args_json, size_t memory_budget_bytes, + const cbm_index_resource_policy_t *resource_policy, + bool single_thread, const char *marker_file, + const char *quarantine_file, cbm_proc_log_cb log_callback, + void *log_context, cbm_index_worker_handle_t **handle_out) { if (handle_out) { *handle_out = NULL; } @@ -699,6 +838,11 @@ int cbm_index_worker_start_with_log(const char *args_json, size_t memory_budget_ if (!handle) { return -1; } + cbm_index_policy_init(&handle->resource_policy); + if (resource_policy) { + handle->resource_policy = *resource_policy; + } + atomic_init(&handle->termination_reason, INDEX_WORKER_TERMINATION_NONE); atomic_init(&handle->terminal, false); handle->log_callback = log_callback; handle->log_context = log_context; @@ -760,15 +904,35 @@ int cbm_index_worker_start_with_log(const char *args_json, size_t memory_budget_ cbm_log_error("index.supervisor.spawn_failed", "action", "fail_closed"); return -1; } + if (handle->resource_policy.max_duration_ms.enabled) { + handle->started_ms = worker_resource_now_ms(); + } *handle_out = handle; return 0; } +int cbm_index_worker_start_with_log(const char *args_json, size_t memory_budget_bytes, + bool single_thread, const char *marker_file, + const char *quarantine_file, cbm_proc_log_cb log_callback, + void *log_context, cbm_index_worker_handle_t **handle_out) { + return worker_start_internal(args_json, memory_budget_bytes, NULL, single_thread, marker_file, + quarantine_file, log_callback, log_context, handle_out); +} + int cbm_index_worker_start(const char *args_json, size_t memory_budget_bytes, bool single_thread, const char *marker_file, const char *quarantine_file, cbm_index_worker_handle_t **handle_out) { - return cbm_index_worker_start_with_log(args_json, memory_budget_bytes, single_thread, - marker_file, quarantine_file, NULL, NULL, handle_out); + return worker_start_internal(args_json, memory_budget_bytes, NULL, single_thread, marker_file, + quarantine_file, NULL, NULL, handle_out); +} + +int cbm_index_worker_start_with_policy(const char *args_json, size_t memory_budget_bytes, + const cbm_index_resource_policy_t *resource_policy, + bool single_thread, const char *marker_file, + const char *quarantine_file, + cbm_index_worker_handle_t **handle_out) { + return worker_start_internal(args_json, memory_budget_bytes, resource_policy, single_thread, + marker_file, quarantine_file, NULL, NULL, handle_out); } cbm_index_worker_poll_t cbm_index_worker_poll(cbm_index_worker_handle_t *handle, @@ -785,10 +949,16 @@ cbm_index_worker_poll_t cbm_index_worker_poll(cbm_index_worker_handle_t *handle, } bool relay_caught_up = true; if (!handle->process_terminal) { + bool rss_probe_failed = worker_check_resource_limits(handle); cbm_proc_result_t process_result; cbm_proc_poll_t state = cbm_subprocess_poll(handle->process, &process_result); relay_caught_up = worker_relay_log(handle); if (state == CBM_PROC_POLL_RUNNING) { + if (rss_probe_failed && cbm_subprocess_root_running(handle->process)) { + worker_request_resource_termination(handle, CBM_INDEX_RESOURCE_RSS_BYTES, 0, + handle->resource_policy.max_rss_bytes.value, + true); + } return CBM_INDEX_WORKER_POLL_RUNNING; } if (state != CBM_PROC_POLL_TERMINAL) { @@ -809,7 +979,12 @@ cbm_index_worker_poll_t cbm_index_worker_poll(cbm_index_worker_handle_t *handle, handle->result.outcome = process_result->outcome; handle->result.exit_code = process_result->exit_code; handle->result.term_signal = process_result->term_signal; - handle->result.cancellation_requested = process_result->cancellation_requested; + int termination_reason = + atomic_load_explicit(&handle->termination_reason, memory_order_acquire); + handle->result.cancellation_requested = + termination_reason == INDEX_WORKER_TERMINATION_CANCEL || + (termination_reason == INDEX_WORKER_TERMINATION_CANCEL_PENDING && + process_result->cancellation_requested); handle->result.forced = process_result->forced; handle->result.tree_quiesced = process_result->tree_quiesced; handle->result.supervision_failed = process_result->supervision_failed; @@ -832,8 +1007,36 @@ cbm_index_worker_poll_t cbm_index_worker_poll(cbm_index_worker_handle_t *handle, } bool cbm_index_worker_request_cancel(cbm_index_worker_handle_t *handle) { - return handle && !atomic_load_explicit(&handle->terminal, memory_order_acquire) && - cbm_subprocess_request_cancel(handle->process); + if (!handle || atomic_load_explicit(&handle->terminal, memory_order_acquire)) { + return false; + } + int expected = INDEX_WORKER_TERMINATION_NONE; + bool claimed = atomic_compare_exchange_strong_explicit( + &handle->termination_reason, &expected, INDEX_WORKER_TERMINATION_CANCEL_PENDING, + memory_order_acq_rel, memory_order_acquire); + if (!claimed) { + if (expected == INDEX_WORKER_TERMINATION_CANCEL) { + return true; + } + if (expected != INDEX_WORKER_TERMINATION_CANCEL_PENDING) { + return false; + } + bool accepted = cbm_subprocess_request_cancel(handle->process); + if (accepted) { + expected = INDEX_WORKER_TERMINATION_CANCEL_PENDING; + (void)atomic_compare_exchange_strong_explicit( + &handle->termination_reason, &expected, INDEX_WORKER_TERMINATION_CANCEL, + memory_order_acq_rel, memory_order_acquire); + } + return accepted; + } + bool accepted = cbm_subprocess_request_cancel(handle->process); + expected = INDEX_WORKER_TERMINATION_CANCEL_PENDING; + (void)atomic_compare_exchange_strong_explicit(&handle->termination_reason, &expected, + accepted ? INDEX_WORKER_TERMINATION_CANCEL + : INDEX_WORKER_TERMINATION_NONE, + memory_order_acq_rel, memory_order_acquire); + return accepted; } const char *cbm_index_worker_response_path(const cbm_index_worker_handle_t *handle) { @@ -858,18 +1061,19 @@ void cbm_index_worker_destroy(cbm_index_worker_handle_t *handle) { free(handle); } -int cbm_index_spawn_worker_with_log_cancel(const char *args_json, bool single_thread, - const char *marker_file, const char *quarantine_file, - cbm_proc_log_cb log_callback, void *log_context, - const atomic_int *cancel_requested, - cbm_index_worker_result_t *result) { +static int worker_spawn_internal(const char *args_json, + const cbm_index_resource_policy_t *resource_policy, + bool single_thread, const char *marker_file, + const char *quarantine_file, cbm_proc_log_cb log_callback, + void *log_context, const atomic_int *cancel_requested, + cbm_index_worker_result_t *result) { if (!result) { return -1; } worker_result_init(result); cbm_index_worker_handle_t *handle = NULL; - if (cbm_index_worker_start_with_log(args_json, 0, single_thread, marker_file, quarantine_file, - log_callback, log_context, &handle) != 0) { + if (worker_start_internal(args_json, 0, resource_policy, single_thread, marker_file, + quarantine_file, log_callback, log_context, &handle) != 0) { return -1; } const cbm_index_worker_result_t *cached = NULL; @@ -895,6 +1099,24 @@ int cbm_index_spawn_worker_with_log_cancel(const char *args_json, bool single_th return 0; } +int cbm_index_spawn_worker_with_log_cancel(const char *args_json, bool single_thread, + const char *marker_file, const char *quarantine_file, + cbm_proc_log_cb log_callback, void *log_context, + const atomic_int *cancel_requested, + cbm_index_worker_result_t *result) { + return worker_spawn_internal(args_json, NULL, single_thread, marker_file, quarantine_file, + log_callback, log_context, cancel_requested, result); +} + +int cbm_index_spawn_worker_with_policy_log_cancel( + const char *args_json, const cbm_index_resource_policy_t *resource_policy, bool single_thread, + const char *marker_file, const char *quarantine_file, cbm_proc_log_cb log_callback, + void *log_context, const atomic_int *cancel_requested, cbm_index_worker_result_t *result) { + return worker_spawn_internal(args_json, resource_policy, single_thread, marker_file, + quarantine_file, log_callback, log_context, cancel_requested, + result); +} + int cbm_index_spawn_worker_with_log(const char *args_json, bool single_thread, const char *marker_file, const char *quarantine_file, cbm_proc_log_cb log_callback, void *log_context, diff --git a/src/mcp/index_supervisor.h b/src/mcp/index_supervisor.h index 42b6bd74a..fdf7b1b44 100644 --- a/src/mcp/index_supervisor.h +++ b/src/mcp/index_supervisor.h @@ -22,8 +22,10 @@ #include #include +#include #include +#include "foundation/index_policy.h" #include "foundation/subprocess.h" /* cbm_proc_outcome_t */ /* Worker-role state, set once from the CLI arg parser (main.c) when this process @@ -148,8 +150,10 @@ typedef struct { bool tree_quiesced; bool supervision_failed; bool response_rejected; /* clean worker exceeded the bounded response protocol */ - char *response; /* worker result only after a contained, uncancelled CLEAN exit; - * borrowed for async polls, caller-owned from the sync wrapper */ + cbm_index_resource_violation_t resource_violation; + bool resource_probe_failed; + char *response; /* worker result only after a contained, uncancelled CLEAN exit; + * borrowed for async polls, caller-owned from the sync wrapper */ } cbm_index_worker_result_t; /* Daemon-owned, nonblocking supervisor for one contained worker process tree. */ @@ -171,6 +175,11 @@ size_t cbm_index_worker_job_memory_limit(size_t memory_budget_bytes); int cbm_index_worker_start(const char *args_json, size_t memory_budget_bytes, bool single_thread, const char *marker_file, const char *quarantine_file, cbm_index_worker_handle_t **handle_out); +int cbm_index_worker_start_with_policy(const char *args_json, size_t memory_budget_bytes, + const cbm_index_resource_policy_t *resource_policy, + bool single_thread, const char *marker_file, + const char *quarantine_file, + cbm_index_worker_handle_t **handle_out); /* Request-scoped variant used by interactive local CLI calls. The callback is * invoked by the owner thread while it polls the contained worker; log_context @@ -234,7 +243,22 @@ int cbm_index_spawn_worker_with_log_cancel(const char *args_json, bool single_th cbm_proc_log_cb log_callback, void *log_context, const atomic_int *cancel_requested, cbm_index_worker_result_t *result); +int cbm_index_spawn_worker_with_policy_log_cancel( + const char *args_json, const cbm_index_resource_policy_t *resource_policy, bool single_thread, + const char *marker_file, const char *quarantine_file, cbm_proc_log_cb log_callback, + void *log_context, const atomic_int *cancel_requested, cbm_index_worker_result_t *result); void cbm_index_worker_result_free(cbm_index_worker_result_t *result); +#ifdef CBM_ENABLE_TEST_SEAMS +typedef uint64_t (*cbm_index_supervisor_clock_fn)(void *context); +typedef cbm_proc_tree_rss_status_t (*cbm_index_supervisor_rss_fn)(cbm_subprocess_t *process, + uint64_t *rss_bytes, + void *context); +void cbm_index_supervisor_set_resource_hooks_for_testing(cbm_index_supervisor_clock_fn clock_fn, + cbm_index_supervisor_rss_fn rss_fn, + void *context); +void cbm_index_supervisor_reset_resource_hooks_for_testing(void); +#endif + #endif /* CBM_INDEX_SUPERVISOR_H */ diff --git a/src/mcp/mcp.c b/src/mcp/mcp.c index 75ff8e396..86d392ece 100644 --- a/src/mcp/mcp.c +++ b/src/mcp/mcp.c @@ -10417,6 +10417,22 @@ static bool build_index_success_response(cbm_mcp_server_t *srv, yyjson_mut_doc * return degraded; } +static bool project_db_is_servable(const char *project, const char *db_path); + +/* Serialize a worker JSON document and release the heap strings it borrowed. + * Extra arguments may be NULL; free(NULL) is a no-op. */ +static char *index_worker_json_result(yyjson_mut_doc *doc, char *repo_path, char *name_override, + char *project_name) { + char *json = yy_doc_to_str(doc); + yyjson_mut_doc_free(doc); + free(repo_path); + free(name_override); + free(project_name); + char *result = cbm_mcp_text_result(json, true); + free(json); + return result; +} + /* Build the response for a worker that crashed/hung/failed without producing a * result. The crash is already contained (this process survived); we report it * rather than dying. Precise skip-and-continue (quarantine the culprit, index the @@ -10449,12 +10465,7 @@ static char *build_worker_failure_response(const char *args, cbm_proc_outcome_t if (repo_path) { yyjson_mut_obj_add_strcpy(doc, root, "repo_path", repo_path); } - char *json = yy_doc_to_str(doc); - yyjson_mut_doc_free(doc); - free(repo_path); - char *result = cbm_mcp_text_result(json, true); - free(json); - return result; + return index_worker_json_result(doc, repo_path, NULL, NULL); } static char *build_worker_unsafe_terminal_response(const char *args, cbm_proc_outcome_t outcome, @@ -10474,12 +10485,50 @@ static char *build_worker_unsafe_terminal_response(const char *args, cbm_proc_ou if (repo_path) { yyjson_mut_obj_add_strcpy(doc, root, "repo_path", repo_path); } - char *json = yy_doc_to_str(doc); - yyjson_mut_doc_free(doc); - free(repo_path); - char *response = cbm_mcp_text_result(json, true); - free(json); - return response; + return index_worker_json_result(doc, repo_path, NULL, NULL); +} + +char *cbm_mcp_index_worker_resource_response(const char *args, + const cbm_index_worker_result_t *worker_result) { + char *repo_path = cbm_mcp_get_string_arg(args, "repo_path"); + char *name_override = cbm_mcp_get_string_arg(args, "name"); + char *project_name = + cbm_project_name_from_path(name_override && name_override[0] ? name_override : repo_path); + char db_path[CBM_SZ_1K] = {0}; + if (project_name) { + project_db_path(project_name, db_path, sizeof(db_path)); + } + + const cbm_index_resource_violation_t *violation = &worker_result->resource_violation; + const char *config_key = cbm_index_resource_config_key(violation->resource); + char message[CBM_SZ_256]; + (void)snprintf(message, sizeof(message), + worker_result->resource_probe_failed + ? "Worker resource measurement failed for %s" + : "Index worker exceeded %s", + config_key); + yyjson_mut_doc *doc = yyjson_mut_doc_new(NULL); + yyjson_mut_val *root = yyjson_mut_obj(doc); + yyjson_mut_doc_set_root(doc, root); + yyjson_mut_obj_add_str(doc, root, "status", "error"); + yyjson_mut_obj_add_str(doc, root, "code", + worker_result->resource_probe_failed ? "resource_probe_failed" + : "resource_limit_exceeded"); + yyjson_mut_obj_add_str(doc, root, "stage", "worker"); + yyjson_mut_obj_add_str(doc, root, "resource", cbm_index_resource_name(violation->resource)); + if (!worker_result->resource_probe_failed) { + yyjson_mut_obj_add_uint(doc, root, "observed", violation->observed); + yyjson_mut_obj_add_uint(doc, root, "limit", violation->limit); + yyjson_mut_obj_add_str(doc, root, "unit", cbm_index_resource_unit(violation->resource)); + } + yyjson_mut_obj_add_bool(doc, root, "retryable", true); + yyjson_mut_obj_add_bool(doc, root, "serving_index_preserved", + project_name && project_db_is_servable(project_name, db_path)); + yyjson_mut_obj_add_strcpy(doc, root, "message", message); + if (repo_path) { + yyjson_mut_obj_add_strcpy(doc, root, "repo_path", repo_path); + } + return index_worker_json_result(doc, repo_path, name_override, project_name); } /* Drop the cached store so the next query reopens whatever the worker wrote (each @@ -10625,6 +10674,9 @@ cbm_mcp_supervised_result_disposition_t cbm_mcp_supervised_result_disposition( worker_result->supervision_failed) { return CBM_MCP_SUPERVISED_RESULT_UNSAFE_TERMINAL; } + if (worker_result->resource_violation.resource != CBM_INDEX_RESOURCE_NONE) { + return CBM_MCP_SUPERVISED_RESULT_RESOURCE_FAILURE; + } if (worker_result->outcome == CBM_PROC_CLEAN) { return worker_result->response ? CBM_MCP_SUPERVISED_RESULT_SUCCESS : CBM_MCP_SUPERVISED_RESULT_FALLBACK; @@ -10632,8 +10684,8 @@ cbm_mcp_supervised_result_disposition_t cbm_mcp_supervised_result_disposition( return CBM_MCP_SUPERVISED_RESULT_CONTAINED_FAILURE; } -static bool index_policy_from_worker_args(const char *args, cbm_index_resource_policy_t *policy, - char *error, size_t error_size) { +bool cbm_mcp_index_policy_from_internal_args(const char *args, cbm_index_resource_policy_t *policy, + char *error, size_t error_size) { yyjson_doc *doc = args ? yyjson_read(args, strlen(args), 0) : NULL; yyjson_val *root = doc ? yyjson_doc_get_root(doc) : NULL; yyjson_val *encoded = @@ -10663,7 +10715,7 @@ static bool index_policy_from_worker_args(const char *args, cbm_index_resource_p static bool load_index_policy(cbm_mcp_server_t *srv, const char *args, cbm_index_resource_policy_t *policy, char *error, size_t error_size) { if (cbm_index_worker_active()) { - return index_policy_from_worker_args(args, policy, error, error_size); + return cbm_mcp_index_policy_from_internal_args(args, policy, error, error_size); } cbm_config_t *owned_config = NULL; cbm_config_t *config = srv ? srv->config : NULL; @@ -10728,13 +10780,14 @@ int cbm_index_restart_cap_for_testing(void) { return index_restart_cap(); } -static char *index_run_supervised(cbm_mcp_server_t *srv, const char *args) { +static char *index_run_supervised(cbm_mcp_server_t *srv, const char *args, + const cbm_index_resource_policy_t *resource_policy) { invalidate_cached_store(srv); /* First attempt: normal parallel run. */ cbm_index_worker_result_t wr; - int rc = cbm_index_spawn_worker_with_log_cancel( - args, false, NULL, NULL, srv ? srv->index_log_callback : NULL, + int rc = cbm_index_spawn_worker_with_policy_log_cancel( + args, resource_policy, false, NULL, NULL, srv ? srv->index_log_callback : NULL, srv ? srv->index_log_context : NULL, srv ? &srv->pipeline_cancel_requested : NULL, &wr); cbm_mcp_supervised_result_disposition_t disposition = cbm_mcp_supervised_result_disposition(rc, &wr); @@ -10752,6 +10805,12 @@ static char *index_run_supervised(cbm_mcp_server_t *srv, const char *args) { invalidate_cached_store(srv); return failure; } + if (disposition == CBM_MCP_SUPERVISED_RESULT_RESOURCE_FAILURE) { + char *failure = cbm_mcp_index_worker_resource_response(args, &wr); + cbm_index_worker_result_free(&wr); + invalidate_cached_store(srv); + return failure; + } if (disposition == CBM_MCP_SUPERVISED_RESULT_SUCCESS) { /* Clean exit → transfer the worker's response (the common path). */ char *resp = wr.response; /* transfer ownership to caller (may be NULL) */ @@ -10801,8 +10860,8 @@ static char *index_run_supervised(cbm_mcp_server_t *srv, const char *args) { bool terminal_cancelled = false; for (int i = 0; i < cap; i++) { cbm_index_worker_result_t wr2; - int rc2 = cbm_index_spawn_worker_with_log_cancel( - args, /*single_thread=*/false, marker_path, quarantine_path, + int rc2 = cbm_index_spawn_worker_with_policy_log_cancel( + args, resource_policy, /*single_thread=*/false, marker_path, quarantine_path, srv ? srv->index_log_callback : NULL, srv ? srv->index_log_context : NULL, srv ? &srv->pipeline_cancel_requested : NULL, &wr2); cbm_mcp_supervised_result_disposition_t recovery_disposition = @@ -10819,6 +10878,11 @@ static char *index_run_supervised(cbm_mcp_server_t *srv, const char *args) { cbm_index_worker_result_free(&wr2); break; } + if (recovery_disposition == CBM_MCP_SUPERVISED_RESULT_RESOURCE_FAILURE) { + resp = cbm_mcp_index_worker_resource_response(args, &wr2); + cbm_index_worker_result_free(&wr2); + break; + } if (recovery_disposition == CBM_MCP_SUPERVISED_RESULT_SUCCESS) { resp = wr2.response; /* transfer ownership to caller */ wr2.response = NULL; @@ -10913,8 +10977,8 @@ static char *index_run_supervised(cbm_mcp_server_t *srv, const char *args) { * so it cannot itself hang. Rare given monotonic progress. */ if (!resp && !unsafe_terminal && quarantined > 0) { cbm_index_worker_result_t wrp; - int rcp = cbm_index_spawn_worker_with_log_cancel( - args, /*single_thread=*/false, NULL, quarantine_path, + int rcp = cbm_index_spawn_worker_with_policy_log_cancel( + args, resource_policy, /*single_thread=*/false, NULL, quarantine_path, srv ? srv->index_log_callback : NULL, srv ? srv->index_log_context : NULL, srv ? &srv->pipeline_cancel_requested : NULL, &wrp); cbm_mcp_supervised_result_disposition_t partial_disposition = @@ -10930,6 +10994,8 @@ static char *index_run_supervised(cbm_mcp_server_t *srv, const char *args) { last_outcome = wrp.outcome; unsafe_terminal = true; terminal_cancelled = wrp.cancellation_requested; + } else if (partial_disposition == CBM_MCP_SUPERVISED_RESULT_RESOURCE_FAILURE) { + resp = cbm_mcp_index_worker_resource_response(args, &wrp); } cbm_index_worker_result_free(&wrp); } @@ -10973,7 +11039,7 @@ static char *index_run_supervised_path(cbm_mcp_server_t *srv, const char *root_p if (!args) { return NULL; } - char *resp = index_run_supervised(srv, args); + char *resp = index_run_supervised(srv, args, &policy); free(args); return resp; } @@ -11195,7 +11261,7 @@ static char *handle_index_repository(cbm_mcp_server_t *srv, const char *args) { index_args_free(repo_path, mode_str, name_override); return cbm_mcp_text_result("failed to prepare supervised index request", true); } - char *supervised = index_run_supervised(srv, worker_args); + char *supervised = index_run_supervised(srv, worker_args, &resource_policy); free(worker_args); if (supervised) { free(mutation_project); diff --git a/src/mcp/mcp.h b/src/mcp/mcp.h index 09befb780..af27e411e 100644 --- a/src/mcp/mcp.h +++ b/src/mcp/mcp.h @@ -226,6 +226,7 @@ char *cbm_mcp_handle_tool(cbm_mcp_server_t *srv, const char *tool_name, const ch typedef enum { CBM_MCP_SUPERVISED_RESULT_FALLBACK = 0, CBM_MCP_SUPERVISED_RESULT_SUCCESS, + CBM_MCP_SUPERVISED_RESULT_RESOURCE_FAILURE, CBM_MCP_SUPERVISED_RESULT_CONTAINED_FAILURE, CBM_MCP_SUPERVISED_RESULT_UNSAFE_TERMINAL, } cbm_mcp_supervised_result_disposition_t; diff --git a/src/mcp/mcp_internal.h b/src/mcp/mcp_internal.h index de170067a..80ab9a720 100644 --- a/src/mcp/mcp_internal.h +++ b/src/mcp/mcp_internal.h @@ -43,6 +43,10 @@ bool cbm_mcp_jsonrpc_response_prepend_notice(char **response_io, const char *not * must remove any untrusted field with the same name before invoking this. */ bool cbm_mcp_index_policy_add_to_args(yyjson_mut_doc *doc, yyjson_mut_val *root, const cbm_index_resource_policy_t *policy); +bool cbm_mcp_index_policy_from_internal_args(const char *args, cbm_index_resource_policy_t *policy, + char *error, size_t error_size); +char *cbm_mcp_index_worker_resource_response(const char *args, + const cbm_index_worker_result_t *worker_result); enum { CBM_MCP_DEFAULT_AUTO_INDEX_LIMIT = 50000 }; diff --git a/src/pipeline/pipeline.c b/src/pipeline/pipeline.c index 6bc0310dd..734a21bf6 100644 --- a/src/pipeline/pipeline.c +++ b/src/pipeline/pipeline.c @@ -499,7 +499,8 @@ const char *cbm_pipeline_repo_path(const cbm_pipeline_t *p) { } const cbm_index_resource_policy_t *cbm_pipeline_resource_policy(const cbm_pipeline_t *p) { - return p && cbm_index_policy_enabled(&p->resource_policy) ? &p->resource_policy : NULL; + return p && cbm_index_policy_discovery_enabled(&p->resource_policy) ? &p->resource_policy + : NULL; } cbm_index_resource_violation_t *cbm_pipeline_resource_violation(cbm_pipeline_t *p) { @@ -2734,7 +2735,7 @@ static int cbm_pipeline_run_staged(cbm_pipeline_t *p) { .ignore_file = NULL, .max_file_size = 0, .resource_policy = - cbm_index_policy_enabled(&p->resource_policy) ? &p->resource_policy : NULL, + cbm_index_policy_discovery_enabled(&p->resource_policy) ? &p->resource_policy : NULL, .resource_violation = &p->resource_violation, }; cbm_file_info_t *files = NULL; diff --git a/tests/test_daemon_application.c b/tests/test_daemon_application.c index 0108e8ec5..7c8b841e0 100644 --- a/tests/test_daemon_application.c +++ b/tests/test_daemon_application.c @@ -1364,6 +1364,8 @@ typedef struct { cbm_project_lock_manager_t *project_locks; bool project_lock_busy_is_error; cbm_proc_outcome_t outcomes[APP_FAKE_MAX_ATTEMPTS]; + cbm_index_resource_violation_t resource_violations[APP_FAKE_MAX_ATTEMPTS]; + bool resource_probe_failures[APP_FAKE_MAX_ATTEMPTS]; const char *responses[APP_FAKE_MAX_ATTEMPTS]; const char *marker_payloads[APP_FAKE_MAX_ATTEMPTS]; char marker_paths[APP_FAKE_MAX_ATTEMPTS][APP_TEST_PATH_CAP]; @@ -1468,7 +1470,9 @@ static int app_fake_worker_start(void *opaque, const char *args_json, size_t mem static bool app_wait_for_atomic_bool(atomic_bool *value, bool expected); static bool app_fake_worker_policy_equals(app_fake_worker_context_t *context, int attempt, - const char *max_files, const char *max_source_mb) { + const char *max_files, const char *max_source_mb, + const char *max_rss_mb, + const char *max_duration_seconds) { if (!context || attempt < 0 || attempt >= APP_FAKE_MAX_ATTEMPTS) { return false; } @@ -1485,9 +1489,18 @@ static bool app_fake_worker_policy_equals(app_fake_worker_context_t *context, in yyjson_val *bytes = policy && yyjson_is_obj(policy) ? yyjson_obj_get(policy, CBM_INDEX_CONFIG_MAX_SOURCE_MB) : NULL; - bool equal = files && bytes && yyjson_is_str(files) && yyjson_is_str(bytes) && + yyjson_val *rss = policy && yyjson_is_obj(policy) + ? yyjson_obj_get(policy, CBM_INDEX_CONFIG_MAX_RSS_MB) + : NULL; + yyjson_val *duration = policy && yyjson_is_obj(policy) + ? yyjson_obj_get(policy, CBM_INDEX_CONFIG_MAX_DURATION_SECONDS) + : NULL; + bool equal = files && bytes && rss && duration && yyjson_is_str(files) && + yyjson_is_str(bytes) && yyjson_is_str(rss) && yyjson_is_str(duration) && strcmp(yyjson_get_str(files), max_files) == 0 && - strcmp(yyjson_get_str(bytes), max_source_mb) == 0; + strcmp(yyjson_get_str(bytes), max_source_mb) == 0 && + strcmp(yyjson_get_str(rss), max_rss_mb) == 0 && + strcmp(yyjson_get_str(duration), max_duration_seconds) == 0; yyjson_doc_free(document); return equal; } @@ -1529,6 +1542,12 @@ static cbm_index_worker_poll_t app_fake_worker_poll(void *opaque, worker->result.outcome = outcome; worker->result.exit_code = outcome == CBM_PROC_CLEAN ? 0 : -1; worker->result.tree_quiesced = true; + if (worker->attempt < APP_FAKE_MAX_ATTEMPTS) { + worker->result.resource_violation = + worker->context->resource_violations[worker->attempt]; + worker->result.resource_probe_failed = + worker->context->resource_probe_failures[worker->attempt]; + } const char *response = worker->attempt < APP_FAKE_MAX_ATTEMPTS && worker->context->responses[worker->attempt] ? worker->context->responses[worker->attempt] @@ -2043,11 +2062,13 @@ TEST(daemon_application_initialize_coalesces_auto_index_for_full_sessions) { bool dirs_ok = cbm_mkdtemp(root) != NULL && cbm_mkdtemp(cache) != NULL; bool cache_set = dirs_ok && cache_saved && cbm_setenv("CBM_CACHE_DIR", cache, 1) == 0; cbm_config_t *stored_config = cache_set ? cbm_config_open(cache) : NULL; - bool config_ready = stored_config && - cbm_config_set(stored_config, CBM_CONFIG_AUTO_INDEX, "true") == 0 && - cbm_config_set(stored_config, CBM_CONFIG_AUTO_WATCH, "false") == 0 && - cbm_config_set(stored_config, CBM_INDEX_CONFIG_MAX_FILES, "3") == 0 && - cbm_config_set(stored_config, CBM_INDEX_CONFIG_MAX_SOURCE_MB, "4") == 0; + bool config_ready = + stored_config && cbm_config_set(stored_config, CBM_CONFIG_AUTO_INDEX, "true") == 0 && + cbm_config_set(stored_config, CBM_CONFIG_AUTO_WATCH, "false") == 0 && + cbm_config_set(stored_config, CBM_INDEX_CONFIG_MAX_FILES, "3") == 0 && + cbm_config_set(stored_config, CBM_INDEX_CONFIG_MAX_SOURCE_MB, "4") == 0 && + cbm_config_set(stored_config, CBM_INDEX_CONFIG_MAX_RSS_MB, "64") == 0 && + cbm_config_set(stored_config, CBM_INDEX_CONFIG_MAX_DURATION_SECONDS, "7") == 0; char canonical_root[APP_TEST_PATH_CAP] = {0}; bool canonical = dirs_ok && cbm_canonical_path(root, canonical_root, sizeof(canonical_root)); char *project = canonical ? cbm_project_name_from_path(canonical_root) : NULL; @@ -2092,7 +2113,8 @@ TEST(daemon_application_initialize_coalesces_auto_index_for_full_sessions) { bool first_owned = first_initialized && project && app_wait_for_subscribers(application, project, 1) && app_wait_for_atomic_int(&fake.starts, 1); - bool auto_policy_propagated = first_owned && app_fake_worker_policy_equals(&fake, 0, "3", "4"); + bool auto_policy_propagated = + first_owned && app_fake_worker_policy_equals(&fake, 0, "3", "4", "64", "7"); bool second_initialized = app_test_initialize_profile(&callbacks, sessions[1], root, CBM_MCP_TOOL_PROFILE_ALL, NULL, NULL); bool coalesced = first_owned && second_initialized && project && @@ -2435,9 +2457,11 @@ TEST(daemon_application_programmatic_index_injects_resource_policy) { ASSERT_NOT_NULL(root); ASSERT_NOT_NULL(cache); cbm_config_t *stored_config = cbm_config_open(cache); - bool configured = stored_config && - cbm_config_set(stored_config, CBM_INDEX_CONFIG_MAX_FILES, "5") == 0 && - cbm_config_set(stored_config, CBM_INDEX_CONFIG_MAX_SOURCE_MB, "6") == 0; + bool configured = + stored_config && cbm_config_set(stored_config, CBM_INDEX_CONFIG_MAX_FILES, "5") == 0 && + cbm_config_set(stored_config, CBM_INDEX_CONFIG_MAX_SOURCE_MB, "6") == 0 && + cbm_config_set(stored_config, CBM_INDEX_CONFIG_MAX_RSS_MB, "64") == 0 && + cbm_config_set(stored_config, CBM_INDEX_CONFIG_MAX_DURATION_SECONDS, "7") == 0; app_fake_worker_context_t fake; app_fake_worker_context_init(&fake); @@ -2457,7 +2481,8 @@ TEST(daemon_application_programmatic_index_injects_resource_policy) { cbm_daemon_application_t *application = configured ? cbm_daemon_application_new(&config) : NULL; int index_rc = application ? cbm_daemon_application_index(application, "policy-programmatic", root) : -1; - bool policy_propagated = index_rc == 0 && app_fake_worker_policy_equals(&fake, 0, "5", "6"); + bool policy_propagated = + index_rc == 0 && app_fake_worker_policy_equals(&fake, 0, "5", "6", "64", "7"); bool stopped = application && cbm_daemon_application_shutdown(application, APP_TEST_TIMEOUT_MS); cbm_daemon_application_free(application); @@ -2469,6 +2494,87 @@ TEST(daemon_application_programmatic_index_injects_resource_policy) { ASSERT_TRUE(stopped); PASS(); } + +TEST(daemon_application_worker_resource_failure_is_structured_and_not_retried) { + char *cache = th_mktempdir("cbm_app_worker_resource_cache"); + cbm_config_t *stored_config = cache ? cbm_config_open(cache) : NULL; + app_fake_worker_context_t fake; + app_fake_worker_context_init(&fake); + atomic_store(&fake.scripted, true); + fake.outcomes[0] = CBM_PROC_KILLED; + fake.resource_violations[0] = (cbm_index_resource_violation_t){ + .resource = CBM_INDEX_RESOURCE_DURATION_MS, .observed = 7001, .limit = 7000}; + cbm_daemon_application_worker_ops_t worker_ops = { + .context = &fake, + .start = app_fake_worker_start, + .poll = app_fake_worker_poll, + .cancel = app_fake_worker_cancel, + .log_path = app_fake_worker_log_path, + .destroy = app_fake_worker_destroy, + }; + cbm_daemon_application_config_t config = { + .config = stored_config, + .worker_ops = &worker_ops, + }; + cbm_daemon_application_t *application = + stored_config ? cbm_daemon_application_new(&config) : NULL; + cbm_daemon_runtime_application_callbacks_t callbacks = + cbm_daemon_application_runtime_callbacks(application); + cbm_daemon_runtime_application_session_t *session = app_test_open(&callbacks, 4188); + char *root = th_mktempdir("cbm_app_worker_resource"); + uint8_t *context = NULL; + uint32_t context_length = 0; + uint8_t *tool = NULL; + uint32_t tool_length = 0; + char args[APP_TEST_PATH_CAP + 96]; + (void)snprintf(args, sizeof(args), + "{\"repo_path\":\"%s\",\"name\":\"DaemonWorkerResourceFixture\"}", + root ? root : ""); + bool setup = cache && stored_config && application && session && root && + app_test_context_request(root, root, &context, &context_length) && + app_test_tool_request("index_repository", args, &tool, &tool_length); + uint8_t *response = NULL; + uint32_t response_length = 0; + if (setup) { + uint8_t *context_response = NULL; + uint32_t context_response_length = 0; + setup = app_test_request(&callbacks, session, context, context_length, &context_response, + &context_response_length) == CBM_DAEMON_RUNTIME_APPLICATION_OK; + free(context_response); + } + cbm_daemon_runtime_application_status_t status = + setup + ? app_test_request(&callbacks, session, tool, tool_length, &response, &response_length) + : CBM_DAEMON_RUNTIME_APPLICATION_TRANSPORT_ERROR; + bool structured = response && strstr((char *)response, "resource_limit_exceeded") && + strstr((char *)response, "\\\"stage\\\":\\\"worker\\\"") && + strstr((char *)response, "\\\"resource\\\":\\\"duration_ms\\\"") && + strstr((char *)response, "\\\"observed\\\":7001") && + strstr((char *)response, "\\\"limit\\\":7000") && + strstr((char *)response, "\\\"unit\\\":\\\"milliseconds\\\""); + if (session) { + callbacks.session_close(callbacks.context, session); + } + bool stopped = application && cbm_daemon_application_shutdown(application, APP_TEST_TIMEOUT_MS); + int starts = atomic_load(&fake.starts); + int destroys = atomic_load(&fake.destroys); + cbm_daemon_application_free(application); + cbm_config_close(stored_config); + free(response); + free(context); + free(tool); + th_cleanup(root); + th_cleanup(cache); + + ASSERT_TRUE(setup); + ASSERT_EQ(status, CBM_DAEMON_RUNTIME_APPLICATION_OK); + ASSERT_TRUE(structured); + ASSERT_EQ(starts, 1); + ASSERT_EQ(destroys, 1); + ASSERT_TRUE(stopped); + PASS(); +} + TEST(daemon_application_auto_index_honors_tracked_file_limit) { app_env_backup_t cache_environment; bool cache_saved = app_env_backup_capture(&cache_environment, "CBM_CACHE_DIR"); @@ -5654,6 +5760,7 @@ SUITE(daemon_application) { RUN_TEST(daemon_application_sensitive_root_blocks_watch_but_preserves_controls); RUN_TEST(daemon_application_index_args_compare_repo_path_separator_equivalently); RUN_TEST(daemon_application_programmatic_index_injects_resource_policy); + RUN_TEST(daemon_application_worker_resource_failure_is_structured_and_not_retried); RUN_TEST(daemon_application_auto_index_honors_tracked_file_limit); RUN_TEST(daemon_application_auto_index_file_count_handles_literal_metacharacter_path); RUN_TEST(daemon_application_auto_index_file_count_supports_non_git_roots); diff --git a/tests/test_index_policy.c b/tests/test_index_policy.c index d866a0631..b33d64cdb 100644 --- a/tests/test_index_policy.c +++ b/tests/test_index_policy.c @@ -21,9 +21,15 @@ TEST(index_policy_defaults_are_disabled) { ASSERT_FALSE(policy.max_files.enabled); ASSERT_FALSE(policy.max_source_bytes.enabled); + ASSERT_FALSE(policy.max_rss_bytes.enabled); + ASSERT_FALSE(policy.max_duration_ms.enabled); ASSERT_FALSE(cbm_index_policy_enabled(&policy)); + ASSERT_FALSE(cbm_index_policy_discovery_enabled(&policy)); + ASSERT_FALSE(cbm_index_policy_worker_enabled(&policy)); ASSERT_STR_EQ(cbm_index_policy_default_value(CBM_INDEX_CONFIG_MAX_FILES), "off"); ASSERT_STR_EQ(cbm_index_policy_default_value(CBM_INDEX_CONFIG_MAX_SOURCE_MB), "off"); + ASSERT_STR_EQ(cbm_index_policy_default_value(CBM_INDEX_CONFIG_MAX_RSS_MB), "off"); + ASSERT_STR_EQ(cbm_index_policy_default_value(CBM_INDEX_CONFIG_MAX_DURATION_SECONDS), "off"); PASS(); } @@ -60,6 +66,38 @@ TEST(index_policy_source_limit_converts_mib_without_overflow) { PASS(); } +TEST(index_policy_worker_limits_validate_and_convert_units) { + cbm_index_resource_policy_t policy; + cbm_index_policy_init(&policy); + char error[256]; + + ASSERT_TRUE( + cbm_index_policy_set(&policy, CBM_INDEX_CONFIG_MAX_RSS_MB, "64", error, sizeof(error))); + ASSERT_TRUE(cbm_index_policy_enabled(&policy)); + ASSERT_TRUE(cbm_index_policy_worker_enabled(&policy)); + ASSERT_FALSE(cbm_index_policy_discovery_enabled(&policy)); + ASSERT_EQ(policy.max_rss_bytes.value, UINT64_C(64) * CBM_INDEX_MIB_BYTES); + ASSERT_TRUE(cbm_index_policy_set(&policy, CBM_INDEX_CONFIG_MAX_RSS_MB, "1048576", error, + sizeof(error))); + ASSERT_EQ(policy.max_rss_bytes.value, UINT64_C(1048576) * CBM_INDEX_MIB_BYTES); + cbm_index_resource_policy_t before = policy; + ASSERT_FALSE( + cbm_index_policy_set(&policy, CBM_INDEX_CONFIG_MAX_RSS_MB, "63", error, sizeof(error))); + ASSERT_EQ(memcmp(&policy, &before, sizeof(policy)), 0); + + ASSERT_TRUE(cbm_index_policy_set(&policy, CBM_INDEX_CONFIG_MAX_DURATION_SECONDS, "1", error, + sizeof(error))); + ASSERT_EQ(policy.max_duration_ms.value, 1000); + ASSERT_TRUE(cbm_index_policy_set(&policy, CBM_INDEX_CONFIG_MAX_DURATION_SECONDS, "86400", error, + sizeof(error))); + ASSERT_EQ(policy.max_duration_ms.value, UINT64_C(86400000)); + before = policy; + ASSERT_FALSE(cbm_index_policy_set(&policy, CBM_INDEX_CONFIG_MAX_DURATION_SECONDS, "86401", + error, sizeof(error))); + ASSERT_EQ(memcmp(&policy, &before, sizeof(policy)), 0); + PASS(); +} + TEST(index_policy_invalid_value_is_rejected_atomically) { static const char *const invalid[] = {"", "0", "-1", "1MB", "1 ", "+1", "10000001", "18446744073709551616"}; @@ -95,6 +133,16 @@ TEST(index_policy_format_round_trips_public_values) { ASSERT_TRUE( cbm_index_policy_format(&policy, CBM_INDEX_CONFIG_MAX_SOURCE_MB, value, sizeof(value))); ASSERT_STR_EQ(value, "7"); + ASSERT_TRUE( + cbm_index_policy_set(&policy, CBM_INDEX_CONFIG_MAX_RSS_MB, "64", error, sizeof(error))); + ASSERT_TRUE( + cbm_index_policy_format(&policy, CBM_INDEX_CONFIG_MAX_RSS_MB, value, sizeof(value))); + ASSERT_STR_EQ(value, "64"); + ASSERT_TRUE(cbm_index_policy_set(&policy, CBM_INDEX_CONFIG_MAX_DURATION_SECONDS, "9", error, + sizeof(error))); + ASSERT_TRUE(cbm_index_policy_format(&policy, CBM_INDEX_CONFIG_MAX_DURATION_SECONDS, value, + sizeof(value))); + ASSERT_STR_EQ(value, "9"); PASS(); } @@ -103,10 +151,18 @@ TEST(index_policy_violation_metadata_is_stable) { ASSERT_STR_EQ(cbm_index_resource_name(CBM_INDEX_RESOURCE_SOURCE_BYTES), "source_bytes"); ASSERT_STR_EQ(cbm_index_resource_unit(CBM_INDEX_RESOURCE_FILES), "files"); ASSERT_STR_EQ(cbm_index_resource_unit(CBM_INDEX_RESOURCE_SOURCE_BYTES), "bytes"); + ASSERT_STR_EQ(cbm_index_resource_name(CBM_INDEX_RESOURCE_RSS_BYTES), "rss_bytes"); + ASSERT_STR_EQ(cbm_index_resource_name(CBM_INDEX_RESOURCE_DURATION_MS), "duration_ms"); + ASSERT_STR_EQ(cbm_index_resource_unit(CBM_INDEX_RESOURCE_RSS_BYTES), "bytes"); + ASSERT_STR_EQ(cbm_index_resource_unit(CBM_INDEX_RESOURCE_DURATION_MS), "milliseconds"); ASSERT_STR_EQ(cbm_index_resource_config_key(CBM_INDEX_RESOURCE_FILES), CBM_INDEX_CONFIG_MAX_FILES); ASSERT_STR_EQ(cbm_index_resource_config_key(CBM_INDEX_RESOURCE_SOURCE_BYTES), CBM_INDEX_CONFIG_MAX_SOURCE_MB); + ASSERT_STR_EQ(cbm_index_resource_config_key(CBM_INDEX_RESOURCE_RSS_BYTES), + CBM_INDEX_CONFIG_MAX_RSS_MB); + ASSERT_STR_EQ(cbm_index_resource_config_key(CBM_INDEX_RESOURCE_DURATION_MS), + CBM_INDEX_CONFIG_MAX_DURATION_SECONDS); PASS(); } @@ -122,9 +178,13 @@ TEST(index_policy_config_loads_defaults_values_and_rejects_corruption) { ASSERT_FALSE(cbm_index_policy_enabled(&policy)); ASSERT_EQ(cbm_config_set(config, CBM_INDEX_CONFIG_MAX_FILES, "9"), 0); ASSERT_EQ(cbm_config_set(config, CBM_INDEX_CONFIG_MAX_SOURCE_MB, "3"), 0); + ASSERT_EQ(cbm_config_set(config, CBM_INDEX_CONFIG_MAX_RSS_MB, "64"), 0); + ASSERT_EQ(cbm_config_set(config, CBM_INDEX_CONFIG_MAX_DURATION_SECONDS, "7"), 0); ASSERT_TRUE(cbm_config_load_index_policy(config, &policy, error, sizeof(error))); ASSERT_EQ(policy.max_files.value, 9); ASSERT_EQ(policy.max_source_bytes.value, UINT64_C(3) * CBM_INDEX_MIB_BYTES); + ASSERT_EQ(policy.max_rss_bytes.value, UINT64_C(64) * CBM_INDEX_MIB_BYTES); + ASSERT_EQ(policy.max_duration_ms.value, UINT64_C(7000)); ASSERT_EQ(cbm_config_set(config, CBM_INDEX_CONFIG_MAX_FILES, "corrupt"), 0); ASSERT_FALSE(cbm_config_load_index_policy(config, &policy, error, sizeof(error))); @@ -135,16 +195,23 @@ TEST(index_policy_config_loads_defaults_values_and_rejects_corruption) { PASS(); } -TEST(index_policy_cli_lists_both_operator_keys) { +TEST(index_policy_cli_lists_all_operator_keys) { bool files_found = false; bool bytes_found = false; + bool rss_found = false; + bool duration_found = false; for (size_t index = 0; index < cbm_cli_config_key_count_for_testing(); index++) { const char *key = cbm_cli_config_key_at_for_testing(index); files_found = files_found || (key && strcmp(key, CBM_INDEX_CONFIG_MAX_FILES) == 0); bytes_found = bytes_found || (key && strcmp(key, CBM_INDEX_CONFIG_MAX_SOURCE_MB) == 0); + rss_found = rss_found || (key && strcmp(key, CBM_INDEX_CONFIG_MAX_RSS_MB) == 0); + duration_found = + duration_found || (key && strcmp(key, CBM_INDEX_CONFIG_MAX_DURATION_SECONDS) == 0); } ASSERT_TRUE(files_found); ASSERT_TRUE(bytes_found); + ASSERT_TRUE(rss_found); + ASSERT_TRUE(duration_found); PASS(); } @@ -391,11 +458,12 @@ SUITE(index_policy) { RUN_TEST(index_policy_defaults_are_disabled); RUN_TEST(index_policy_file_limit_accepts_off_and_exact_range); RUN_TEST(index_policy_source_limit_converts_mib_without_overflow); + RUN_TEST(index_policy_worker_limits_validate_and_convert_units); RUN_TEST(index_policy_invalid_value_is_rejected_atomically); RUN_TEST(index_policy_format_round_trips_public_values); RUN_TEST(index_policy_violation_metadata_is_stable); RUN_TEST(index_policy_config_loads_defaults_values_and_rejects_corruption); - RUN_TEST(index_policy_cli_lists_both_operator_keys); + RUN_TEST(index_policy_cli_lists_all_operator_keys); RUN_TEST(index_policy_cli_set_rejects_invalid_value_without_overwrite); RUN_TEST(index_policy_cli_set_reports_a_failed_write); RUN_TEST(index_policy_worker_rejects_missing_parent_policy); diff --git a/tests/test_index_supervisor.c b/tests/test_index_supervisor.c index d69354d76..d1227dd06 100644 --- a/tests/test_index_supervisor.c +++ b/tests/test_index_supervisor.c @@ -941,6 +941,269 @@ TEST(index_supervisor_job_memory_limit_has_floor_headroom_and_no_overflow) { PASS(); } +typedef struct { + uint64_t now_ms; + uint64_t rss_values[4]; + int rss_value_count; + int rss_calls; + int clock_calls; + cbm_proc_tree_rss_status_t rss_status; +} index_supervisor_resource_fake_t; + +static uint64_t index_supervisor_fake_clock(void *context) { + index_supervisor_resource_fake_t *fake = context; + fake->clock_calls++; + return fake->now_ms; +} + +static cbm_proc_tree_rss_status_t index_supervisor_fake_rss(cbm_subprocess_t *process, + uint64_t *rss_bytes, void *context) { + (void)process; + index_supervisor_resource_fake_t *fake = context; + int value_index = + fake->rss_calls < fake->rss_value_count ? fake->rss_calls : fake->rss_value_count - 1; + fake->rss_calls++; + if (fake->rss_status == CBM_PROC_TREE_RSS_OK && rss_bytes && value_index >= 0) { + *rss_bytes = fake->rss_values[value_index]; + } + return fake->rss_status; +} + +static cbm_index_resource_policy_t index_supervisor_test_worker_policy(uint64_t rss_bytes, + uint64_t duration_ms) { + cbm_index_resource_policy_t policy; + cbm_index_policy_init(&policy); + policy.max_rss_bytes = (cbm_index_limit_u64_t){.enabled = rss_bytes > 0, .value = rss_bytes}; + policy.max_duration_ms = + (cbm_index_limit_u64_t){.enabled = duration_ms > 0, .value = duration_ms}; + return policy; +} + +TEST(index_supervisor_disabled_limits_do_not_probe) { + index_supervisor_resource_fake_t fake = { + .now_ms = 100, + .rss_values = {UINT64_MAX}, + .rss_value_count = 1, + .rss_status = CBM_PROC_TREE_RSS_OK, + }; + cbm_index_supervisor_set_resource_hooks_for_testing(index_supervisor_fake_clock, + index_supervisor_fake_rss, &fake); + cbm_index_resource_policy_t policy = index_supervisor_test_worker_policy(0, 0); + cbm_index_worker_handle_t *handle = NULL; + int start_rc = cbm_index_worker_start_with_policy("{\"__cbm_test_worker\":\"clean\"}", 0, + &policy, false, NULL, NULL, &handle); + const cbm_index_worker_result_t *result = NULL; + bool terminal = handle && index_supervisor_test_poll_terminal( + handle, INDEX_SUPERVISOR_TEST_TERMINAL_MS, &result); + cbm_index_resource_t resource = + result ? result->resource_violation.resource : CBM_INDEX_RESOURCE_RSS_BYTES; + bool probe_failed = result && result->resource_probe_failed; + if (terminal) { + cbm_index_worker_destroy(handle); + } else { + index_supervisor_test_cleanup_handle(handle); + } + cbm_index_supervisor_reset_resource_hooks_for_testing(); + + ASSERT_EQ(start_rc, 0); + ASSERT_TRUE(terminal); + ASSERT_EQ(fake.rss_calls, 0); + ASSERT_EQ(fake.clock_calls, 0); + ASSERT_EQ(resource, CBM_INDEX_RESOURCE_NONE); + ASSERT_FALSE(probe_failed); + PASS(); +} + +TEST(index_supervisor_rss_equality_runs_then_excess_terminates_tree) { + const uint64_t limit = UINT64_C(64) * CBM_INDEX_MIB_BYTES; + index_supervisor_resource_fake_t fake = { + .now_ms = 100, + .rss_values = {limit, limit + 1}, + .rss_value_count = 2, + .rss_status = CBM_PROC_TREE_RSS_OK, + }; + cbm_index_supervisor_set_resource_hooks_for_testing(index_supervisor_fake_clock, + index_supervisor_fake_rss, &fake); + cbm_index_resource_policy_t policy = index_supervisor_test_worker_policy(limit, 0); + cbm_index_worker_handle_t *handle = NULL; + int start_rc = cbm_index_worker_start_with_policy("{\"__cbm_test_worker\":\"hang-tree\"}", 0, + &policy, false, NULL, NULL, &handle); + const cbm_index_worker_result_t *result = NULL; + cbm_index_worker_poll_t equal_state = + handle ? cbm_index_worker_poll(handle, &result) : CBM_INDEX_WORKER_POLL_ERROR; + fake.now_ms += 250; + bool terminal = handle && index_supervisor_test_poll_terminal( + handle, INDEX_SUPERVISOR_TEST_TERMINAL_MS, &result); + const cbm_index_worker_result_t *cached = NULL; + bool cached_terminal = + terminal && cbm_index_worker_poll(handle, &cached) == CBM_INDEX_WORKER_POLL_TERMINAL && + cached == result; + bool limited = terminal && result && + result->resource_violation.resource == CBM_INDEX_RESOURCE_RSS_BYTES && + result->resource_violation.observed == limit + 1 && + result->resource_violation.limit == limit && !result->cancellation_requested && + result->tree_quiesced && !result->supervision_failed; + if (terminal) { + cbm_index_worker_destroy(handle); + } else { + index_supervisor_test_cleanup_handle(handle); + } + cbm_index_supervisor_reset_resource_hooks_for_testing(); + + ASSERT_EQ(start_rc, 0); + ASSERT_EQ(equal_state, CBM_INDEX_WORKER_POLL_RUNNING); + ASSERT_TRUE(limited); + ASSERT_TRUE(cached_terminal); + PASS(); +} + +TEST(index_supervisor_duration_is_total_time_not_quiet_timeout) { + index_supervisor_resource_fake_t fake = { + .now_ms = 100, + .rss_status = CBM_PROC_TREE_RSS_EMPTY, + }; + cbm_index_supervisor_set_resource_hooks_for_testing(index_supervisor_fake_clock, + index_supervisor_fake_rss, &fake); + cbm_index_resource_policy_t policy = index_supervisor_test_worker_policy(0, 1000); + cbm_index_worker_handle_t *handle = NULL; + int start_rc = cbm_index_worker_start_with_policy("{\"__cbm_test_worker\":\"hang-tree\"}", 0, + &policy, false, NULL, NULL, &handle); + fake.now_ms = 1100; + const cbm_index_worker_result_t *result = NULL; + cbm_index_worker_poll_t equal_state = + handle ? cbm_index_worker_poll(handle, &result) : CBM_INDEX_WORKER_POLL_ERROR; + fake.now_ms = 1101; + bool terminal = handle && index_supervisor_test_poll_terminal( + handle, INDEX_SUPERVISOR_TEST_TERMINAL_MS, &result); + bool limited = terminal && result && + result->resource_violation.resource == CBM_INDEX_RESOURCE_DURATION_MS && + result->resource_violation.observed == 1001 && + result->resource_violation.limit == 1000 && result->outcome != CBM_PROC_HANG && + result->tree_quiesced; + if (terminal) { + cbm_index_worker_destroy(handle); + } else { + index_supervisor_test_cleanup_handle(handle); + } + cbm_index_supervisor_reset_resource_hooks_for_testing(); + + ASSERT_EQ(start_rc, 0); + ASSERT_EQ(equal_state, CBM_INDEX_WORKER_POLL_RUNNING); + ASSERT_TRUE(limited); + PASS(); +} + +TEST(index_supervisor_quiet_timeout_remains_hang_with_duration_enabled) { + const char *saved_timeout = getenv("CBM_INDEX_WORKER_TIMEOUT_S"); + char *saved_timeout_copy = saved_timeout ? cbm_strdup(saved_timeout) : NULL; + (void)cbm_setenv("CBM_INDEX_WORKER_TIMEOUT_S", "1", 1); + cbm_index_supervisor_reset_resource_hooks_for_testing(); + cbm_index_resource_policy_t policy = index_supervisor_test_worker_policy(0, 60000); + cbm_index_worker_handle_t *handle = NULL; + int start_rc = cbm_index_worker_start_with_policy("{\"__cbm_test_worker\":\"hang-tree\"}", 0, + &policy, false, NULL, NULL, &handle); + const cbm_index_worker_result_t *result = NULL; + bool terminal = handle && index_supervisor_test_poll_terminal( + handle, INDEX_SUPERVISOR_TEST_TERMINAL_MS, &result); + bool remained_hang = terminal && result && result->outcome == CBM_PROC_HANG && + result->resource_violation.resource == CBM_INDEX_RESOURCE_NONE && + !result->resource_probe_failed; + if (terminal) { + cbm_index_worker_destroy(handle); + } else { + index_supervisor_test_cleanup_handle(handle); + } + index_supervisor_test_restore_env("CBM_INDEX_WORKER_TIMEOUT_S", saved_timeout_copy); + + ASSERT_EQ(start_rc, 0); + ASSERT_TRUE(terminal); + ASSERT_TRUE(remained_hang); + PASS(); +} + +TEST(index_supervisor_cancel_precedes_resource_probe) { + const uint64_t limit = UINT64_C(64) * CBM_INDEX_MIB_BYTES; + index_supervisor_resource_fake_t fake = { + .now_ms = 100, + .rss_values = {limit + 1}, + .rss_value_count = 1, + .rss_status = CBM_PROC_TREE_RSS_OK, + }; + cbm_index_supervisor_set_resource_hooks_for_testing(index_supervisor_fake_clock, + index_supervisor_fake_rss, &fake); + cbm_index_resource_policy_t policy = index_supervisor_test_worker_policy(limit, 1000); + cbm_index_worker_handle_t *handle = NULL; + int start_rc = cbm_index_worker_start_with_policy("{\"__cbm_test_worker\":\"hang-tree\"}", 0, + &policy, false, NULL, NULL, &handle); + fake.now_ms = 2000; + bool cancel_accepted = handle && cbm_index_worker_request_cancel(handle); + const cbm_index_worker_result_t *result = NULL; + bool terminal = handle && index_supervisor_test_poll_terminal( + handle, INDEX_SUPERVISOR_TEST_TERMINAL_MS, &result); + bool cancelled = terminal && result && result->cancellation_requested && + result->resource_violation.resource == CBM_INDEX_RESOURCE_NONE && + !result->resource_probe_failed && result->tree_quiesced; + if (terminal) { + cbm_index_worker_destroy(handle); + } else { + index_supervisor_test_cleanup_handle(handle); + } + cbm_index_supervisor_reset_resource_hooks_for_testing(); + + ASSERT_EQ(start_rc, 0); + ASSERT_TRUE(cancel_accepted); + ASSERT_TRUE(cancelled); + ASSERT_EQ(fake.rss_calls, 0); + PASS(); +} + +TEST(index_supervisor_three_failed_rss_probes_fail_closed) { + const uint64_t limit = UINT64_C(64) * CBM_INDEX_MIB_BYTES; + index_supervisor_resource_fake_t fake = { + .now_ms = 100, + .rss_status = CBM_PROC_TREE_RSS_ERROR, + }; + cbm_index_supervisor_set_resource_hooks_for_testing(index_supervisor_fake_clock, + index_supervisor_fake_rss, &fake); + cbm_index_resource_policy_t policy = index_supervisor_test_worker_policy(limit, 0); + cbm_index_worker_handle_t *handle = NULL; + int start_rc = cbm_index_worker_start_with_policy("{\"__cbm_test_worker\":\"hang-tree\"}", 0, + &policy, false, NULL, NULL, &handle); + const cbm_index_worker_result_t *result = NULL; + cbm_index_worker_poll_t first = + handle ? cbm_index_worker_poll(handle, &result) : CBM_INDEX_WORKER_POLL_ERROR; + cbm_index_worker_poll_t throttled = + handle ? cbm_index_worker_poll(handle, &result) : CBM_INDEX_WORKER_POLL_ERROR; + int calls_before_interval = fake.rss_calls; + fake.now_ms += 250; + cbm_index_worker_poll_t second = + handle ? cbm_index_worker_poll(handle, &result) : CBM_INDEX_WORKER_POLL_ERROR; + fake.now_ms += 250; + cbm_index_worker_poll_t third = + handle ? cbm_index_worker_poll(handle, &result) : CBM_INDEX_WORKER_POLL_ERROR; + bool terminal = handle && index_supervisor_test_poll_terminal( + handle, INDEX_SUPERVISOR_TEST_TERMINAL_MS, &result); + bool failed_closed = terminal && result && result->resource_probe_failed && + result->resource_violation.resource == CBM_INDEX_RESOURCE_RSS_BYTES && + !result->cancellation_requested && result->tree_quiesced; + if (terminal) { + cbm_index_worker_destroy(handle); + } else { + index_supervisor_test_cleanup_handle(handle); + } + cbm_index_supervisor_reset_resource_hooks_for_testing(); + + ASSERT_EQ(start_rc, 0); + ASSERT_EQ(first, CBM_INDEX_WORKER_POLL_RUNNING); + ASSERT_EQ(throttled, CBM_INDEX_WORKER_POLL_RUNNING); + ASSERT_EQ(calls_before_interval, 1); + ASSERT_EQ(second, CBM_INDEX_WORKER_POLL_RUNNING); + ASSERT_EQ(third, CBM_INDEX_WORKER_POLL_RUNNING); + ASSERT_TRUE(failed_closed); + ASSERT_EQ(fake.rss_calls, 3); + PASS(); +} + SUITE(index_supervisor) { RUN_TEST(index_supervisor_job_memory_limit_has_floor_headroom_and_no_overflow); RUN_TEST(index_supervisor_worker_argv_requires_exact_build_bound_grammar); @@ -951,4 +1214,10 @@ SUITE(index_supervisor) { RUN_TEST(index_supervisor_drains_terminal_backlog_into_request_progress_callback); RUN_TEST(index_supervisor_oversized_response_is_contained_and_log_is_retained); RUN_TEST(index_supervisor_killed_worker_log_is_never_empty_and_names_the_run); + RUN_TEST(index_supervisor_disabled_limits_do_not_probe); + RUN_TEST(index_supervisor_rss_equality_runs_then_excess_terminates_tree); + RUN_TEST(index_supervisor_duration_is_total_time_not_quiet_timeout); + RUN_TEST(index_supervisor_quiet_timeout_remains_hang_with_duration_enabled); + RUN_TEST(index_supervisor_cancel_precedes_resource_probe); + RUN_TEST(index_supervisor_three_failed_rss_probes_fail_closed); } diff --git a/tests/test_mcp.c b/tests/test_mcp.c index ffa9ea3a3..9b418ac79 100644 --- a/tests/test_mcp.c +++ b/tests/test_mcp.c @@ -18100,6 +18100,17 @@ TEST(index_supervisor_unsafe_clean_is_never_fallback_or_recovery) { CBM_MCP_SUPERVISED_RESULT_UNSAFE_TERMINAL); result.supervision_failed = false; + result.tree_quiesced = true; + result.outcome = CBM_PROC_KILLED; + result.resource_violation.resource = CBM_INDEX_RESOURCE_RSS_BYTES; + ASSERT_EQ(cbm_mcp_supervised_result_disposition(0, &result), + CBM_MCP_SUPERVISED_RESULT_RESOURCE_FAILURE); + result.tree_quiesced = false; + ASSERT_EQ(cbm_mcp_supervised_result_disposition(0, &result), + CBM_MCP_SUPERVISED_RESULT_UNSAFE_TERMINAL); + + result.resource_violation = (cbm_index_resource_violation_t){0}; + result.tree_quiesced = true; result.outcome = CBM_PROC_CRASH; result.response = NULL; ASSERT_EQ(cbm_mcp_supervised_result_disposition(0, &result), @@ -18109,6 +18120,102 @@ TEST(index_supervisor_unsafe_clean_is_never_fallback_or_recovery) { PASS(); } +#ifndef _WIN32 +typedef struct { + int calls; +} index_worker_resource_clock_t; + +static uint64_t index_worker_resource_clock(void *context) { + index_worker_resource_clock_t *clock = context; + clock->calls++; + return clock->calls == 1 ? 100 : 1101; +} + +enum { + IDXRESOURCE_OK = 0, + IDXRESOURCE_SETUP = 71, + IDXRESOURCE_NO_RESPONSE = 72, + IDXRESOURCE_BAD_CONTRACT = 73, + IDXRESOURCE_RETRIED = 74, +}; + +static int index_worker_resource_response_check(const char *repo, const char *cache) { + (void)cbm_setenv("CBM_CACHE_DIR", cache, 1); + cbm_unsetenv("CBM_INDEX_SUPERVISOR"); + cbm_index_supervisor_mark_host(); + cbm_config_t *config = cbm_config_open(cache); + cbm_mcp_server_t *server = cbm_mcp_server_new(NULL); + if (!config || !server || + cbm_config_set(config, CBM_INDEX_CONFIG_MAX_DURATION_SECONDS, "1") != 0) { + cbm_mcp_server_free(server); + cbm_config_close(config); + return IDXRESOURCE_SETUP; + } + cbm_mcp_server_set_config(server, config); + index_worker_resource_clock_t clock = {0}; + cbm_index_supervisor_set_resource_hooks_for_testing(index_worker_resource_clock, NULL, &clock); + int before = cbm_index_supervisor_spawn_count(); + char args[CBM_SZ_2K]; + (void)snprintf(args, sizeof(args), + "{\"repo_path\":\"%s\",\"mode\":\"fast\"," + "\"_cbm_index_policy\":{\"index_max_files\":\"off\"," + "\"index_max_source_mb\":\"off\",\"index_max_rss_mb\":\"off\"," + "\"index_max_duration_seconds\":\"off\"}}", + repo); + char *response = cbm_mcp_handle_tool(server, "index_repository", args); + int after = cbm_index_supervisor_spawn_count(); + cbm_index_supervisor_reset_resource_hooks_for_testing(); + int result = IDXRESOURCE_OK; + if (!response) { + result = IDXRESOURCE_NO_RESPONSE; + } else if (!strstr(response, "resource_limit_exceeded") || + !strstr(response, "\\\"stage\\\":\\\"worker\\\"") || + !strstr(response, "\\\"resource\\\":\\\"duration_ms\\\"") || + !strstr(response, "\\\"observed\\\":1001") || + !strstr(response, "\\\"limit\\\":1000") || + !strstr(response, "\\\"unit\\\":\\\"milliseconds\\\"")) { + result = IDXRESOURCE_BAD_CONTRACT; + } else if (after - before != 1) { + result = IDXRESOURCE_RETRIED; + } + free(response); + cbm_mcp_server_free(server); + cbm_config_close(config); + return result; +} +#endif + +TEST(index_worker_resource_limit_is_trusted_structured_and_not_retried) { +#ifdef _WIN32 + SKIP_PLATFORM("fork-isolated MCP host harness; supervisor state machine is cross-platform"); +#else + char *repo = th_mktempdir("cbm_worker_resource_repo"); + char *cache = th_mktempdir("cbm_worker_resource_cache"); + ASSERT_NOT_NULL(repo); + ASSERT_NOT_NULL(cache); + ASSERT_EQ(th_write_file(TH_PATH(repo, "main.c"), "int main(void) { return 0; }\n"), 0); + fflush(NULL); + pid_t pid = fork(); + if (pid == 0) { + alarm(60); + _exit(index_worker_resource_response_check(repo, cache)); + } + int status = 0; + bool waited = waitpid(pid, &status, 0) == pid; + int exit_code = waited && WIFEXITED(status) ? WEXITSTATUS(status) : -1; + int signal_code = waited && WIFSIGNALED(status) ? WTERMSIG(status) : 0; + th_cleanup(repo); + th_cleanup(cache); + if (exit_code != IDXRESOURCE_OK) { + printf(" worker resource child exit=%d signal=%d\n", exit_code, signal_code); + } + ASSERT_TRUE(waited); + ASSERT_EQ(signal_code, 0); + ASSERT_EQ(exit_code, IDXRESOURCE_OK); + PASS(); +#endif +} + /* Child-side check: index a tiny fixture and verify it ran IN-PROCESS. * Distinct exit codes so the parent can report the exact failure mode. */ enum { @@ -20682,6 +20789,7 @@ SUITE(mcp) { RUN_TEST(index_repository_cli_name_override_issue823); RUN_TEST(index_repository_over_budget_reports_named_reason); RUN_TEST(index_supervisor_unsafe_clean_is_never_fallback_or_recovery); + RUN_TEST(index_worker_resource_limit_is_trusted_structured_and_not_retried); RUN_TEST(index_supervisor_gate_requires_marked_host_issue845); RUN_TEST(index_supervisor_start_failure_is_fail_closed_in_real_host); RUN_TEST(index_bg_paths_route_through_supervisor_issue832); diff --git a/tests/test_subprocess.c b/tests/test_subprocess.c index 308d6b584..cf3a6b53a 100644 --- a/tests/test_subprocess.c +++ b/tests/test_subprocess.c @@ -86,6 +86,12 @@ TEST(subprocess_outcome_str) { PASS(); } +TEST(subprocess_tree_rss_sum_saturates_on_overflow) { + const uint64_t values[] = {UINT64_MAX - 4, 4, 1}; + ASSERT_EQ(cbm_subprocess_rss_sum_for_testing(values, 3), UINT64_MAX); + PASS(); +} + /* ── Layer 2: real spawn/reap (POSIX) ─────────────────────────────────────── */ #ifndef _WIN32 @@ -490,6 +496,45 @@ TEST(subprocess_spawn_returns_while_child_is_running) { #endif } +TEST(subprocess_tree_rss_measures_contained_descendants) { +#ifdef _WIN32 + SKIP_PLATFORM("native Job Object process-tree RSS runs on Windows CI"); +#else + char pid_path[64]; + ASSERT_TRUE(make_tree_pid_path(pid_path)); + cbm_subprocess_t *process = NULL; + ASSERT_EQ(spawn_ignoring_tree(pid_path, 0, 100, &process), 0); + pid_t parent_pid = 0; + pid_t grandchild_pid = 0; + bool ready = wait_for_tree_pids(pid_path, process, &parent_pid, &grandchild_pid, 3000); + uint64_t rss_bytes = 0; + cbm_proc_tree_rss_status_t rss_status = + ready ? cbm_subprocess_tree_rss_bytes(process, &rss_bytes) : CBM_PROC_TREE_RSS_ERROR; + bool cancel_accepted = cbm_subprocess_request_cancel(process); + cbm_proc_result_t result; + bool terminal = poll_until_terminal(process, 5000, &result); + uint64_t final_rss = UINT64_MAX; + cbm_proc_tree_rss_status_t final_rss_status = + terminal ? cbm_subprocess_tree_rss_bytes(process, &final_rss) : CBM_PROC_TREE_RSS_ERROR; + if (terminal) { + cbm_subprocess_destroy(process); + } else { + force_probe_cleanup(parent_pid, grandchild_pid); + } + (void)unlink(pid_path); + + ASSERT_TRUE(ready); + ASSERT_EQ(rss_status, CBM_PROC_TREE_RSS_OK); + ASSERT_TRUE(rss_bytes > 0); + ASSERT_TRUE(cancel_accepted); + ASSERT_TRUE(terminal); + ASSERT_TRUE(result.tree_quiesced); + ASSERT_EQ(final_rss_status, CBM_PROC_TREE_RSS_EMPTY); + ASSERT_EQ(final_rss, 0); + PASS(); +#endif +} + TEST(subprocess_natural_completion_is_cached_across_polls) { #ifdef _WIN32 SKIP_PLATFORM("POSIX /bin/sh completion-cache probe; native Job Object coverage pending"); @@ -1264,6 +1309,7 @@ SUITE(subprocess) { RUN_TEST(subprocess_classify_non_fault_signal_is_killed); RUN_TEST(subprocess_classify_timeout_dominates); RUN_TEST(subprocess_outcome_str); + RUN_TEST(subprocess_tree_rss_sum_saturates_on_overflow); RUN_TEST(subprocess_run_clean); RUN_TEST(subprocess_run_exit_nonzero); RUN_TEST(subprocess_run_resolves_literal_binary_name_from_path); @@ -1275,6 +1321,7 @@ SUITE(subprocess) { RUN_TEST(subprocess_run_spawn_failure); RUN_TEST(subprocess_run_null_bin_rejected); RUN_TEST(subprocess_spawn_returns_while_child_is_running); + RUN_TEST(subprocess_tree_rss_measures_contained_descendants); RUN_TEST(subprocess_natural_completion_is_cached_across_polls); RUN_TEST(subprocess_cancel_is_idempotent_and_kills_ignoring_tree); RUN_TEST(subprocess_quiet_timeout_kills_ignoring_tree); From 7d014c89cd2c3b9c64feca0b8ea6e39b421fb403 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=88=98=E5=86=B2?= Date: Tue, 22 Sep 2026 01:28:50 +0800 Subject: [PATCH 2/3] fix(test): send the complete worker policy from the session-scope fixture MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit PR 2257 still named only the two discovery keys in the worker-scope script. The worker now fail-closes unless every known policy key is present, so the fixture must stand in for the supervisor with cbm_test_index_worker_policy_json. Signed-off-by: 刘冲 --- tests/test_worker_session_scope.sh | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/tests/test_worker_session_scope.sh b/tests/test_worker_session_scope.sh index ab018ba76..1f32c4682 100644 --- a/tests/test_worker_session_scope.sh +++ b/tests/test_worker_session_scope.sh @@ -73,9 +73,9 @@ run_worker() { # A worker never reads the resource policy from config or environment: the # supervisor resolves it and sends it inside the request, and a request # without one is refused ("missing or incomplete trusted worker policy"). -# This script plays the supervisor, so it sends what the supervisor sends -- -# both limits off, the default. -policy='"_cbm_index_policy":{"index_max_files":"off","index_max_source_mb":"off"}' +# This script plays the supervisor, so it sends the same complete object +# cbm_test_index_worker_policy_json produces -- every known key, all off. +policy="$(cbm_test_index_worker_policy_json)" response="${tmpdir}/scoped.response" if ! run_worker "{\"repo_path\":\"${repo}\",\"mode\":\"fast\",${policy}}" "${response}" \ From a8da9a366353bb0a67806b40f840500444ac963f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=88=98=E5=86=B2?= Date: Tue, 22 Sep 2026 03:05:13 +0800 Subject: [PATCH 3/3] fix(index): charge macOS phys_footprint and pin duration across recovery MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Watchdogs now sample the OS-charged footprint instead of resident size on Darwin, and index_max_duration_seconds covers the whole request so a crash/hang respawn cannot reset the clock. Signed-off-by: 刘冲 --- README.md | 4 +- docs/CONFIGURATION.md | 4 +- docs/INDEX_RESOURCE_LIMITS.md | 28 ++++++---- src/cli/cli.c | 5 +- src/daemon/application.c | 17 +++--- src/daemon/application.h | 2 +- src/foundation/subprocess.c | 12 ++++- src/foundation/subprocess.h | 10 ++-- src/mcp/index_supervisor.c | 33 ++++++++---- src/mcp/index_supervisor.h | 9 +++- src/mcp/mcp.c | 8 +-- tests/test_daemon_application.c | 2 + tests/test_index_supervisor.c | 94 ++++++++++++++++++++++++++------- 13 files changed, 164 insertions(+), 64 deletions(-) diff --git a/README.md b/README.md index a70aa7805..434b74c18 100644 --- a/README.md +++ b/README.md @@ -738,8 +738,8 @@ codebase-memory-mcp config set auto_watch false # don't register backgr codebase-memory-mcp config set watcher_enabled false # stop the watcher thread entirely (default: true) codebase-memory-mcp config set index_max_files 250000 # optional per-index source-file limit codebase-memory-mcp config set index_max_source_mb 16384 # optional per-index source-size limit -codebase-memory-mcp config set index_max_rss_mb 8192 # optional worker-tree current RSS limit -codebase-memory-mcp config set index_max_duration_seconds 3600 # optional total worker duration +codebase-memory-mcp config set index_max_rss_mb 8192 # optional worker-tree charged-memory limit +codebase-memory-mcp config set index_max_duration_seconds 3600 # optional per-request duration codebase-memory-mcp config reset auto_index # reset to default ``` diff --git a/docs/CONFIGURATION.md b/docs/CONFIGURATION.md index 79b74e7e1..d47f19bb5 100644 --- a/docs/CONFIGURATION.md +++ b/docs/CONFIGURATION.md @@ -91,8 +91,8 @@ Current keys: | `watcher_enabled` | `true` | Master switch for the background watcher subsystem. Set `false` to stop the watcher from starting at all — no poll thread and no project registration. Reindex manually with `index_repository` when disabled. | | `index_max_files` | `off` | Optional maximum number of accepted source files in one discovery run. | | `index_max_source_mb` | `off` | Optional maximum accepted source size in MiB in one discovery run. | -| `index_max_rss_mb` | `off` | Optional maximum current RSS in MiB for the complete contained index-worker process tree (`64..1048576`). | -| `index_max_duration_seconds` | `off` | Optional maximum total worker duration in seconds (`1..86400`). | +| `index_max_rss_mb` | `off` | Optional maximum charged memory in MiB for the complete contained index-worker process tree (`64..1048576`). On macOS this is phys_footprint, not resident size. | +| `index_max_duration_seconds` | `off` | Optional maximum wall-clock duration in seconds for the whole index request, including crash/hang recovery (`1..86400`). | > **`watcher_enabled` vs `auto_watch`.** `watcher_enabled` controls whether the > watcher *subsystem* starts at all (the background poll thread). `auto_watch` is diff --git a/docs/INDEX_RESOURCE_LIMITS.md b/docs/INDEX_RESOURCE_LIMITS.md index 66be4ecf4..152599360 100644 --- a/docs/INDEX_RESOURCE_LIMITS.md +++ b/docs/INDEX_RESOURCE_LIMITS.md @@ -27,24 +27,30 @@ ranges are rejected without changing the stored value. | Key | Default | Accepted value | Protects | |---|---:|---:|---| -| `index_max_rss_mb` | `off` | `off` or `64..1048576` | Current RSS of the complete worker process tree | -| `index_max_duration_seconds` | `off` | `off` or `1..86400` | Total worker wall-clock duration | +| `index_max_rss_mb` | `off` | `off` or `64..1048576` | Charged memory of the complete worker process tree | +| `index_max_duration_seconds` | `off` | `off` or `1..86400` | Wall-clock duration of the whole index request | ```bash codebase-memory-mcp config set index_max_rss_mb 8192 codebase-memory-mcp config set index_max_duration_seconds 3600 ``` -RSS is the current resident memory of the contained worker and every descendant, -not the worker's allocation budget and not peak memory. This hard watchdog is +RSS is the charged memory of the contained worker and every descendant, not +the worker's allocation budget and not peak memory. On macOS that quantity is +`phys_footprint` (the same number `cbm_mem_charged()` enforces), not +`resident_size`, which still counts pages the allocator has already handed +back. Linux and Windows use RSS / working set. This hard watchdog is separate from the internal `CBM_MEM_BUDGET_MB` soft budget. The supervisor -samples RSS at most once every 250 milliseconds so the watchdog does not turn -full process-table enumeration into a busy loop. - -Duration uses a monotonic clock from successful spawn. It is independent of the -existing 15-minute quiet timeout: continuous log progress does not reset total -duration, while the quiet timeout continues to identify a worker that stops -making progress. +samples the tree at most once every 250 milliseconds so the watchdog does not +turn full process-table enumeration into a busy loop. + +Duration is per request: the clock starts at the first worker spawn and is +not reset when crash/hang recovery starts a later attempt. Continuous log +progress does not reset it. It is independent of the existing 15-minute quiet +timeout, which still identifies a worker that stops making progress. A +duration limit shorter than that quiet timeout kills a hung worker before hang +quarantine can name the in-flight file, so the next attempt may hang on the +same file. Equality is allowed. The first RSS or elapsed-duration observation above its limit starts the existing graceful-to-force process-tree shutdown. CBM reports diff --git a/src/cli/cli.c b/src/cli/cli.c index 66817c902..ab3fb5037 100644 --- a/src/cli/cli.c +++ b/src/cli/cli.c @@ -7449,8 +7449,9 @@ static const config_key_def_t CONFIG_KEYS[] = { {CBM_CONFIG_UI_PORT, "9749", "Port for the graph UI listener when enabled"}, {CBM_INDEX_CONFIG_MAX_FILES, "off", "Max accepted source files per index, or off"}, {CBM_INDEX_CONFIG_MAX_SOURCE_MB, "off", "Max accepted source MiB per index, or off"}, - {CBM_INDEX_CONFIG_MAX_RSS_MB, "off", "Max worker process-tree RSS MiB, or off"}, - {CBM_INDEX_CONFIG_MAX_DURATION_SECONDS, "off", "Max worker duration in seconds, or off"}, + {CBM_INDEX_CONFIG_MAX_RSS_MB, "off", "Max worker process-tree charged memory MiB, or off"}, + {CBM_INDEX_CONFIG_MAX_DURATION_SECONDS, "off", + "Max wall-clock seconds for the whole index request, or off"}, }; /* #1558: ui_enabled and ui_port were reachable ONLY by hand-editing diff --git a/src/daemon/application.c b/src/daemon/application.c index f8922e36f..5e5576edf 100644 --- a/src/daemon/application.c +++ b/src/daemon/application.c @@ -147,6 +147,7 @@ struct cbm_daemon_application_job { bool cancelled; bool cancel_requested; bool supervision_failed; + uint64_t request_started_ms; cbm_daemon_application_job_t *next; }; @@ -303,6 +304,7 @@ static void application_project_lock_release_fully(cbm_project_lock_lease_t **le static int application_worker_start_default(void *context, const char *args_json, size_t memory_budget_bytes, const char *marker_file, const char *quarantine_file, + uint64_t duration_origin_ms, cbm_daemon_application_worker_t *worker_out) { (void)context; cbm_index_resource_policy_t resource_policy; @@ -314,9 +316,9 @@ static int application_worker_start_default(void *context, const char *args_json return -1; } cbm_index_worker_handle_t *worker = NULL; - int result = - cbm_index_worker_start_with_policy(args_json, memory_budget_bytes, &resource_policy, false, - marker_file, quarantine_file, &worker); + int result = cbm_index_worker_start_with_policy(args_json, memory_budget_bytes, + &resource_policy, false, marker_file, + quarantine_file, duration_origin_ms, &worker); *worker_out = worker; return result; } @@ -1124,10 +1126,13 @@ static application_attempt_status_t application_job_run_attempt(cbm_daemon_appli } cbm_daemon_application_worker_t worker = NULL; + if (job->request_started_ms == 0) { + job->request_started_ms = cbm_index_worker_now_ms(); + } application_tmp_lock(); - int start_result = - application->worker_ops.start(application->worker_ops.context, job->args_json, - memory_budget_bytes, marker_path, quarantine_path, &worker); + int start_result = application->worker_ops.start( + application->worker_ops.context, job->args_json, memory_budget_bytes, marker_path, + quarantine_path, job->request_started_ms, &worker); application_tmp_unlock(); if (start_result != 0 || !worker) { return application_job_cancel_requested(job) ? APPLICATION_ATTEMPT_CANCELLED diff --git a/src/daemon/application.h b/src/daemon/application.h index 933186aa7..426bef3db 100644 --- a/src/daemon/application.h +++ b/src/daemon/application.h @@ -31,7 +31,7 @@ typedef void *cbm_daemon_application_update_worker_t; typedef struct { void *context; int (*start)(void *context, const char *args_json, size_t memory_budget_bytes, - const char *marker_file, const char *quarantine_file, + const char *marker_file, const char *quarantine_file, uint64_t duration_origin_ms, cbm_daemon_application_worker_t *worker_out); cbm_index_worker_poll_t (*poll)(void *context, cbm_daemon_application_worker_t worker, const cbm_index_worker_result_t **result_out); diff --git a/src/foundation/subprocess.c b/src/foundation/subprocess.c index b14a74fbc..813e8f461 100644 --- a/src/foundation/subprocess.c +++ b/src/foundation/subprocess.c @@ -495,6 +495,12 @@ static cbm_proc_tree_rss_status_t cbm_subprocess_tree_rss_platform(cbm_subproces bool root_failed = false; for (DWORD index = 0; index < processes->NumberOfProcessIdsInList; index++) { DWORD pid = (DWORD)processes->ProcessIdList[index]; + /* GetProcessMemoryInfo documents PROCESS_QUERY_INFORMATION | + * PROCESS_VM_READ. PROCESS_QUERY_LIMITED_INFORMATION is a lesser + * alternative on newer Windows; the Job Object membership list + * already closes most of the pid-reuse window between listing and + * opening. A member we cannot open is a probe failure when it is + * the root. */ HANDLE member = OpenProcess(PROCESS_QUERY_INFORMATION | PROCESS_VM_READ, FALSE, pid); if (!member) { root_failed = root_failed || pid == process->process_id; @@ -574,7 +580,11 @@ static cbm_proc_tree_rss_status_t cbm_subprocess_tree_rss_platform(cbm_subproces root_failed = root_failed || pids[index] == process->pid; continue; } - cbm_rss_add_saturated(&total, usage.ri_resident_size); + /* phys_footprint is what the OS charges, matching cbm_mem_charged(). + * ri_resident_size still counts pages mimalloc already marked + * MADV_FREE_REUSABLE, so a Mac worker can look like 17 GB RSS while + * the kernel holds it to 5.5 GB. */ + cbm_rss_add_saturated(&total, usage.ri_phys_footprint); measured++; } cbm_free(CBM_MEM_CLASS_OTHER, pids); diff --git a/src/foundation/subprocess.h b/src/foundation/subprocess.h index 1179d97c4..876ca871e 100644 --- a/src/foundation/subprocess.h +++ b/src/foundation/subprocess.h @@ -129,10 +129,12 @@ int cbm_subprocess_spawn(const cbm_proc_opts_t *opts, cbm_subprocess_t **out); * tree when cancel_grace_ms elapses. Callers must keep polling to make progress. */ cbm_proc_poll_t cbm_subprocess_poll(cbm_subprocess_t *process, cbm_proc_result_t *out); -/* Read the current resident-set size of the complete contained process tree. - * OK returns an overflow-safe byte total, EMPTY means the owned tree currently - * has no observable members, and ERROR means no trustworthy measurement could - * be obtained. Individual processes that exit during enumeration are ignored. */ +/* Read the current charged memory of the complete contained process tree. + * On macOS this is phys_footprint (what the OS holds against the process); + * elsewhere it is RSS. OK returns an overflow-safe byte total, EMPTY means + * the owned tree currently has no observable members, and ERROR means no + * trustworthy measurement could be obtained. Individual processes that exit + * during enumeration are ignored. */ cbm_proc_tree_rss_status_t cbm_subprocess_tree_rss_bytes(cbm_subprocess_t *process, uint64_t *rss_bytes); bool cbm_subprocess_root_running(const cbm_subprocess_t *process); diff --git a/src/mcp/index_supervisor.c b/src/mcp/index_supervisor.c index 93d3bef46..0ee29f8ec 100644 --- a/src/mcp/index_supervisor.c +++ b/src/mcp/index_supervisor.c @@ -806,7 +806,8 @@ static int worker_start_internal(const char *args_json, size_t memory_budget_byt const cbm_index_resource_policy_t *resource_policy, bool single_thread, const char *marker_file, const char *quarantine_file, cbm_proc_log_cb log_callback, - void *log_context, cbm_index_worker_handle_t **handle_out) { + void *log_context, uint64_t duration_origin_ms, + cbm_index_worker_handle_t **handle_out) { if (handle_out) { *handle_out = NULL; } @@ -905,34 +906,42 @@ static int worker_start_internal(const char *args_json, size_t memory_budget_byt return -1; } if (handle->resource_policy.max_duration_ms.enabled) { - handle->started_ms = worker_resource_now_ms(); + /* Duration is per request, not per spawn: a crash/hang recovery must + * not reset the clock. Callers pass the request origin; 0 means now. */ + uint64_t origin = duration_origin_ms != 0 ? duration_origin_ms : worker_resource_now_ms(); + handle->started_ms = origin; } *handle_out = handle; return 0; } +uint64_t cbm_index_worker_now_ms(void) { + return worker_resource_now_ms(); +} + int cbm_index_worker_start_with_log(const char *args_json, size_t memory_budget_bytes, bool single_thread, const char *marker_file, const char *quarantine_file, cbm_proc_log_cb log_callback, void *log_context, cbm_index_worker_handle_t **handle_out) { return worker_start_internal(args_json, memory_budget_bytes, NULL, single_thread, marker_file, - quarantine_file, log_callback, log_context, handle_out); + quarantine_file, log_callback, log_context, 0, handle_out); } int cbm_index_worker_start(const char *args_json, size_t memory_budget_bytes, bool single_thread, const char *marker_file, const char *quarantine_file, cbm_index_worker_handle_t **handle_out) { return worker_start_internal(args_json, memory_budget_bytes, NULL, single_thread, marker_file, - quarantine_file, NULL, NULL, handle_out); + quarantine_file, NULL, NULL, 0, handle_out); } int cbm_index_worker_start_with_policy(const char *args_json, size_t memory_budget_bytes, const cbm_index_resource_policy_t *resource_policy, bool single_thread, const char *marker_file, - const char *quarantine_file, + const char *quarantine_file, uint64_t duration_origin_ms, cbm_index_worker_handle_t **handle_out) { return worker_start_internal(args_json, memory_budget_bytes, resource_policy, single_thread, - marker_file, quarantine_file, NULL, NULL, handle_out); + marker_file, quarantine_file, NULL, NULL, duration_origin_ms, + handle_out); } cbm_index_worker_poll_t cbm_index_worker_poll(cbm_index_worker_handle_t *handle, @@ -1066,14 +1075,15 @@ static int worker_spawn_internal(const char *args_json, bool single_thread, const char *marker_file, const char *quarantine_file, cbm_proc_log_cb log_callback, void *log_context, const atomic_int *cancel_requested, - cbm_index_worker_result_t *result) { + uint64_t duration_origin_ms, cbm_index_worker_result_t *result) { if (!result) { return -1; } worker_result_init(result); cbm_index_worker_handle_t *handle = NULL; if (worker_start_internal(args_json, 0, resource_policy, single_thread, marker_file, - quarantine_file, log_callback, log_context, &handle) != 0) { + quarantine_file, log_callback, log_context, duration_origin_ms, + &handle) != 0) { return -1; } const cbm_index_worker_result_t *cached = NULL; @@ -1105,16 +1115,17 @@ int cbm_index_spawn_worker_with_log_cancel(const char *args_json, bool single_th const atomic_int *cancel_requested, cbm_index_worker_result_t *result) { return worker_spawn_internal(args_json, NULL, single_thread, marker_file, quarantine_file, - log_callback, log_context, cancel_requested, result); + log_callback, log_context, cancel_requested, 0, result); } int cbm_index_spawn_worker_with_policy_log_cancel( const char *args_json, const cbm_index_resource_policy_t *resource_policy, bool single_thread, const char *marker_file, const char *quarantine_file, cbm_proc_log_cb log_callback, - void *log_context, const atomic_int *cancel_requested, cbm_index_worker_result_t *result) { + void *log_context, const atomic_int *cancel_requested, uint64_t duration_origin_ms, + cbm_index_worker_result_t *result) { return worker_spawn_internal(args_json, resource_policy, single_thread, marker_file, quarantine_file, log_callback, log_context, cancel_requested, - result); + duration_origin_ms, result); } int cbm_index_spawn_worker_with_log(const char *args_json, bool single_thread, diff --git a/src/mcp/index_supervisor.h b/src/mcp/index_supervisor.h index fdf7b1b44..61b406e2c 100644 --- a/src/mcp/index_supervisor.h +++ b/src/mcp/index_supervisor.h @@ -178,9 +178,13 @@ int cbm_index_worker_start(const char *args_json, size_t memory_budget_bytes, bo int cbm_index_worker_start_with_policy(const char *args_json, size_t memory_budget_bytes, const cbm_index_resource_policy_t *resource_policy, bool single_thread, const char *marker_file, - const char *quarantine_file, + const char *quarantine_file, uint64_t duration_origin_ms, cbm_index_worker_handle_t **handle_out); +/* Supervisor clock used to stamp a request-scoped duration origin. Tests may + * replace it; production is cbm_now_ms(). */ +uint64_t cbm_index_worker_now_ms(void); + /* Request-scoped variant used by interactive local CLI calls. The callback is * invoked by the owner thread while it polls the contained worker; log_context * remains caller-owned until terminal. No process-global sink is installed. @@ -246,7 +250,8 @@ int cbm_index_spawn_worker_with_log_cancel(const char *args_json, bool single_th int cbm_index_spawn_worker_with_policy_log_cancel( const char *args_json, const cbm_index_resource_policy_t *resource_policy, bool single_thread, const char *marker_file, const char *quarantine_file, cbm_proc_log_cb log_callback, - void *log_context, const atomic_int *cancel_requested, cbm_index_worker_result_t *result); + void *log_context, const atomic_int *cancel_requested, uint64_t duration_origin_ms, + cbm_index_worker_result_t *result); void cbm_index_worker_result_free(cbm_index_worker_result_t *result); diff --git a/src/mcp/mcp.c b/src/mcp/mcp.c index 86d392ece..d474b9508 100644 --- a/src/mcp/mcp.c +++ b/src/mcp/mcp.c @@ -10783,12 +10783,14 @@ int cbm_index_restart_cap_for_testing(void) { static char *index_run_supervised(cbm_mcp_server_t *srv, const char *args, const cbm_index_resource_policy_t *resource_policy) { invalidate_cached_store(srv); + uint64_t request_started_ms = cbm_index_worker_now_ms(); /* First attempt: normal parallel run. */ cbm_index_worker_result_t wr; int rc = cbm_index_spawn_worker_with_policy_log_cancel( args, resource_policy, false, NULL, NULL, srv ? srv->index_log_callback : NULL, - srv ? srv->index_log_context : NULL, srv ? &srv->pipeline_cancel_requested : NULL, &wr); + srv ? srv->index_log_context : NULL, srv ? &srv->pipeline_cancel_requested : NULL, + request_started_ms, &wr); cbm_mcp_supervised_result_disposition_t disposition = cbm_mcp_supervised_result_disposition(rc, &wr); @@ -10863,7 +10865,7 @@ static char *index_run_supervised(cbm_mcp_server_t *srv, const char *args, int rc2 = cbm_index_spawn_worker_with_policy_log_cancel( args, resource_policy, /*single_thread=*/false, marker_path, quarantine_path, srv ? srv->index_log_callback : NULL, srv ? srv->index_log_context : NULL, - srv ? &srv->pipeline_cancel_requested : NULL, &wr2); + srv ? &srv->pipeline_cancel_requested : NULL, request_started_ms, &wr2); cbm_mcp_supervised_result_disposition_t recovery_disposition = cbm_mcp_supervised_result_disposition(rc2, &wr2); if (recovery_disposition == CBM_MCP_SUPERVISED_RESULT_FALLBACK) { @@ -10980,7 +10982,7 @@ static char *index_run_supervised(cbm_mcp_server_t *srv, const char *args, int rcp = cbm_index_spawn_worker_with_policy_log_cancel( args, resource_policy, /*single_thread=*/false, NULL, quarantine_path, srv ? srv->index_log_callback : NULL, srv ? srv->index_log_context : NULL, - srv ? &srv->pipeline_cancel_requested : NULL, &wrp); + srv ? &srv->pipeline_cancel_requested : NULL, request_started_ms, &wrp); cbm_mcp_supervised_result_disposition_t partial_disposition = cbm_mcp_supervised_result_disposition(rcp, &wrp); if (partial_disposition == CBM_MCP_SUPERVISED_RESULT_SUCCESS) { diff --git a/tests/test_daemon_application.c b/tests/test_daemon_application.c index 7c8b841e0..d2c0670cb 100644 --- a/tests/test_daemon_application.c +++ b/tests/test_daemon_application.c @@ -1417,7 +1417,9 @@ static void app_fake_worker_read_file(const char *path, char *out, size_t out_si static int app_fake_worker_start(void *opaque, const char *args_json, size_t memory_budget_bytes, const char *marker_file, const char *quarantine_file, + uint64_t duration_origin_ms, cbm_daemon_application_worker_t *worker_out) { + (void)duration_origin_ms; app_fake_worker_context_t *context = opaque; app_fake_worker_t *worker = calloc(1, sizeof(*worker)); if (!worker) { diff --git a/tests/test_index_supervisor.c b/tests/test_index_supervisor.c index d1227dd06..a9855b04b 100644 --- a/tests/test_index_supervisor.c +++ b/tests/test_index_supervisor.c @@ -687,8 +687,8 @@ TEST(index_supervisor_worker_keeps_default_info_liveness_heartbeat) { (void)cbm_unsetenv("CBM_LOG_LEVEL"); cbm_index_worker_handle_t *handle = NULL; - int start_rc = cbm_index_worker_start("{\"__cbm_test_worker\":\"heartbeat\"}", 0, false, - NULL, NULL, &handle); + int start_rc = cbm_index_worker_start("{\"__cbm_test_worker\":\"heartbeat\"}", 0, false, NULL, + NULL, &handle); char log_path[INDEX_SUPERVISOR_TEST_PATH_CAP] = {0}; if (handle) { (void)snprintf(log_path, sizeof(log_path), "%s", cbm_index_worker_log_path(handle)); @@ -696,8 +696,8 @@ TEST(index_supervisor_worker_keeps_default_info_liveness_heartbeat) { bool ready = log_path[0] && index_supervisor_test_wait_file_text( log_path, "async worker heartbeat probe ready", INDEX_SUPERVISOR_TEST_READY_MS); - bool heartbeat = ready && index_supervisor_test_wait_file_text( - log_path, "msg=pipeline.discover", 1000); + bool heartbeat = + ready && index_supervisor_test_wait_file_text(log_path, "msg=pipeline.discover", 1000); const cbm_index_worker_result_t *result = NULL; bool terminal = handle && index_supervisor_test_poll_terminal( handle, INDEX_SUPERVISOR_TEST_TERMINAL_MS, &result); @@ -956,8 +956,9 @@ static uint64_t index_supervisor_fake_clock(void *context) { return fake->now_ms; } -static cbm_proc_tree_rss_status_t index_supervisor_fake_rss(cbm_subprocess_t *process, - uint64_t *rss_bytes, void *context) { +static cbm_proc_tree_rss_status_t index_supervisor_fake_charged_rss(cbm_subprocess_t *process, + uint64_t *rss_bytes, + void *context) { (void)process; index_supervisor_resource_fake_t *fake = context; int value_index = @@ -987,11 +988,11 @@ TEST(index_supervisor_disabled_limits_do_not_probe) { .rss_status = CBM_PROC_TREE_RSS_OK, }; cbm_index_supervisor_set_resource_hooks_for_testing(index_supervisor_fake_clock, - index_supervisor_fake_rss, &fake); + index_supervisor_fake_charged_rss, &fake); cbm_index_resource_policy_t policy = index_supervisor_test_worker_policy(0, 0); cbm_index_worker_handle_t *handle = NULL; int start_rc = cbm_index_worker_start_with_policy("{\"__cbm_test_worker\":\"clean\"}", 0, - &policy, false, NULL, NULL, &handle); + &policy, false, NULL, NULL, 0, &handle); const cbm_index_worker_result_t *result = NULL; bool terminal = handle && index_supervisor_test_poll_terminal( handle, INDEX_SUPERVISOR_TEST_TERMINAL_MS, &result); @@ -1014,7 +1015,7 @@ TEST(index_supervisor_disabled_limits_do_not_probe) { PASS(); } -TEST(index_supervisor_rss_equality_runs_then_excess_terminates_tree) { +TEST(index_supervisor_charged_rss_equality_runs_then_excess_terminates_tree) { const uint64_t limit = UINT64_C(64) * CBM_INDEX_MIB_BYTES; index_supervisor_resource_fake_t fake = { .now_ms = 100, @@ -1023,11 +1024,11 @@ TEST(index_supervisor_rss_equality_runs_then_excess_terminates_tree) { .rss_status = CBM_PROC_TREE_RSS_OK, }; cbm_index_supervisor_set_resource_hooks_for_testing(index_supervisor_fake_clock, - index_supervisor_fake_rss, &fake); + index_supervisor_fake_charged_rss, &fake); cbm_index_resource_policy_t policy = index_supervisor_test_worker_policy(limit, 0); cbm_index_worker_handle_t *handle = NULL; int start_rc = cbm_index_worker_start_with_policy("{\"__cbm_test_worker\":\"hang-tree\"}", 0, - &policy, false, NULL, NULL, &handle); + &policy, false, NULL, NULL, 0, &handle); const cbm_index_worker_result_t *result = NULL; cbm_index_worker_poll_t equal_state = handle ? cbm_index_worker_poll(handle, &result) : CBM_INDEX_WORKER_POLL_ERROR; @@ -1063,11 +1064,11 @@ TEST(index_supervisor_duration_is_total_time_not_quiet_timeout) { .rss_status = CBM_PROC_TREE_RSS_EMPTY, }; cbm_index_supervisor_set_resource_hooks_for_testing(index_supervisor_fake_clock, - index_supervisor_fake_rss, &fake); + index_supervisor_fake_charged_rss, &fake); cbm_index_resource_policy_t policy = index_supervisor_test_worker_policy(0, 1000); cbm_index_worker_handle_t *handle = NULL; int start_rc = cbm_index_worker_start_with_policy("{\"__cbm_test_worker\":\"hang-tree\"}", 0, - &policy, false, NULL, NULL, &handle); + &policy, false, NULL, NULL, 0, &handle); fake.now_ms = 1100; const cbm_index_worker_result_t *result = NULL; cbm_index_worker_poll_t equal_state = @@ -1093,6 +1094,60 @@ TEST(index_supervisor_duration_is_total_time_not_quiet_timeout) { PASS(); } +TEST(index_supervisor_duration_spans_recovery_attempts) { + index_supervisor_resource_fake_t fake = { + .now_ms = 100, + .rss_status = CBM_PROC_TREE_RSS_EMPTY, + }; + cbm_index_supervisor_set_resource_hooks_for_testing(index_supervisor_fake_clock, + index_supervisor_fake_charged_rss, &fake); + cbm_index_resource_policy_t policy = index_supervisor_test_worker_policy(0, 1000); + cbm_index_worker_handle_t *first = NULL; + int first_rc = cbm_index_worker_start_with_policy("{\"__cbm_test_worker\":\"hang-tree\"}", 0, + &policy, false, NULL, NULL, 100, &first); + fake.now_ms = 500; + const cbm_index_worker_result_t *first_result = NULL; + cbm_index_worker_poll_t mid = + first ? cbm_index_worker_poll(first, &first_result) : CBM_INDEX_WORKER_POLL_ERROR; + bool cancel_accepted = first && cbm_index_worker_request_cancel(first); + bool first_terminal = first && index_supervisor_test_poll_terminal( + first, INDEX_SUPERVISOR_TEST_TERMINAL_MS, &first_result); + if (first_terminal) { + cbm_index_worker_destroy(first); + first = NULL; + } else { + index_supervisor_test_cleanup_handle(first); + first = NULL; + } + + fake.now_ms = 1101; + cbm_index_worker_handle_t *second = NULL; + int second_rc = cbm_index_worker_start_with_policy("{\"__cbm_test_worker\":\"hang-tree\"}", 0, + &policy, false, NULL, NULL, 100, &second); + const cbm_index_worker_result_t *second_result = NULL; + bool second_terminal = second && index_supervisor_test_poll_terminal( + second, INDEX_SUPERVISOR_TEST_TERMINAL_MS, &second_result); + bool limited = second_terminal && second_result && + second_result->resource_violation.resource == CBM_INDEX_RESOURCE_DURATION_MS && + second_result->resource_violation.observed == 1001 && + second_result->resource_violation.limit == 1000 && + second_result->outcome != CBM_PROC_HANG && second_result->tree_quiesced; + if (second_terminal) { + cbm_index_worker_destroy(second); + } else { + index_supervisor_test_cleanup_handle(second); + } + cbm_index_supervisor_reset_resource_hooks_for_testing(); + + ASSERT_EQ(first_rc, 0); + ASSERT_EQ(mid, CBM_INDEX_WORKER_POLL_RUNNING); + ASSERT_TRUE(cancel_accepted); + ASSERT_TRUE(first_terminal); + ASSERT_EQ(second_rc, 0); + ASSERT_TRUE(limited); + PASS(); +} + TEST(index_supervisor_quiet_timeout_remains_hang_with_duration_enabled) { const char *saved_timeout = getenv("CBM_INDEX_WORKER_TIMEOUT_S"); char *saved_timeout_copy = saved_timeout ? cbm_strdup(saved_timeout) : NULL; @@ -1101,7 +1156,7 @@ TEST(index_supervisor_quiet_timeout_remains_hang_with_duration_enabled) { cbm_index_resource_policy_t policy = index_supervisor_test_worker_policy(0, 60000); cbm_index_worker_handle_t *handle = NULL; int start_rc = cbm_index_worker_start_with_policy("{\"__cbm_test_worker\":\"hang-tree\"}", 0, - &policy, false, NULL, NULL, &handle); + &policy, false, NULL, NULL, 0, &handle); const cbm_index_worker_result_t *result = NULL; bool terminal = handle && index_supervisor_test_poll_terminal( handle, INDEX_SUPERVISOR_TEST_TERMINAL_MS, &result); @@ -1130,11 +1185,11 @@ TEST(index_supervisor_cancel_precedes_resource_probe) { .rss_status = CBM_PROC_TREE_RSS_OK, }; cbm_index_supervisor_set_resource_hooks_for_testing(index_supervisor_fake_clock, - index_supervisor_fake_rss, &fake); + index_supervisor_fake_charged_rss, &fake); cbm_index_resource_policy_t policy = index_supervisor_test_worker_policy(limit, 1000); cbm_index_worker_handle_t *handle = NULL; int start_rc = cbm_index_worker_start_with_policy("{\"__cbm_test_worker\":\"hang-tree\"}", 0, - &policy, false, NULL, NULL, &handle); + &policy, false, NULL, NULL, 0, &handle); fake.now_ms = 2000; bool cancel_accepted = handle && cbm_index_worker_request_cancel(handle); const cbm_index_worker_result_t *result = NULL; @@ -1164,11 +1219,11 @@ TEST(index_supervisor_three_failed_rss_probes_fail_closed) { .rss_status = CBM_PROC_TREE_RSS_ERROR, }; cbm_index_supervisor_set_resource_hooks_for_testing(index_supervisor_fake_clock, - index_supervisor_fake_rss, &fake); + index_supervisor_fake_charged_rss, &fake); cbm_index_resource_policy_t policy = index_supervisor_test_worker_policy(limit, 0); cbm_index_worker_handle_t *handle = NULL; int start_rc = cbm_index_worker_start_with_policy("{\"__cbm_test_worker\":\"hang-tree\"}", 0, - &policy, false, NULL, NULL, &handle); + &policy, false, NULL, NULL, 0, &handle); const cbm_index_worker_result_t *result = NULL; cbm_index_worker_poll_t first = handle ? cbm_index_worker_poll(handle, &result) : CBM_INDEX_WORKER_POLL_ERROR; @@ -1215,8 +1270,9 @@ SUITE(index_supervisor) { RUN_TEST(index_supervisor_oversized_response_is_contained_and_log_is_retained); RUN_TEST(index_supervisor_killed_worker_log_is_never_empty_and_names_the_run); RUN_TEST(index_supervisor_disabled_limits_do_not_probe); - RUN_TEST(index_supervisor_rss_equality_runs_then_excess_terminates_tree); + RUN_TEST(index_supervisor_charged_rss_equality_runs_then_excess_terminates_tree); RUN_TEST(index_supervisor_duration_is_total_time_not_quiet_timeout); + RUN_TEST(index_supervisor_duration_spans_recovery_attempts); RUN_TEST(index_supervisor_quiet_timeout_remains_hang_with_duration_enabled); RUN_TEST(index_supervisor_cancel_precedes_resource_probe); RUN_TEST(index_supervisor_three_failed_rss_probes_fail_closed);