MCPcopy Create free account
hub / github.com/DeusData/codebase-memory-mcp / application_job_subscribe_locked

Function application_job_subscribe_locked

src/daemon/application.c:1582–1643  ·  view source on GitHub ↗

Caller holds application->mutex. Keeping watcher ownership validation and * this admission in the same critical section closes the unwatch race. */

Source from the content-addressed store, hash-verified

1580/* Caller holds application->mutex. Keeping watcher ownership validation and
1581 * this admission in the same critical section closes the unwatch race. */
1582static cbm_daemon_application_job_t *application_job_subscribe_locked(
1583 cbm_daemon_application_t *application, const char *project_key, const char *root_path,
1584 const char *args_json, application_job_subscribe_status_t *status_out) {
1585 *status_out = APPLICATION_JOB_SUBSCRIBE_UNAVAILABLE;
1586 if (application->stopping) {
1587 return NULL;
1588 }
1589 cbm_daemon_application_job_t *job =
1590 application_find_active_job_locked(application, project_key);
1591 if (job) {
1592 if (job->cancel_requested) {
1593 *status_out = APPLICATION_JOB_SUBSCRIBE_CANCELLING;
1594 return NULL;
1595 }
1596 if (!application_index_args_equal(job->args_json, args_json)) {
1597 *status_out = APPLICATION_JOB_SUBSCRIBE_OPTIONS_CONFLICT;
1598 return NULL;
1599 }
1600 job->subscribers++;
1601 *status_out = APPLICATION_JOB_SUBSCRIBE_OK;
1602 return job;
1603 }
1604
1605 if (application_active_job_count_locked(application) >= application->physical_job_limit) {
1606 char limit[32];
1607 (void)snprintf(limit, sizeof(limit), "%zu", application->physical_job_limit);
1608 cbm_log_warn("daemon.index.admission_busy", "limit", limit, "project", project_key);
1609 *status_out = APPLICATION_JOB_SUBSCRIBE_BUSY;
1610 return NULL;
1611 }
1612
1613 job = calloc(1, sizeof(*job));
1614 if (job) {
1615 job->project_key = strdup(project_key);
1616 job->root_path = strdup(root_path);
1617 job->args_json = strdup(args_json);
1618 }
1619 if (!job || !job->project_key || !job->root_path || !job->args_json) {
1620 application_job_free(job);
1621 *status_out = APPLICATION_JOB_SUBSCRIBE_ALLOCATION_FAILED;
1622 return NULL;
1623 }
1624 job->application = application;
1625 job->subscribers = 1;
1626 job->next = application->jobs;
1627 application->jobs = job;
1628 if (application_job_thread_create(&job->thread, job) == 0) {
1629 job->thread_started = true;
1630 } else {
1631 /* The job was linked only so a concurrently started thread could
1632 * observe its reservation. No thread exists on this path, so roll the
1633 * reservation back synchronously and let background callers retry. */
1634 application->jobs = job->next;
1635 job->next = NULL;
1636 application_job_free(job);
1637 *status_out = APPLICATION_JOB_SUBSCRIBE_UNAVAILABLE;
1638 cbm_log_warn("daemon.index.thread_start_failed", "action", "retry");
1639 return NULL;

Tested by

no test coverage detected