From 28b41c35b71cb039988f39f72689dd657e30522d Mon Sep 17 00:00:00 2001 From: Martin Vogel Date: Fri, 25 Sep 2026 23:03:56 +0200 Subject: [PATCH] feat(mcp): opt-in async index_repository with status polling (#2144) index_repository blocks until the whole index completes. MCP clients with a per-call deadline (Copilot for IntelliJ) give up on a large repository; their cancel drops the daemon job's last subscriber, the daemon cancels the worker, and every retry starts over, so the index never finishes. Add an opt-in async mode and a status query to the same tool; the default synchronous behaviour is unchanged: - async: true starts the project's index job in the daemon, or joins the one already running for it, and returns at once with its state. The application holds one subscriber reference for an async job until its terminal publish, so the request returning, the client cancelling or the starting session closing no longer cancels it. Only daemon shutdown (or the final session of a non-permanent generation) ends it. Async requests never queue behind the physical job limit; a full daemon answers busy. async/status are stripped from the worker args, so an async request coalesces with an identical synchronous one. - status: true reports the running or last job for the project (queued/running/cancelling/succeeded/failed/cancelled, started_at, finished_at, error summary) from a small per-project record kept in the daemon's job registry after the job itself is reaped. An unknown project is an error; a project indexed before this daemon generation reports idle. No freshness key is added: index_status keeps the #1561 freshness object. - async + status together, non-boolean values, and async with cross-repo-intelligence are refused. In-process servers (index worker, embedders) refuse both modes: nothing there outlives the call. - On a temporary (non-permanent) daemon, async is refused when the request arrives on the one-shot `cli` tool channel and that session is the daemon's only live session: the daemon would stop and cancel the job the moment the command exits. The error points at `daemon start` (or an open MCP session). A long-lived MCP session stays a valid host even when it is the only one, which is the IDE case this exists for; status stays allowed everywhere. The daemon lifecycle is unchanged. - New allocations go through the memory core (cbm_alloc/cbm_calloc/ cbm_mem_strdup/cbm_free); blocks returned by non-core APIs are released with safe_free, so no file's raw allocator count grows. - A cut-short synchronous call now offers the async alternative: in the cancellation text the daemon or frontend delivers (including the -32800 reply to a cancelled index_repository), and, because a timed-out client never reads that reply, as a notice on the next index_repository or status call for the project. - The tool description and schema, the CLI help and the README explain the flags, the deadline problem and the polling pattern. No server-side timeout or budget is introduced. Tests (deterministic, held fake worker, bounded observation waits): async returns before completion and status observes running then succeeded; an async job survives its session closing while a sync job in the same shape is still cancelled; validation errors; the cancellation reply and the next-call notice carry the advice; a lone one-shot client of a temporary daemon is refused async while status and an MCP session on the same daemon are allowed; in-process refusal; schema; notice helper. Refs #2031 Fixes #2144 Signed-off-by: Martin Vogel --- README.md | 30 +- src/daemon/application.c | 445 ++++++++++++++++++++++++++-- src/daemon/frontend.c | 20 +- src/daemon/frontend.h | 6 + src/mcp/mcp.c | 207 ++++++++++++- src/mcp/mcp.h | 31 +- src/ui/http_server.c | 3 +- tests/test_daemon_application.c | 496 ++++++++++++++++++++++++++++++++ tests/test_daemon_frontend.c | 22 ++ tests/test_mcp.c | 63 ++++ 10 files changed, 1295 insertions(+), 28 deletions(-) diff --git a/README.md b/README.md index 9c82c18aed..4752006fd6 100644 --- a/README.md +++ b/README.md @@ -669,12 +669,39 @@ JSON arguments can also be piped on stdin, for tools that take arguments. A tool | Tool | Description | |------|-------------| -| `index_repository` | Index a repository into the graph. Auto-sync keeps it fresh after that. | +| `index_repository` | Index a repository into the graph. Auto-sync keeps it fresh after that. Waits for the whole index by default; pass `async: true` to start it in the daemon and return at once, then poll with `status: true` (see below). | | `list_projects` | List all indexed projects with node/edge counts. | | `delete_project` | Remove a project and all its graph data. | | `index_status` | Check indexing status of a project. | | `check_index_coverage` | Check whether exact paths or a scope are indexed and fresh. A clean result means no recorded gap, not proof of completeness. | +**Long indexes and client call deadlines.** A synchronous `index_repository` on a large +repository can take longer than an MCP client allows one tool call (some IDE clients give up +after a fixed deadline). When a client cancels or disconnects, the daemon cancels an index that +nobody else is waiting for, so retrying the same blocking call never finishes. Use the async +mode instead: + +1. `index_repository(repo_path="/abs/path", async: true)` starts the index in the daemon (or + joins the one already running for that project) and returns immediately with + `state` (`queued`/`running`). The job keeps running even if the client cancels, times out or + disconnects; only stopping the daemon ends it. +2. `index_repository(repo_path="/abs/path", status: true)` reports `state` + (`queued`, `running`, `cancelling`, `succeeded`, `failed`, `cancelled`), `started_at`, + `finished_at` and an `error` summary. Poll it until the state is `succeeded`, `failed` or + `cancelled`. Pass the same `repo_path` (and `name`, if the index call used one). + +`async` and `status` are exclusive; `async` does not apply to `cross-repo-intelligence`. Both +need the daemon-backed server (the default `codebase-memory-mcp` entry point). A temporary +daemon (started on demand rather than by `codebase-memory-mcp daemon start`) stops, and +cancels its jobs, when its last client disconnects. A connected MCP session keeps it alive, so +async from your editor works. A one-shot `codebase-memory-mcp cli index_repository --async` +that is the daemon's only client is refused with a clear error, because the job would die the +moment the command exits: run `codebase-memory-mcp daemon start` first, keep an MCP session +open, or call without `--async`. `status` works everywhere. When a synchronous call was cut +short, the next +`index_repository` or `status` call for that project carries a `notice` suggesting the async +mode. `index_status` keeps describing the published graph and its freshness. + ### Querying | Tool | Description | @@ -799,6 +826,7 @@ SQLite databases stored at `~/.cache/codebase-memory-mcp/`. Persists across rest |---------|-----| | `/mcp` doesn't show the server | Check `.mcp.json` path is absolute. Restart agent. Test: `echo '{}' \| /path/to/binary` should output JSON. | | `index_repository` fails | Pass absolute path: `index_repository(repo_path="/absolute/path")` | +| `index_repository` times out in the client | Start it with `async: true`, then poll with `status: true` (see [Indexing](#indexing)). | | `trace_path` returns 0 results | Use `search_graph(name_pattern=".*PartialName.*")` first to find the exact name. | | Queries return wrong project results | Add `project="name"` parameter. Use `list_projects` to see names. | | Binary not found after install | Add to PATH: `export PATH="$HOME/.local/bin:$PATH"` | diff --git a/src/daemon/application.c b/src/daemon/application.c index 748b8bb757..22c62029b4 100644 --- a/src/daemon/application.c +++ b/src/daemon/application.c @@ -10,6 +10,7 @@ #include "foundation/compat_thread.h" #include "foundation/log.h" #include "foundation/mem.h" +#include "foundation/mem_core.h" #include "foundation/platform.h" #include "foundation/secure_random.h" #include "foundation/sha256.h" @@ -32,6 +33,7 @@ #include #include #include +#include #ifdef _WIN32 #ifndef WIN32_LEAN_AND_MEAN @@ -127,6 +129,9 @@ struct cbm_daemon_application_session { bool update_notice_delivered; bool pending_background_initialize; bool pending_update_notice; + /* #2144: the current request arrived on the one-shot tool channel (the + * `cli` command path), whose client closes right after the reply. */ + bool one_shot_request; cbm_daemon_application_session_t *next; }; @@ -147,9 +152,32 @@ struct cbm_daemon_application_job { bool cancelled; bool cancel_requested; bool supervision_failed; + /* #2144: an async index_repository holds one subscriber reference on + * the application's behalf, released only at terminal publish, so the + * job outlives the request and session that started it. */ + bool async_owned; + bool worker_started; cbm_daemon_application_job_t *next; }; +/* #2144: the last known index outcome per project, kept after the terminal + * job itself is reaped so index_repository(status: true) can answer after + * completion. sync_cut_short records that a synchronous caller gave up (client + * cancel, deadline, disconnect) before its index finished; the next + * index_repository or status call for that project carries the async advice. + * This is job bookkeeping, not graph freshness (index_status owns that). */ +typedef struct cbm_daemon_application_index_record cbm_daemon_application_index_record_t; +struct cbm_daemon_application_index_record { + char *project_key; + const char *state; /* static string: queued/succeeded/failed/cancelled */ + bool async; + bool sync_cut_short; + char started_at[32]; + char finished_at[32]; + char *error; + cbm_daemon_application_index_record_t *next; +}; + /* A watcher-triggered physical job is owned by the exact live sessions that * currently subscribe to its project/root watch. The callback waiting for the * job is only a storage waiter; it is deliberately not an ownership @@ -177,6 +205,7 @@ struct cbm_daemon_application { cbm_daemon_application_job_t *jobs; cbm_daemon_application_watch_job_subscription_t *watch_job_subscriptions; cbm_daemon_application_mutation_t *mutations; + cbm_daemon_application_index_record_t *index_records; cbm_daemon_application_worker_ops_t worker_ops; cbm_daemon_application_update_ops_t update_ops; cbm_project_lock_manager_t *project_locks; @@ -556,6 +585,131 @@ static void application_job_free(cbm_daemon_application_job_t *job) { free(job); } +enum { APPLICATION_RECORD_ERROR_CAP = 512 }; + +static void application_utc_now(char out[32]) { + time_t now = time(NULL); + struct tm parts; + out[0] = '\0'; + if (cbm_gmtime_r(&now, &parts)) { + (void)strftime(out, 32, "%Y-%m-%dT%H:%M:%SZ", &parts); + } +} + +static cbm_daemon_application_index_record_t *application_index_record_find_locked( + cbm_daemon_application_t *application, const char *project_key) { + for (cbm_daemon_application_index_record_t *record = application->index_records; record; + record = record->next) { + if (strcmp(record->project_key, project_key) == 0) { + return record; + } + } + return NULL; +} + +/* Find or create; NULL only on allocation failure (bookkeeping is then + * skipped, never the index itself). */ +static cbm_daemon_application_index_record_t *application_index_record_get_locked( + cbm_daemon_application_t *application, const char *project_key) { + cbm_daemon_application_index_record_t *record = + application_index_record_find_locked(application, project_key); + if (record) { + return record; + } + record = cbm_calloc(CBM_MEM_CLASS_OTHER, sizeof(*record)); + if (!record || !(record->project_key = cbm_mem_strdup(CBM_MEM_CLASS_OTHER, project_key))) { + cbm_free(CBM_MEM_CLASS_OTHER, record); + return NULL; + } + record->state = "queued"; + record->next = application->index_records; + application->index_records = record; + return record; +} + +static void application_index_records_free(cbm_daemon_application_index_record_t *record) { + while (record) { + cbm_daemon_application_index_record_t *next = record->next; + cbm_free(CBM_MEM_CLASS_OTHER, record->project_key); + cbm_free(CBM_MEM_CLASS_OTHER, record->error); + cbm_free(CBM_MEM_CLASS_OTHER, record); + record = next; + } +} + +/* A new physical job for the project begins a fresh attempt record. The + * cut-short marker deliberately survives: it describes the caller, and the + * advice stays pending until a later call has carried it. */ +static void application_index_record_begin_locked(cbm_daemon_application_job_t *job) { + cbm_daemon_application_index_record_t *record = + application_index_record_get_locked(job->application, job->project_key); + if (!record) { + return; + } + record->state = "queued"; + record->async = false; + application_utc_now(record->started_at); + record->finished_at[0] = '\0'; + cbm_free(CBM_MEM_CLASS_OTHER, record->error); + record->error = NULL; +} + +/* The first content text of a tool result, bounded for a status summary. */ +static char *application_tool_result_summary(const char *response) { + yyjson_doc *doc = response ? yyjson_read(response, strlen(response), 0) : NULL; + yyjson_val *root = doc ? yyjson_doc_get_root(doc) : NULL; + yyjson_val *content = root && yyjson_is_obj(root) ? yyjson_obj_get(root, "content") : NULL; + yyjson_val *first = content && yyjson_is_arr(content) ? yyjson_arr_get_first(content) : NULL; + yyjson_val *text = first && yyjson_is_obj(first) ? yyjson_obj_get(first, "text") : NULL; + char *summary = NULL; + if (text && yyjson_is_str(text)) { + const char *value = yyjson_get_str(text); + size_t length = strlen(value); + if (length >= APPLICATION_RECORD_ERROR_CAP) { + length = APPLICATION_RECORD_ERROR_CAP - 1U; + /* Never split a UTF-8 sequence. */ + while (length > 0 && ((unsigned char)value[length] & 0xC0U) == 0x80U) { + length--; + } + } + summary = cbm_alloc(CBM_MEM_CLASS_OTHER, length + 1U); + if (summary) { + memcpy(summary, value, length); + summary[length] = '\0'; + } + } + yyjson_doc_free(doc); + return summary; +} + +static void application_index_record_finish_locked(cbm_daemon_application_job_t *job) { + cbm_daemon_application_index_record_t *record = + application_index_record_get_locked(job->application, job->project_key); + if (!record) { + return; + } + record->state = job->successful ? "succeeded" : job->cancelled ? "cancelled" : "failed"; + application_utc_now(record->finished_at); + cbm_free(CBM_MEM_CLASS_OTHER, record->error); + record->error = job->successful ? NULL : application_tool_result_summary(job->response); +} + +/* A synchronous caller stopped waiting before its index finished. */ +static void application_index_record_mark_cut_short_locked(cbm_daemon_application_t *application, + const char *project_key) { + cbm_daemon_application_index_record_t *record = + application_index_record_get_locked(application, project_key); + if (record) { + record->sync_cut_short = true; + } +} + +static void application_index_record_cut_short_locked(cbm_daemon_application_job_t *job) { + if (job && !job->terminal) { + application_index_record_mark_cut_short_locked(job->application, job->project_key); + } +} + /* Reap completed job threads only after every logical demand subscription and * watcher callback storage waiter has released the job. Exactly one caller * removes a job under the mutex. */ @@ -1127,6 +1281,7 @@ static application_attempt_status_t application_job_run_attempt(cbm_daemon_appli cbm_mutex_lock(&application->mutex); job->worker = worker; + job->worker_started = true; bool cancel_now = job->cancel_requested || application->stopping; cbm_mutex_unlock(&application->mutex); if (cancel_now) { @@ -1187,7 +1342,9 @@ static char *application_job_failure_response(const cbm_index_worker_result_t *r char message[1024]; if (result && result->cancellation_requested) { (void)snprintf(message, sizeof(message), - "index operation cancelled after its final owning session disconnected"); + "index operation cancelled before completion (its last waiting " + "client disconnected or the daemon stopped). %s", + CBM_MCP_INDEX_ASYNC_HINT); } else if (result && (result->supervision_failed || !result->tree_quiesced)) { (void)snprintf(message, sizeof(message), "index worker containment failed (%s); inspect log: %s", @@ -1495,6 +1652,12 @@ static void application_job_publish(cbm_daemon_application_job_t *job, !execution->last_result.cancellation_requested); job->terminal = true; job->thread_done = true; + application_index_record_finish_locked(job); + if (job->async_owned) { + /* The async reference ends with the job; terminal, so no cancel. */ + job->async_owned = false; + application_job_unsubscribe_locked(job); + } for (cbm_daemon_application_session_t *session = application->sessions; session; session = session->next) { if (session->auto_index_job != job || !session->auto_index_subscribed) { @@ -1733,6 +1896,8 @@ static cbm_daemon_application_job_t *application_job_subscribe_locked( application->jobs = job; if (application_job_thread_create(&job->thread, job) == 0) { job->thread_started = true; + /* Still under the mutex, so this precedes the thread's publish. */ + application_index_record_begin_locked(job); } else { /* The job was linked only so a concurrently started thread could * observe its reservation. No thread exists on this path, so roll the @@ -2286,7 +2451,8 @@ static char *application_job_wait_for_session(cbm_daemon_application_session_t * cbm_mutex_lock(&application->mutex); if (session->active_job != job || !session->active_job_subscribed) { cbm_mutex_unlock(&application->mutex); - return cbm_mcp_text_result("index operation cancelled for this session", true); + return cbm_mcp_text_result( + "index operation cancelled for this session. " CBM_MCP_INDEX_ASYNC_HINT, true); } if (job->terminal) { char *response = job->response ? strdup(job->response) : NULL; @@ -2303,12 +2469,20 @@ static char *application_job_wait_for_session(cbm_daemon_application_session_t * } } -static char *application_index_execute(void *context, const char *root_path, - const char *args_json) { - cbm_daemon_application_session_t *session = context; - if (!session || !root_path || !args_json) { - return NULL; +static const char *application_index_admission_error(application_job_subscribe_status_t status) { + if (status == APPLICATION_JOB_SUBSCRIBE_OPTIONS_CONFLICT) { + return "another index operation for this project is active with different options"; } + if (status == APPLICATION_JOB_SUBSCRIBE_ALLOCATION_FAILED) { + return "daemon index coordinator could not allocate an index job"; + } + return "daemon index coordinator is stopping or unavailable"; +} + +/* The synchronous contract: wait for the whole index (queueing behind the + * physical job limit); a cancelled request drops only its own reference. */ +static char *application_index_execute_sync(cbm_daemon_application_session_t *session, + const char *root_path, const char *args_json) { char *project_key = application_index_project_key(root_path, args_json); if (!project_key) { return cbm_mcp_text_result("failed to derive index project identity", true); @@ -2332,28 +2506,28 @@ static char *application_index_execute(void *context, const char *root_path, memory_order_release); cbm_mutex_lock(&session->application->mutex); bool queued_cancelled = application_request_cancelled_locked(session); + if (queued_cancelled) { + application_index_record_mark_cut_short_locked(session->application, project_key); + } cbm_mutex_unlock(&session->application->mutex); if (queued_cancelled) { free(project_key); - return cbm_mcp_text_result("index operation cancelled for this session", true); + return cbm_mcp_text_result( + "index operation cancelled for this session. " CBM_MCP_INDEX_ASYNC_HINT, true); } cbm_usleep(APPLICATION_JOB_POLL_US); } free(project_key); if (!job) { - const char *message = "daemon index coordinator is stopping or unavailable"; - if (subscribe_status == APPLICATION_JOB_SUBSCRIBE_OPTIONS_CONFLICT) { - message = "another index operation for this project is active with different options"; - } else if (subscribe_status == APPLICATION_JOB_SUBSCRIBE_ALLOCATION_FAILED) { - message = "daemon index coordinator could not allocate an index job"; - } - return cbm_mcp_text_result(message, true); + return cbm_mcp_text_result(application_index_admission_error(subscribe_status), true); } cbm_mutex_lock(&session->application->mutex); if (application_request_cancelled_locked(session)) { + application_index_record_cut_short_locked(job); application_job_unsubscribe_locked(job); cbm_mutex_unlock(&session->application->mutex); - return cbm_mcp_text_result("index operation cancelled for this session", true); + return cbm_mcp_text_result( + "index operation cancelled for this session. " CBM_MCP_INDEX_ASYNC_HINT, true); } if (session->active_job) { application_job_unsubscribe_locked(job); @@ -2366,6 +2540,236 @@ static char *application_index_execute(void *context, const char *root_path, return application_job_wait_for_session(session, job); } +/* Take (and clear) the project's pending cut-short advice. Taken at the START + * of a synchronous call: if this call is cut short too, its own cancel path + * sets the marker again, so the advice is never consumed by a reply nobody + * reads. */ +static bool application_index_take_cut_short(cbm_daemon_application_t *application, + const char *project_key) { + cbm_mutex_lock(&application->mutex); + cbm_daemon_application_index_record_t *record = + application_index_record_find_locked(application, project_key); + bool pending = record && record->sync_cut_short; + if (pending) { + record->sync_cut_short = false; + } + cbm_mutex_unlock(&application->mutex); + return pending; +} + +#define APPLICATION_CUT_SHORT_NOTICE \ + "The previous synchronous index_repository call for this project was cut short before the " \ + "index finished (client cancel, call deadline or disconnect). " CBM_MCP_INDEX_ASYNC_HINT + +static char *application_index_run_sync(cbm_daemon_application_session_t *session, + const char *root_path, const char *args_json) { + char *project_key = application_index_project_key(root_path, args_json); + bool notice = + project_key && application_index_take_cut_short(session->application, project_key); + safe_free(project_key); + char *response = application_index_execute_sync(session, root_path, args_json); + return notice ? cbm_mcp_tool_result_add_notice(response, APPLICATION_CUT_SHORT_NOTICE) + : response; +} + +static char *application_index_async_response(const char *project_key, const char *state, + bool joined) { + yyjson_mut_doc *doc = yyjson_mut_doc_new(NULL); + yyjson_mut_val *root = doc ? yyjson_mut_obj(doc) : NULL; + char *json = NULL; + if (root) { + yyjson_mut_doc_set_root(doc, root); + yyjson_mut_obj_add_strcpy(doc, root, "project", project_key); + yyjson_mut_obj_add_str(doc, root, "state", state); + yyjson_mut_obj_add_bool(doc, root, "async", true); + yyjson_mut_obj_add_bool(doc, root, "joined", joined); + yyjson_mut_obj_add_str(doc, root, "poll", + "call index_repository with the same repo_path and status: true " + "until state is succeeded, failed or cancelled"); + json = yyjson_mut_write(doc, 0, NULL); + } + yyjson_mut_doc_free(doc); + char *result = cbm_mcp_text_result(json ? json : "{}", false); + safe_free(json); + return result; +} + +/* #2144 user decision: refuse async where the job would certainly die with its + * requester. A temporary (non-permanent) daemon stops and cancels every job + * when its final live session ends. The one-shot tool channel (`cli ...`) + * closes its session right after the reply, so if it is the only live session + * the job is doomed. A long-lived MCP session - the IDE case this exists for - + * stays connected and is a valid host even when it is the only one, as is any + * other live session and every permanent (`daemon start`) generation. */ +static bool application_async_host_missing(cbm_daemon_application_session_t *session) { + cbm_daemon_application_t *application = session->application; + cbm_mutex_lock(&application->mutex); + bool other_live_session = false; + for (cbm_daemon_application_session_t *other = application->sessions; + other && !other_live_session; other = other->next) { + other_live_session = other != session && !other->session_cancelled; + } + bool missing = !application->permanent && session->one_shot_request && !other_live_session; + cbm_mutex_unlock(&application->mutex); + return missing; +} + +/* index_repository(async: true) (#2144): start or join the project's job and + * return at once. The application keeps one subscriber reference until the + * terminal publish, so neither this request returning, nor its client + * cancelling, nor this session closing cancels the job: only daemon shutdown + * (or the final live session of a non-permanent generation) does. Never + * queues behind the physical job limit: a full daemon answers busy. */ +#define APPLICATION_ASYNC_NEEDS_HOST \ + "async indexing needs a daemon that outlives this command: this one-shot call is the only " \ + "client of a temporary daemon, which stops (cancelling the job) when the command exits. Run " \ + "`codebase-memory-mcp daemon start` (or keep an MCP session open), then retry; or call " \ + "without async" + +static char *application_index_start_async(cbm_daemon_application_session_t *session, + const char *root_path, const char *args_json) { + cbm_daemon_application_t *application = session->application; + if (application_async_host_missing(session)) { + return cbm_mcp_text_result(APPLICATION_ASYNC_NEEDS_HOST, true); + } + char *project_key = application_index_project_key(root_path, args_json); + if (!project_key) { + return cbm_mcp_text_result("failed to derive index project identity", true); + } + application_jobs_reap_completed(application); + application_job_subscribe_status_t status = APPLICATION_JOB_SUBSCRIBE_UNAVAILABLE; + cbm_mutex_lock(&application->mutex); + /* Subscribe and claim in one critical section: the job cannot publish in + * between, so the async reference is always released by its publish. */ + cbm_daemon_application_job_t *job = + application_job_subscribe_locked(application, project_key, root_path, args_json, &status); + bool joined = false; + const char *state = "queued"; + if (job) { + joined = job->async_owned || job->subscribers > 1; + if (job->async_owned) { + application_job_unsubscribe_locked(job); /* one async reference per job */ + } else { + job->async_owned = true; + } + state = job->worker_started ? "running" : "queued"; + cbm_daemon_application_index_record_t *record = + application_index_record_get_locked(application, project_key); + if (record) { + record->async = true; + record->sync_cut_short = false; /* the advice was taken */ + } + } + cbm_mutex_unlock(&application->mutex); + char *result = NULL; + if (job) { + cbm_log_info("daemon.index.async", "project", project_key, "joined", + joined ? "true" : "false"); + result = application_index_async_response(project_key, state, joined); + } else if (status == APPLICATION_JOB_SUBSCRIBE_BUSY) { + result = cbm_mcp_text_result( + "daemon index job limit reached; async requests do not queue. Retry async later, or " + "call without async to wait for a slot", + true); + } else if (status == APPLICATION_JOB_SUBSCRIBE_CANCELLING) { + result = cbm_mcp_text_result( + "the previous index of this project is still being cancelled; retry async shortly", + true); + } else { + result = cbm_mcp_text_result(application_index_admission_error(status), true); + } + safe_free(project_key); + return result; +} + +static char *application_index_execute(void *context, const char *root_path, const char *args_json, + bool async) { + cbm_daemon_application_session_t *session = context; + if (!session || !root_path || !args_json) { + return NULL; + } + return async ? application_index_start_async(session, root_path, args_json) + : application_index_run_sync(session, root_path, args_json); +} + +static void application_status_add_record(yyjson_mut_doc *doc, yyjson_mut_val *root, + const cbm_daemon_application_index_record_t *record) { + yyjson_mut_obj_add_bool(doc, root, "async", record->async); + if (record->started_at[0]) { + yyjson_mut_obj_add_strcpy(doc, root, "started_at", record->started_at); + } + if (record->finished_at[0]) { + yyjson_mut_obj_add_strcpy(doc, root, "finished_at", record->finished_at); + } + if (record->error) { + yyjson_mut_obj_add_strcpy(doc, root, "error", record->error); + } + if (record->sync_cut_short) { + yyjson_mut_obj_add_str(doc, root, "notice", APPLICATION_CUT_SHORT_NOTICE); + } +} + +/* index_repository(status: true) (#2144). A live job reports queued, running + * or cancelling; otherwise the last recorded outcome. A project with neither a + * job record nor a database is an error; one indexed before this daemon + * generation reports idle and points at index_status for graph freshness. */ +static char *application_index_status(void *context, const char *project) { + cbm_daemon_application_session_t *session = context; + if (!session || !project) { + return NULL; + } + cbm_daemon_application_t *application = session->application; + bool db_exists = application_regular_db_exists(project); + yyjson_mut_doc *doc = yyjson_mut_doc_new(NULL); + yyjson_mut_val *root = doc ? yyjson_mut_obj(doc) : NULL; + if (!root) { + yyjson_mut_doc_free(doc); + return NULL; + } + yyjson_mut_doc_set_root(doc, root); + yyjson_mut_obj_add_strcpy(doc, root, "project", project); + cbm_mutex_lock(&application->mutex); + cbm_daemon_application_job_t *job = application_find_active_job_locked(application, project); + cbm_daemon_application_index_record_t *record = + application_index_record_find_locked(application, project); + bool known = job || record || db_exists; + if (job) { + yyjson_mut_obj_add_str(doc, root, "state", + job->cancel_requested ? "cancelling" + : job->worker_started ? "running" + : "queued"); + } else if (record) { + yyjson_mut_obj_add_str(doc, root, "state", record->state); + } else { + yyjson_mut_obj_add_str(doc, root, "state", "idle"); + yyjson_mut_obj_add_str(doc, root, "detail", + "no index job has run for this project in this daemon generation; " + "index_status describes the published graph"); + } + if (record) { + application_status_add_record(doc, root, record); + } + if (job && !job->cancel_requested) { + yyjson_mut_obj_add_str(doc, root, "poll", + "call again with status: true until state is succeeded, failed or " + "cancelled"); + } + cbm_mutex_unlock(&application->mutex); + char *json = known ? yyjson_mut_write(doc, 0, NULL) : NULL; + yyjson_mut_doc_free(doc); + if (!known) { + char message[APPLICATION_PATH_CAP]; + (void)snprintf(message, sizeof(message), + "no index job or index is known for project '%s'; start one with " + "index_repository(repo_path=..., async: true)", + project); + return cbm_mcp_text_result(message, true); + } + char *result = cbm_mcp_text_result(json ? json : "{}", false); + safe_free(json); + return result; +} + static cbm_daemon_runtime_application_session_t *application_session_open( void *context, cbm_daemon_client_id_t client_id, uint64_t authenticated_process_id) { cbm_daemon_application_t *application = context; @@ -2386,6 +2790,7 @@ static cbm_daemon_runtime_application_session_t *application_session_open( cbm_mcp_server_set_background_tasks(session->mcp, false); cbm_mcp_server_set_config(session->mcp, application->config); cbm_mcp_server_set_index_executor(session->mcp, application_index_execute, session); + cbm_mcp_server_set_index_status_provider(session->mcp, application_index_status, session); cbm_mcp_server_set_project_mutation_guard(session->mcp, application_session_mutation_begin, application_session_mutation_end, session); cbm_mcp_server_set_project_mutation_try_guard(session->mcp, @@ -2617,7 +3022,9 @@ static cbm_daemon_runtime_application_status_t application_tool_request( free(args); return CBM_DAEMON_RUNTIME_APPLICATION_REJECTED; } + session->one_shot_request = true; /* only this request thread reads it */ char *response = cbm_mcp_handle_tool(session->mcp, tool, args); + session->one_shot_request = false; free(tool); free(args); if (!response) { @@ -2849,6 +3256,7 @@ static void application_request_cancel(void *context, cbm_daemon_application_job_t *job = session->active_job; session->active_job = NULL; session->active_job_subscribed = false; + application_index_record_cut_short_locked(job); application_job_unsubscribe_locked(job); } if (active_match) { @@ -2881,6 +3289,7 @@ static void application_session_cancel(void *context, cbm_daemon_application_job_t *job = session->active_job; session->active_job = NULL; session->active_job_subscribed = false; + application_index_record_cut_short_locked(job); application_job_unsubscribe_locked(job); } join_auto_index = application_auto_index_release_locked(session); @@ -2937,6 +3346,7 @@ static void application_session_close(void *context, *cursor = session->next; } if (session->active_job && session->active_job_subscribed) { + application_index_record_cut_short_locked(session->active_job); application_job_unsubscribe_locked(session->active_job); session->active_job = NULL; session->active_job_subscribed = false; @@ -3148,7 +3558,10 @@ bool cbm_daemon_application_free_with_timeout(cbm_daemon_application_t *applicat application->watch_job_subscriptions = NULL; cbm_daemon_application_mutation_t *mutations = application->mutations; application->mutations = NULL; + cbm_daemon_application_index_record_t *index_records = application->index_records; + application->index_records = NULL; cbm_mutex_unlock(&application->mutex); + application_index_records_free(index_records); while (sessions) { cbm_daemon_application_session_t *next = sessions->next; cbm_mcp_server_free(sessions->mcp); diff --git a/src/daemon/frontend.c b/src/daemon/frontend.c index 2dd9869011..cbe4b282e3 100644 --- a/src/daemon/frontend.c +++ b/src/daemon/frontend.c @@ -382,8 +382,26 @@ static bool frontend_write_response(FILE *out, const uint8_t *response, uint32_t return fflush(out) == 0 && written; } +const char *cbm_daemon_frontend_cancelled_error_message(const char *request_message) { + static const char plain[] = "Request cancelled"; + static const char index_call[] = "Request cancelled. " CBM_MCP_INDEX_ASYNC_HINT; + cbm_jsonrpc_request_t request = {0}; + if (!request_message || cbm_jsonrpc_parse(request_message, &request) != 0) { + return plain; + } + char *tool = request.method && strcmp(request.method, "tools/call") == 0 && request.params_raw + ? cbm_mcp_get_string_arg(request.params_raw, "name") + : NULL; + bool index_repository = tool && strcmp(tool, "index_repository") == 0; + safe_free(tool); + cbm_jsonrpc_request_free(&request); + return index_repository ? index_call : plain; +} + static bool frontend_write_cancelled_response(FILE *out, const frontend_item_t *item) { - static const char cancelled_error[] = "{\"code\":-32800,\"message\":\"Request cancelled\"}"; + char cancelled_error[512]; + (void)snprintf(cancelled_error, sizeof(cancelled_error), "{\"code\":-32800,\"message\":\"%s\"}", + cbm_daemon_frontend_cancelled_error_message(item->message)); cbm_jsonrpc_response_t response = { .id = item->id, .id_str = item->id_str, diff --git a/src/daemon/frontend.h b/src/daemon/frontend.h index f23395774a..9d172606e0 100644 --- a/src/daemon/frontend.h +++ b/src/daemon/frontend.h @@ -29,6 +29,12 @@ bool cbm_daemon_frontend_is_cancellation_notification(const char *message); bool cbm_daemon_frontend_cancellation_matches_request(const char *message, int64_t active_id, const char *active_id_str); +/* The JSON-RPC error message answered for a cancelled request (#2144). A + * cancelled index_repository call is usually a client deadline on a long + * index, so its reply names the async alternative; every other request keeps + * the plain message. Returns a static string free of JSON metacharacters. */ +const char *cbm_daemon_frontend_cancelled_error_message(const char *request_message); + /* Start a temporary observer for a one-shot local CLI command or physical * supervised worker. manager and cancel_context are borrowed until stop. On * maintenance intent the observer invokes cancel once, permits a fixed bounded diff --git a/src/mcp/mcp.c b/src/mcp/mcp.c index d6c7bd4686..0d8cbf910f 100644 --- a/src/mcp/mcp.c +++ b/src/mcp/mcp.c @@ -428,6 +428,60 @@ char *cbm_mcp_text_result(const char *text, bool is_error) { return out; } +/* Rebuild the payload text with the notice attached; NULL when the payload + * cannot be rewritten. The result is a memory-core block (class other) on + * both branches, so the caller has one release path. */ +static char *tool_result_payload_with_notice(const char *payload, const char *notice) { + yyjson_doc *payload_doc = yyjson_read(payload, strlen(payload), 0); + yyjson_val *payload_root = payload_doc ? yyjson_doc_get_root(payload_doc) : NULL; + char *rewritten = NULL; + if (payload_root && yyjson_is_obj(payload_root)) { + yyjson_mut_doc *copy = yyjson_doc_mut_copy(payload_doc, NULL); + yyjson_mut_val *copy_root = copy ? yyjson_mut_doc_get_root(copy) : NULL; + char *written = copy_root && yyjson_mut_obj_add_strcpy(copy, copy_root, "notice", notice) + ? yy_doc_to_str(copy) + : NULL; + rewritten = written ? cbm_mem_strdup(CBM_MEM_CLASS_OTHER, written) : NULL; + safe_free(written); /* yyjson's writer block, not a core block */ + yyjson_mut_doc_free(copy); + } else { + size_t payload_len = strlen(payload); + size_t notice_len = strlen(notice); + rewritten = cbm_alloc(CBM_MEM_CLASS_OTHER, payload_len + notice_len + 3U); + if (rewritten) { + memcpy(rewritten, payload, payload_len); + memcpy(rewritten + payload_len, "\n\n", 2U); + memcpy(rewritten + payload_len + 2U, notice, notice_len + 1U); + } + } + yyjson_doc_free(payload_doc); + return rewritten; +} + +char *cbm_mcp_tool_result_add_notice(char *result, const char *notice) { + if (!result || !notice || !notice[0]) { + return result; + } + yyjson_doc *doc = yyjson_read(result, strlen(result), 0); + yyjson_val *root = doc ? yyjson_doc_get_root(doc) : NULL; + yyjson_val *content = root && yyjson_is_obj(root) ? yyjson_obj_get(root, "content") : NULL; + yyjson_val *first = content && yyjson_is_arr(content) ? yyjson_arr_get_first(content) : NULL; + yyjson_val *text = first && yyjson_is_obj(first) ? yyjson_obj_get(first, "text") : NULL; + yyjson_val *is_error = root && yyjson_is_obj(root) ? yyjson_obj_get(root, "isError") : NULL; + char *payload = text && yyjson_is_str(text) + ? tool_result_payload_with_notice(yyjson_get_str(text), notice) + : NULL; + char *rebuilt = + payload ? cbm_mcp_text_result(payload, is_error && yyjson_get_bool(is_error)) : NULL; + cbm_free(CBM_MEM_CLASS_OTHER, payload); + yyjson_doc_free(doc); + if (!rebuilt) { + return result; + } + safe_free(result); /* the caller's tool result, not a core block */ + return rebuilt; +} + bool cbm_mcp_cancel_request_matches(const char *params_json, int64_t active_id, const char *active_id_str) { if (!params_json) { @@ -466,7 +520,11 @@ typedef struct { static const tool_def_t TOOLS[] = { {"index_repository", "Index a repository. full/moderate add semantics; fast omits them; cross-repo-intelligence " - "links services. Reports coverage gaps.", + "links services. Reports coverage gaps. Waits for the whole index by default; a large repo " + "can exceed a client's per-call deadline. async:true starts (or joins) the index in the " + "daemon and returns at once; poll with status:true (same repo_path) until state is " + "succeeded, failed or cancelled. async needs a daemon that outlives the call: an MCP session " + "or `daemon start`.", "{\"type\":\"object\",\"properties\":{\"repo_path\":{\"type\":\"string\",\"description\":" "\"Repository path\"}," "\"mode\":{\"type\":\"string\"," @@ -478,7 +536,14 @@ static const tool_def_t TOOLS[] = { "\"name\":{\"type\":\"string\",\"description\":" "\"Name override; Non-ASCII bytes are encoded; unsafe characters normalized.\"}," "\"persistence\":{\"type\":\"boolean\",\"default\":false,\"description\":" - "\"Write .codebase-memory/graph.db.zst.\"}" + "\"Write .codebase-memory/graph.db.zst.\"}," + "\"async\":{\"type\":\"boolean\",\"default\":false,\"description\":" + "\"Start the index in the background and return immediately; it keeps running if this " + "call is cancelled or times out. Not with status or cross-repo-intelligence. Refused for a " + "one-shot cli call that is the only client of a temporary daemon (run daemon start).\"}," + "\"status\":{\"type\":\"boolean\",\"default\":false,\"description\":" + "\"Report the running or last index job (queued/running/succeeded/failed/cancelled) " + "instead of indexing. Pass the same repo_path (and name, if any) as the index call.\"}" "},\"required\":[\"repo_path\"]}"}, {"search_graph", @@ -1660,6 +1725,8 @@ struct cbm_mcp_server { struct cbm_config *config; /* external config ref (not owned) */ cbm_mcp_index_executor_fn index_executor; void *index_executor_context; + cbm_mcp_index_status_fn index_status_provider; /* #2144; NULL outside the daemon */ + void *index_status_context; cbm_proc_log_cb index_log_callback; void *index_log_context; cbm_mcp_project_mutation_begin_fn mutation_begin; @@ -1815,6 +1882,14 @@ void cbm_mcp_server_set_index_executor(cbm_mcp_server_t *srv, cbm_mcp_index_exec } } +void cbm_mcp_server_set_index_status_provider(cbm_mcp_server_t *srv, + cbm_mcp_index_status_fn provider, void *context) { + if (srv) { + srv->index_status_provider = provider; + srv->index_status_context = context; + } +} + void cbm_mcp_server_set_index_log_callback(cbm_mcp_server_t *srv, cbm_proc_log_cb callback, void *context) { if (srv) { @@ -10468,7 +10543,8 @@ static char *build_worker_unsafe_terminal_response(const char *args, cbm_proc_ou yyjson_mut_obj_add_str( doc, root, "hint", cancellation_requested - ? "Indexing worker was cancelled. No in-process retry was started." + ? "Indexing worker was cancelled. No in-process retry was " + "started. " CBM_MCP_INDEX_ASYNC_HINT : "Indexing worker process-tree containment failed. No in-process retry was " "started; inspect daemon logs."); if (repo_path) { @@ -11044,6 +11120,16 @@ static char *index_args_with_repo_path(const char *args, const char *canonical_r while (yyjson_mut_obj_get(copy_root, "_cbm_index_policy")) { (void)yyjson_mut_obj_remove_key(copy_root, "_cbm_index_policy"); } + /* #2144: async/status choose how the caller waits, not what is indexed. + * They never reach the worker, and an async request coalesces with an + * identical synchronous one instead of being refused as an options + * conflict. */ + static const char *const call_mode_keys[] = {"async", "status"}; + for (size_t i = 0; i < sizeof(call_mode_keys) / sizeof(call_mode_keys[0]); i++) { + while (yyjson_mut_obj_get(copy_root, call_mode_keys[i])) { + (void)yyjson_mut_obj_remove_key(copy_root, call_mode_keys[i]); + } + } if (!yyjson_mut_obj_add_strcpy(copy, copy_root, "repo_path", canonical_repo_path) || !cbm_mcp_index_policy_add_to_args(copy, copy_root, policy)) { yyjson_mut_doc_free(copy); @@ -11107,7 +11193,107 @@ static void index_args_free(char *repo_path, char *mode_str, char *name_override free(name_override); } +/* #2144: async and status are strict booleans. A present non-boolean value is + * refused rather than read as false, so "async":"true" never silently blocks + * for the whole index. */ +static bool index_call_mode_arg(const char *args, const char *key, bool *value_out) { + *value_out = false; + yyjson_doc *doc = args ? yyjson_read(args, strlen(args), 0) : NULL; + yyjson_val *root = doc ? yyjson_doc_get_root(doc) : NULL; + yyjson_val *value = root && yyjson_is_obj(root) ? yyjson_obj_get(root, key) : NULL; + bool valid = !value || yyjson_is_bool(value); + *value_out = valid && value && yyjson_get_bool(value); + yyjson_doc_free(doc); + return valid; +} + +/* The project key a status query names: an explicit name override wins (the + * key the index call used), then repo_path resolved exactly as for indexing + * (session root, canonical form, workspace boundary), then a known project + * alias. NULL with *error_out set when none identifies a project. */ +static char *index_status_project_key(cbm_mcp_server_t *srv, const char *args, char *error_out, + size_t error_size) { + char *name = cbm_mcp_get_string_arg(args, "name"); + if (name && name[0]) { + char *key = cbm_project_name_from_path(name); + safe_free(name); + return key; + } + safe_free(name); + char *repo_path = cbm_mcp_get_string_arg(args, "repo_path"); + if (!repo_path) { + char *project = get_project_arg(args); + if (!project) { + (void)snprintf(error_out, error_size, + "status needs the repo_path (or name) the index call used"); + } + return project; + } + cbm_normalize_path_sep(repo_path); + if (!resolve_session_repo_path(srv, &repo_path)) { + safe_free(repo_path); + (void)snprintf(error_out, error_size, "failed to resolve repo_path"); + return NULL; + } + repo_path = canonicalize_repo_path_if_exists(repo_path); + const char *allowed_root = + srv->allowed_root_policy_set ? srv->allowed_root : getenv("CBM_ALLOWED_ROOT"); + if (repo_path && repo_path[0] && + !cbm_workspace_root_allowed(repo_path, cbm_workspace_home_dir(), cbm_workspace_cache_dir(), + allowed_root, error_out, error_size)) { + safe_free(repo_path); + return NULL; + } + char *key = repo_path ? cbm_project_name_from_path(repo_path) : NULL; + safe_free(repo_path); + if (!key) { + (void)snprintf(error_out, error_size, "could not resolve index project name"); + } + return key; +} + +/* index_repository(status: true): report the project's running or most recent + * index job without starting one. Job state lives in the daemon's job + * registry; a server without one (CLI index worker, embedder) has no job that + * could outlive a call, so it refuses instead of guessing. */ +static char *handle_index_repository_status(cbm_mcp_server_t *srv, const char *args) { + if (!srv->index_status_provider) { + return cbm_mcp_text_result("index_repository status needs the daemon-backed MCP server; " + "this in-process server tracks no index jobs", + true); + } + char error[CBM_SZ_1K] = {0}; + char *project = index_status_project_key(srv, args, error, sizeof(error)); + if (!project) { + return cbm_mcp_text_result(error[0] ? error : "could not resolve index project name", true); + } + char *result = srv->index_status_provider(srv->index_status_context, project); + safe_free(project); + return result ? result : cbm_mcp_text_result("index job status is unavailable", true); +} + static char *handle_index_repository(cbm_mcp_server_t *srv, const char *args) { + bool async_mode = false; + bool status_mode = false; + if (!index_call_mode_arg(args, "async", &async_mode) || + !index_call_mode_arg(args, "status", &status_mode)) { + return cbm_mcp_text_result("async and status must be booleans", true); + } + if (async_mode && status_mode) { + return cbm_mcp_text_result( + "async and status are exclusive: start with async: true, then poll with status: true", + true); + } + if (status_mode) { + return handle_index_repository_status(srv, args); + } + if (async_mode && !srv->index_executor) { + return cbm_mcp_text_result("index_repository async needs the daemon-backed MCP server: " + "this in-process server has no process that outlives the call; " + "call without async", + true); + } + char *repo_path = cbm_mcp_get_string_arg(args, "repo_path"); char *mode_str = cbm_mcp_get_string_arg(args, "mode"); char *name_override = cbm_mcp_get_string_arg(args, "name"); @@ -11149,6 +11335,12 @@ static char *handle_index_repository(cbm_mcp_server_t *srv, const char *args) { } if (mode_str && strcmp(mode_str, "cross-repo-intelligence") == 0) { + if (async_mode) { + index_args_free(repo_path, mode_str, name_override); + return cbm_mcp_text_result( + "async is not supported for mode cross-repo-intelligence (it runs no index job)", + true); + } char *result = handle_cross_repo_mode(srv, repo_path, name_override, args); index_args_free(repo_path, mode_str, name_override); return result; @@ -11165,9 +11357,9 @@ static char *handle_index_repository(cbm_mcp_server_t *srv, const char *args) { * registry only after path canonicalization and workspace authorization. */ if (srv->index_executor) { char *worker_args = index_args_with_repo_path(args, repo_path, &resource_policy); - char *coordinated = - worker_args ? srv->index_executor(srv->index_executor_context, repo_path, worker_args) - : NULL; + char *coordinated = worker_args ? srv->index_executor(srv->index_executor_context, + repo_path, worker_args, async_mode) + : NULL; free(worker_args); index_args_free(repo_path, mode_str, name_override); return coordinated ? coordinated @@ -11223,7 +11415,8 @@ static char *handle_index_repository(cbm_mcp_server_t *srv, const char *args) { mcp_project_mutation_end(srv, mutation_project); free(mutation_project); index_args_free(repo_path, mode_str, name_override); - return cbm_mcp_text_result("index operation cancelled for this request", true); + return cbm_mcp_text_result( + "index operation cancelled for this request. " CBM_MCP_INDEX_ASYNC_HINT, true); } cbm_index_mode_t mode = CBM_MODE_FULL; diff --git a/src/mcp/mcp.h b/src/mcp/mcp.h index 09befb7809..31ffd69b04 100644 --- a/src/mcp/mcp.h +++ b/src/mcp/mcp.h @@ -56,6 +56,22 @@ char *cbm_jsonrpc_format_error(int64_t id, int code, const char *message); /* Format an MCP tool result with text content. Returns heap-allocated JSON. */ char *cbm_mcp_text_result(const char *text, bool is_error); +/* Attach a notice to a complete MCP tool result: a JSON-object payload gains a + * "notice" key (content text and structuredContent stay in step), a text + * payload gains a trailing paragraph. Takes + * ownership of result and returns the rewritten result, or result unchanged + * when it is not a tool-result object (never NULL for a non-NULL result). */ +char *cbm_mcp_tool_result_add_notice(char *result, const char *notice); + +/* #2144: the advice every cut-short synchronous index_repository surfaces. A + * long index can outlive a client's per-call deadline; the async flag lets the + * daemon finish it while the client polls. */ +#define CBM_MCP_INDEX_ASYNC_HINT \ + "Long indexes can exceed an MCP client's per-call deadline: retry with " \ + "index_repository(repo_path=..., async: true), then poll " \ + "index_repository(repo_path=..., status: true) until state is succeeded, failed or " \ + "cancelled." + /* Return true when notifications/cancelled params target the active request. */ bool cbm_mcp_cancel_request_matches(const char *params_json, int64_t active_id, const char *active_id_str); @@ -128,9 +144,17 @@ bool cbm_mcp_tool_profile_allows_http(cbm_mcp_tool_profile_t profile); * already authorized against this session; args_json contains that canonical * path plus all caller options. The callback returns a complete malloc-owned * MCP tool result. A NULL result is reported as an error and never degrades to - * an uncoordinated in-process index. */ + * an uncoordinated in-process index. async (#2144) asks the executor to start + * or join the project's job and return at once; the job must then outlive the + * request and its session. */ typedef char *(*cbm_mcp_index_executor_fn)(void *context, const char *repo_path, - const char *args_json); + const char *args_json, bool async); + +/* Optional daemon-owned index job status (#2144): the state of the running or + * most recent index job for one project key. Returns a malloc-owned tool + * result. Servers without one (in-process CLI workers, embedders) refuse + * index_repository's async and status modes. */ +typedef char *(*cbm_mcp_index_status_fn)(void *context, const char *project); /* Daemon-owned exclusive lease for operations that mutate a published project * database. begin may wait, but must remain cancellable by its session owner; @@ -180,6 +204,9 @@ void cbm_mcp_server_set_background_tasks(cbm_mcp_server_t *srv, bool enabled); void cbm_mcp_server_set_index_executor(cbm_mcp_server_t *srv, cbm_mcp_index_executor_fn executor, void *context); +void cbm_mcp_server_set_index_status_provider(cbm_mcp_server_t *srv, + cbm_mcp_index_status_fn provider, void *context); + /* Relay supervised worker logs to one local request (for CLI progress). The * callback/context are borrowed until the synchronous tool call returns. */ void cbm_mcp_server_set_index_log_callback(cbm_mcp_server_t *srv, cbm_proc_log_cb callback, diff --git a/src/ui/http_server.c b/src/ui/http_server.c index f47ac59bee..e4fe6de3e0 100644 --- a/src/ui/http_server.c +++ b/src/ui/http_server.c @@ -2065,10 +2065,11 @@ static void dispatch_request(cbm_http_server_t *srv, cbm_http_conn_t *c, /* ── Public API ───────────────────────────────────────────────── */ static char *http_read_only_index_rejected(void *context, const char *repo_path, - const char *args_json) { + const char *args_json, bool async) { (void)context; (void)repo_path; (void)args_json; + (void)async; return cbm_mcp_text_result("UI RPC indexing is disabled; use the coordinated /api/index route", true); } diff --git a/tests/test_daemon_application.c b/tests/test_daemon_application.c index 0108e8ec50..d2b479a395 100644 --- a/tests/test_daemon_application.c +++ b/tests/test_daemon_application.c @@ -5634,6 +5634,495 @@ TEST(daemon_application_index_args_compare_repo_path_separator_equivalently) { PASS(); } +/* ── #2144: async index_repository + status polling ───────────────────────── + * + * A client with a per-call deadline (Copilot for IntelliJ) gives up on a long + * synchronous index; its cancel makes the daemon drop the job's last + * subscriber and cancel the worker, so every retry starts over and a large + * repository never finishes. async:true starts the job detached from the + * request and session; status:true polls it. Every wait below is a bounded + * observation of a held fake worker (never a timing assertion). */ + +typedef struct { + app_fake_worker_context_t fake; + cbm_daemon_application_worker_ops_t ops; + cbm_daemon_application_t *application; + cbm_daemon_runtime_application_callbacks_t callbacks; + char root[APP_TEST_PATH_CAP]; + char *project; + uint8_t *context; + uint32_t context_length; +} app_async_fixture_t; + +static bool app_async_fixture_init(app_async_fixture_t *fixture, const char *label, + bool permanent) { + memset(fixture, 0, sizeof(*fixture)); + app_fake_worker_context_init(&fixture->fake); + fixture->ops = (cbm_daemon_application_worker_ops_t){ + .context = &fixture->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 = {.worker_ops = &fixture->ops}; + fixture->application = cbm_daemon_application_new(&config); + if (!fixture->application) { + return false; + } + /* permanent mirrors `daemon start`; false is the temporary generation an + * MCP client or one-shot command spawns. */ + cbm_daemon_application_set_permanent(fixture->application, permanent); + fixture->callbacks = cbm_daemon_application_runtime_callbacks(fixture->application); + (void)snprintf(fixture->root, sizeof(fixture->root), "%s/cbm-app-%s-XXXXXX", cbm_tmpdir(), + label); + if (!cbm_mkdtemp(fixture->root)) { + fixture->root[0] = '\0'; + return false; + } + fixture->project = cbm_project_name_from_path(fixture->root); + return fixture->project && + app_test_context_request(fixture->root, fixture->root, &fixture->context, + &fixture->context_length); +} + +static cbm_daemon_runtime_application_session_t *app_async_session(app_async_fixture_t *fixture, + cbm_daemon_client_id_t id) { + cbm_daemon_runtime_application_session_t *session = app_test_open(&fixture->callbacks, id); + uint8_t *empty = NULL; + uint32_t empty_length = 0; + bool ok = session && app_test_request(&fixture->callbacks, session, fixture->context, + fixture->context_length, &empty, + &empty_length) == CBM_DAEMON_RUNTIME_APPLICATION_OK; + free(empty); + if (!ok && session) { + fixture->callbacks.session_close(fixture->callbacks.context, session); + session = NULL; + } + return session; +} + +/* One index_repository request on its own thread, so a call that (wrongly) + * blocks for the whole index cannot hang the suite: the caller observes + * `done` within the hang-detector budget, then always releases the worker + * before joining. */ +typedef struct { + app_request_thread_t request; + uint8_t *tool; + cbm_thread_t thread; + bool started; +} app_async_call_t; + +static bool app_async_call_start(app_async_fixture_t *fixture, app_async_call_t *call, + cbm_daemon_runtime_application_session_t *session, + cbm_daemon_runtime_application_token_t token, + const char *extra_args) { + memset(call, 0, sizeof(*call)); + char args[APP_TEST_PATH_CAP + 128]; + (void)snprintf(args, sizeof(args), "{\"repo_path\":\"%s\"%s}", fixture->root, + extra_args ? extra_args : ""); + uint32_t tool_length = 0; + if (!app_test_tool_request("index_repository", args, &call->tool, &tool_length)) { + return false; + } + call->request.callbacks = fixture->callbacks; + call->request.session = session; + call->request.request_token = token; + call->request.request = call->tool; + call->request.request_length = tool_length; + atomic_init(&call->request.done, false); + call->started = cbm_thread_create(&call->thread, 0, app_request_thread, &call->request) == 0; + return call->started; +} + +static const char *app_async_call_response(app_async_call_t *call) { + return call->request.response ? (const char *)call->request.response : ""; +} + +static void app_async_call_join(app_async_call_t *call) { + if (call->started) { + (void)cbm_thread_join(&call->thread); + call->started = false; + } +} + +static void app_async_call_free(app_async_call_t *call) { + app_async_call_join(call); + free(call->request.response); + free(call->tool); + memset(call, 0, sizeof(*call)); +} + +/* A request that must return promptly (status, validation, async start). */ +static char *app_async_call_now(app_async_fixture_t *fixture, + cbm_daemon_runtime_application_session_t *session, + const char *extra_args) { + app_async_call_t call; + memset(&call, 0, sizeof(call)); + char *response = NULL; + if (app_async_call_start(fixture, &call, session, CBM_DAEMON_RUNTIME_APPLICATION_TOKEN_INVALID, + extra_args) && + app_wait_for_atomic_bool(&call.request.done, true)) { + response = strdup(app_async_call_response(&call)); + } + /* Never strand a request thread that blocked on the held worker. */ + if (call.started && !atomic_load(&call.request.done)) { + atomic_store(&fixture->fake.allow_completion, true); + } + app_async_call_free(&call); + return response; +} + +static bool app_async_fixture_finish(app_async_fixture_t *fixture) { + atomic_store(&fixture->fake.allow_completion, true); + bool stopped = fixture->application && + cbm_daemon_application_shutdown(fixture->application, APP_TEST_TIMEOUT_MS); + cbm_daemon_application_free(fixture->application); + free(fixture->context); + free(fixture->project); + if (fixture->root[0]) { + (void)cbm_rmdir(fixture->root); + } + return stopped; +} + +static bool app_async_contains(const char *text, const char *needle) { + return text && strstr(text, needle) != NULL; +} + +TEST(daemon_application_async_index_returns_before_completion_and_status_polls) { + app_async_fixture_t fixture; + bool setup = app_async_fixture_init(&fixture, "async-poll", true); + cbm_daemon_runtime_application_session_t *session = + setup ? app_async_session(&fixture, 7101) : NULL; + + /* The worker is held: a blocking call cannot return until it is released. */ + char *started = session ? app_async_call_now(&fixture, session, ",\"async\":true") : NULL; + bool worker_running = started && app_wait_for_atomic_int(&fixture.fake.starts, 1); + char *running = + worker_running ? app_async_call_now(&fixture, session, ",\"status\":true") : NULL; + atomic_store(&fixture.fake.allow_completion, true); + bool drained = running && app_wait_for_active_jobs(fixture.application, 0); + char *finished = drained ? app_async_call_now(&fixture, session, ",\"status\":true") : NULL; + int cancels = atomic_load(&fixture.fake.cancels); + if (session) { + fixture.callbacks.session_close(fixture.callbacks.context, session); + } + bool stopped = app_async_fixture_finish(&fixture); + + ASSERT_TRUE(setup); + ASSERT_NOT_NULL(session); + ASSERT_NOT_NULL(started); + ASSERT_TRUE(app_async_contains(started, "\"async\":true")); + ASSERT_FALSE(app_async_contains(started, "\"isError\":true")); + ASSERT_TRUE(worker_running); + ASSERT_NOT_NULL(running); + ASSERT_TRUE(app_async_contains(running, "\"state\":\"running\"")); + ASSERT_TRUE(drained); + ASSERT_NOT_NULL(finished); + ASSERT_TRUE(app_async_contains(finished, "\"state\":\"succeeded\"")); + ASSERT_TRUE(app_async_contains(finished, "\"finished_at\"")); + ASSERT_EQ(cancels, 0); + ASSERT_EQ(atomic_load(&fixture.fake.starts), 1); + ASSERT_TRUE(stopped); + free(started); + free(running); + free(finished); + PASS(); +} + +/* The async job belongs to the daemon, not to the session that started it: + * closing that session (the client gave up and went away) leaves the worker + * running while another session still keeps this daemon generation alive. */ +TEST(daemon_application_async_index_survives_requesting_session_close) { + app_async_fixture_t fixture; + bool setup = app_async_fixture_init(&fixture, "async-survive", false); + cbm_daemon_runtime_application_session_t *starter = + setup ? app_async_session(&fixture, 7201) : NULL; + cbm_daemon_runtime_application_session_t *observer = + starter ? app_async_session(&fixture, 7202) : NULL; + char *started = observer ? app_async_call_now(&fixture, starter, ",\"async\":true") : NULL; + bool worker_running = started && app_wait_for_atomic_int(&fixture.fake.starts, 1); + if (worker_running) { + fixture.callbacks.session_cancel(fixture.callbacks.context, starter); + fixture.callbacks.session_close(fixture.callbacks.context, starter); + starter = NULL; + } + int cancels_after_close = atomic_load(&fixture.fake.cancels); + size_t jobs_after_close = cbm_daemon_application_active_jobs(fixture.application); + atomic_store(&fixture.fake.allow_completion, true); + bool drained = worker_running && app_wait_for_active_jobs(fixture.application, 0); + char *finished = drained ? app_async_call_now(&fixture, observer, ",\"status\":true") : NULL; + if (starter) { + fixture.callbacks.session_close(fixture.callbacks.context, starter); + } + if (observer) { + fixture.callbacks.session_close(fixture.callbacks.context, observer); + } + bool stopped = app_async_fixture_finish(&fixture); + + ASSERT_TRUE(setup); + ASSERT_NOT_NULL(observer); + ASSERT_TRUE(app_async_contains(started, "\"async\":true")); + ASSERT_TRUE(worker_running); + ASSERT_EQ(cancels_after_close, 0); + ASSERT_EQ(jobs_after_close, 1); + ASSERT_TRUE(drained); + ASSERT_TRUE(app_async_contains(finished, "\"state\":\"succeeded\"")); + ASSERT_EQ(atomic_load(&fixture.fake.cancels), 0); + ASSERT_TRUE(stopped); + free(started); + free(finished); + PASS(); +} + +/* Regression guard for the unchanged synchronous contract, in the exact + * shape of the test above: the same close with a SYNC caller still cancels + * the worker once its last subscriber is gone, even though another session + * keeps the daemon generation alive. */ +TEST(daemon_application_sync_index_still_cancels_on_last_subscriber_close) { + app_async_fixture_t fixture; + bool setup = app_async_fixture_init(&fixture, "sync-cancel", false); + cbm_daemon_runtime_application_session_t *starter = + setup ? app_async_session(&fixture, 7301) : NULL; + cbm_daemon_runtime_application_session_t *observer = + starter ? app_async_session(&fixture, 7302) : NULL; + app_async_call_t call; + memset(&call, 0, sizeof(call)); + bool call_started = + observer && app_async_call_start(&fixture, &call, starter, UINT64_C(7303), NULL); + bool worker_running = call_started && app_wait_for_atomic_int(&fixture.fake.starts, 1); + if (worker_running) { + fixture.callbacks.session_cancel(fixture.callbacks.context, starter); + } + bool cancelled = worker_running && app_wait_for_atomic_int(&fixture.fake.cancels, 1); + if (call_started) { + if (!atomic_load(&call.request.done)) { + atomic_store(&fixture.fake.allow_completion, true); + } + app_async_call_join(&call); + } + cbm_daemon_runtime_application_status_t status = call.request.status; + if (call_started) { + app_async_call_free(&call); + } + if (starter) { + fixture.callbacks.session_close(fixture.callbacks.context, starter); + } + if (observer) { + fixture.callbacks.session_close(fixture.callbacks.context, observer); + } + bool stopped = app_async_fixture_finish(&fixture); + + ASSERT_TRUE(setup); + ASSERT_TRUE(call_started); + ASSERT_TRUE(worker_running); + ASSERT_TRUE(cancelled); + ASSERT_EQ(status, CBM_DAEMON_RUNTIME_APPLICATION_CANCELLED); + ASSERT_TRUE(stopped); + PASS(); +} + +TEST(daemon_application_async_and_status_validation_errors) { + app_async_fixture_t fixture; + bool setup = app_async_fixture_init(&fixture, "async-validate", false); + cbm_daemon_runtime_application_session_t *session = + setup ? app_async_session(&fixture, 7401) : NULL; + char *both = + session ? app_async_call_now(&fixture, session, ",\"async\":true,\"status\":true") : NULL; + char *not_bool = session ? app_async_call_now(&fixture, session, ",\"async\":\"yes\"") : NULL; + char *cross = session + ? app_async_call_now(&fixture, session, + ",\"async\":true,\"mode\":\"cross-repo-intelligence\"") + : NULL; + /* Nothing was ever indexed here and no job ran: status must say so, not + * report an empty success. */ + char *unknown = session ? app_async_call_now(&fixture, session, ",\"status\":true") : NULL; + int starts = atomic_load(&fixture.fake.starts); + if (session) { + fixture.callbacks.session_close(fixture.callbacks.context, session); + } + bool stopped = app_async_fixture_finish(&fixture); + + ASSERT_TRUE(setup); + ASSERT_TRUE(app_async_contains(both, "\"isError\":true")); + ASSERT_TRUE(app_async_contains(both, "exclusive")); + ASSERT_TRUE(app_async_contains(not_bool, "\"isError\":true")); + ASSERT_TRUE(app_async_contains(not_bool, "must be booleans")); + ASSERT_TRUE(app_async_contains(cross, "\"isError\":true")); + ASSERT_TRUE(app_async_contains(cross, "cross-repo-intelligence")); + ASSERT_TRUE(app_async_contains(unknown, "\"isError\":true")); + ASSERT_TRUE(app_async_contains(unknown, "no index job or index is known")); + ASSERT_EQ(starts, 0); + ASSERT_TRUE(stopped); + free(both); + free(not_bool); + free(cross); + free(unknown); + PASS(); +} + +/* Cut-short advice (a): a synchronous caller whose job is cancelled out from + * under it (here: the daemon stops) receives the async alternative in the + * cancellation text it is actually delivered. */ +TEST(daemon_application_cancelled_sync_index_reply_offers_async) { + app_async_fixture_t fixture; + bool setup = app_async_fixture_init(&fixture, "sync-hint", false); + cbm_daemon_runtime_application_session_t *session = + setup ? app_async_session(&fixture, 7501) : NULL; + app_async_call_t call; + memset(&call, 0, sizeof(call)); + bool call_started = + session && app_async_call_start(&fixture, &call, session, UINT64_C(7502), NULL); + bool worker_running = call_started && app_wait_for_atomic_int(&fixture.fake.starts, 1); + /* Shutdown cancels the job but not the waiting request, so the waiter is + * handed the job's cancellation reply. */ + bool stopped = + worker_running && cbm_daemon_application_shutdown(fixture.application, APP_TEST_TIMEOUT_MS); + if (call_started && !atomic_load(&call.request.done)) { + atomic_store(&fixture.fake.allow_completion, true); + } + if (call_started) { + app_async_call_join(&call); + } + cbm_daemon_runtime_application_status_t status = call.request.status; + char *reply = call_started ? strdup(app_async_call_response(&call)) : NULL; + if (call_started) { + app_async_call_free(&call); + } + if (session) { + fixture.callbacks.session_close(fixture.callbacks.context, session); + } + (void)app_async_fixture_finish(&fixture); + + ASSERT_TRUE(setup); + ASSERT_TRUE(worker_running); + ASSERT_TRUE(stopped); + ASSERT_EQ(status, CBM_DAEMON_RUNTIME_APPLICATION_OK); + ASSERT_TRUE(app_async_contains(reply, "cancelled")); + ASSERT_TRUE(app_async_contains(reply, "async: true")); + ASSERT_TRUE(app_async_contains(reply, "status: true")); + free(reply); + PASS(); +} + +/* Cut-short advice (b): the Copilot-deadline case. The client cancels its + * synchronous call, which never sees a reply; the NEXT status call and the + * NEXT index_repository call for that project carry the advice, once. */ +TEST(daemon_application_cut_short_sync_index_advice_reaches_next_call) { + app_async_fixture_t fixture; + bool setup = app_async_fixture_init(&fixture, "sync-next", false); + cbm_daemon_runtime_application_session_t *session = + setup ? app_async_session(&fixture, 7601) : NULL; + app_async_call_t call; + memset(&call, 0, sizeof(call)); + cbm_daemon_runtime_application_token_t token = UINT64_C(7602); + bool call_started = session && app_async_call_start(&fixture, &call, session, token, NULL); + bool worker_running = call_started && app_wait_for_atomic_int(&fixture.fake.starts, 1); + if (worker_running) { + fixture.callbacks.request_cancel(fixture.callbacks.context, session, token); + } + bool cancelled_returned = worker_running && app_wait_for_atomic_bool(&call.request.done, true); + if (call_started && !atomic_load(&call.request.done)) { + atomic_store(&fixture.fake.allow_completion, true); + } + if (call_started) { + app_async_call_free(&call); + } + bool drained = cancelled_returned && app_wait_for_active_jobs(fixture.application, 0); + char *status = drained ? app_async_call_now(&fixture, session, ",\"status\":true") : NULL; + /* Every later worker completes at once. */ + atomic_store(&fixture.fake.allow_completion, true); + char *next = status ? app_async_call_now(&fixture, session, NULL) : NULL; + char *after = next ? app_async_call_now(&fixture, session, NULL) : NULL; + if (session) { + fixture.callbacks.session_close(fixture.callbacks.context, session); + } + bool stopped = app_async_fixture_finish(&fixture); + + ASSERT_TRUE(setup); + ASSERT_TRUE(worker_running); + ASSERT_TRUE(cancelled_returned); + ASSERT_TRUE(drained); + ASSERT_TRUE(app_async_contains(status, "\"state\":\"cancelled\"")); + ASSERT_TRUE(app_async_contains(status, "\"notice\"")); + ASSERT_TRUE(app_async_contains(status, "async: true")); + ASSERT_TRUE(app_async_contains(next, "indexed")); + ASSERT_TRUE(app_async_contains(next, "\"notice\"")); + ASSERT_TRUE(app_async_contains(next, "cut short")); + ASSERT_TRUE(app_async_contains(after, "indexed")); + ASSERT_FALSE(app_async_contains(after, "\"notice\"")); + ASSERT_TRUE(stopped); + free(status); + free(next); + free(after); + PASS(); +} + +/* #2144 user decision: on a temporary daemon, a one-shot `cli` request that is + * the daemon's only live session is refused async (the daemon would stop and + * cancel the job as soon as the command exits), while status stays allowed and + * a long-lived MCP session on the same temporary daemon remains a valid host. */ +TEST(daemon_application_async_refused_for_lone_one_shot_client_of_temporary_daemon) { + app_async_fixture_t fixture; + bool setup = app_async_fixture_init(&fixture, "async-refuse", false); + cbm_daemon_runtime_application_session_t *session = + setup ? app_async_session(&fixture, 7701) : NULL; + char *refused = session ? app_async_call_now(&fixture, session, ",\"async\":true") : NULL; + int starts_after_refusal = atomic_load(&fixture.fake.starts); + char *status = refused ? app_async_call_now(&fixture, session, ",\"status\":true") : NULL; + + /* Same temporary daemon, same lone session, but the MCP channel. */ + char message[APP_TEST_PATH_CAP + 256]; + (void)snprintf(message, sizeof(message), + "{\"jsonrpc\":\"2.0\",\"id\":7702,\"method\":\"tools/call\",\"params\":{" + "\"name\":\"index_repository\",\"arguments\":{\"repo_path\":\"%s\"," + "\"async\":true}}}", + fixture.root); + uint8_t *mcp_request = NULL; + uint32_t mcp_request_length = 0; + uint8_t *mcp_response = NULL; + uint32_t mcp_response_length = 0; + cbm_daemon_runtime_application_status_t mcp_status = + status && app_test_text_request(CBM_DAEMON_APPLICATION_REQUEST_MCP, message, &mcp_request, + &mcp_request_length) + ? app_test_request(&fixture.callbacks, session, mcp_request, mcp_request_length, + &mcp_response, &mcp_response_length) + : CBM_DAEMON_RUNTIME_APPLICATION_TRANSPORT_ERROR; + bool worker_started = mcp_status == CBM_DAEMON_RUNTIME_APPLICATION_OK && + app_wait_for_atomic_int(&fixture.fake.starts, 1); + int cancels_while_hosted = atomic_load(&fixture.fake.cancels); + atomic_store(&fixture.fake.allow_completion, true); + bool drained = worker_started && app_wait_for_active_jobs(fixture.application, 0); + if (session) { + fixture.callbacks.session_close(fixture.callbacks.context, session); + } + bool stopped = app_async_fixture_finish(&fixture); + char *mcp_text = mcp_response ? strdup((const char *)mcp_response) : NULL; + free(mcp_response); + free(mcp_request); + + ASSERT_TRUE(setup); + ASSERT_NOT_NULL(session); + ASSERT_TRUE(app_async_contains(refused, "\"isError\":true")); + ASSERT_TRUE(app_async_contains(refused, "daemon start")); + ASSERT_EQ(starts_after_refusal, 0); + ASSERT_NOT_NULL(status); + ASSERT_FALSE(app_async_contains(status, "daemon start")); + ASSERT_EQ(mcp_status, CBM_DAEMON_RUNTIME_APPLICATION_OK); + ASSERT_TRUE(app_async_contains(mcp_text, "\\\"async\\\":true")); + ASSERT_FALSE(app_async_contains(mcp_text, "daemon start")); + ASSERT_TRUE(worker_started); + ASSERT_EQ(cancels_while_hosted, 0); + ASSERT_TRUE(drained); + ASSERT_TRUE(stopped); + free(refused); + free(status); + free(mcp_text); + PASS(); +} + SUITE(daemon_application) { RUN_TEST(daemon_application_oversized_reply_is_a_jsonrpc_error_not_a_death); RUN_TEST(daemon_application_new_session_does_not_retain_initial_store); @@ -5653,6 +6142,13 @@ SUITE(daemon_application) { RUN_TEST(daemon_application_sensitive_root_blocks_auto_index_but_preserves_controls); 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_async_index_returns_before_completion_and_status_polls); + RUN_TEST(daemon_application_async_index_survives_requesting_session_close); + RUN_TEST(daemon_application_sync_index_still_cancels_on_last_subscriber_close); + RUN_TEST(daemon_application_async_and_status_validation_errors); + RUN_TEST(daemon_application_cancelled_sync_index_reply_offers_async); + RUN_TEST(daemon_application_cut_short_sync_index_advice_reaches_next_call); + RUN_TEST(daemon_application_async_refused_for_lone_one_shot_client_of_temporary_daemon); RUN_TEST(daemon_application_programmatic_index_injects_resource_policy); RUN_TEST(daemon_application_auto_index_honors_tracked_file_limit); RUN_TEST(daemon_application_auto_index_file_count_handles_literal_metacharacter_path); diff --git a/tests/test_daemon_frontend.c b/tests/test_daemon_frontend.c index ed0e868e90..b91705f157 100644 --- a/tests/test_daemon_frontend.c +++ b/tests/test_daemon_frontend.c @@ -1274,6 +1274,27 @@ TEST(daemon_frontend_correlates_cancellation_to_exact_request) { PASS(); } +/* #2144 (a): the -32800 reply to a cancelled index_repository call is the one + * text a deadline-bound client can still read, so it names the async mode; + * other requests and unparseable input keep the plain message. */ +TEST(daemon_frontend_cancelled_index_reply_offers_async_issue2144) { + const char *index_reply = cbm_daemon_frontend_cancelled_error_message( + "{\"jsonrpc\":\"2.0\",\"id\":5,\"method\":\"tools/call\"," + "\"params\":{\"name\":\"index_repository\",\"arguments\":{\"repo_path\":\"/r\"}}}"); + ASSERT_NOT_NULL(strstr(index_reply, "Request cancelled")); + ASSERT_NOT_NULL(strstr(index_reply, "async: true")); + ASSERT_NOT_NULL(strstr(index_reply, "status: true")); + ASSERT_NULL(strchr(index_reply, '"')); + ASSERT_NULL(strchr(index_reply, '\\')); + ASSERT_STR_EQ(cbm_daemon_frontend_cancelled_error_message( + "{\"jsonrpc\":\"2.0\",\"id\":6,\"method\":\"tools/call\"," + "\"params\":{\"name\":\"search_graph\",\"arguments\":{}}}"), + "Request cancelled"); + ASSERT_STR_EQ(cbm_daemon_frontend_cancelled_error_message("not json"), "Request cancelled"); + ASSERT_STR_EQ(cbm_daemon_frontend_cancelled_error_message(NULL), "Request cancelled"); + PASS(); +} + TEST(daemon_frontend_ignores_cancellation_text_in_string_content) { ASSERT_FALSE(cbm_daemon_frontend_is_cancellation_notification( "{\"jsonrpc\":\"2.0\",\"id\":1,\"method\":\"tools/call\"," @@ -1659,6 +1680,7 @@ TEST(daemon_local_participant_monitor_joins_before_manager_teardown) { SUITE(daemon_frontend) { RUN_TEST(daemon_frontend_recognizes_exact_cancellation_notification); RUN_TEST(daemon_frontend_correlates_cancellation_to_exact_request); + RUN_TEST(daemon_frontend_cancelled_index_reply_offers_async_issue2144); RUN_TEST(daemon_frontend_ignores_cancellation_text_in_string_content); RUN_TEST(daemon_frontend_rejects_non_notification_cancellation_shapes); #if defined(CBM_ENABLE_TEST_SEAMS) && CBM_ENABLE_TEST_SEAMS diff --git a/tests/test_mcp.c b/tests/test_mcp.c index ffa9ea3a3b..849cdbc5b8 100644 --- a/tests/test_mcp.c +++ b/tests/test_mcp.c @@ -20397,7 +20397,70 @@ TEST(bm25_searches_legacy_four_column_fts_without_error_issue518) { PASS(); } +/* #2144: async/status need the daemon's job registry. An in-process server + * (index worker, embedder) has no process that outlives the call, so it must + * refuse both modes instead of silently running a blocking index. */ +TEST(index_repository_async_and_status_refused_without_daemon_issue2144) { + cbm_mcp_server_t *srv = cbm_mcp_server_new(NULL); + ASSERT_NOT_NULL(srv); + char *async_reply = cbm_mcp_handle_tool(srv, "index_repository", + "{\"repo_path\":\"/nonexistent-2144\"," + "\"async\":true}"); + char *status_reply = cbm_mcp_handle_tool(srv, "index_repository", + "{\"repo_path\":\"/nonexistent-2144\"," + "\"status\":true}"); + cbm_mcp_server_free(srv); + ASSERT_NOT_NULL(async_reply); + ASSERT_NOT_NULL(strstr(async_reply, "\"isError\":true")); + ASSERT_NOT_NULL(strstr(async_reply, "daemon")); + ASSERT_NOT_NULL(status_reply); + ASSERT_NOT_NULL(strstr(status_reply, "\"isError\":true")); + ASSERT_NOT_NULL(strstr(status_reply, "daemon")); + free(async_reply); + free(status_reply); + PASS(); +} + +/* #2144: the tool schema advertises both flags and the description explains + * the deadline problem and the polling pattern. */ +TEST(index_repository_schema_documents_async_polling_issue2144) { + cbm_mcp_server_t *srv = cbm_mcp_server_new(NULL); + ASSERT_NOT_NULL(srv); + char *resp = + cbm_mcp_server_handle(srv, "{\"jsonrpc\":\"2.0\",\"id\":2144,\"method\":\"tools/list\"}"); + cbm_mcp_server_free(srv); + ASSERT_NOT_NULL(resp); + ASSERT_NOT_NULL(strstr(resp, "\"async\":{\"type\":\"boolean\"")); + ASSERT_NOT_NULL(strstr(resp, "\"status\":{\"type\":\"boolean\"")); + ASSERT_NOT_NULL(strstr(resp, "per-call deadline")); + ASSERT_NOT_NULL(strstr(resp, "poll with status:true")); + free(resp); + PASS(); +} + +/* #2144: a notice rides on a JSON payload as a key (content text and + * structuredContent together) and on a text payload as a paragraph. */ +TEST(tool_result_add_notice_keeps_payload_shape_issue2144) { + char *object = cbm_mcp_tool_result_add_notice( + cbm_mcp_text_result("{\"status\":\"indexed\"}", false), "retry async"); + char *text = + cbm_mcp_tool_result_add_notice(cbm_mcp_text_result("plain failure", true), "retry async"); + ASSERT_NOT_NULL(object); + ASSERT_NOT_NULL(strstr(object, "\"structuredContent\":{\"status\":\"indexed\"," + "\"notice\":\"retry async\"}")); + ASSERT_NOT_NULL(strstr(object, "\"isError\":false")); + ASSERT_NOT_NULL(text); + ASSERT_NOT_NULL(strstr(text, "plain failure\\n\\nretry async")); + ASSERT_NOT_NULL(strstr(text, "\"isError\":true")); + free(object); + free(text); + PASS(); +} + SUITE(mcp) { + RUN_TEST(index_repository_async_and_status_refused_without_daemon_issue2144); + RUN_TEST(index_repository_schema_documents_async_polling_issue2144); + RUN_TEST(tool_result_add_notice_keeps_payload_shape_issue2144); /* #518/#519 — BM25 prose search */ RUN_TEST(bm25_finds_section_by_its_prose_issue518); RUN_TEST(bm25_finds_module_by_promoted_description_issue519);