From e147f476ee24d803ab263732eca1e6470688e66d Mon Sep 17 00:00:00 2001 From: Sounkou Mahamane Toure <56392505+sounkou-bioinfo@users.noreply.github.com> Date: Thu, 17 Sep 2026 01:53:59 +0400 Subject: [PATCH] fix(coordinator): reap the real task schema and strengthen regression coverage Align native maintenance with the existing reaping transition. Replace the synthetic test table with an installed-package fixture, make poll control flow explicit, and correct NULL handling and partial-load cleanup. Preserve the current CI matrices and install this checkout for the native integration test. --- .github/workflows/R-CMD-check.yaml | 6 +- tools/canard-coordinator/README.md | 46 ++-- .../src/canard_coordinator.c | 56 +++-- tools/canard-coordinator/test/test-v1.R | 218 +++++++++++++----- 4 files changed, 230 insertions(+), 96 deletions(-) diff --git a/.github/workflows/R-CMD-check.yaml b/.github/workflows/R-CMD-check.yaml index f8852e1..3a72e9f 100644 --- a/.github/workflows/R-CMD-check.yaml +++ b/.github/workflows/R-CMD-check.yaml @@ -92,9 +92,13 @@ jobs: - name: Install runtime test dependencies shell: Rscript {0} run: | - install.packages(c("DBI", "duckdb")) + install.packages(c("bit64", "DBI", "duckdb", "S7")) stopifnot(packageVersion("duckdb") >= "1.5.5") + - name: Install CanardAbsurd from this checkout + shell: bash + run: R CMD INSTALL . + - name: Fetch DuckDB 1.5.5 shell: bash run: git clone --branch v1.5.5 --depth 1 https://github.com/duckdb/duckdb.git artifacts/upstream/duckdb diff --git a/tools/canard-coordinator/README.md b/tools/canard-coordinator/README.md index b9a73ac..6da7e00 100644 --- a/tools/canard-coordinator/README.md +++ b/tools/canard-coordinator/README.md @@ -9,26 +9,26 @@ tasks that have exhausted `max_failures`. It exposes: ```sql SELECT ca_coordinator_start(1000, 64); -- poll milliseconds, row limit -SELECT (ca_coordinator_status()).*; +SELECT unnest(ca_coordinator_status()); SELECT ca_coordinator_stop(); ``` +Reaping applies the transition in `inst/sql/reap.sql` across all queues: mark +an exhausted task failed, increment its failure count, record the lease-expiry +error, and clear its worker, token, and lease deadline. It does not change +retryable expired tasks, live leases, or terminal tasks. + The extension opens its dedicated connection while its load entrypoint still has a valid borrowed database handle. `ca_coordinator_stop()` interrupts active database work, joins the thread, and closes that connection. The host must call -it before closing the database. The coordinator cannot restart on that database -handle after the connection closes. - -C API v1 can drive the maintenance query, but it does not provide a loaded C -extension with a database-shutdown callback or an owned clone of the host -database handle. The one-shot lifecycle is the narrow compatible contract; -a permanently autonomous loadable-extension lifecycle is not supplied by v1. +it before closing the database, even if the extension was loaded but never +started. The coordinator cannot restart on that database handle after the +connection closes. The maintenance statement uses `duckdb_pending_prepared()` and repeatedly calls `duckdb_pending_execute_task()`. It counts a mutation only after -`duckdb_execute_pending()` succeeds. This is the v1 equivalent needed for this -bounded operation; v2's chunked result state machine and structured errors are -not required here. +`duckdb_execute_pending()` succeeds. Status reports both successful and failed +polls in `poll_count`; `last_error` describes the latest poll. ## Build against DuckDB 1.5.5 @@ -45,16 +45,26 @@ cmake --build build/canard --target canard_coordinator_loadable_extension The output is under `build/canard/extension/canard_coordinator/canard_coordinator.duckdb_extension`. -Run the compatibility test with the R `duckdb` 1.5.5 package: + +Install this CanardAbsurd checkout before testing so the test uses the package's +schema and task-input encoding, not a substitute table definition: + +```sh +R CMD INSTALL /absolute/path/to/CanardAbsurd +``` + +From `tools/canard-coordinator`, run the compatibility test with the R `duckdb` +1.5.5 package: ```sh Rscript test/test-v1.R /absolute/path/to/canard_coordinator.duckdb_extension ``` -The test enables unsigned extensions only on its temporary development -connection. Production artifacts need the normal DuckDB extension signing and -distribution path. +The test checks exhausted tasks in multiple queues, unchanged retryable/live/ +terminal tasks, ownership cleanup, argument validation, and the one-shot +lifecycle. Cleanup also runs when an assertion fails. Unsigned extensions are +enabled only on its temporary development connection; production artifacts +need the normal DuckDB extension signing and distribution path. -This is a maintenance coordinator, not an external-job supervisor. Submission -identity, backend closure, resource reservations, and immutable artifact -publication remain separate contracts. +This is a maintenance coordinator. It does not dispatch jobs, allocate physical +resources, or execute R handlers. diff --git a/tools/canard-coordinator/src/canard_coordinator.c b/tools/canard-coordinator/src/canard_coordinator.c index b4fe641..95d6666 100644 --- a/tools/canard-coordinator/src/canard_coordinator.c +++ b/tools/canard-coordinator/src/canard_coordinator.c @@ -12,17 +12,21 @@ extern duckdb_ext_api_v1 duckdb_ext_api; +/* Same transition as inst/sql/reap.sql, applied across all queues. */ static const char *CANARD_REAP_SQL = - "UPDATE canard_absurd.tasks " - "SET state = 'failed', completed_at = current_timestamp, updated_at = current_timestamp, " - "lease_until = NULL, lease_token = NULL " - "WHERE id IN (" - " SELECT id FROM canard_absurd.tasks " - " WHERE state = 'running' AND lease_until <= current_timestamp " - " AND failures + 1 >= max_failures " - " AND (SELECT version FROM canard_absurd.schema_version) = 1 " - " ORDER BY lease_until, id LIMIT ?" - ") RETURNING id"; + "UPDATE canard_absurd.tasks\n" + "SET state = 'failed', failures = failures + 1,\n" + " error = 'worker lease expired',\n" + " worker = NULL, token = NULL, lease_until = NULL,\n" + " updated_at = current_timestamp\n" + "WHERE id IN (\n" + " SELECT id FROM canard_absurd.tasks\n" + " WHERE state = 'running'\n" + " AND lease_until <= current_timestamp AND failures + 1 >= max_failures\n" + " AND (SELECT version FROM canard_absurd.schema_version) = 1\n" + " ORDER BY lease_until, id LIMIT ?\n" + ")\n" + "RETURNING 1 AS changed, id;"; typedef enum { CANARD_COORDINATOR_READY = 0, @@ -170,16 +174,19 @@ static void canard_coordinator_main(void *argument) { while (!canard_stop_requested(runtime)) { char error[CANARD_ERROR_CAPACITY] = {0}; uint64_t reaped = 0; - if (!statement && canard_prepare_reaper(runtime, &statement, error) == DuckDBError) { - if (!canard_stop_requested(runtime)) { - canard_record_poll(runtime, 0, error); - } - } else if (statement && canard_execute_reaper(runtime, statement, &reaped, error) == DuckDBError) { + duckdb_state poll_result = DuckDBSuccess; + if (!statement) { + poll_result = canard_prepare_reaper(runtime, &statement, error); + } + if (poll_result == DuckDBSuccess) { + poll_result = canard_execute_reaper(runtime, statement, &reaped, error); + } + if (poll_result == DuckDBError) { duckdb_destroy_prepare(&statement); if (!canard_stop_requested(runtime)) { canard_record_poll(runtime, 0, error); } - } else if (statement) { + } else { canard_record_poll(runtime, reaped, NULL); } if (canard_wait(runtime, runtime->poll_milliseconds)) { @@ -204,6 +211,7 @@ static bool canard_runtime_start(canard_runtime *runtime, uint64_t poll_millisec canard_mutex_unlock(&runtime->mutex); return false; } + runtime->poll_milliseconds = poll_milliseconds; runtime->reap_limit = reap_limit; runtime->stop_requested = false; @@ -310,8 +318,10 @@ static void canard_start_function(duckdb_function_info info, duckdb_data_chunk i if (!canard_valid_control_call(info, input, 2)) { return; } - uint64_t poll_milliseconds = ((uint64_t *)duckdb_vector_get_data(duckdb_data_chunk_get_vector(input, 0)))[0]; - uint64_t reap_limit = ((uint64_t *)duckdb_vector_get_data(duckdb_data_chunk_get_vector(input, 1)))[0]; + duckdb_vector poll_vector = duckdb_data_chunk_get_vector(input, 0); + duckdb_vector limit_vector = duckdb_data_chunk_get_vector(input, 1); + uint64_t poll_milliseconds = ((uint64_t *)duckdb_vector_get_data(poll_vector))[0]; + uint64_t reap_limit = ((uint64_t *)duckdb_vector_get_data(limit_vector))[0]; if (poll_milliseconds < 1 || poll_milliseconds > CANARD_MAX_POLL_MILLISECONDS) { duckdb_scalar_function_set_error(info, "poll_milliseconds must be between 1 and 86400000"); return; @@ -339,7 +349,6 @@ static void canard_stop_function(duckdb_function_info info, duckdb_data_chunk in } static void canard_status_function(duckdb_function_info info, duckdb_data_chunk input, duckdb_vector output) { - (void)input; canard_runtime *runtime = (canard_runtime *)duckdb_scalar_function_get_extra_info(info); canard_status status; canard_status_read(runtime, &status); @@ -372,6 +381,8 @@ static bool canard_register_start(duckdb_connection connection, canard_runtime * duckdb_scalar_function_add_parameter(function, unsigned_type); duckdb_scalar_function_set_return_type(function, boolean_type); duckdb_scalar_function_set_volatile(function); + /* Let the callback reject NULL instead of silently returning SQL NULL. */ + duckdb_scalar_function_set_special_handling(function); canard_runtime_retain(runtime); duckdb_scalar_function_set_extra_info(function, runtime, canard_runtime_release); duckdb_scalar_function_set_function(function, canard_start_function); @@ -453,10 +464,11 @@ bool canard_coordinator_load(duckdb_connection connection, duckdb_extension_info bool registered = canard_register_start(connection, runtime) && canard_register_stop(connection, runtime) && canard_register_status(connection, runtime); - canard_runtime_release(runtime); if (!registered) { + /* Earlier registrations can still retain the runtime after LOAD fails. */ + (void)canard_runtime_stop(runtime); access->set_error(info, "Unable to register the Canard coordinator SQL functions"); - return false; } - return true; + canard_runtime_release(runtime); + return registered; } diff --git a/tools/canard-coordinator/test/test-v1.R b/tools/canard-coordinator/test/test-v1.R index 2a958d5..ffdd894 100644 --- a/tools/canard-coordinator/test/test-v1.R +++ b/tools/canard-coordinator/test/test-v1.R @@ -1,56 +1,164 @@ -args <- commandArgs(trailingOnly = TRUE) -if (length(args) != 1L) stop("Usage: Rscript test-v1.R /path/to/canard_coordinator.duckdb_extension") -extension <- normalizePath(args[[1L]], mustWork = TRUE) - -driver <- duckdb::duckdb(tempfile(fileext = ".duckdb"), config = list( - allow_unsigned_extensions = "true", autoinstall_known_extensions = "false")) -con <- DBI::dbConnect(driver) -on.exit(DBI::dbDisconnect(con, shutdown = TRUE), add = TRUE) -DBI::dbExecute(con, paste("LOAD", DBI::dbQuoteString(con, extension))) -DBI::dbExecute(con, "CREATE SCHEMA canard_absurd") -DBI::dbExecute(con, "CREATE TABLE canard_absurd.schema_version (version INTEGER PRIMARY KEY)") -DBI::dbExecute(con, "INSERT INTO canard_absurd.schema_version VALUES (1)") -DBI::dbExecute(con, " - CREATE TABLE canard_absurd.tasks ( - id VARCHAR PRIMARY KEY, - state VARCHAR NOT NULL, - failures INTEGER NOT NULL, - max_failures INTEGER NOT NULL, - lease_until TIMESTAMPTZ, - lease_token UUID, - completed_at TIMESTAMPTZ, - updated_at TIMESTAMPTZ NOT NULL - )") -DBI::dbExecute(con, " - INSERT INTO canard_absurd.tasks - VALUES ('expired', 'running', 0, 1, TIMESTAMPTZ '2000-01-01 00:00:00+00', - NULL, NULL, current_timestamp)") - -started <- DBI::dbGetQuery(con, "SELECT ca_coordinator_start(10, 16) AS started")$started[[1L]] -stopifnot(isTRUE(started)) -deadline <- proc.time()[["elapsed"]] + 10 -repeat { - state <- DBI::dbGetQuery(con, - "SELECT state FROM canard_absurd.tasks WHERE id = 'expired'")$state[[1L]] - if (identical(state, "failed")) break - if (proc.time()[["elapsed"]] >= deadline) stop("Coordinator did not reap the expired task") - Sys.sleep(0.01) +assert_query_error <- function(con, sql, message) { + error <- tryCatch(DBI::dbGetQuery(con, sql), error = identity) + stopifnot(inherits(error, "error")) + stopifnot(grepl(message, conditionMessage(error), fixed = TRUE)) } -status <- DBI::dbGetQuery(con, " - SELECT (ca_coordinator_status()).state AS state, - (ca_coordinator_status()).poll_count AS poll_count, - (ca_coordinator_status()).reaped_count AS reaped_count, - (ca_coordinator_status()).last_error AS last_error") -stopifnot(identical(status$state[[1L]], "running")) -stopifnot(as.numeric(status$poll_count[[1L]]) >= 1) -stopifnot(as.numeric(status$reaped_count[[1L]]) == 1) -stopifnot(identical(status$last_error[[1L]], "")) -stopifnot(identical( - DBI::dbGetQuery(con, "SELECT ca_coordinator_start(10, 16) AS started")$started[[1L]], - FALSE)) -stopifnot(isTRUE(DBI::dbGetQuery(con, "SELECT ca_coordinator_stop() AS stopped")$stopped[[1L]])) -stopifnot(identical( - DBI::dbGetQuery(con, "SELECT (ca_coordinator_status()).state AS state")$state[[1L]], - "closed")) -restart <- try(DBI::dbGetQuery(con, "SELECT ca_coordinator_start(10, 16)"), silent = TRUE) -stopifnot(inherits(restart, "try-error"), grepl("cannot restart", as.character(restart), fixed = TRUE)) + +read_status <- function(con) { + DBI::dbGetQuery(con, " + WITH value AS MATERIALIZED (SELECT ca_coordinator_status() AS status) + SELECT status.state, status.poll_count, status.reaped_count, status.last_error + FROM value + ") +} + +create_fixture <- function(path) { + # Use the installed package's schema and native input encoding. + db <- CanardAbsurd::ca_open(path) + on.exit(CanardAbsurd::ca_close(db)) + + tasks <- data.frame( + id = c("expired-a", "expired-b", "expired-c", "retryable", "live", + "ready", "completed", "failed", "cancelled"), + queue = c("alpha", "beta", "alpha", rep("other", 6)), + max_failures = c(1L, 3L, 2L, 3L, rep(1L, 5)) + ) + for (index in seq_len(nrow(tasks))) { + CanardAbsurd::ca_spawn( + db, + name = "coordinator-test", + input = list(sample = tasks$id[[index]]), + id = tasks$id[[index]], + queue = tasks$queue[[index]], + max_failures = tasks$max_failures[[index]] + ) + } + + # Inject expired/live claims without sleeping through actual lease durations. + DBI::dbExecute(db@con, " + UPDATE canard_absurd.tasks AS task + SET state = 'running', attempt = fixture.failures + 1, + failures = fixture.failures, worker = 'test-worker', token = uuid(), + lease_until = fixture.deadline + FROM (VALUES + ('expired-a', 0, TIMESTAMPTZ '2000-01-01 00:00:00+00'), + ('expired-b', 2, TIMESTAMPTZ '2000-01-02 00:00:00+00'), + ('expired-c', 1, TIMESTAMPTZ '2000-01-03 00:00:00+00'), + ('retryable', 0, TIMESTAMPTZ '2000-01-01 00:00:00+00'), + ('live', 0, current_timestamp + INTERVAL '1 hour') + ) AS fixture(id, failures, deadline) + WHERE task.id = fixture.id + ") + DBI::dbExecute(db@con, " + UPDATE canard_absurd.tasks + SET state = id, attempt = 1, + failures = CASE WHEN id = 'failed' THEN 1 ELSE 0 END + WHERE id IN ('completed', 'failed', 'cancelled') + ") +} + +test_reaper <- function(extension) { + path <- tempfile(fileext = ".duckdb") + create_fixture(path) + driver <- duckdb::duckdb(path, config = list( + allow_unsigned_extensions = "true", + autoinstall_known_extensions = "false" + )) + con <- DBI::dbConnect(driver) + loaded <- FALSE + on.exit({ + tryCatch({ + if (loaded) { + DBI::dbGetQuery(con, "SELECT ca_coordinator_stop()") + } + }, finally = { + DBI::dbDisconnect(con, shutdown = TRUE) + }) + }) + + DBI::dbExecute(con, " + CREATE TEMP TABLE before_poll AS SELECT * FROM canard_absurd.tasks + ") + DBI::dbExecute(con, paste("LOAD", DBI::dbQuoteString(con, extension))) + loaded <- TRUE + + assert_query_error(con, "SELECT ca_coordinator_start(NULL, 1)", "cannot be NULL") + assert_query_error(con, "SELECT ca_coordinator_start(10, NULL)", "cannot be NULL") + assert_query_error(con, "SELECT ca_coordinator_start(0, 1)", "poll_milliseconds") + assert_query_error(con, "SELECT ca_coordinator_start(10, 0)", "reap_limit") + + started <- DBI::dbGetQuery(con, + "SELECT ca_coordinator_start(10, 1) AS started")$started[[1L]] + stopifnot(isTRUE(started)) + + deadline <- proc.time()[["elapsed"]] + 10 + repeat { + status <- read_status(con) + if (nzchar(status$last_error[[1L]])) { + stop("Coordinator poll failed: ", status$last_error[[1L]]) + } + if (as.numeric(status$reaped_count[[1L]]) == 3) { + break + } + if (proc.time()[["elapsed"]] >= deadline) { + stop("Coordinator did not reap the three exhausted tasks") + } + Sys.sleep(0.01) + } + stopifnot(identical(status$state[[1L]], "running")) + stopifnot(as.numeric(status$poll_count[[1L]]) >= 3) + + started_again <- DBI::dbGetQuery(con, + "SELECT ca_coordinator_start(10, 1) AS started")$started[[1L]] + stopifnot(identical(started_again, FALSE)) + stopped <- DBI::dbGetQuery(con, + "SELECT ca_coordinator_stop() AS stopped")$stopped[[1L]] + stopifnot(isTRUE(stopped)) + + reaped <- DBI::dbGetQuery(con, " + SELECT task.id, task.state, task.failures, task.attempt, task.error, + task.worker IS NULL AND task.token IS NULL + AND task.lease_until IS NULL AS ownership_cleared, + task.input IS NOT DISTINCT FROM before.input AS input_preserved, + task.input_rtype IS NOT DISTINCT FROM before.input_rtype AS type_preserved, + task.created_at = before.created_at AS created_at_preserved + FROM canard_absurd.tasks AS task + JOIN before_poll AS before USING (id) + WHERE task.id LIKE 'expired-%' + ORDER BY task.id + ") + stopifnot(identical(reaped$id, c("expired-a", "expired-b", "expired-c"))) + stopifnot(all(reaped$state == "failed")) + stopifnot(identical(reaped$failures, c(1L, 3L, 2L))) + stopifnot(identical(reaped$attempt, c(1L, 3L, 2L))) + stopifnot(all(reaped$error == "worker lease expired")) + stopifnot(all(reaped$ownership_cleared)) + stopifnot(all(reaped$input_preserved), all(reaped$type_preserved)) + stopifnot(all(reaped$created_at_preserved)) + + changed <- DBI::dbGetQuery(con, " + SELECT * FROM canard_absurd.tasks WHERE id NOT LIKE 'expired-%' + EXCEPT + SELECT * FROM before_poll WHERE id NOT LIKE 'expired-%' + ") + stopifnot(nrow(changed) == 0L) + status <- read_status(con) + stopifnot(identical(status$state[[1L]], "closed")) + stopifnot(as.numeric(status$reaped_count[[1L]]) == 3) + stopifnot(identical(status$last_error[[1L]], "")) + assert_query_error(con, "SELECT ca_coordinator_start(10, 1)", "cannot restart") + stopped_again <- DBI::dbGetQuery(con, + "SELECT ca_coordinator_stop() AS stopped")$stopped[[1L]] + stopifnot(identical(stopped_again, FALSE)) +} + +main <- function() { + args <- commandArgs(trailingOnly = TRUE) + if (length(args) != 1L) { + stop("Usage: Rscript test-v1.R /path/to/canard_coordinator.duckdb_extension") + } + extension <- normalizePath(args[[1L]], mustWork = TRUE) + test_reaper(extension) +} + +main()