Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 18 additions & 4 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -467,7 +467,7 @@ $(HARDWARE_I2C_O): src/hardware/hardware_i2c.c src/hardware/hardware_i2c.h
$(HARDWARE_CAMERA_O): src/hardware/hardware_camera.c src/hardware/hardware_camera.h src/hardware/board_detect.h
$(CC) $(CFLAGS) $(INC) -c -o $@ src/hardware/hardware_camera.c

$(CRON_O): src/tools/cron.c src/tools/cron.h src/crypto/crypto.h src/core/memory.h src/channels/channel.h src/core/config.h
$(CRON_O): src/tools/cron.c src/tools/cron.h src/crypto/crypto.h src/core/agent.h src/core/memory.h src/channels/channel.h src/core/config.h
$(CC) $(CFLAGS) $(INC) -c -o $@ src/tools/cron.c

$(ASAP_INVOKE_O): src/tools/asap_invoke.c src/tools/asap_invoke.h src/tools/tool.h src/asap/envelope.h src/asap/client.h src/asap/registry.h src/asap/ulid.h src/asap/asap_version.h src/core/config.h vendor/cJSON/cJSON.h
Expand Down Expand Up @@ -631,9 +631,10 @@ test_web_search: tests/test_web_search.c $(WEBSEARCH_O) $(CONFIG_O) $(TOML_O) $(
$(CC) $(CFLAGS) $(LDFLAGS) $(INC) -o $(BINDIR)/$@ tests/test_web_search.c $(WEBSEARCH_O) $(CONFIG_O) $(TOML_O) $(CJSON_O) $(LDLIBS)
$(DSYM_SCRIPT)

test_cron: tests/test_cron.c $(CRON_O) $(CRYPTO_LINK) $(MEMORY_O) $(SQLITE3_O) $(CHANNEL_COMMON_O) $(CJSON_O)
test_cron: tests/test_cron.c $(CRON_O) $(CRYPTO_LINK) $(MEMORY_O) $(SQLITE3_O) $(CHANNEL_COMMON_O) $(CJSON_O) $(AGENT_O) $(SKILL_O) $(CONFIG_O) $(TOML_O) $(PROVIDER_COMMON_O)
@mkdir -p $(BINDIR)
$(CC) $(CFLAGS) $(LDFLAGS) $(INC) -o $(BINDIR)/$@ tests/test_cron.c $(CRON_O) $(CRYPTO_LINK) $(MEMORY_O) $(SQLITE3_O) $(CHANNEL_COMMON_O) $(CJSON_O) $(LDLIBS)
$(CC) $(CFLAGS) $(LDFLAGS) $(INC) -o $(BINDIR)/$@ tests/test_cron.c $(CRON_O) $(CRYPTO_LINK) $(MEMORY_O) $(SQLITE3_O) $(CHANNEL_COMMON_O) $(CJSON_O) $(AGENT_O) $(SKILL_O) $(CONFIG_O) $(TOML_O) $(PROVIDER_COMMON_O) $(LDLIBS) -pthread
$(DSYM_SCRIPT)

test_ws: tests/test_ws.c $(WS_TEST_O)
@mkdir -p $(BINDIR)
Expand Down Expand Up @@ -735,6 +736,18 @@ test_gateway_http: shellclaw tests/test_gateway_http.c $(AUTH_O) $(CONFIG_O) $(T
$(CC) $(CFLAGS) $(LDFLAGS) $(INC) -DSHELLCLAW_GATEWAY -pthread -o $(BINDIR)/$@ tests/test_gateway_http.c $(AUTH_O) $(CRYPTO_LINK) $(CONFIG_O) $(TOML_O) $(CJSON_O) $(LDLIBS) -pthread
$(DSYM_SCRIPT)

GATEWAY_DB_LOCK_OBJS := $(filter-out src/core/main.o,$(SHELLCLAW_OBJS))

test_gateway_db_lock: tests/test_gateway_db_lock.c $(GATEWAY_DB_LOCK_OBJS)
@mkdir -p $(BINDIR)
@if [ "$(GATEWAY)" != "1" ]; then \
if [ "$(CI)" = "true" ]; then echo "test_gateway_db_lock: GATEWAY=0 in CI — install libwebsockets-dev"; exit 1; fi; \
echo "test_gateway_db_lock: skipped (GATEWAY=0)"; \
else \
$(CC) $(CFLAGS) $(LDFLAGS) $(INC) -DSHELLCLAW_GATEWAY -pthread -o $(BINDIR)/$@ tests/test_gateway_db_lock.c $(GATEWAY_DB_LOCK_OBJS) $(LDLIBS) $(GATEWAY_LDLIBS) $(LIBGPIOD_LDLIBS) -pthread; \
fi
$(DSYM_SCRIPT)

test_routes_hardware: tests/test_routes_hardware.c tests/test_routes_json_stub.c $(ROUTES_HARDWARE_O) $(HARDWARE_GPIO_SNAPSHOT_O) $(HARDWARE_TEGRASTATS_O) $(HARDWARE_INIT_O) $(HARDWARE_STUB_O) $(BOARD_DETECT_O) $(HARDWARE_I2C_O) $(HARDWARE_CAMERA_O) $(HARDWARE_LIBGPIOD_O) $(CONFIG_O) $(TOML_O) $(CJSON_O)
@if [ "$(GATEWAY)" != "1" ]; then \
if [ "$(CI)" = "true" ]; then echo "test_routes_hardware: GATEWAY=0 in CI — install libwebsockets-dev"; exit 1; fi; \
Expand Down Expand Up @@ -889,6 +902,7 @@ test: test_config test_config_patch test_memory test_skill test_provider test_an
$(MAKE) test_static && $(BINDIR)/test_static
@if [ "$(GATEWAY)" = "1" ]; then \
$(BINDIR)/test_routes_hardware && \
$(MAKE) test_gateway_db_lock && $(BINDIR)/test_gateway_db_lock && \
$(MAKE) test_asap_http_body && $(BINDIR)/test_asap_http_body && \
$(MAKE) test_gateway_http && $(BINDIR)/test_gateway_http; \
elif [ "$(CI)" = "true" ]; then echo "GATEWAY=0 in CI — install libwebsockets-dev"; exit 1; fi
Expand All @@ -912,6 +926,6 @@ clean: clean-root-dsym
rm -f $(OBJS) $(PROVIDER_COMMON_O) $(STUB_O) $(ANTHROPIC_O) $(OPENAI_COMPAT_O) $(OPENAI_O) $(LOCAL_O) $(ROUTER_O) $(CJSON_O) $(TWEETNACL_O) $(ANTHROPIC_TEST_O) $(OPENAI_TEST_O) $(LOCAL_TEST_O) $(CONTEXT_TEST_OBJS) $(HEARTBEAT_TEST_O) $(CHANNEL_TG_TEST_O) $(CHANNEL_COMMON_O) $(CHANNEL_STUB_O) $(CHANNEL_CLI_O) $(CHANNEL_TG_O) $(CHANNEL_DISCORD_O) $(DISCORD_HELPERS_O) $(CHANNEL_HEARTBEAT_O) $(CHANNEL_WEBCHAT_O) $(AUTH_O) $(STATIC_O) $(HTTP_O) $(HTTP_LWS_O) $(ASAP_HTTP_BODY_O) $(ROUTES_O) $(ROUTES_HARDWARE_O) $(WS_O) $(MANIFEST_O) $(MANIFEST_PROFILES_O) $(MANIFEST_BUILD_O) $(MANIFEST_SIGN_O) $(MANIFEST_KEYS_O) $(ENVELOPE_O) $(ULID_O) $(CLIENT_O) $(ASAP_REGISTRY_O) $(SERVER_O) $(ASAP_LOG_O) $(RATE_LIMIT_O) $(SHELL_O) $(WEBSEARCH_O) $(FILE_O) $(REGISTRY_O) $(CONTEXT_O) $(CONTEXT_CACHE_O) $(CONTEXT_HTTP_O) $(CONTEXT_GEO_O) $(CRYPTO_O) $(JCS_O) $(HARDWARE_STUB_O) $(HARDWARE_INIT_O) $(HARDWARE_GPIO_SNAPSHOT_O) $(HARDWARE_TEGRASTATS_O) $(HARDWARE_TOOLS_O) $(BOARD_DETECT_O) src/hardware/hardware_libgpiod.o $(HARDWARE_I2C_O) $(HARDWARE_CAMERA_O) $(CRON_O) $(ASAP_INVOKE_O) $(SANDBOX_CORE_O) $(SANDBOX_LANDLOCK_O) $(ALLOWLIST_SCAN_O) $(ALLOWLIST_PATH_O)
rm -f src/gateway/ui_assets.h
find . -name '*.gcno' -o -name '*.gcda' -o -name '*.gcov' | xargs rm -f 2>/dev/null || true
rm -f $(WS_TEST_O) $(BINDIR)/asap_registry_test.o $(BINDIR)/asap_invoke_test.o $(CONTEXT_TEST_OBJS) $(HEARTBEAT_TEST_O) $(BINDIR)/shellclaw $(BINDIR)/test_tweetnacl_smoke $(BINDIR)/test_config $(BINDIR)/test_config_patch $(BINDIR)/test_memory $(BINDIR)/test_skill $(BINDIR)/test_provider $(BINDIR)/test_anthropic $(BINDIR)/test_openai $(BINDIR)/test_local_provider $(BINDIR)/test_router $(BINDIR)/test_heartbeat $(BINDIR)/test_crypto $(BINDIR)/test_hardware_stub $(BINDIR)/test_board_detect $(BINDIR)/test_hardware_libgpiod $(BINDIR)/test_hardware_i2c $(BINDIR)/test_hardware_camera $(BINDIR)/test_pin_tables $(BINDIR)/test_hardware_init $(BINDIR)/test_hardware_tools $(BINDIR)/test_registry $(BINDIR)/test_ws $(BINDIR)/test_agent $(BINDIR)/test_channel $(BINDIR)/test_cli $(BINDIR)/test_shell $(BINDIR)/test_file $(BINDIR)/test_telegram $(BINDIR)/test_discord_helpers $(BINDIR)/test_web_search $(BINDIR)/test_cron $(BINDIR)/test_context $(BINDIR)/test_manifest_build $(BINDIR)/test_manifest_keys $(BINDIR)/test_jcs $(BINDIR)/test_asap_envelope $(BINDIR)/test_asap_ulid $(BINDIR)/test_asap_client $(BINDIR)/test_asap_registry $(BINDIR)/test_asap_server $(BINDIR)/test_asap_invoke $(BINDIR)/test_asap_log $(BINDIR)/test_auth $(BINDIR)/test_gateway_http $(BINDIR)/test_static $(BINDIR)/test_sandbox $(BINDIR)/test_allowlist $(BINDIR)/test_rate_limit
rm -f $(WS_TEST_O) $(BINDIR)/asap_registry_test.o $(BINDIR)/asap_invoke_test.o $(CONTEXT_TEST_OBJS) $(HEARTBEAT_TEST_O) $(BINDIR)/shellclaw $(BINDIR)/test_tweetnacl_smoke $(BINDIR)/test_config $(BINDIR)/test_config_patch $(BINDIR)/test_memory $(BINDIR)/test_skill $(BINDIR)/test_provider $(BINDIR)/test_anthropic $(BINDIR)/test_openai $(BINDIR)/test_local_provider $(BINDIR)/test_router $(BINDIR)/test_heartbeat $(BINDIR)/test_crypto $(BINDIR)/test_hardware_stub $(BINDIR)/test_board_detect $(BINDIR)/test_hardware_libgpiod $(BINDIR)/test_hardware_i2c $(BINDIR)/test_hardware_camera $(BINDIR)/test_pin_tables $(BINDIR)/test_hardware_init $(BINDIR)/test_hardware_tools $(BINDIR)/test_registry $(BINDIR)/test_ws $(BINDIR)/test_agent $(BINDIR)/test_channel $(BINDIR)/test_cli $(BINDIR)/test_shell $(BINDIR)/test_file $(BINDIR)/test_telegram $(BINDIR)/test_discord_helpers $(BINDIR)/test_web_search $(BINDIR)/test_cron $(BINDIR)/test_context $(BINDIR)/test_manifest_build $(BINDIR)/test_manifest_keys $(BINDIR)/test_jcs $(BINDIR)/test_asap_envelope $(BINDIR)/test_asap_ulid $(BINDIR)/test_asap_client $(BINDIR)/test_asap_registry $(BINDIR)/test_asap_server $(BINDIR)/test_asap_invoke $(BINDIR)/test_asap_log $(BINDIR)/test_auth $(BINDIR)/test_gateway_http $(BINDIR)/test_gateway_db_lock $(BINDIR)/test_static $(BINDIR)/test_sandbox $(BINDIR)/test_allowlist $(BINDIR)/test_rate_limit
rm -rf $(BINDIR)/*.dSYM $(DSYMDIR)
rm -f $(BOOTSTRAP_DISPATCH_STUB_O) $(TOOL_RELOAD_STUB_O) $(RELOAD_CHANNEL_STUB_O) $(HTTP_RELOAD_STUB_O)
9 changes: 6 additions & 3 deletions src/core/agent.h
Original file line number Diff line number Diff line change
Expand Up @@ -58,14 +58,17 @@ int agent_run(const config_t *cfg, const char *session_id, const char *user_mess
* cron_ack_delivery() (SQLite amalgamation is SQLITE_THREADSAFE=0).
* Inbound mcp.tool_call must hold it around tool execute (see #60).
* Inbound state.query must hold it around memory_get_row_counts() (see #60).
* Release before channel I/O (ch->send). The mutex is not recursive.
* Dashboard /api/memory, /api/sessions list, and /api/cron handlers must hold
* it around g_db. cron_poll must hold it around due-job SQL, and must release
* it before the re-offer sleep. The mutex is not recursive.
* Release before channel I/O (ch->send).
*/
void agent_lock(void);

/**
* Release the global agent mutex after agent_run(), a locked session_delete(),
* inbound mcp.tool_call execute, a locked state.query memory read, or a
* locked cron_ack_delivery().
* inbound mcp.tool_call execute, a locked state.query memory read, a locked
* cron_ack_delivery(), a dashboard g_db handler, or cron_poll's due-job SQL.
*/
void agent_unlock(void);

Expand Down
20 changes: 20 additions & 0 deletions src/core/memory.c
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
static sqlite3 *g_db;
static void (*g_session_delete_hook_for_test)(const char *session_id);
static void (*g_memory_get_row_counts_hook_for_test)(void);
static void (*g_sqlite_entry_hook_for_test)(void);

void session_delete_set_hook_for_test(void (*hook)(const char *session_id))
{
Expand All @@ -28,6 +29,17 @@ void memory_get_row_counts_set_hook_for_test(void (*hook)(void))
g_memory_get_row_counts_hook_for_test = hook;
}

void memory_sqlite_entry_set_hook_for_test(void (*hook)(void))
{
g_sqlite_entry_hook_for_test = hook;
}

static void memory_sqlite_entry_probe(void)
{
if (g_sqlite_entry_hook_for_test)
g_sqlite_entry_hook_for_test();
}

static const char *SCHEMA_MEMORIES =
"CREATE TABLE IF NOT EXISTS memories ("
" rowid INTEGER PRIMARY KEY,"
Expand Down Expand Up @@ -202,6 +214,7 @@ int memory_save(const char *key, const char *content, const char *metadata)
int memory_recall(const char *query, char *results, size_t max_len, int limit)
{
if (!g_db || !results || max_len == 0 || !query) return -1;
memory_sqlite_entry_probe();
results[0] = '\0';
if (limit <= 0) limit = 10;
const char *sql = "SELECT content FROM memories_fts WHERE memories_fts MATCH ?1 ORDER BY rank LIMIT ?2";
Expand Down Expand Up @@ -291,6 +304,7 @@ int session_delete(const char *session_id)
int session_list(char **session_ids_out, int max_count)
{
if (!g_db || !session_ids_out || max_count <= 0) return -1;
memory_sqlite_entry_probe();
const char *sql = "SELECT id FROM sessions ORDER BY updated_at DESC";
sqlite3_stmt *stmt = NULL;
if (sqlite3_prepare_v2(g_db, sql, -1, &stmt, NULL) != SQLITE_OK) return -1;
Expand Down Expand Up @@ -398,6 +412,7 @@ int cron_job_create(const char *id, const char *schedule, const char *message,
const char *channel, const char *recipient, long long next_run, int enabled)
{
if (!g_db || !id || !schedule || !message) return -1;
memory_sqlite_entry_probe();
/* fill_cron_job_row truncates id[128]; a longer id breaks the due cursor. */
if (strlen(id) >= sizeof(((cron_job_row_t *)0)->id)) return -1;
if (strlen(schedule) > (size_t)CRON_JOB_TEXT_MAX) return -1;
Expand All @@ -423,6 +438,7 @@ int cron_job_create(const char *id, const char *schedule, const char *message,
int cron_job_delete(const char *id)
{
if (!g_db || !id) return -1;
memory_sqlite_entry_probe();
const char *sql = "DELETE FROM cron_jobs WHERE id = ?1";
sqlite3_stmt *stmt = NULL;
if (sqlite3_prepare_v2(g_db, sql, -1, &stmt, NULL) != SQLITE_OK) return -1;
Expand All @@ -435,6 +451,7 @@ int cron_job_delete(const char *id)
int cron_job_toggle(const char *id)
{
if (!g_db || !id) return -1;
memory_sqlite_entry_probe();
const char *sql = "UPDATE cron_jobs SET enabled = 1 - enabled WHERE id = ?1";
sqlite3_stmt *stmt = NULL;
if (sqlite3_prepare_v2(g_db, sql, -1, &stmt, NULL) != SQLITE_OK) return -1;
Expand All @@ -460,6 +477,7 @@ int cron_job_update_next_run(const char *id, long long next_run)
int cron_job_list(cron_job_row_t *out, int max_count)
{
if (!g_db || !out || max_count <= 0) return -1;
memory_sqlite_entry_probe();
const char *sql = "SELECT id, schedule, message, channel, recipient, next_run, enabled "
"FROM cron_jobs ORDER BY next_run ASC, id ASC";
sqlite3_stmt *stmt = NULL;
Expand All @@ -485,6 +503,7 @@ int cron_job_list_due(long long now, cron_due_key_t *out, int max_count)
int count = 0;

if (!g_db || !out || max_count <= 0) return -1;
memory_sqlite_entry_probe();
sql = "SELECT id, next_run FROM cron_jobs WHERE next_run <= ?1 AND enabled = 1 "
"ORDER BY next_run ASC, id ASC";
if (sqlite3_prepare_v2(g_db, sql, -1, &stmt, NULL) != SQLITE_OK) return -1;
Expand Down Expand Up @@ -551,6 +570,7 @@ int cron_job_get_by_id(const char *id, cron_job_row_t *out)
int ret = 0;

if (!g_db || !id || !out) return -1;
memory_sqlite_entry_probe();
sql = "SELECT id, schedule, message, channel, recipient, next_run, enabled FROM cron_jobs WHERE id = ?1";
if (sqlite3_prepare_v2(g_db, sql, -1, &stmt, NULL) != SQLITE_OK) return -1;
sqlite3_bind_text(stmt, 1, id, -1, SQLITE_TRANSIENT);
Expand Down
9 changes: 9 additions & 0 deletions src/core/memory.h
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,15 @@ void session_delete_set_hook_for_test(void (*hook)(const char *session_id));
*/
void memory_get_row_counts_set_hook_for_test(void (*hook)(void));

/**
* Test-only: invoke @p hook when a dashboard or cron-poll SQL entry runs.
* Pass NULL to clear. Used to assert those callers hold agent_lock()
* (SQLite is SQLITE_THREADSAFE=0).
*
* Example: memory_sqlite_entry_set_hook_for_test(probe); session_list(ids, 8);
*/
void memory_sqlite_entry_set_hook_for_test(void (*hook)(void));

/**
* List session IDs from the database.
*
Expand Down
37 changes: 30 additions & 7 deletions src/gateway/routes.c
Original file line number Diff line number Diff line change
Expand Up @@ -458,8 +458,13 @@ static void handle_skill_delete(const config_t *cfg, const char *name,
static void handle_memory_get(const char *query, int limit, char *buf, size_t size, int *status)
{
char *results = malloc(MEMORY_RESULTS_MAX);
int recall_rc;
if (!results) { json_error(buf, size, status, 500, "Out of memory"); return; }
if (memory_recall(query ? query : "", results, MEMORY_RESULTS_MAX, limit > 0 ? limit : 20) != 0) {
agent_lock();
recall_rc = memory_recall(query ? query : "", results, MEMORY_RESULTS_MAX,
limit > 0 ? limit : 20);
agent_unlock();
if (recall_rc != 0) {
free(results);
json_error(buf, size, status, 500, "Memory recall failed");
return;
Expand All @@ -481,7 +486,10 @@ static void handle_memory_get(const char *query, int limit, char *buf, size_t si
static void handle_sessions_list(char *buf, size_t size, int *status)
{
char *ids[64];
int n = session_list(ids, 64);
int n;
agent_lock();
n = session_list(ids, 64);
agent_unlock();
cJSON *arr = cJSON_CreateArray();
if (!arr) { json_error(buf, size, status, 500, "Internal error"); return; }
for (int i = 0; i < n; i++) {
Expand Down Expand Up @@ -522,7 +530,9 @@ static void handle_cron_list(char *buf, size_t size, int *status)
int n;
int i;
if (!rows) { json_error(buf, size, status, 500, "Out of memory"); return; }
agent_lock();
n = cron_job_list(rows, 64);
agent_unlock();
if (n < 0) {
free(rows);
json_error(buf, size, status, 500, "Internal error");
Expand Down Expand Up @@ -580,7 +590,10 @@ static void handle_cron_create(const char *body, size_t body_len, char *buf, siz
id[0] = '\0';
const char *ch = (channel && cJSON_IsString(channel)) ? channel->valuestring : NULL;
const char *rec = (recipient && cJSON_IsString(recipient)) ? recipient->valuestring : NULL;
int ret = cron_create_job(schedule->valuestring, message->valuestring, ch, rec, id, sizeof(id));
int ret;
agent_lock();
ret = cron_create_job(schedule->valuestring, message->valuestring, ch, rec, id, sizeof(id));
agent_unlock();
cJSON_Delete(root);
if (ret != 0) {
json_error(buf, size, status, 400, "Failed to create job");
Expand All @@ -600,7 +613,11 @@ static void handle_cron_create(const char *body, size_t body_len, char *buf, siz

static void handle_cron_delete(const char *id, char *buf, size_t size, int *status)
{
if (cron_job_delete(id) != 0) {
int rc;
agent_lock();
rc = cron_job_delete(id);
agent_unlock();
if (rc != 0) {
json_error(buf, size, status, 404, "Job not found");
return;
}
Expand All @@ -610,7 +627,11 @@ static void handle_cron_delete(const char *id, char *buf, size_t size, int *stat

static void handle_cron_toggle(const char *id, char *buf, size_t size, int *status)
{
if (cron_job_toggle(id) != 0) {
int rc;
agent_lock();
rc = cron_job_toggle(id);
agent_unlock();
if (rc != 0) {
json_error(buf, size, status, 404, "Job not found");
return;
}
Expand Down Expand Up @@ -964,8 +985,10 @@ int dispatch_route(http_server_ctx_t *ctx, struct lws *wsi, int method,
if (method != HTTP_GET) { json_error(buf, size, status, 405, "Method not allowed"); return 0; }
char qbuf[256] = {0};
char lbuf[32] = {0};
lws_get_urlarg_by_name_safe(wsi, "q", qbuf, sizeof(qbuf));
lws_get_urlarg_by_name_safe(wsi, "limit", lbuf, sizeof(lbuf));
if (wsi) {
lws_get_urlarg_by_name_safe(wsi, "q", qbuf, sizeof(qbuf));
lws_get_urlarg_by_name_safe(wsi, "limit", lbuf, sizeof(lbuf));
}
int limit = parse_memory_limit_arg(lbuf[0] ? lbuf : NULL);
handle_memory_get(qbuf[0] ? qbuf : NULL, limit, buf, size, status);
return 0;
Expand Down
38 changes: 25 additions & 13 deletions src/tools/cron.c
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@

#include "tools/cron.h"
#include "crypto/crypto.h"
#include "core/agent.h"
#include "core/memory.h"
#include "channels/channel.h"
#include "core/config.h"
Expand Down Expand Up @@ -320,28 +321,39 @@ static int cron_pick_due_row(long long now, int timeout_ms, cron_job_row_t *row)
{
cron_due_key_t keys[CRON_DUE_SCAN_MAX];
struct timespec mono;
long age = 0;
int n;
int i;
int wait_for_first = 0;

if (clock_gettime(CLOCK_MONOTONIC, &mono) != 0)
memset(&mono, 0, sizeof(mono));
/* HTTP dashboard routes use this same connection (SQLITE_THREADSAFE=0).
* Drop the lock before the re-offer sleep so that thread is not stalled. */
agent_lock();
n = cron_job_list_due(now, keys, CRON_DUE_SCAN_MAX);
if (n <= 0)
return 0;
for (i = 0; i < n; i++) {
if (!cron_due_is_returnable(keys[i].id, timeout_ms, &mono))
continue;
if (cron_job_get_by_id(keys[i].id, row) == 1)
return 1;
if (n > 0) {
for (i = 0; i < n; i++) {
if (!cron_due_is_returnable(keys[i].id, timeout_ms, &mono))
continue;
if (cron_job_get_by_id(keys[i].id, row) == 1) {
agent_unlock();
return 1;
}
}
if (cron_job_get_by_id(keys[0].id, row) == 1)
wait_for_first = 1;
}
if (cron_job_get_by_id(keys[0].id, row) != 1)
agent_unlock();
if (!wait_for_first)
return 0;
/* An id that never fit in the table has no age, so wait_remaining would return immediately. */
if (cron_offer_age_ms(row->id, &mono, &age))
cron_wait_remaining(row->id, timeout_ms);
else
cron_sleep_ms((long)timeout_ms);
{
long age = 0;
if (cron_offer_age_ms(row->id, &mono, &age))
cron_wait_remaining(row->id, timeout_ms);
else
cron_sleep_ms((long)timeout_ms);
}
return 1;
}

Expand Down
Loading
Loading