Caller holds application->mutex. Keeping watcher ownership validation and * this admission in the same critical section closes the unwatch race. */
| 1580 | /* Caller holds application->mutex. Keeping watcher ownership validation and |
| 1581 | * this admission in the same critical section closes the unwatch race. */ |
| 1582 | static 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; |
no test coverage detected