diff --git a/NEWS.md b/NEWS.md index 6c3b108..7c45512 100644 --- a/NEWS.md +++ b/NEWS.md @@ -1,5 +1,24 @@ # uscogdata 0.4.0 +## DuckDB's resource budget is configurable + +`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`. + +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 `cog_spending()`, `cog_revenue()` and `cog_balances()` gain optional `state` diff --git a/R/config.R b/R/config.R index 37354a3..8c1c4eb 100644 --- a/R/config.R +++ b/R/config.R @@ -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) +} diff --git a/R/session.R b/R/session.R index 468e980..26f2943 100644 --- a/R/session.R +++ b/R/session.R @@ -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) || diff --git a/README.md b/README.md index 7590d76..1578400 100644 --- a/README.md +++ b/README.md @@ -110,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 diff --git a/tests/testthat/test-duckdb-limits.R b/tests/testthat/test-duckdb-limits.R new file mode 100644 index 0000000..d6fd803 --- /dev/null +++ b/tests/testthat/test-duckdb-limits.R @@ -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)) +})