Skip to content
Closed
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
26 changes: 16 additions & 10 deletions tools/canard-coordinator/src/canard_coordinator.c
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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)) {
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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);
Expand Down
184 changes: 130 additions & 54 deletions tools/canard-coordinator/test/test-v1.R
Original file line number Diff line number Diff line change
@@ -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)
Loading