Compare commits

..
7 Commits
Author SHA1 Message Date
jared d2caa6de97 Merge pull request 'docs: re-measure the corpus-access table against the published corpus (#56)' (#65) from docs/readme-perf-remeasure-56 into main
Mirror to GitHub / mirror (push) Successful in 9s
R-CMD-check / check (push) Successful in 3m25s
2026-08-10 19:36:42 -04:00
jared 303aa59b07 docs: order the 0.4.0 NEWS sections by user impact, not merge order
R-CMD-check / check (push) Successful in 3m58s
R-CMD-check / check (pull_request) Successful in 3m38s
The three 0.4.0 features landed in the order their PRs merged, which buried the
headline change (cohort predicates, 4.8x) below an operator config knob and a
documentation note. Reordered to: cohorts, pagination, DuckDB budget, docs,
Fixes -- with Fixes last, where it was already.

The stacked PRs each appended their own '## Fixes' heading, so resolving the
conflicts also folded two of them into the single section that belongs there.

Content is byte-identical to what merged; only section order changed. Verified
by diffing the sorted non-blank lines against the previous commit.
2026-08-10 19:29:43 -04:00
jared 81f72321ee docs: re-measure the corpus-access table against the published corpus (#56)
figures predate the row-group rechunk (cog_pipeline#93, published 2026-08-09)
and reported the mirrored column as 'local speed' with no number -- hiding the
largest difference available to a user.

Measured 2026-08-10, fresh R session per arm, against the live corpus at
pipeline_commit 3d28ddd. Madison WI, 16-core Linux workstation.

Three findings the old table could not express:

- A local mirror is 60-80x faster. A one-off question is ~12 s end to end
  remotely against ~0.15 s mirrored. Stated outright now, because it is a
  bigger and cheaper win for users than anything in the R code.

- Opening the session is the LARGEST remote cost (~7.5 s), bigger than any
  individual query, and it lands on the user's first query rather than on
  library(). The old table accounted for it nowhere, so every per-query figure
  was quietly missing it.

- 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 (verified by running the arms
  in both orders). This is why #93's 1.4-1.7x, measured through cog-api against
  a local mount, does not show up on the remote path -- there, network latency
  swamps scan time.

Corpus size corrected to ~201 MB: row-group chunking added ~3.4%, and 190.6 was
ambiguous between MB and MiB besides. Measured from the manifest and on disk.
The 0.3.0 NEWS section keeps 190.6 -- it was correct for that release.

Also documented HTTP 429: a burst of remote queries gets rate-limited by the
host. Hit while taking these measurements.
2026-08-10 19:29:05 -04:00
jared 5cb83d8f1d Merge pull request 'feat: limit/offset on cog_gov_search() and cog_balances() (#57)' (#63) from feat/pagination-search-balances-57 into main
Mirror to GitHub / mirror (push) Successful in 12s
R-CMD-check / check (push) Successful in 4m26s
2026-08-10 19:28:09 -04:00
jared 2bb9d41d72 feat: limit/offset on cog_gov_search() and cog_balances() (#57)
R-CMD-check / check (pull_request) Successful in 4m22s
R-CMD-check / check (push) Successful in 4m28s
verbs were left materializing everything and slicing in R -- the pattern behind
the 2026-08-06 production incident. cog_gov_search() had no LIMIT at all, so an
unfiltered call returns the entire 40,336-row crosswalk.

Extracted the #39 machinery into R/pagination.R first (.validate_pagination(),
.paginate_sql(), .take_pagination_total()) rather than growing a third inline
copy: three definitions of what total_rows means is three places for it to
drift. Conflict refusals stay at the call sites because each verb's conflict
set differs. .verb_spendrev() now uses the shared helpers and is unchanged in
behaviour.

The empty-page fallback query is now passed as a thunk, so the unpaginated SQL
is only BUILT when an offset actually lands past the end instead of on every
paged call.

Two things #57 did not anticipate:

- cog_gov_search()'s ORDER BY was not a total order. population_acs DESC NULLS
  LAST leaves ties -- and the whole NULL block -- in scan order, so two requests
  can order them differently and a paged sweep duplicates one row while dropping
  another. Added canonical_govid as tiebreaker. Unpaginated output changes only
  in the relative order of already-tied rows.

- Basket mode returns one resolved row per requested name plus a sidecar
  covering all of them, so a page of it is not a page of anything the caller
  asked for. Refused with uscogdata_basket_pagination_conflict rather than
  silently ignoring the arguments.

Both default to NULL, so cog-api adopts them behind its existing formals()
probe with no lockstep deploy.

Suite: 1067 passed, 0 failed, 0 warnings (2 pre-existing live-corpus skips).
2026-08-10 19:27:11 -04:00
jared 700ae93c9c Merge pull request 'feat: cog_open() honours a DuckDB thread and memory budget (#60)' (#62) from feat/duckdb-threads-60 into main
Mirror to GitHub / mirror (push) Successful in 10s
R-CMD-check / check (push) Successful in 4m6s
Reviewed-on: #62
2026-08-10 19:23:14 -04:00
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
12 changed files with 712 additions and 84 deletions
+73 -24
View File
@@ -1,29 +1,5 @@
# uscogdata 0.4.0 # uscogdata 0.4.0
## Documentation: the corpus-access table is re-measured and honest
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:
* **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.
## Cohorts can be named by predicate, not just by id ## Cohorts can be named by predicate, not just by id
`cog_spending()`, `cog_revenue()` and `cog_balances()` gain optional `state` `cog_spending()`, `cog_revenue()` and `cog_balances()` gain optional `state`
@@ -70,8 +46,81 @@ When the cohort is named by predicate there is no id list to report, so
`provenance$scope$cohort` carries `state`, `type` and `n_governments` instead. `provenance$scope$cohort` carries `state`, `type` and `n_governments` instead.
A `govid`-named cohort's provenance is unchanged. A `govid`-named cohort's provenance is unchanged.
## `cog_gov_search()` and `cog_balances()` gain `limit`/`offset`
Pagination arrived on `cog_spending()`/`cog_revenue()` in 0.3.0; the other two
verbs were left materializing everything and slicing in R. Both now take
`limit`/`offset` with the same semantics: `NULL` default, the page applied in
SQL behind a deterministic `ORDER BY`, and the unpaginated count returned as a
`total_rows` attribute computed by `COUNT(*) OVER()` in the same scan rather
than a second query.
`cog_gov_search()` had no `LIMIT` at all, which made it the one verb that
returns the entire 40,336-row crosswalk when called with no filter.
Two refusals rather than silent surprises:
* `cog_balances(recipe = , limit = )` aborts with class
`uscogdata_recipe_pagination_conflict` -- a recipe's result comes from a
separate query that pagination is not wired into.
* `cog_gov_search()` in basket mode (`length(name) > 1`) aborts with class
`uscogdata_basket_pagination_conflict`. Basket mode returns one resolved row
per requested name with a sidecar covering all of them; a page of that is not
a page of anything the caller asked for.
## 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.
## Documentation: the corpus-access table is re-measured and honest
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:
* **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.
## Fixes ## Fixes
* `cog_gov_search()` now orders by `population_acs DESC NULLS LAST,
canonical_govid`. **`population_acs` alone is not a total order** -- ties, and
the entire `NULLS LAST` block, came back in whatever order the scan produced.
That was invisible while every call returned the full result set, but it makes
a paged sweep unsound: two requests can order tied rows differently, so a row
is duplicated on one page and missing from the next. Unpaginated results are
unchanged except for the relative order of rows that were already tied.
* An unknown `state` abbreviation now aborts with "Unknown state abbreviation" * An unknown `state` abbreviation now aborts with "Unknown state abbreviation"
(class `uscogdata_unknown_state`) instead of base R's "subscript out of (class `uscogdata_unknown_state`) instead of base R's "subscript out of
bounds". `.state_abbrev_to_fips` is a named character vector, so `[[` on an bounds". `.state_abbrev_to_fips` is a named character vector, so `[[` on an
+41 -2
View File
@@ -47,6 +47,12 @@
#' @param recipe Optional harmonization recipe id (see [cog_recipes()]). #' @param recipe Optional harmonization recipe id (see [cog_recipes()]).
#' `"cash_securities_z77_wide"` and `"cash_securities_z78_wide"` bridge the #' `"cash_securities_z77_wide"` and `"cash_securities_z78_wide"` bridge the
#' wide era to the modern one. #' wide era to the modern one.
#' @param limit Maximum number of result rows to return, pushed into the SQL
#' rather than applied after materializing every row. `NULL` (default)
#' returns everything. Cannot be combined with `recipe` -- see `offset` and
#' `total_rows`.
#' @param offset Rows to skip before `limit` starts counting (0-based).
#' Ignored if `limit` is `NULL`; defaults to `0L` when `limit` is set.
#' #'
#' @return Tibble with columns `year`, `canonical_govid`, `gov_name`, #' @return Tibble with columns `year`, `canonical_govid`, `gov_name`,
#' `balance_subtype`, `category`, `amt_nominal`, `codes_included`, #' `balance_subtype`, `category`, `amt_nominal`, `codes_included`,
@@ -64,11 +70,16 @@
#' and `truncated` (the observed subtypes whose coverage falls short of the #' and `truncated` (the observed subtypes whose coverage falls short of the
#' requested years). `expenditure_concept`/`revenue_concept` are `NA` -- #' requested years). `expenditure_concept`/`revenue_concept` are `NA` --
#' holdings are a stock, not a flow, so neither concept vocabulary applies. #' holdings are a stock, not a flow, so neither concept vocabulary applies.
#'
#' When `limit` is set, also carries a `total_rows` attribute: the full
#' unpaginated row count, computed by the same query (`COUNT(*) OVER()`)
#' rather than a second scan.
#' @export #' @export
cog_balances <- function(govid = NULL, years, category = NULL, cog_balances <- function(govid = NULL, years, category = NULL,
per_capita = FALSE, adjust_to_year = NULL, per_capita = FALSE, adjust_to_year = NULL,
basis = c("harmonized", "raw"), recipe = NULL, basis = c("harmonized", "raw"), recipe = NULL,
state = NULL, type = NULL) { state = NULL, type = NULL,
limit = NULL, offset = NULL) {
call <- match.call() call <- match.call()
basis <- match.arg(basis, c("harmonized", "raw")) basis <- match.arg(basis, c("harmonized", "raw"))
# Coerce FIRST, validate second: .validate_verb_inputs() asserts # Coerce FIRST, validate second: .validate_verb_inputs() asserts
@@ -91,6 +102,22 @@ cog_balances <- function(govid = NULL, years, category = NULL,
# validator's own doc comment for the incident that made that matter. # validator's own doc comment for the incident that made that matter.
.validate_verb_inputs(govid, years, category, per_capita, adjust_to_year, .validate_verb_inputs(govid, years, category, per_capita, adjust_to_year,
recipe) recipe)
# Same semantics as the money verbs (R/pagination.R). Only the `recipe`
# conflict applies here: cog_balances() has no `complete` argument, and a
# recipe's result comes from .run_recipe()'s own query, which pagination is
# not wired into.
paging <- .validate_pagination(limit, offset)
limit <- paging$limit
offset <- paging$offset
if (!is.null(limit) && !is.null(recipe)) {
cli::cli_abort(c(
"`limit`/`offset` cannot be combined with `recipe`.",
"i" = "A recipe's result comes from a separate query (`.run_recipe()`) that pagination is not wired into yet.",
"*" = "Drop `limit`/`offset`, or drop `recipe`."
), class = "uscogdata_recipe_pagination_conflict")
}
years <- as.integer(years) years <- as.integer(years)
if (!is.null(adjust_to_year)) adjust_to_year <- as.integer(adjust_to_year) if (!is.null(adjust_to_year)) adjust_to_year <- as.integer(adjust_to_year)
@@ -108,6 +135,7 @@ cog_balances <- function(govid = NULL, years, category = NULL,
manifest <- .uscogdata_env$manifest manifest <- .uscogdata_env$manifest
recipe_block <- NULL recipe_block <- NULL
category_for_prov <- category category_for_prov <- category
total_rows <- NULL # set below only when limit is non-NULL (non-recipe path)
if (!is.null(recipe)) { if (!is.null(recipe)) {
.require_schema_v5(con, manifest, "recipe =") .require_schema_v5(con, manifest, "recipe =")
@@ -125,8 +153,18 @@ cog_balances <- function(govid = NULL, years, category = NULL,
} else { } else {
sql <- .build_verb_sql("balance_annotated", "balance_subtype", sql <- .build_verb_sql("balance_annotated", "balance_subtype",
cohort, years, category, cohort, years, category,
ig_view = NULL, subtype_scope = NULL) ig_view = NULL, subtype_scope = NULL,
limit = limit, offset = offset)
result <- tibble::as_tibble(DBI::dbGetQuery(con, sql)) result <- tibble::as_tibble(DBI::dbGetQuery(con, sql))
if (!is.null(limit)) {
paged <- .take_pagination_total(result, con, function() {
.build_verb_sql("balance_annotated", "balance_subtype",
cohort, years, category,
ig_view = NULL, subtype_scope = NULL)
})
result <- paged$result
total_rows <- paged$total_rows
}
} }
# Order matters (matches .verb_spendrev()): per-capita first, so # Order matters (matches .verb_spendrev()): per-capita first, so
@@ -158,6 +196,7 @@ cog_balances <- function(govid = NULL, years, category = NULL,
.emit_balance_caveats(prov$balance_caveats) .emit_balance_caveats(prov$balance_caveats)
attr(result, "provenance") <- prov attr(result, "provenance") <- prov
if (!is.null(limit)) attr(result, "total_rows") <- total_rows
result result
} }
+60 -1
View File
@@ -18,7 +18,15 @@
# Nextcloud share or a local copy made by cog_mirror(). # Nextcloud share or a local copy made by cog_mirror().
url = "https://huggingface.co/datasets/civilytics/us-cog-finance/resolve/main/", url = "https://huggingface.co/datasets/civilytics/us-cog-finance/resolve/main/",
cache_dir = NULL, 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 #' Resolve a config value: env var > option > default
@@ -61,3 +69,54 @@
v <- .cfg("cache_dir") v <- .cfg("cache_dir")
if (is.null(v)) tools::R_user_dir("uscogdata", "cache") else v 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)
}
+81
View File
@@ -0,0 +1,81 @@
# R/pagination.R
#
# Shared limit/offset machinery. #39 established the semantics inside
# .verb_spendrev(); #57 extends them to cog_gov_search() and cog_balances(),
# which is what made a single definition worth having: three inline copies of
# "coerce, refuse, unwrap the count" would be three places for the meaning of
# `total_rows` to drift.
#
# The SQL side stays in .build_verb_sql() (R/spending.R) -- it already wraps
# the aggregate in an outer SELECT so COUNT(*) OVER() sees the post-GROUP-BY
# row count rather than the pre-aggregation one, and that is the subtle part
# worth not duplicating either.
#' Coerce and check a limit/offset pair.
#'
#' Returns the coerced pair, or NULL for `limit` when no page was requested.
#' `offset` defaults to 0 whenever `limit` is set, so a caller can supply just
#' `limit` and get the first page.
#'
#' Conflicts with other arguments are deliberately NOT checked here: they
#' differ per verb (`complete`/`recipe` for the money verbs, basket mode for
#' `cog_gov_search()`, `recipe` alone for `cog_balances()`), and a shared
#' function taking a list of conflict flags would be harder to read than the
#' three explicit refusals at the call sites.
#' @noRd
.validate_pagination <- function(limit, offset) {
if (is.null(limit)) {
return(list(limit = NULL, offset = NULL))
}
limit <- as.integer(limit)
if (length(limit) != 1L || is.na(limit) || limit < 0L) {
cli::cli_abort("`limit` must be a single non-negative integer.",
class = "uscogdata_invalid_pagination")
}
offset <- if (is.null(offset)) 0L else as.integer(offset)
if (length(offset) != 1L || is.na(offset) || offset < 0L) {
cli::cli_abort("`offset` must be a single non-negative integer.",
class = "uscogdata_invalid_pagination")
}
list(limit = limit, offset = offset)
}
#' Wrap a query so one page comes back carrying the unpaginated total.
#'
#' `COUNT(*) OVER()` rides along as an ordinary column, so the caller gets the
#' true total from the SAME scan instead of a second round trip. The outer
#' `SELECT *` matters: appending LIMIT/OFFSET directly to a grouped query would
#' have the window function count pre-aggregation rows.
#' @noRd
.paginate_sql <- function(base_sql, limit, offset) {
if (is.null(limit)) return(base_sql)
sprintf(
"SELECT *, COUNT(*) OVER() AS pagination_total_rows
FROM (%s) AS _paged
LIMIT %d OFFSET %d",
base_sql, limit, offset
)
}
#' Strip the count column back out and report the unpaginated total.
#'
#' Returns `list(result = , total_rows = )`.
#'
#' An empty page -- an offset past the end -- carries no row to read the window
#' function off, so that one case falls back to a second, unpaginated
#' `COUNT(*)` rather than reporting a wrong zero. `unpaged_sql` is passed as a
#' function so the fallback query is only BUILT when it is actually needed;
#' every caller's unpaginated SQL is otherwise constructed on every paged call
#' and thrown away.
#' @noRd
.take_pagination_total <- function(result, con, unpaged_sql) {
if (nrow(result) > 0L) {
total <- result$pagination_total_rows[[1]]
result$pagination_total_rows <- NULL
return(list(result = result, total_rows = as.integer(total)))
}
count_sql <- sprintf("SELECT COUNT(*) AS n FROM (%s) AS _uncounted",
if (is.function(unpaged_sql)) unpaged_sql() else unpaged_sql)
list(result = result,
total_rows = as.integer(DBI::dbGetQuery(con, count_sql)$n[[1]]))
}
+50 -7
View File
@@ -44,10 +44,22 @@
#' in basket mode (recycles from length 1). Excluded types `4`/`5` (or #' in basket mode (recycles from length 1). Excluded types `4`/`5` (or
#' `"special_district"` / `"school_district"`) trigger an explanatory #' `"special_district"` / `"school_district"`) trigger an explanatory
#' message and an empty result. #' message and an empty result.
#' @param limit Maximum number of rows to return, applied in SQL. `NULL`
#' (default) returns every match -- which, with no other filter, is the
#' entire crosswalk. Utility mode only: pagination has no meaning in basket
#' mode, where the result is one resolved row per requested name in input
#' order, and is refused there with class
#' `uscogdata_basket_pagination_conflict`.
#' @param offset Rows to skip before `limit` starts counting (0-based).
#' Ignored if `limit` is `NULL`; defaults to `0L` when `limit` is set.
#' @return A tibble of `canonical_fips_xwalk` rows. In utility mode, all #' @return A tibble of `canonical_fips_xwalk` rows. In utility mode, all
#' matches sorted by `population_acs` desc. In basket mode, resolved #' matches sorted by `population_acs` desc, ties broken by
#' rows in input order, with `attr(., "resolution")` set to the #' `canonical_govid`. In basket mode, resolved rows in input order, with
#' sidecar tibble. #' `attr(., "resolution")` set to the sidecar tibble.
#'
#' When `limit` is set, carries a `total_rows` attribute: the full
#' unpaginated match count, computed by the same query (`COUNT(*) OVER()`)
#' rather than a second scan.
#' @seealso [cog_basket_resolution()], [cog_basket_unresolved()], #' @seealso [cog_basket_resolution()], [cog_basket_unresolved()],
#' [cog_spending()], [cog_revenue()]. #' [cog_spending()], [cog_revenue()].
#' @examples #' @examples
@@ -82,7 +94,12 @@
#' ) #' )
#' } #' }
#' @export #' @export
cog_gov_search <- function(name = NULL, state = NULL, type = NULL) { cog_gov_search <- function(name = NULL, state = NULL, type = NULL,
limit = NULL, offset = NULL) {
paging <- .validate_pagination(limit, offset)
limit <- paging$limit
offset <- paging$offset
if (!is.null(type) && length(type) == 1L && .is_excluded_type(type)) { if (!is.null(type) && length(type) == 1L && .is_excluded_type(type)) {
cli::cli_inform(c( cli::cli_inform(c(
i = "v0.1 covers gov_types 0-3 (state/county/city/township) only.", i = "v0.1 covers gov_types 0-3 (state/county/city/township) only.",
@@ -94,6 +111,18 @@ cog_gov_search <- function(name = NULL, state = NULL, type = NULL) {
con <- .ensure_session() con <- .ensure_session()
if (length(name) > 1L) { if (length(name) > 1L) {
# Basket mode returns one resolved row per requested name, in input order,
# with a resolution sidecar describing how each was matched. A page of that
# is not a page of anything the caller asked for -- the sidecar would still
# describe every name -- so refuse rather than silently ignoring the
# arguments. Same shape as the recipe/complete refusals in .verb_spendrev().
if (!is.null(limit)) {
cli::cli_abort(c(
"`limit`/`offset` cannot be combined with basket mode.",
"i" = "Basket mode ({.code length(name) > 1}) returns one resolved row per requested name, in input order, with a resolution sidecar covering all of them.",
"*" = "Drop `limit`/`offset`, or search one name at a time."
), class = "uscogdata_basket_pagination_conflict")
}
return(.resolve_basket(name = name, state = state, type = type, con = con)) return(.resolve_basket(name = name, state = state, type = type, con = con))
} }
@@ -123,12 +152,26 @@ cog_gov_search <- function(name = NULL, state = NULL, type = NULL) {
} }
where <- if (length(preds) == 0L) "" else paste("WHERE", paste(preds, collapse = " AND ")) where <- if (length(preds) == 0L) "" else paste("WHERE", paste(preds, collapse = " AND "))
sql <- paste( # canonical_govid breaks ties. population_acs alone is NOT a total order --
# governments sharing a population, and the whole NULLS LAST block, came back
# in whatever order the scan produced. That was invisible while every call
# returned the full result set, but it makes a paged sweep unsound: two
# requests can order the tied rows differently, so a row is duplicated on one
# page and missing from the next. Any pagination has to sit on a total order.
base_sql <- paste(
"SELECT * FROM canonical_fips_xwalk", "SELECT * FROM canonical_fips_xwalk",
where, where,
"ORDER BY population_acs DESC NULLS LAST" "ORDER BY population_acs DESC NULLS LAST, canonical_govid"
) )
tibble::as_tibble(DBI::dbGetQuery(con, sql)) result <- tibble::as_tibble(
DBI::dbGetQuery(con, .paginate_sql(base_sql, limit, offset))
)
if (is.null(limit)) return(result)
paged <- .take_pagination_total(result, con, base_sql)
out <- paged$result
attr(out, "total_rows") <- paged$total_rows
out
} }
#' @noRd #' @noRd
+31 -1
View File
@@ -2,13 +2,23 @@
#' Internal: open session, register views, cache manifest. #' Internal: open session, register views, cache manifest.
#' Not exported. Called lazily by verbs via .ensure_session(). #' 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 #' @noRd
cog_open <- function(url = .resolve_url(), 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) .check_url_configured(url)
if (!dir.exists(cache_dir)) dir.create(cache_dir, recursive = TRUE) if (!dir.exists(cache_dir)) dir.create(cache_dir, recursive = TRUE)
con <- DBI::dbConnect(duckdb::duckdb()) 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;") DBI::dbExecute(con, "INSTALL httpfs; LOAD httpfs;")
manifest <- .fetch_or_cache_manifest(url, cache_dir) manifest <- .fetch_or_cache_manifest(url, cache_dir)
@@ -25,6 +35,26 @@ cog_open <- function(url = .resolve_url(),
invisible(con) 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 #' @noRd
.ensure_session <- function() { .ensure_session <- function() {
if (is.null(.uscogdata_env$con) || if (is.null(.uscogdata_env$con) ||
+14 -44
View File
@@ -398,17 +398,10 @@ cog_spending <- function(govid = NULL, years, category = NULL,
# up front rather than silently ignored: complete = TRUE fills a grid over # up front rather than silently ignored: complete = TRUE fills a grid over
# the FULL requested (year, category) space, and a recipe's result comes # the FULL requested (year, category) space, and a recipe's result comes
# from .run_recipe()'s own query, which this function does not touch. # from .run_recipe()'s own query, which this function does not touch.
paging <- .validate_pagination(limit, offset)
limit <- paging$limit
offset <- paging$offset
if (!is.null(limit)) { if (!is.null(limit)) {
limit <- as.integer(limit)
if (length(limit) != 1L || is.na(limit) || limit < 0L) {
cli::cli_abort("`limit` must be a single non-negative integer.",
class = "uscogdata_invalid_pagination")
}
offset <- if (is.null(offset)) 0L else as.integer(offset)
if (length(offset) != 1L || is.na(offset) || offset < 0L) {
cli::cli_abort("`offset` must be a single non-negative integer.",
class = "uscogdata_invalid_pagination")
}
if (complete) { if (complete) {
cli::cli_abort(c( cli::cli_abort(c(
"`limit`/`offset` cannot be combined with `complete = TRUE`.", "`limit`/`offset` cannot be combined with `complete = TRUE`.",
@@ -466,24 +459,14 @@ cog_spending <- function(govid = NULL, years, category = NULL,
limit = limit, offset = offset) limit = limit, offset = offset)
result <- tibble::as_tibble(DBI::dbGetQuery(con, sql)) result <- tibble::as_tibble(DBI::dbGetQuery(con, sql))
if (!is.null(limit)) { if (!is.null(limit)) {
# COUNT(*) OVER() rides along as an ordinary column so the total comes paged <- .take_pagination_total(result, con, function() {
# from the same scan when this page has any rows -- see .build_verb_sql(view, subtype_col, cohort, years,
# .build_verb_sql(). An empty page (offset past the end) carries no if (all_categories) NULL else category,
# such row to read it from, so that one case falls back to a second, ig_view, subtype_scope,
# unpaginated COUNT(*) query rather than reporting a wrong zero. all_categories = all_categories)
if (nrow(result) > 0L) { })
total_rows <- result$pagination_total_rows[[1]] result <- paged$result
result$pagination_total_rows <- NULL total_rows <- paged$total_rows
} else {
count_sql <- sprintf(
"SELECT COUNT(*) AS n FROM (%s) AS _uncounted",
.build_verb_sql(view, subtype_col, cohort, years,
if (all_categories) NULL else category,
ig_view, subtype_scope,
all_categories = all_categories)
)
total_rows <- as.integer(DBI::dbGetQuery(con, count_sql)$n[[1]])
}
} }
} }
@@ -873,22 +856,9 @@ cog_spending <- function(govid = NULL, years, category = NULL,
# matching row across the network only to slice and discard most of it # matching row across the network only to slice and discard most of it
# afterward (the pattern behind the 2026-08-06 production incident: a # afterward (the pattern behind the 2026-08-06 production incident: a
# 193,105-row/194-page sweep re-ran the full query and re-listified every # 193,105-row/194-page sweep re-ran the full query and re-listified every
# row on EVERY page). COUNT(*) OVER() rides along as an ordinary column so # row on EVERY page). See .paginate_sql() in R/pagination.R for why the
# the caller gets the true total from this same scan -- see the call site # wrapping is an outer SELECT rather than a bare LIMIT on base_sql.
# in .verb_spendrev(), which reads it off row 1 and strips it back out. .paginate_sql(base_sql, limit, offset)
# The outer SELECT * wrapping (rather than appending LIMIT/OFFSET directly
# to base_sql) is what makes COUNT(*) OVER() see the post-GROUP-BY row
# count, not the pre-aggregation one.
if (is.null(limit)) {
base_sql
} else {
sprintf(
"SELECT *, COUNT(*) OVER() AS pagination_total_rows
FROM (%s) AS _paged
LIMIT %d OFFSET %d",
base_sql, limit, offset
)
}
} }
#' Join population onto a result and derive the per-capita columns. #' Join population onto a result and derive the per-capita columns.
+11
View File
@@ -134,6 +134,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_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_CACHE_DIR` — where the manifest is cached (default: user cache dir)
- `USCOGDATA_MANIFEST_TTL_SECS` — manifest re-fetch interval (default 3600) - `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 ## Amounts are in full US dollars
+15 -1
View File
@@ -13,7 +13,9 @@ cog_balances(
basis = c("harmonized", "raw"), basis = c("harmonized", "raw"),
recipe = NULL, recipe = NULL,
state = NULL, state = NULL,
type = NULL type = NULL,
limit = NULL,
offset = NULL
) )
} }
\arguments{ \arguments{
@@ -76,6 +78,14 @@ wide era to the modern one.}
and `govids_missing` are empty -- there is no id list to report against -- and `govids_missing` are empty -- there is no id list to report against --
and `provenance$scope$cohort` carries `state`, `type` and and `provenance$scope$cohort` carries `state`, `type` and
`n_governments` instead. A `govid`-named cohort reports exactly as before.} `n_governments` instead. A `govid`-named cohort reports exactly as before.}
\item{limit}{Maximum number of result rows to return, pushed into the SQL
rather than applied after materializing every row. `NULL` (default)
returns everything. Cannot be combined with `recipe` -- see `offset` and
`total_rows`.}
\item{offset}{Rows to skip before `limit` starts counting (0-based).
Ignored if `limit` is `NULL`; defaults to `0L` when `limit` is set.}
} }
\value{ \value{
Tibble with columns `year`, `canonical_govid`, `gov_name`, Tibble with columns `year`, `canonical_govid`, `gov_name`,
@@ -94,6 +104,10 @@ Tibble with columns `year`, `canonical_govid`, `gov_name`,
and `truncated` (the observed subtypes whose coverage falls short of the and `truncated` (the observed subtypes whose coverage falls short of the
requested years). `expenditure_concept`/`revenue_concept` are `NA` -- requested years). `expenditure_concept`/`revenue_concept` are `NA` --
holdings are a stock, not a flow, so neither concept vocabulary applies. holdings are a stock, not a flow, so neither concept vocabulary applies.
When `limit` is set, also carries a `total_rows` attribute: the full
unpaginated row count, computed by the same query (`COUNT(*) OVER()`)
rather than a second scan.
} }
\description{ \description{
Returns Census cash-and-security holdings (`category_type = "balance"`): Returns Census cash-and-security holdings (`category_type = "balance"`):
+24 -4
View File
@@ -4,7 +4,13 @@
\alias{cog_gov_search} \alias{cog_gov_search}
\title{Search for governments by name, state, and/or type} \title{Search for governments by name, state, and/or type}
\usage{ \usage{
cog_gov_search(name = NULL, state = NULL, type = NULL) cog_gov_search(
name = NULL,
state = NULL,
type = NULL,
limit = NULL,
offset = NULL
)
} }
\arguments{ \arguments{
\item{name}{Character vector of place name(s). Length 1 = utility mode; \item{name}{Character vector of place name(s). Length 1 = utility mode;
@@ -19,12 +25,26 @@ all entries; otherwise must match `length(name)`.}
in basket mode (recycles from length 1). Excluded types `4`/`5` (or in basket mode (recycles from length 1). Excluded types `4`/`5` (or
`"special_district"` / `"school_district"`) trigger an explanatory `"special_district"` / `"school_district"`) trigger an explanatory
message and an empty result.} message and an empty result.}
\item{limit}{Maximum number of rows to return, applied in SQL. `NULL`
(default) returns every match -- which, with no other filter, is the
entire crosswalk. Utility mode only: pagination has no meaning in basket
mode, where the result is one resolved row per requested name in input
order, and is refused there with class
`uscogdata_basket_pagination_conflict`.}
\item{offset}{Rows to skip before `limit` starts counting (0-based).
Ignored if `limit` is `NULL`; defaults to `0L` when `limit` is set.}
} }
\value{ \value{
A tibble of `canonical_fips_xwalk` rows. In utility mode, all A tibble of `canonical_fips_xwalk` rows. In utility mode, all
matches sorted by `population_acs` desc. In basket mode, resolved matches sorted by `population_acs` desc, ties broken by
rows in input order, with `attr(., "resolution")` set to the `canonical_govid`. In basket mode, resolved rows in input order, with
sidecar tibble. `attr(., "resolution")` set to the sidecar tibble.
When `limit` is set, carries a `total_rows` attribute: the full
unpaginated match count, computed by the same query (`COUNT(*) OVER()`)
rather than a second scan.
} }
\description{ \description{
Resolves human-readable place names into rows of `canonical_fips_xwalk`, Resolves human-readable place names into rows of `canonical_fips_xwalk`,
+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))
})
@@ -0,0 +1,178 @@
# tests/testthat/test-search-balances-pagination.R
#
# uscogdata#57. cog_spending()/cog_revenue() gained limit/offset in #39;
# cog_gov_search() and cog_balances() did not, so every consumer of those two
# was back to materialize-then-slice -- the exact pattern that wedged the
# production API for hours on 2026-08-06.
#
# cog_gov_search() was also the one verb with no LIMIT at all, so an
# unfiltered call returns the entire 40,336-row crosswalk by accident.
# --- cog_gov_search() -------------------------------------------------------
test_that("cog_gov_search() limit returns the first page of the unpaginated result", {
skip_if_no_corpus()
full <- cog_gov_search(state = "WI", type = "city")
skip_if(nrow(full) < 12L, "fixture has too few WI cities to page")
page <- cog_gov_search(state = "WI", type = "city", limit = 5L)
expect_equal(nrow(page), 5L)
expect_equal(page$canonical_govid, full$canonical_govid[1:5])
})
test_that("cog_gov_search() offset skips ahead without gaps or overlap", {
skip_if_no_corpus()
full <- cog_gov_search(state = "WI", type = "city")
skip_if(nrow(full) < 12L, "fixture has too few WI cities to page")
p1 <- cog_gov_search(state = "WI", type = "city", limit = 5L)
p2 <- cog_gov_search(state = "WI", type = "city", limit = 5L, offset = 5L)
expect_equal(p2$canonical_govid, full$canonical_govid[6:10])
expect_length(intersect(p1$canonical_govid, p2$canonical_govid), 0L)
})
test_that("walking every page reconstructs the unpaginated search exactly", {
skip_if_no_corpus()
full <- cog_gov_search(state = "WI", type = "city")
n <- nrow(full)
limit <- 7L
pages <- list()
offset <- 0L
repeat {
p <- cog_gov_search(state = "WI", type = "city", limit = limit, offset = offset)
if (nrow(p) == 0L) break
pages[[length(pages) + 1L]] <- p
offset <- offset + limit
if (offset > n + limit) stop("test runaway: paging did not terminate")
}
walked <- dplyr::bind_rows(pages)
expect_equal(nrow(walked), n)
expect_equal(walked$canonical_govid, full$canonical_govid)
})
test_that("cog_gov_search() total_rows reports the full unpaginated count", {
skip_if_no_corpus()
full <- cog_gov_search(state = "WI", type = "city")
page <- cog_gov_search(state = "WI", type = "city", limit = 3L)
expect_equal(attr(page, "total_rows"), nrow(full))
})
test_that("cog_gov_search() offset past the end reports the true total, not zero", {
skip_if_no_corpus()
full <- cog_gov_search(state = "WI", type = "city")
# No row survives to carry COUNT(*) OVER(), so this is the branch that has
# to fall back to a second count rather than reporting 0 rows out of 0.
page <- cog_gov_search(state = "WI", type = "city",
limit = 5L, offset = nrow(full) + 50L)
expect_equal(nrow(page), 0L)
expect_equal(attr(page, "total_rows"), nrow(full))
})
test_that("cog_gov_search() bounds an otherwise-unfiltered crosswalk sweep", {
skip_if_no_corpus()
# The reason this verb needed a limit most: with no filter it returns the
# whole crosswalk.
page <- cog_gov_search(limit = 10L)
expect_equal(nrow(page), 10L)
expect_gt(attr(page, "total_rows"), 10L)
})
test_that("cog_gov_search() orders by a total order, not population alone", {
skip_if_no_corpus()
# population_acs is not unique -- NA in particular repeats across many rows
# -- so paging on it alone can duplicate a row on one page and drop it from
# the next. The tiebreaker is what makes the sequence reproducible.
full <- cog_gov_search(state = "WI")
skip_if(nrow(full) < 5L, "fixture has too few WI governments")
expect_equal(cog_gov_search(state = "WI")$canonical_govid,
full$canonical_govid)
ties <- full[is.na(full$population_acs), ]
skip_if(nrow(ties) < 2L, "no tied rows in the fixture to order")
expect_false(is.unsorted(ties$canonical_govid))
})
test_that("cog_gov_search() refuses pagination in basket mode", {
skip_if_no_corpus()
expect_error(
cog_gov_search(name = c("MADISON CITY", "MILWAUKEE CITY"),
state = c("WI", "WI"), limit = 1L),
class = "uscogdata_basket_pagination_conflict"
)
})
test_that("cog_gov_search() rejects a malformed limit or offset", {
skip_if_no_corpus()
expect_error(cog_gov_search(state = "WI", limit = -1L),
class = "uscogdata_invalid_pagination")
expect_error(cog_gov_search(state = "WI", limit = 5L, offset = -1L),
class = "uscogdata_invalid_pagination")
})
# --- cog_balances() ---------------------------------------------------------
test_that("cog_balances() limit/offset walk the unpaginated result exactly", {
skip_if_no_corpus()
full <- cog_balances(years = 2019:2020, state = "WI", type = "city")
skip_if(nrow(full) < 6L, "fixture has too few WI city balance rows to page")
key <- c("year", "canonical_govid", "balance_subtype", "amt_nominal")
p1 <- cog_balances(years = 2019:2020, state = "WI", type = "city", limit = 3L)
p2 <- cog_balances(years = 2019:2020, state = "WI", type = "city",
limit = 3L, offset = 3L)
expect_equal(nrow(p1), 3L)
expect_equal(p1[key], full[1:3, key], ignore_attr = TRUE)
expect_equal(p2[key], full[4:6, key], ignore_attr = TRUE)
# The window-function column is an implementation detail and must not reach
# the caller's data frame.
expect_false("pagination_total_rows" %in% names(p1))
})
test_that("cog_balances() total_rows reports the full unpaginated count", {
skip_if_no_corpus()
full <- cog_balances(years = 2019:2020, state = "WI", type = "city")
page <- cog_balances(years = 2019:2020, state = "WI", type = "city", limit = 2L)
expect_equal(attr(page, "total_rows"), nrow(full))
})
test_that("cog_balances() offset past the end reports the true total", {
skip_if_no_corpus()
full <- cog_balances(years = 2019:2020, state = "WI", type = "city")
page <- cog_balances(years = 2019:2020, state = "WI", type = "city",
limit = 5L, offset = nrow(full) + 50L)
expect_equal(nrow(page), 0L)
expect_equal(attr(page, "total_rows"), nrow(full))
})
test_that("cog_balances() refuses pagination alongside a recipe", {
skip_if_no_corpus()
expect_error(
cog_balances(years = 2011, state = "WI", type = "city",
recipe = "cash_securities_z77_wide", limit = 5L),
class = "uscogdata_recipe_pagination_conflict"
)
})
test_that("cog_balances() rejects a malformed limit or offset", {
skip_if_no_corpus()
expect_error(cog_balances(years = 2019, state = "WI", type = "city", limit = -1L),
class = "uscogdata_invalid_pagination")
expect_error(cog_balances(years = 2019, state = "WI", type = "city",
limit = 5L, offset = -1L),
class = "uscogdata_invalid_pagination")
})
# --- Unchanged without the arguments ----------------------------------------
test_that("both verbs are unchanged when limit is not supplied", {
skip_if_no_corpus()
# The adoption contract for cog-api: NULL default, so a formals() probe can
# feature-detect without any call site changing behaviour.
s <- cog_gov_search(state = "WI", type = "city")
b <- cog_balances(years = 2019, state = "WI", type = "city")
expect_null(attr(s, "total_rows"))
expect_null(attr(b, "total_rows"))
expect_true(all(c("limit", "offset") %in% names(formals(cog_gov_search))))
expect_true(all(c("limit", "offset") %in% names(formals(cog_balances))))
})