Skip to content
Merged
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
6 changes: 5 additions & 1 deletion .github/workflows/R-CMD-check.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
46 changes: 28 additions & 18 deletions tools/canard-coordinator/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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.
56 changes: 34 additions & 22 deletions tools/canard-coordinator/src/canard_coordinator.c
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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)) {
Expand All @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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;
}
Loading