From 1e0488ec66e03953759ca99e6f3e622517f7a430 Mon Sep 17 00:00:00 2001 From: sounkou-bioinfo Date: Sat, 19 Sep 2026 14:17:36 +0200 Subject: [PATCH] sync --- README.Rmd | 37 +++++++++------- README.md | 41 +++++++++++------- inst/tinytest/test_laws.R | 5 ++- inst/tinytest/test_payloads.R | 5 ++- inst/tinytest/test_protocol.R | 5 ++- inst/tinytest/test_quack.R | 5 ++- inst/tinytest/test_values.R | 15 +++++-- vignettes/articles/quack-server.Rmd | 28 ++++++++++++ vignettes/durability.Rmd | 66 ++++++++++++++++++++++++----- vignettes/getting-started.Rmd | 25 +++++++++++ vignettes/native-values.Rmd | 16 +++++++ 11 files changed, 200 insertions(+), 48 deletions(-) diff --git a/README.Rmd b/README.Rmd index d0dd1e4..b93b330 100644 --- a/README.Rmd +++ b/README.Rmd @@ -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
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
ca_serve()"] --> SharedDB[("DuckDB file")] + Maintenance["optional native
lease maintenance"] --> SharedDB + end end ``` @@ -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 diff --git a/README.md b/README.md index f9cc32b..726200a 100644 --- a/README.md +++ b/README.md @@ -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
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
ca_serve()"] --> SharedDB[("DuckDB file")] + Maintenance["optional native
lease maintenance"] --> SharedDB + end end ``` @@ -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 diff --git a/inst/tinytest/test_laws.R b/inst/tinytest/test_laws.R index 4209af7..2461fa0 100644 --- a/inst/tinytest/test_laws.R +++ b/inst/tinytest/test_laws.R @@ -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") diff --git a/inst/tinytest/test_payloads.R b/inst/tinytest/test_payloads.R index 110b9f6..32529c6 100644 --- a/inst/tinytest/test_payloads.R +++ b/inst/tinytest/test_payloads.R @@ -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, diff --git a/inst/tinytest/test_protocol.R b/inst/tinytest/test_protocol.R index 4c00813..2de2583 100644 --- a/inst/tinytest/test_protocol.R +++ b/inst/tinytest/test_protocol.R @@ -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")), diff --git a/inst/tinytest/test_quack.R b/inst/tinytest/test_quack.R index 5495fb2..adbeef3 100644 --- a/inst/tinytest/test_quack.R +++ b/inst/tinytest/test_quack.R @@ -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()) diff --git a/inst/tinytest/test_values.R b/inst/tinytest/test_values.R index db0b27e..9d1bae0 100644 --- a/inst/tinytest/test_values.R +++ b/inst/tinytest/test_values.R @@ -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"), @@ -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") } @@ -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) diff --git a/vignettes/articles/quack-server.Rmd b/vignettes/articles/quack-server.Rmd index 40d0552..25f386f 100644 --- a/vignettes/articles/quack-server.Rmd +++ b/vignettes/articles/quack-server.Rmd @@ -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
ca_work()"] -->|"ca_connect()"| Quack + subgraph Owner["one database-owner process"] + Quack["Quack listener
ca_serve()"] --> Database[("DuckDB file")] + Maintenance["optional native
lease maintenance"] --> Database + end +``` + ## Install the runtime Install DuckDB's Quack extension explicitly before serving or connecting. From @@ -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 diff --git a/vignettes/durability.Rmd b/vignettes/durability.Rmd index 732ade0..b6ff7e8 100644 --- a/vignettes/durability.Rmd +++ b/vignettes/durability.Rmd @@ -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
polls all queues"] --> Reap + Reap --> DB[("DuckDB task rows")] + Poll --> Claim["claim eligible work
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 @@ -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 @@ -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 diff --git a/vignettes/getting-started.Rmd b/vignettes/getting-started.Rmd index bac95eb..f815f91 100644 --- a/vignettes/getting-started.Rmd +++ b/vignettes/getting-started.Rmd @@ -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
align → call gVCF"] --> Saved1[("DuckDB
sample-101 completed
gVCF location")] + Sample2["sample-102 task
align → call gVCF"] --> Saved2[("DuckDB
sample-102 completed
gVCF location")] + Saved1 --> Controller["application or targets
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. diff --git a/vignettes/native-values.Rmd b/vignettes/native-values.Rmd index 03a910c..296e3b6 100644 --- a/vignettes/native-values.Rmd +++ b/vignettes/native-values.Rmd @@ -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