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
37 changes: 22 additions & 15 deletions README.Rmd
Original file line number Diff line number Diff line change
Expand Up @@ -105,9 +105,12 @@ flowchart LR
LocalR["your R process"] -->|"ca_open()"| LocalDB[("DuckDB")]
end
subgraph Shared["separate producer and worker processes"]
Producer["producer"] -->|"ca_connect() over Quack"| Server["server<br/>ca_serve()"]
Producer["producer"] -->|"ca_connect() over Quack"| Server
Worker["worker"] -->|"ca_connect() over Quack"| Server
Server -->|"only file owner"| SharedDB[("DuckDB file")]
subgraph Owner["database-owner process"]
Server["Quack server<br/>ca_serve()"] --> SharedDB[("DuckDB file")]
Maintenance["optional native<br/>lease maintenance"] --> SharedDB
end
end
```

Expand All @@ -117,19 +120,23 @@ in workers, not in the server.

## Optional native lease maintenance

The optional `canard_coordinator` extension runs bounded expired-lease
maintenance in the database-owner process, independently of worker claims. It
targets stable DuckDB C API v1.2 and uses a dedicated connection with pending
query execution. `ca_coordinator_start()` starts the loop,
`ca_coordinator_status()` reports its state and counters, and
`ca_coordinator_stop()` joins the native thread and releases its connection.
`ca_close()` performs that shutdown automatically.

The v1 coordinator is one-shot for each database handle: it cannot restart after
its dedicated connection closes. It does not supervise external jobs, reserve
CPU or memory, or publish artifacts. See the
[native extension instructions](tools/canard-coordinator/README.md) for the
build, signing, and lifecycle contract.
If a worker disappears on its last allowed attempt, its task must eventually be
marked failed. Workers perform this cleanup when they poll a queue. The optional
`canard_coordinator` extension performs the same cleanup across all queues,
including queues with no active workers.

It runs a service thread and dedicated connection in the database-owner process.
It does **not** execute handlers, renew worker leases, check dependencies, or
launch downstream jobs. Workers still claim and run tasks through `ca_work()`.

Start it explicitly with `ca_coordinator_start()`, inspect progress and errors
with `ca_coordinator_status()`, and stop it with `ca_coordinator_stop()`.
`ca_close()` stops a coordinator started through that handle. A stopped instance
cannot restart; start a fresh database instance instead. The extension is built
separately; package loading does not install or start it.

See [lease maintenance and recovery](https://rgenomicsetl.github.io/CanardAbsurd/articles/durability.html#native-lease-maintenance)
and the [extension build instructions](https://github.com/RGenomicsETL/CanardAbsurd/tree/main/tools/canard-coordinator).

## Where repeat work can happen

Expand Down
41 changes: 26 additions & 15 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -117,9 +117,12 @@ flowchart LR
LocalR["your R process"] -->|"ca_open()"| LocalDB[("DuckDB")]
end
subgraph Shared["separate producer and worker processes"]
Producer["producer"] -->|"ca_connect() over Quack"| Server["server<br/>ca_serve()"]
Producer["producer"] -->|"ca_connect() over Quack"| Server
Worker["worker"] -->|"ca_connect() over Quack"| Server
Server -->|"only file owner"| SharedDB[("DuckDB file")]
subgraph Owner["database-owner process"]
Server["Quack server<br/>ca_serve()"] --> SharedDB[("DuckDB file")]
Maintenance["optional native<br/>lease maintenance"] --> SharedDB
end
end
```

Expand All @@ -129,19 +132,27 @@ through Quack. Handlers run in workers, not in the server.

## Optional native lease maintenance

The optional `canard_coordinator` extension runs bounded expired-lease
maintenance in the database-owner process, independently of worker
claims. It targets stable DuckDB C API v1.2 and uses a dedicated
connection with pending query execution. `ca_coordinator_start()` starts
the loop, `ca_coordinator_status()` reports its state and counters, and
`ca_coordinator_stop()` joins the native thread and releases its
connection. `ca_close()` performs that shutdown automatically.

The v1 coordinator is one-shot for each database handle: it cannot
restart after its dedicated connection closes. It does not supervise
external jobs, reserve CPU or memory, or publish artifacts. See the
[native extension instructions](tools/canard-coordinator/README.md) for
the build, signing, and lifecycle contract.
If a worker disappears on its last allowed attempt, its task must
eventually be marked failed. Workers perform this cleanup when they poll
a queue. The optional `canard_coordinator` extension performs the same
cleanup across all queues, including queues with no active workers.

It runs a service thread and dedicated connection in the database-owner
process. It does **not** execute handlers, renew worker leases, check
dependencies, or launch downstream jobs. Workers still claim and run
tasks through `ca_work()`.

