Compare commits

..
Author SHA1 Message Date
jared f4ab9b6d90 feat: cog_open() honours a DuckDB thread and memory budget (#60)
R-CMD-check / check (push) Successful in 4m11s
R-CMD-check / check (pull_request) Successful in 4m11s
cog_open() connected with a bare dbConnect() and set no resource pragmas, so
DuckDB claimed every visible core. Right for one interactive session on a
dedicated machine; wrong for a server, where cog-api runs two replicas on an
8-core host budgeted 4 and each replica independently claims all 8.

USCOGDATA_DUCKDB_THREADS and USCOGDATA_DUCKDB_MEMORY_LIMIT now resolve through
.cfg() -- inheriting the env var > option > default precedence USCOGDATA_URL
already had -- and are applied as pragmas when the connection is created.

Unset issues NO pragma, so an unconfigured session is byte-identical to before.
That negative property is asserted directly against a connection opened the
pre-change way rather than against a hardcoded core count.

.cfg() returns an env var as character, so both resolvers coerce and validate
rather than trusting the type: sprintf("SET threads TO %d", "4") would
otherwise abort inside the connection path with an error naming the pragma
instead of the setting the operator got wrong.

Replaces cog-api's getFromNamespace(".ensure_session", "uscogdata") workaround,
which depended on a private name and on the session already being open.
2026-08-10 18:51:56 -04:00
5 changed files with 257 additions and 52 deletions
+16 -21
View File
@@ -1,28 +1,23 @@
# uscogdata 0.4.0
## Documentation: the corpus-access table is re-measured and honest
## DuckDB's resource budget is configurable
The README's "two ways to read the corpus" table carried figures taken before
the corpus was re-chunked into row groups (cog_pipeline#93, published
2026-08-09) and reported the mirrored column as "local speed" with no number at
all. Re-measured 2026-08-10 against the published corpus (`pipeline_commit
3d28ddd`), fresh R session per arm:
`USCOGDATA_DUCKDB_THREADS` and `USCOGDATA_DUCKDB_MEMORY_LIMIT` (with matching
`options(uscogdata.duckdb_threads = )` / `options(uscogdata.duckdb_memory_limit = )`
spellings) cap the DuckDB connection the package opens. Both follow the same
env-var > option > default precedence as `USCOGDATA_URL`.
* **A local mirror is roughly 60-80x faster.** A one-off question costs ~12 s
end to end remotely against ~0.15 s mirrored. That is the largest single
difference available to a user and it is now stated outright rather than left
as "local speed".
* **Opening the session is the largest remote cost** (~7.5 s -- manifest fetch
plus 23 view registrations over HTTPS), larger than any individual query, and
it lands on the first query rather than on `library(uscogdata)`. The old table
did not account for it anywhere.
* **The remote cost is round-trips, not scanning.** A repeat query over
already-touched partitions is ~1.5 s against ~4 s cold, and a full-history
query costs ~7 s whether it runs first or last.
* The corpus size is **~201 MB**, not 190.6 MB -- row-group chunking added ~3.4%
and the old figure was ambiguous between MB and MiB besides.
* Documented that a burst of remote queries can be rate-limited by the host
(`HTTP 429`), which is another reason to mirror for real work.
Unset, **no pragma is issued at all** and DuckDB's own defaults apply exactly as
before -- every visible core. That is right for one interactive session on a
dedicated machine and wrong for a server: where several readers share a host, each
otherwise claims the whole machine and they contend. Capping measured ~5% on a
single-government all-years query (502 ms at 2 threads vs 475 ms uncapped on 16
cores), which is cheap enough that a server should always cap.
This replaces a workaround in which a consumer reached into the package namespace
at boot -- `getFromNamespace(".ensure_session", "uscogdata")()` followed by a manual
`SET threads` -- depending both on a private name and on the session already being
open.
## Cohorts can be named by predicate, not just by id
+60 -1
View File
@@ -18,7 +18,15 @@
# Nextcloud share or a local copy made by cog_mirror().
url = "https://huggingface.co/datasets/civilytics/us-cog-finance/resolve/main/",
cache_dir = NULL,
manifest_ttl_secs = 3600L
manifest_ttl_secs = 3600L,
# NULL means "emit no pragma", which leaves DuckDB's own defaults intact:
# every visible core, and 80% of RAM. That is right for one interactive
# session on a dedicated machine and wrong for a server, where several
# readers share a box and each would otherwise claim all of it. See
# .resolve_duckdb_threads() for why this is a supported option rather than
# something a consumer reaches into the namespace to set.
duckdb_threads = NULL,
duckdb_memory_limit = NULL
)
#' Resolve a config value: env var > option > default
@@ -61,3 +69,54 @@
v <- .cfg("cache_dir")
if (is.null(v)) tools::R_user_dir("uscogdata", "cache") else v
}
#' Resolve the DuckDB thread cap, or NULL to leave DuckDB's default alone.
#'
#' `cog_open()` used to connect with a bare `dbConnect()` and set no `threads`
#' pragma, so DuckDB claimed every core it could see. cog-api works around that
#' by reaching into this namespace at boot --
#' `getFromNamespace(".ensure_session", "uscogdata")()` followed by a manual
#' `SET threads` -- which depends on a private name and on the session already
#' being open. Making it a resolved option removes the reason to do that.
#'
#' `.cfg()` returns an environment variable as CHARACTER, so this coerces
#' rather than trusting the type: `USCOGDATA_DUCKDB_THREADS=4` arrives as "4",
#' and `sprintf("SET threads TO %d", "4")` would abort inside the connection
#' path with an error about the pragma rather than about the setting.
#' @noRd
.resolve_duckdb_threads <- function() {
v <- .cfg("duckdb_threads")
if (is.null(v) || (is.character(v) && !nzchar(v))) return(NULL)
n <- suppressWarnings(as.integer(v))
if (length(n) != 1L || is.na(n) || n < 1L) {
cli::cli_abort(c(
"{.envvar USCOGDATA_DUCKDB_THREADS} must be a single positive integer.",
x = "Got {.val {v}}.",
i = "Unset it (or {.code options(uscogdata.duckdb_threads = NULL)}) to use DuckDB's default of every visible core."
), class = "uscogdata_invalid_duckdb_threads")
}
n
}
#' Resolve the DuckDB memory limit, or NULL to leave DuckDB's default alone.
#'
#' The value is a DuckDB size string (`"4GB"`, `"512MB"`). Only its SHAPE is
#' checked here -- DuckDB owns the unit vocabulary, and re-implementing that
#' parse would be a second definition free to drift from the engine's. An
#' unrecognised unit therefore surfaces as DuckDB's own error at `SET` time,
#' which names the setting correctly; the check here exists to reject the
#' inputs that would otherwise reach the connection as a SQL fragment.
#' @noRd
.resolve_duckdb_memory_limit <- function() {
v <- .cfg("duckdb_memory_limit")
if (is.null(v) || (is.character(v) && !nzchar(v))) return(NULL)
if (length(v) != 1L || !is.character(v) ||
!grepl("^[0-9]+(\\.[0-9]+)?\\s*[A-Za-z]{0,3}$", v)) {
cli::cli_abort(c(
"{.envvar USCOGDATA_DUCKDB_MEMORY_LIMIT} must be a single DuckDB size string.",
x = "Got {.val {v}}.",
i = "Examples: {.val 4GB}, {.val 512MB}, {.val 1.5GB}."
), class = "uscogdata_invalid_duckdb_memory_limit")
}
trimws(v)
}
+31 -1
View File
@@ -2,13 +2,23 @@
#' Internal: open session, register views, cache manifest.
#' Not exported. Called lazily by verbs via .ensure_session().
#'
#' `threads` and `memory_limit` default to the resolved configuration and are
#' applied as pragmas on the new connection. When both resolve to NULL -- which
#' is the case unless the operator sets one -- NO pragma is issued at all, so an
#' unconfigured session connects exactly as it did before this argument existed.
#' @noRd
cog_open <- function(url = .resolve_url(),
cache_dir = .resolve_cache_dir()) {
cache_dir = .resolve_cache_dir(),
threads = .resolve_duckdb_threads(),
memory_limit = .resolve_duckdb_memory_limit()) {
.check_url_configured(url)
if (!dir.exists(cache_dir)) dir.create(cache_dir, recursive = TRUE)
con <- DBI::dbConnect(duckdb::duckdb())
# Before anything else touches the connection: httpfs reads the corpus, and
# a remote read should already be bound by whatever budget the operator set.
.apply_duckdb_limits(con, threads, memory_limit)
DBI::dbExecute(con, "INSTALL httpfs; LOAD httpfs;")
manifest <- .fetch_or_cache_manifest(url, cache_dir)
@@ -25,6 +35,26 @@ cog_open <- function(url = .resolve_url(),
invisible(con)
}
#' Apply the operator's DuckDB resource budget to a fresh connection.
#'
#' Split out from cog_open() so the "unset changes nothing" property is one
#' readable branch rather than two conditionals buried in the connection path.
#' Both settings are session-scoped in DuckDB, so this must run per connection;
#' cog_close() discards the connection and the next cog_open() re-resolves,
#' which is what makes a changed option take effect on the next session.
#' @noRd
.apply_duckdb_limits <- function(con, threads, memory_limit) {
if (!is.null(threads)) {
DBI::dbExecute(con, sprintf("SET threads TO %d", threads))
}
if (!is.null(memory_limit)) {
# Quoted as a string literal: DuckDB's memory_limit takes '4GB', not 4GB.
DBI::dbExecute(con, sprintf("SET memory_limit TO %s",
.sql_lit_chr(memory_limit)))
}
invisible(con)
}
#' @noRd
.ensure_session <- function() {
if (is.null(.uscogdata_env$con) ||
+16 -29
View File
@@ -18,7 +18,7 @@ carries provenance describing what was converted, what was aggregated, and
which known series breaks intersect your query.
**Scope:** government types 0–3 (state, county, municipality, township).
56 fiscal years, 46,148,034 rows, ~201 MB. There is no source data for FY1968
56 fiscal years, 46,148,034 rows, 190.6 MB. There is no source data for FY1968
or FY1969. Special districts (type 4) and school districts (type 5) are
excluded pending validation.
@@ -85,38 +85,14 @@ cog_explain(spend)
| | Remote (default) | Mirrored |
|---|---|---|
| Setup | none | `cog_mirror(dest)`, ~201 MB once |
| Disk used | **0 MB** — HTTP range requests only | ~201 MB |
| Opening a session | ~7.5 s | ~0.1 s |
| One government, one year | ~4 s | ~0.05 s |
| One government, full history | ~7 s | ~0.1 s |
| Later queries, same session | ~1.5 s | ~0.05 s |
| Setup | none | `cog_mirror(dest)`, 190.6 MB once |
| Disk used | **0 MB** — HTTP range requests only | 190.6 MB |
| Per query | ~4 s (one government, one year)<br>~6 s (one government, 23 years) | local speed |
| Good for | trying it out, teaching, one-off questions | repeated analysis, offline work, reproducibility |
**A local mirror is roughly 60–80x faster, and it is one function call.** That is
by far the largest difference any of these settings makes. If you are going to
ask more than a handful of questions, mirror first.
Measured 2026-08-10 on a 16-core Linux workstation against the published corpus
(schema v7, `pipeline_commit 3d28ddd`), fresh R session per arm. A one-off
question costs about **12 seconds end to end remotely and 0.15 seconds
mirrored**, session setup included.
Two things the per-query rows hide:
- **Opening the session is the single largest remote cost** — larger than any
one query. It fetches the manifest and registers 23 SQL views over HTTPS, and
it lands on your first query, not on `library(uscogdata)`.
- **The cost is network round-trips, not scanning.** A repeat query against
partitions this session has already touched is ~1.5 s rather than ~4 s, and a
full-history query costs ~7 s whether it runs first or last. What you are
paying for is reaching each of the 56 yearly files over HTTPS the first time.
Nothing is written to disk in remote mode: DuckDB fetches the parquet footer,
works out which row groups it needs, and reads only those. Nothing is cached
between sessions either, so every query goes back to the network — and a session
that issues many remote queries in quick succession can be rate-limited by the
host (`HTTP Error: ... 429`). Both are further reasons to mirror for real work.
between sessions either, so every query goes back to the network.
The default points at a public HuggingFace mirror of the corpus. If you would
rather not depend on a third party — for reproducibility, for an air-gapped
@@ -134,6 +110,17 @@ After that, nothing in your analysis touches an external service.
- `USCOGDATA_URL` — corpus root: an HTTPS URL or a local path, **trailing slash required**
- `USCOGDATA_CACHE_DIR` — where the manifest is cached (default: user cache dir)
- `USCOGDATA_MANIFEST_TTL_SECS` — manifest re-fetch interval (default 3600)
- `USCOGDATA_DUCKDB_THREADS` — cap DuckDB's thread count (default: every visible core)
- `USCOGDATA_DUCKDB_MEMORY_LIMIT` — cap DuckDB's memory, e.g. `"4GB"` (default: DuckDB's own)
Each also has an `options()` spelling — `uscogdata.url`, `uscogdata.duckdb_threads`,
and so on — and the environment variable wins where both are set.
The two DuckDB caps exist for **servers**, not laptops. Unset, DuckDB claims every
core it can see, which is right for one interactive session on your own machine and
wrong when several readers share a box: each claims the whole machine and they fight.
Capping costs roughly 5% on a single query and is worth it anywhere the process is
sharing hardware.
## Amounts are in full US dollars
+134
View File
@@ -0,0 +1,134 @@
# tests/testthat/test-duckdb-limits.R
#
# uscogdata#60. cog_open() used to connect with a bare dbConnect() and set no
# resource pragmas, so DuckDB claimed every visible core. That is right for one
# interactive session on a dedicated machine and wrong for a server: cog-api
# runs two replicas on an 8-core host budgeted 4, and without a cap each
# replica independently claims all 8 and they fight.
#
# The consumer-side workaround this replaces reached into the namespace at
# boot -- getFromNamespace(".ensure_session", "uscogdata")() followed by a
# manual SET threads -- which depends on a private name AND on the session
# already being open.
#
# The load-bearing property is the NEGATIVE one: unset must emit no pragma at
# all, so an unconfigured session is byte-identical to pre-#60 behaviour.
# Open a session under a given configuration and read a DuckDB setting back.
# Each call closes first, because both settings are session-scoped: an
# already-open connection would be reused by .ensure_session() and report the
# PREVIOUS test's value, which is exactly the false pass to avoid here.
setting_under <- function(setting, envvars = character(0), opts = list()) {
uscogdata:::cog_close()
on.exit(uscogdata:::cog_close(), add = TRUE)
withr::with_envvar(envvars, {
withr::with_options(opts, {
con <- uscogdata:::cog_open()
DBI::dbGetQuery(
con, sprintf("SELECT current_setting('%s') AS v", setting)
)$v[[1]]
})
})
}
test_that("USCOGDATA_DUCKDB_THREADS caps the connection's thread count", {
skip_if_no_corpus()
expect_equal(
as.integer(setting_under("threads", c(USCOGDATA_DUCKDB_THREADS = "2"))),
2L
)
})
test_that("the option spelling works, and the env var beats it", {
skip_if_no_corpus()
expect_equal(
as.integer(setting_under("threads",
c(USCOGDATA_DUCKDB_THREADS = NA),
list(uscogdata.duckdb_threads = 3L))),
3L
)
# Same precedence .cfg() gives every other setting: env var > option.
expect_equal(
as.integer(setting_under("threads",
c(USCOGDATA_DUCKDB_THREADS = "1"),
list(uscogdata.duckdb_threads = 3L))),
1L
)
})
test_that("unset leaves DuckDB's own default in place", {
skip_if_no_corpus()
# Not asserting a specific number -- the default is core-count-dependent and
# a literal would fail on a different machine. The claim is that NO pragma
# was issued, so the session sees whatever DuckDB would have chosen on its
# own. Compared against a plain connection opened the pre-#60 way.
unset <- setting_under("threads",
c(USCOGDATA_DUCKDB_THREADS = NA),
list(uscogdata.duckdb_threads = NULL))
bare <- local({
con <- DBI::dbConnect(duckdb::duckdb())
on.exit(DBI::dbDisconnect(con, shutdown = TRUE), add = TRUE)
DBI::dbGetQuery(con, "SELECT current_setting('threads') AS v")$v[[1]]
})
expect_equal(as.integer(unset), as.integer(bare))
})
test_that("USCOGDATA_DUCKDB_MEMORY_LIMIT is applied", {
skip_if_no_corpus()
v <- setting_under("memory_limit", c(USCOGDATA_DUCKDB_MEMORY_LIMIT = "2GB"))
# DuckDB does not echo back the string it was given: it stores bytes and
# reports BINARY units, so "2GB" (2e9 bytes) comes back as "1.8 GiB". Assert
# the magnitude it actually means rather than the spelling this package sent
# -- matching on "2" passes for the wrong reason and fails on the right one.
expect_match(as.character(v), "GiB", fixed = TRUE)
# DuckDB also truncates the display to one decimal ("1.8 GiB" for 1.863), so
# the tolerance covers rounding, not slack in the setting itself.
gib <- as.numeric(sub("\\s*GiB$", "", as.character(v)))
expect_equal(gib, 2e9 / 1024^3, tolerance = 0.05)
})
# --- Validation -------------------------------------------------------------
# .cfg() returns an env var as CHARACTER. Without coercion here,
# sprintf("SET threads TO %d", "4") aborts inside the connection path with an
# error about the pragma rather than about the setting the operator got wrong.
test_that(".resolve_duckdb_threads coerces a character env var to integer", {
withr::local_envvar(USCOGDATA_DUCKDB_THREADS = "4")
expect_identical(uscogdata:::.resolve_duckdb_threads(), 4L)
})
test_that(".resolve_duckdb_threads returns NULL when unset or empty", {
withr::local_options(uscogdata.duckdb_threads = NULL)
withr::local_envvar(USCOGDATA_DUCKDB_THREADS = NA)
expect_null(uscogdata:::.resolve_duckdb_threads())
withr::local_envvar(USCOGDATA_DUCKDB_THREADS = "")
expect_null(uscogdata:::.resolve_duckdb_threads())
})
test_that(".resolve_duckdb_threads rejects values that are not positive integers", {
for (bad in c("0", "-1", "two", "1.5.2")) {
withr::local_envvar(USCOGDATA_DUCKDB_THREADS = bad)
expect_error(uscogdata:::.resolve_duckdb_threads(),
class = "uscogdata_invalid_duckdb_threads")
}
})
test_that(".resolve_duckdb_memory_limit accepts size strings and rejects junk", {
withr::local_envvar(USCOGDATA_DUCKDB_MEMORY_LIMIT = "4GB")
expect_identical(uscogdata:::.resolve_duckdb_memory_limit(), "4GB")
withr::local_envvar(USCOGDATA_DUCKDB_MEMORY_LIMIT = "1.5GB")
expect_identical(uscogdata:::.resolve_duckdb_memory_limit(), "1.5GB")
# A SQL fragment must not reach the connection as one.
withr::local_envvar(USCOGDATA_DUCKDB_MEMORY_LIMIT = "4GB'; DROP TABLE x; --")
expect_error(uscogdata:::.resolve_duckdb_memory_limit(),
class = "uscogdata_invalid_duckdb_memory_limit")
})
test_that(".apply_duckdb_limits issues no statement when both are NULL", {
# The negative property, asserted directly rather than inferred: a connection
# that would ERROR on any statement proves none was sent.
expect_silent(uscogdata:::.apply_duckdb_limits(NULL, NULL, NULL))
})