diff --git a/tools/canard-coordinator/src/canard_coordinator.c b/tools/canard-coordinator/src/canard_coordinator.c index b4fe641..736cf59 100644 --- a/tools/canard-coordinator/src/canard_coordinator.c +++ b/tools/canard-coordinator/src/canard_coordinator.c @@ -12,17 +12,20 @@ extern duckdb_ext_api_v1 duckdb_ext_api; +/* Same failure 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 " + "SET state = 'failed', failures = failures + 1, " + "error = 'worker lease expired', " + "worker = NULL, token = NULL, lease_until = NULL, " + "updated_at = current_timestamp " "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"; + ") RETURNING 1 AS changed, id"; typedef enum { CANARD_COORDINATOR_READY = 0, @@ -170,16 +173,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_state = DuckDBSuccess; + if (!statement) { + poll_state = canard_prepare_reaper(runtime, &statement, error); + } + if (poll_state == DuckDBSuccess) { + poll_state = canard_execute_reaper(runtime, statement, &reaped, error); + } + if (poll_state == 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)) { @@ -339,7 +345,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 +377,7 @@ 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); + 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); diff --git a/tools/canard-coordinator/test/test-v1.R b/tools/canard-coordinator/test/test-v1.R index e624fb8..848bde2 100644 --- a/tools/canard-coordinator/test/test-v1.R +++ b/tools/canard-coordinator/test/test-v1.R @@ -1,56 +1,132 @@ -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) +# Run from any working directory; the schema belongs to the package, not this test. +script_argument <- grep("^--file=", commandArgs(), value = TRUE) +script_path <- normalizePath(sub("^--file=", "", script_argument), mustWork = TRUE) +package_root <- normalizePath(file.path(dirname(script_path), "../../.."), mustWork = TRUE) -driver <- duckdb::duckdb(tempfile(fileext = ".duckdb"), shared_home = FALSE, - 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) +coordinator_status <- function(con) { + DBI::dbGetQuery(con, " + WITH snapshot AS MATERIALIZED (SELECT ca_coordinator_status() AS status) + SELECT status.state, status.poll_count, status.reaped_count, status.last_error + FROM snapshot") } -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)) + +expect_sql_error <- function(con, sql, message) { + error <- tryCatch(DBI::dbGetQuery(con, sql), error = identity) + stopifnot(inherits(error, "error")) + stopifnot(grepl(message, conditionMessage(error), fixed = TRUE)) +} + +test_coordinator <- function(extension, package_root) { + driver <- duckdb::duckdb(tempfile(fileext = ".duckdb"), shared_home = FALSE, + config = list(allow_unsigned_extensions = "true", autoinstall_known_extensions = "false")) + con <- DBI::dbConnect(driver) + on.exit(DBI::dbDisconnect(con, shutdown = TRUE), add = TRUE) + + schema <- readLines(file.path(package_root, "inst/sql/schema.sql"), warn = FALSE) + statements <- strsplit(paste(schema, collapse = "\n"), ";", fixed = TRUE)[[1L]] + for (statement in statements) { + if (nzchar(trimws(statement))) { + DBI::dbExecute(con, statement) + } + } + + DBI::dbExecute(con, paste("LOAD", DBI::dbQuoteString(con, extension))) + # Stop before disconnecting, including when an assertion below fails. + on.exit(DBI::dbGetQuery(con, "SELECT ca_coordinator_stop()"), add = TRUE, after = FALSE) + + expect_sql_error(con, "SELECT ca_coordinator_start(NULL, 16)", "cannot be NULL") + expect_sql_error(con, "SELECT ca_coordinator_start(10, NULL)", "cannot be NULL") + expect_sql_error(con, "SELECT ca_coordinator_start(0, 16)", "poll_milliseconds must be") + expect_sql_error(con, "SELECT ca_coordinator_start(10, 0)", "reap_limit must be") + stopifnot(identical(coordinator_status(con)$state[[1L]], "ready")) + + # Only the two exhausted, expired claims should be reaped. Use actual lease + # ownership, counters, and native payloads so the package constraints apply. + DBI::dbExecute(con, " + INSERT INTO canard_absurd.tasks + (id, queue, name, input, input_rtype, state, attempt, failures, max_failures, + worker, token, lease_until, updated_at) + SELECT id, queue, 'example', 17::VARIANT, + struct_pack(kind := 'integer', length := 1)::VARIANT, + 'running', failures + 1, failures, max_failures, 'test-worker', uuid(), + lease_until, TIMESTAMPTZ '1999-01-01 00:00:00+00' + FROM (VALUES + ('expired-default', 'default', 0, 1, TIMESTAMPTZ '2000-01-01 00:00:00+00'), + ('expired-other', 'other', 2, 3, TIMESTAMPTZ '2000-01-01 00:00:00+00'), + ('retryable', 'default', 0, 3, TIMESTAMPTZ '2000-01-01 00:00:00+00'), + ('live', 'default', 0, 1, TIMESTAMPTZ '2100-01-01 00:00:00+00') + ) AS fixture(id, queue, failures, max_failures, lease_until)") + DBI::dbExecute(con, " + INSERT INTO canard_absurd.tasks (id, queue, name, input_rtype, state, max_failures) + SELECT state, 'default', 'example', + struct_pack(kind := 'NULL', length := 0)::VARIANT, state, 1 + FROM (VALUES ('ready'), ('completed'), ('failed'), ('cancelled')) AS fixture(state)") + DBI::dbExecute(con, " + UPDATE canard_absurd.tasks + SET checkpoints = map(['saved'], [struct_pack( + kind := 'step', value := 42::VARIANT, + rtype := struct_pack(kind := 'integer', length := 1)::VARIANT)]) + WHERE id LIKE 'expired-%'") + DBI::dbExecute(con, "CREATE TABLE before_reaping AS SELECT * FROM canard_absurd.tasks") + + started <- DBI::dbGetQuery(con, "SELECT ca_coordinator_start(10, 1) AS started")$started[[1L]] + stopifnot(isTRUE(started)) + deadline <- proc.time()[["elapsed"]] + 10 + repeat { + status <- coordinator_status(con) + if (nzchar(status$last_error[[1L]])) { + stop("Coordinator poll failed: ", status$last_error[[1L]]) + } + if (as.numeric(status$reaped_count[[1L]]) == 2) { + break + } + if (proc.time()[["elapsed"]] >= deadline) { + stop("Coordinator did not reap the exhausted claims") + } + Sys.sleep(0.01) + } + stopifnot(identical(status$state[[1L]], "running")) + stopifnot(as.numeric(status$poll_count[[1L]]) >= 2) + stopifnot(identical( + DBI::dbGetQuery(con, "SELECT ca_coordinator_start(10, 1) AS started")$started[[1L]], + FALSE)) + stopifnot(isTRUE(DBI::dbGetQuery(con, "SELECT ca_coordinator_stop() AS stopped")$stopped[[1L]])) + + # Match the complete inst/sql/reap.sql transition, not merely its state label. + reaped <- DBI::dbGetQuery(con, " + SELECT current.id, + current.state = 'failed' AS failed, + current.failures = previous.failures + 1 AS failure_counted, + current.error = 'worker lease expired' AS reason_recorded, + current.worker IS NULL AND current.token IS NULL + AND current.lease_until IS NULL AS ownership_cleared, + current.updated_at > previous.updated_at AS timestamp_updated, + current.attempt = previous.attempt + AND current.input IS NOT DISTINCT FROM previous.input + AND current.input_rtype IS NOT DISTINCT FROM previous.input_rtype + AND current.checkpoints IS NOT DISTINCT FROM previous.checkpoints AS work_preserved + FROM canard_absurd.tasks AS current + JOIN before_reaping AS previous USING (id) + WHERE current.id LIKE 'expired-%'") + stopifnot(nrow(reaped) == 2L) + stopifnot(all(as.matrix(reaped[, -1L, drop = FALSE]))) + + changed_unrelated <- DBI::dbGetQuery(con, " + SELECT * FROM canard_absurd.tasks WHERE id NOT LIKE 'expired-%' + EXCEPT + SELECT * FROM before_reaping WHERE id NOT LIKE 'expired-%'") + stopifnot(nrow(changed_unrelated) == 0L) + stopifnot(identical(coordinator_status(con)$state[[1L]], "closed")) + stopifnot(as.numeric(coordinator_status(con)$reaped_count[[1L]]) == 2) + stopifnot(identical( + DBI::dbGetQuery(con, "SELECT ca_coordinator_stop() AS stopped")$stopped[[1L]], + FALSE)) + expect_sql_error(con, "SELECT ca_coordinator_start(10, 1)", "cannot restart") +} + +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_coordinator(extension, package_root)