Start it explicitly with `ca_coordinator_start()`, inspect progress and
errors with `ca_coordinator_status()`, and stop it with
`ca_coordinator_stop()`. `ca_close()` stops a coordinator started
through that handle. A stopped instance cannot restart; start a fresh
database instance instead. The extension is built separately; package
loading does not install or start it.

See [lease maintenance and
recovery](https://rgenomicsetl.github.io/CanardAbsurd/articles/durability.html#native-lease-maintenance)
and the [extension build
instructions](https://github.com/RGenomicsETL/CanardAbsurd/tree/main/tools/canard-coordinator).

## Where repeat work can happen

Expand Down
5 changes: 4 additions & 1 deletion inst/tinytest/test_laws.R
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,10 @@ local({
current <- ca_claim(db)
}
rejected <- vapply(stale, function(task) {
tryCatch({ ca_complete(task, "stale"); FALSE },
tryCatch({
ca_complete(task, "stale")
FALSE
},
canard_lease_lost = function(e) TRUE)
}, logical(1L))
ca_complete(current, "current")
Expand Down
5 changes: 4 additions & 1 deletion inst/tinytest/test_payloads.R
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,10 @@ local({
id <- ca_spawn(db, "work", value)
task <- ca_claim(db)
calls <- 0L
fn <- function() { calls <<- calls + 1L; value }
fn <- function() {
calls <<- calls + 1L
value
}
expect_identical(task@input, value)
expect_identical(ca_step(task, key, fn), value)
stored <- DBI::dbGetQuery(db@con,
Expand Down
5 changes: 4 additions & 1 deletion inst/tinytest/test_protocol.R
Original file line number Diff line number Diff line change
Expand Up @@ -91,7 +91,10 @@ if (requireNamespace("s7contract", quietly = TRUE)) local({
tryCatch({
ca_spawn(db, "work", id = "protocol", max_failures = 100L)
db
}, error = function(error) { ca_close(db); stop(error) })
}, error = function(error) {
ca_close(db)
stop(error)
})
},
teardown = ca_close,
classify = function(sequence) unique(vapply(sequence, `[[`, character(1L), "command")),
Expand Down
5 changes: 4 additions & 1 deletion inst/tinytest/test_quack.R
Original file line number Diff line number Diff line change
Expand Up @@ -152,7 +152,10 @@ local({
task <- ca_claim(db, lease_seconds = 0.5)
file.create(file.path(directory, "old-claim"))
while (!file.exists(file.path(directory, "release"))) Sys.sleep(0.01)
tryCatch({ ca_complete(task, "late"); "accepted" },
tryCatch({
ca_complete(task, "late")
"accepted"
},
canard_lease_lost = function(e) "fenced")
}, args = list(fixture$uri, fixture$directory), libpath = .libPaths(), supervise = TRUE)
withr::defer(if (late$is_alive()) late$kill())
Expand Down
15 changes: 12 additions & 3 deletions inst/tinytest/test_values.R
Original file line number Diff line number Diff line change
Expand Up @@ -102,7 +102,10 @@ local({
db <- local_database()
query <- db@query
queries <- 0L
db@query <- function(sql) { queries <<- queries + 1L; query(sql) }
db@query <- function(sql) {
queries <<- queries + 1L
query(sql)
}
for (value in list(identity, globalenv(), quote(x + y), matrix(1:4, 2L),
structure(1, class = "custom"), setNames(list(1, 2), c("x", "x")),
as.Date(NaN, origin = "1970-01-01"),
Expand All @@ -126,7 +129,10 @@ for (remote in c(FALSE, TRUE)) local({
# Case-colliding fields fail admission without issuing submission SQL.
query <- db@query
queries <- 0L
db@query <- function(sql) { queries <<- queries + 1L; query(sql) }
db@query <- function(sql) {
queries <<- queries + 1L
query(sql)
}
for (value in case_collisions) {
expect_error(ca_spawn(db, "invalid", value), class = "canard_value_error")
}
Expand All @@ -138,7 +144,10 @@ for (remote in c(FALSE, TRUE)) local({
id <- ca_spawn(db, "invalid", max_failures = 1L)
calls <- 0L
outcome <- ca_run(ca_claim(db), function(input, task) {
produce <- function() { calls <<- calls + 1L; value }
produce <- function() {
calls <<- calls + 1L
value
}
if (operation == "step") ca_step(task, "invalid", produce) else produce()
})
expect_identical(calls, 1L)
Expand Down
28 changes: 28 additions & 0 deletions vignettes/articles/quack-server.Rmd
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,16 @@ endpoint. Production deployment separates the server and
workers as described below. This network-dependent article is built by pkgdown,
while the package vignettes run offline.

```mermaid
flowchart LR
Producer["API / producer"] -->|"ca_connect()"| Quack
Worker["R worker<br/>ca_work()"] -->|"ca_connect()"| Quack
subgraph Owner["one database-owner process"]
Quack["Quack listener<br/>ca_serve()"] --> Database[("DuckDB file")]
Maintenance["optional native<br/>lease maintenance"] --> Database
end
```

## Install the runtime

Install DuckDB's Quack extension explicitly before serving or connecting. From
Expand All @@ -40,6 +50,24 @@ The extensions are stored in DuckDB's shared extension directory so later R
processes using the same DuckDB runtime can load them. Package loading and
connections never download extensions.

## Optional maintenance belongs on the owner

The `canard_coordinator` extension is separate from Quack. If it is available in
the deployment, start it explicitly with `ca_coordinator_start()` on the owner
handle returned by `ca_serve()`. A `ca_connect()` client is not its lifecycle
owner. The ordinary server/worker path below does not require the extension.

The native loop marks expired, exhausted attempts failed across all queues, even
when no worker is polling them. It neither runs R handlers nor monitors external
tools or downstream dependencies. Inspect progress and errors with
`ca_coordinator_status()`. Stop it before database teardown;
`ca_close()` does so for a coordinator started through that handle. Stopping is
one-shot: a fresh database instance is required to start another loop.

See [durability and recovery](durability.html#native-lease-maintenance) for the
maintenance contract and [extension instructions](https://github.com/RGenomicsETL/CanardAbsurd/tree/main/tools/canard-coordinator)
for building and signing. Installing Quack does not install this extension.

## Run the server/client path

The demonstration uses a disposable localhost token. Use a securely supplied
Expand Down
66 changes: 55 additions & 11 deletions vignettes/durability.Rmd
Original file line number Diff line number Diff line change
Expand Up @@ -69,9 +69,47 @@ explicit decision, the error propagates.
Recovery replaces an expired lease when the failure budget allows another
attempt. The `attempt` counter counts claims, including sleep resumptions. The
`failures` counter counts reported failures and expired leases. `max_failures`
is the threshold for terminal failure. Reaping happens on the claim path, so an
expired last attempt is marked failed when workers next poll its queue.
`reap_limit` controls the size of that maintenance batch; its default is 64 rows.
is the threshold for terminal failure. Before claiming, `ca_claim()` and
`ca_work()` mark a bounded batch of expired, exhausted tasks as failed in their
selected queue. `reap_limit` defaults to 64. An expired task with failure budget
remaining can instead be reclaimed by a worker; expiry does not by itself start
another attempt.

## Native lease maintenance

A queue with no polling workers still needs cleanup when its final attempt
expires. The optional `canard_coordinator` extension performs that cleanup across
all queues, on a dedicated connection and native thread in the database-owner
process. It increments the failure count, records `Lease expired`, clears the
worker/token/deadline, and marks an exhausted task failed in one SQL statement.

```mermaid
flowchart LR
Poll["worker polls its queue"] --> Reap["mark expired final attempts failed"]
Native["optional native maintenance<br/>polls all queues"] --> Reap
Reap --> DB[("DuckDB task rows")]
Poll --> Claim["claim eligible work<br/>or reclaim an expired retryable task"]
Claim --> Handler["R worker runs the handler"]
```

Use `ca_coordinator_start()` on the `ca_open()` or `ca_serve()` owner handle, not
a remote `ca_connect()` client. `poll_milliseconds` sets the maintenance interval;
`reap_limit` bounds the tasks changed by each poll, not the amount of table data
DuckDB may scan. The native loop does not replace worker polling and does not
renew a worker's lease, run a handler, enforce CPU/memory admission, or decide
which downstream tasks are ready.

`ca_coordinator_status()` reports state, poll/reaped/error counters, and the last
error. Starting the loop is not proof that maintenance SQL succeeds; inspect
those counters and errors. `ca_coordinator_stop()` joins the thread and releases
its connection. `ca_close()` does this automatically when the handle started the
coordinator. The instance is one-shot: after stop, use a fresh database instance
to start another coordinator.

The extension is built separately against the stable DuckDB C extension API.
Loading and starting it are explicit; opening an ordinary database does neither.
See the [extension instructions](https://github.com/RGenomicsETL/CanardAbsurd/tree/main/tools/canard-coordinator)
for build, signing, and executable native integration tests.

## A replaced claim cannot write

Expand Down Expand Up @@ -170,12 +208,16 @@ objects. See `?ca_conditions` for the condition fields and retry notifications.

## Error metadata compatibility

Two upstream paths discard structured information:
Error preservation depends on the installed R driver and Quack build. Affected
R wrappers discard the original condition's fields and cause chain; affected
Quack bind paths turn the transported structured error into an
`InvalidInputException`.

- [DuckDB R #2711](https://github.com/duckdb/duckdb-r/issues/2711): the DBI rethrow
discards the original error fields and parent.
- [Quack PR #212](https://github.com/duckdb/duckdb-quack/pull/212): the proposed
protocol change preserves the server's error type and transferable metadata.
The R wrapper report is [duckdb-r #2711](https://github.com/duckdb/duckdb-r/issues/2711).
The fixes in [duckdb-r #2714](https://github.com/duckdb/duckdb-r/pull/2714) and
[duckdb-quack #212](https://github.com/duckdb/duckdb-quack/pull/212) are merged.
Their merge status alone does not establish that an installed driver/extension
pair includes both fixes.

CanardAbsurd first uses a structured DuckDB transaction error when available.
Otherwise, one compatibility adapter recognizes exact tested write-conflict
Expand All @@ -188,9 +230,11 @@ supplied it. Structured non-transaction types take precedence over the text.
This is a workaround for released drivers, **not a stable error protocol**.
Changed or unknown message forms propagate without a retry restart, as do
transport failures. An unhandled recognized conflict also propagates: retry
limits and delays still belong to the caller. Revalidate against the DuckDB/Quack pair in deployment when upgrading; the integration suite uses installed packages
and independent R processes. Tests exercise DuckDB R 1.5.3 and 1.5.5 and the
community Quack builds available for them.
limits and delays still belong to the caller. Revalidate the DuckDB/Quack pair
when upgrading; the integration suite uses installed packages and independent
R processes. The package requires DuckDB R 1.5.5 or newer. Historical probes on
1.5.3 explain the compatibility adapter, but that version is not a supported
package runtime.

## External effects still need idempotency

Expand Down
25 changes: 25 additions & 0 deletions vignettes/getting-started.Rmd
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,31 @@ A zero-second sleep demonstrates the state transition without delaying the
render. Positive durations use the database clock. Sleeps do not consume the
failure budget.

## Several samples, then a cohort

For genomics work, one sample task could have named steps for alignment and gVCF
calling. Two samples are two tasks that separate workers can run concurrently.
Joint genotyping is another task, submitted when the required sample outputs are
ready.

```mermaid
flowchart LR
Sample1["sample-101 task<br/>align → call gVCF"] --> Saved1[("DuckDB<br/>sample-101 completed<br/>gVCF location")]
Sample2["sample-102 task<br/>align → call gVCF"] --> Saved2[("DuckDB<br/>sample-102 completed<br/>gVCF location")]
Saved1 --> Controller["application or targets<br/>checks cohort readiness"]
Saved2 --> Controller
Controller --> Joint["submit joint genotyping task"]
```

These arrows describe application orchestration, not a dependency API supplied
by CanardAbsurd. Sequential `ca_step()` calls inside one sample handler are not
independently scheduled jobs. A stored gVCF path also does not store the file:
workers need access to the outputs and compatible execution environments.

The optional native coordinator handles expired final attempts. It does not
supply the cohort-readiness check in this diagram. See
[native lease maintenance](durability.html#native-lease-maintenance).

## Choose a worker lifetime

- `max_tasks = 1` executes one claimed attempt, including a failure or suspension.
Expand Down
16 changes: 16 additions & 0 deletions vignettes/native-values.Rmd
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,22 @@ compatibility v1.5.0. This is separate from CanardAbsurd's own
schema version, which remains 1. These unreleased schema-1 files require matching
package code; no schema migration is provided.

## Driver validation limits

DuckDB R 1.5.5 can return malformed type descriptors under garbage-collection
pressure while materializing VARIANT values. A local Linux probe reproduces
corruption before CanardAbsurd restores the value. The missing-descriptor error
in this [Windows package check](https://github.com/RGenomicsETL/CanardAbsurd/actions/runs/35221419601/job/105202242080)
is consistent with that defect; the Windows run itself has not been reproduced
locally.

This can make inspection fail even though a task's SQL transition committed.
The typed projections described below address nested-NULL conversion, but type
descriptors still pass through the driver's generic VARIANT reader. Neither
those projections nor a successful example establishes that this GC-sensitive
path is safe. A driver-side allocation-protection fix is being tested; the
minimum package version alone does not establish its presence.

## Field-name admission

Named lists and data frame columns require nonempty names that are unique
Expand Down
Loading