Compare commits

..
Author SHA1 Message Date
jared 8bd16bb085 docs: re-measure the corpus-access table against the published corpus (#56)
R-CMD-check / check (pull_request) Successful in 3m42s
R-CMD-check / check (push) Successful in 3m39s
#56 step 4 asked for the README table to be re-measured after the pass. The old
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:16:03 -04:00
21 changed files with 87 additions and 1012 deletions
-12
View File
@@ -39,17 +39,5 @@ jobs:
echo "PAT_GH is unset -- add it under Settings > Actions > Secrets." >&2
exit 1
fi
# On a tag-triggered run, checkout materializes refs/tags/<tag> as a
# LIGHTWEIGHT tag at the commit SHA -- the annotated tag object Gitea
# holds is never fetched. Mirroring that strips the annotation, and the
# NEXT run on main (which does fetch the real object) is then rejected
# with "already exists" trying to correct it, because git will not
# clobber an existing tag. That is why v0.4.0 failed to mirror.
#
# Re-fetch canonical tag objects from Gitea first. --force here rewrites
# LOCAL tag refs only; it is not a force push and does not weaken the
# non-force guarantee on main documented above.
git fetch --tags --force origin
git push "https://x-access-token:${PAT_GH}@github.com/civilytics/uscogdata.git" \
HEAD:refs/heads/main --tags
-64
View File
@@ -1,64 +0,0 @@
# Explain the mirror contribution flow on every incoming pull request.
#
# This repository is a MIRROR. A PR opened here is landed on the canonical Gitea
# repository and syncs back; because the merge preserves the contributor's
# commits at their original SHAs, GitHub marks the PR "Merged" on its own as
# soon as the mirror syncs -- with nobody visibly clicking Merge.
#
# Without this comment, that reads as a rejection: the contributor sees their PR
# close with no review, no merge button pressed, and no explanation. It is
# actually the successful outcome. Say so up front, before it happens.
#
# WHY pull_request_target AND NOT pull_request:
# a `pull_request` run from a fork gets a read-only token, so it cannot post a
# comment -- which is exactly the case this workflow exists to serve.
# `pull_request_target` runs in the context of the BASE repo and gets a writable
# token. That is only safe because this job never checks out or executes the
# contributor's code; it posts a fixed string. Do not add a checkout of
# `github.event.pull_request.head.sha` here -- that combination is the standard
# pull_request_target privilege-escalation hole.
name: Explain the mirror flow
on:
pull_request_target:
types: [opened]
permissions:
pull-requests: write
jobs:
comment:
runs-on: ubuntu-latest
steps:
- name: Post the contribution-flow explainer
uses: actions/github-script@v7
with:
script: |
const body = [
"Thanks for this — and one thing worth knowing before it happens.",
"",
"**This repository is a mirror.** Development happens on Gitea at",
"`gitea.civilytics.org/Civilytics/uscogdata`. Your pull request will be fetched",
"from here, landed there, and synced back.",
"",
"Because that merge preserves your commits at their original SHAs, **GitHub will",
"mark this pull request \"Merged\" on its own** — without anyone visibly clicking",
"the Merge button, and possibly without a review comment on this page first.",
"",
"> If your pull request closes as \"Merged\" and nobody appears to have merged it,",
"> that is the normal, successful outcome — not a rejection.",
"",
"If it is *not* going to be merged, you will get an actual reply saying so.",
"",
"Substantial contributions get a `ctb` entry in `DESCRIPTION`, which surfaces in",
"`citation(\"uscogdata\")`. There is no CLA and no DCO sign-off.",
"",
"Full details: [CONTRIBUTING.md](https://github.com/civilytics/uscogdata/blob/main/CONTRIBUTING.md).",
].join("\n");
await github.rest.issues.createComment({
owner: context.repo.owner,
repo: context.repo.repo,
issue_number: context.payload.pull_request.number,
body,
});
+1 -14
View File
@@ -4,17 +4,7 @@
.Ruserdata
*.Rproj
inst/doc
# pkgdown output. Listed as children rather than `docs/` so compass's
# docs/pm/ and docs/decisions/ can be re-included -- git cannot re-include
# anything beneath an excluded directory.
#
# Note the anchoring change this forces: a bare `docs/` matches a directory of
# that name at ANY depth, while `/docs/*` matches only at the repo root. The
# fixture corpus's own docs/ therefore needs its own rule to stay excluded.
/docs/*
!/docs/pm/
!/docs/decisions/
inst/extdata/fixture_corpus/docs/
docs/
/doc/
/Meta/
.DS_Store
@@ -22,6 +12,3 @@ inst/extdata/fixture_corpus/docs/
# SDD working artifacts (ledger, briefs, review packages) — plans/ stays tracked
.superpowers/sdd/
.compass-cache/
# roborev snapshots
/.roborev/
-56
View File
@@ -1,56 +0,0 @@
# roborev configuration, initialised by compass.
# Reviews are queued to a background daemon -- they never block a commit.
post_commit_review = 'commit'
excluded_commit_patterns = ['WIP', 'chore:', 'docs:', 'Merge ']
review_guidelines = '''
# --- compass:begin (generated -- edit the sources, not this) ---
- Prefer returning new values to mutating arguments in place. A function that edits
its caller's object is a bug waiting for a second caller.
- Validate at system boundaries -- user input, API responses, file contents, config.
Fail fast with a message naming the field and the file.
- Never swallow an error. Handle it or let it propagate; a bare catch that continues
is worse than a crash.
- No hardcoded secrets, tokens, or credentials, and no secrets in log output or error
messages.
- Parameterise every query. String-built SQL is a defect even when the input looks safe.
- Keep functions under roughly 50 lines and files under roughly 400. Flag nesting
deeper than four levels.
- No magic numbers or hardcoded paths -- name them as constants or read them from config.
- New behaviour needs a test. A bug fix needs a test that fails without the fix.
- Use the native pipe `|>`, not magrittr `%>%`.
- snake_case for objects and functions; UPPER_SNAKE for constants. Never use `.` as a
word separator in a function name -- it collides with S3 dispatch.
- Validate arguments at the top of exported functions with `stopifnot()` or an explicit
check, and say which argument was wrong.
- Never `setDT()`, `set()`, or otherwise modify by reference a data.table the caller
still owns. `as.data.table()` copies; use it.
- Prefer `vapply()` to `sapply()` -- `sapply()` silently returns a list when the type
varies, which turns a type error into a downstream mystery.
- Use `seq_len(n)` / `seq_along(x)`, never `1:n`, which iterates backwards when n is 0.
- Compare strings with `==` only after checking for NA; use `identical()` for scalars
where NA would be wrong.
- Do not call `library()` inside package or module files; attach packages in scripts and
test helpers only.
- Namespace-qualify calls into other packages (`stats::sd`) in code that is sourced.
- Every exported function needs roxygen with `@param` for each argument (type, meaning,
and why the default is what it is) and `@return`. Add `@examples` for exported API.
- Declare dependencies in DESCRIPTION. Prefer base R or an existing dependency over
adding a new one; a package with zero hard deps is worth keeping that way.
- Signal errors with `stop()` carrying a condition class, so callers can catch the kind
rather than matching on message text.
- Keep internals internal. Export only what a user needs; an accidentally exported
helper becomes an API you have to keep.
- Tests use testthat edition 3. Each test is self-sufficient -- no reliance on state
left by an earlier test or on a fixture built elsewhere in the file.
- Prefer duplication in tests over a helper that hides what is being asserted.
- Every verb calls .ensure_session() first, then queries via DBI::dbGetQuery().
- A verb's return value is always a tbl_df carrying a provenance attribute.
- govid inputs always go through .coerce_govid_input(); it accepts a character vector or a data frame.
- SQL has two layers: view definitions are numbered .sql files in inst/sql/ registered by .register_views(); query construction is inline sprintf() in R. Add a view as a file; build a query in R.
- No arrow dependency -- DuckDB reads parquet natively.
- withr is Suggests-only and must appear in tests alone.
- Tests must pass offline against the bundled fixture; tests/testthat/setup.R sets USCOGDATA_URL for that.
# --- compass:end ---
'''
+2 -10
View File
@@ -101,13 +101,5 @@ but "usually" is not a release gate.
variables set. This is the only check that catches a
corpus-unreachable defect, and its absence is why 0.3.0 needed fixing.
8. Bump `Version` and add a `NEWS.md` section.
9. Tag on **Gitea** (`git tag -a vX.Y.Z && git push origin vX.Y.Z`). The mirror
workflow carries tags to GitHub on its own — confirm the tag appears at
`github.com/civilytics/uscogdata/tags` before continuing.
10. Update the r-universe registry pin at
`github.com/civilytics/civilytics.r-universe.dev` — edit `packages.json`'s
`branch` to the new tag. **r-universe will not pick up a release until this
is edited**: the pin is a tag, deliberately, so a mid-refactor `main` is
never published as a release. `"branch": "*release"` would track releases
automatically, but it needs a GitHub *Release* object and the mirror pushes
tags only — so it would silently never update.
9. Tag, then update the r-universe registry pin at
`github.com/civilytics/civilytics.r-universe.dev`.
+24 -73
View File
@@ -1,5 +1,29 @@
# 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
`cog_spending()`, `cog_revenue()` and `cog_balances()` gain optional `state`
@@ -46,81 +70,8 @@ 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.
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
* `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"
(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
+2 -41
View File
@@ -47,12 +47,6 @@
#' @param recipe Optional harmonization recipe id (see [cog_recipes()]).
#' `"cash_securities_z77_wide"` and `"cash_securities_z78_wide"` bridge the
#' 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`,
#' `balance_subtype`, `category`, `amt_nominal`, `codes_included`,
@@ -70,16 +64,11 @@
#' and `truncated` (the observed subtypes whose coverage falls short of the
#' requested years). `expenditure_concept`/`revenue_concept` are `NA` --
#' 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
cog_balances <- function(govid = NULL, years, category = NULL,
per_capita = FALSE, adjust_to_year = NULL,
basis = c("harmonized", "raw"), recipe = NULL,
state = NULL, type = NULL,
limit = NULL, offset = NULL) {
state = NULL, type = NULL) {
call <- match.call()
basis <- match.arg(basis, c("harmonized", "raw"))
# Coerce FIRST, validate second: .validate_verb_inputs() asserts
@@ -102,22 +91,6 @@ cog_balances <- function(govid = NULL, years, category = NULL,
# validator's own doc comment for the incident that made that matter.
.validate_verb_inputs(govid, years, category, per_capita, adjust_to_year,
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)
if (!is.null(adjust_to_year)) adjust_to_year <- as.integer(adjust_to_year)
@@ -135,7 +108,6 @@ cog_balances <- function(govid = NULL, years, category = NULL,
manifest <- .uscogdata_env$manifest
recipe_block <- NULL
category_for_prov <- category
total_rows <- NULL # set below only when limit is non-NULL (non-recipe path)
if (!is.null(recipe)) {
.require_schema_v5(con, manifest, "recipe =")
@@ -153,18 +125,8 @@ cog_balances <- function(govid = NULL, years, category = NULL,
} else {
sql <- .build_verb_sql("balance_annotated", "balance_subtype",
cohort, years, category,
ig_view = NULL, subtype_scope = NULL,
limit = limit, offset = offset)
ig_view = NULL, subtype_scope = NULL)
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
@@ -196,7 +158,6 @@ cog_balances <- function(govid = NULL, years, category = NULL,
.emit_balance_caveats(prov$balance_caveats)
attr(result, "provenance") <- prov
if (!is.null(limit)) attr(result, "total_rows") <- total_rows
result
}
+1 -60
View File
@@ -18,15 +18,7 @@
# 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,
# 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
manifest_ttl_secs = 3600L
)
#' Resolve a config value: env var > option > default
@@ -69,54 +61,3 @@
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)
}
-81
View File
@@ -1,81 +0,0 @@
# 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]]))
}
+7 -50
View File
@@ -44,22 +44,10 @@
#' in basket mode (recycles from length 1). Excluded types `4`/`5` (or
#' `"special_district"` / `"school_district"`) trigger an explanatory
#' 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
#' matches sorted by `population_acs` desc, ties broken by
#' `canonical_govid`. In basket mode, resolved rows in input order, with
#' `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.
#' matches sorted by `population_acs` desc. In basket mode, resolved
#' rows in input order, with `attr(., "resolution")` set to the
#' sidecar tibble.
#' @seealso [cog_basket_resolution()], [cog_basket_unresolved()],
#' [cog_spending()], [cog_revenue()].
#' @examples
@@ -94,12 +82,7 @@
#' )
#' }
#' @export
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
cog_gov_search <- function(name = NULL, state = NULL, type = NULL) {
if (!is.null(type) && length(type) == 1L && .is_excluded_type(type)) {
cli::cli_inform(c(
i = "v0.1 covers gov_types 0-3 (state/county/city/township) only.",
@@ -111,18 +94,6 @@ cog_gov_search <- function(name = NULL, state = NULL, type = NULL,
con <- .ensure_session()
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))
}
@@ -152,26 +123,12 @@ cog_gov_search <- function(name = NULL, state = NULL, type = NULL,
}
where <- if (length(preds) == 0L) "" else paste("WHERE", paste(preds, collapse = " AND "))
# 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(
sql <- paste(
"SELECT * FROM canonical_fips_xwalk",
where,
"ORDER BY population_acs DESC NULLS LAST, canonical_govid"
"ORDER BY population_acs DESC NULLS LAST"
)
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
tibble::as_tibble(DBI::dbGetQuery(con, sql))
}
#' @noRd
+1 -31
View File
@@ -2,23 +2,13 @@
#' 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(),
threads = .resolve_duckdb_threads(),
memory_limit = .resolve_duckdb_memory_limit()) {
cache_dir = .resolve_cache_dir()) {
.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)
@@ -35,26 +25,6 @@ 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) ||
+44 -14
View File
@@ -398,10 +398,17 @@ cog_spending <- function(govid = NULL, years, category = NULL,
# up front rather than silently ignored: complete = TRUE fills a grid over
# the FULL requested (year, category) space, and a recipe's result comes
# 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)) {
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) {
cli::cli_abort(c(
"`limit`/`offset` cannot be combined with `complete = TRUE`.",
@@ -459,14 +466,24 @@ cog_spending <- function(govid = NULL, years, category = NULL,
limit = limit, offset = offset)
result <- tibble::as_tibble(DBI::dbGetQuery(con, sql))
if (!is.null(limit)) {
paged <- .take_pagination_total(result, con, function() {
.build_verb_sql(view, subtype_col, cohort, years,
if (all_categories) NULL else category,
ig_view, subtype_scope,
all_categories = all_categories)
})
result <- paged$result
total_rows <- paged$total_rows
# COUNT(*) OVER() rides along as an ordinary column so the total comes
# from the same scan when this page has any rows -- see
# .build_verb_sql(). An empty page (offset past the end) carries no
# such row to read it from, so that one case falls back to a second,
# unpaginated COUNT(*) query rather than reporting a wrong zero.
if (nrow(result) > 0L) {
total_rows <- result$pagination_total_rows[[1]]
result$pagination_total_rows <- NULL
} 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]])
}
}
}
@@ -856,9 +873,22 @@ cog_spending <- function(govid = NULL, years, category = NULL,
# matching row across the network only to slice and discard most of it
# 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
# row on EVERY page). See .paginate_sql() in R/pagination.R for why the
# wrapping is an outer SELECT rather than a bare LIMIT on base_sql.
.paginate_sql(base_sql, limit, offset)
# row on EVERY page). COUNT(*) OVER() rides along as an ordinary column so
# the caller gets the true total from this same scan -- see the call site
# in .verb_spendrev(), which reads it off row 1 and strips it back out.
# 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.
-12
View File
@@ -1,7 +1,6 @@
# uscogdata
<!-- badges: start -->
[![R-CMD-check](https://github.com/civilytics/uscogdata/actions/workflows/R-CMD-check.yaml/badge.svg)](https://github.com/civilytics/uscogdata/actions/workflows/R-CMD-check.yaml)
[![r-universe](https://civilytics.r-universe.dev/badges/uscogdata)](https://civilytics.r-universe.dev/uscogdata)
[![License: MIT](https://img.shields.io/badge/license-MIT-blue.svg)](LICENSE.md)
<!-- badges: end -->
@@ -135,17 +134,6 @@ 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
-12
View File
@@ -1,12 +0,0 @@
# Decisions
One file per decision, numbered and immutable. A decision that changes is superseded
by a new record, never edited in place — the old reasoning is the point.
The table below is **generated** by `compass:decide`. Do not hand-edit it.
<!-- compass:begin decisions -->
| # | Date | Decision | Status |
|---|---|---|---|
| — | — | *No decisions recorded yet.* | — |
<!-- compass:end decisions -->
-14
View File
@@ -1,14 +0,0 @@
# Project journal
Append-only, newest first. **Entries are never edited** — the value of this file is
that it records what was believed at the time, including the parts that turned out
wrong. Where things stand *today* is in `STATUS.md`, which is generated.
Four lines per entry. The analysis belongs in the issue or the decision record; this
file carries the reasoning and the pointers.
- **Why** — the driver. The one line git cannot reconstruct later.
- **Obligates** — issues this change created elsewhere. Numbers, not prose.
- **Refs** — commits, issues, decision records.
---
-73
View File
@@ -1,73 +0,0 @@
# Project status
> Between the compass markers is generated. Edit the sources, not this.
<!-- compass:begin -->
<!-- compass:board -->
## Where this stands
uscogdata is at 0.4.0 and its public surface is settled: the query verbs, the cohort
predicates added in this release, and the provenance contract every verb returns.
The six open issues split cleanly. Two are API work carried out of the #9 review pass
and deliberately deferred there rather than fixed in that branch. Three concern the
corpus layer, and the largest of them, partition-level caching, was named the single
highest-leverage change on the remote path before being deferred. One, the
data-correction intake, is a decision rather than a task: it was parked during the
0.3.0 design and it gates the API announcement, because without it the corpus cannot
make the "traceable and correctable" claim that most distinguishes it from Census's
own files.
Nothing here is blocked on anything else, so the ordering is a judgement about value
rather than a dependency graph.
## Ready to work on next
- **#34** cog_revenue() offers expenditure recipes as suggestions: scope the candidate query by category_type · `ws/api` — nothing is blocking it; something is currently wrong
- **#36** n_units_reporting is category-conditional and cannot be read as a response rate · `ws/corpus` — nothing is blocking it; owed work from an earlier change
- **#2** Extend population data to be households as an alternate spending denominator · `ws/corpus` — nothing is blocking it
- **#33** Decompose .build_suggestions() (106 lines) into named helpers · `ws/api` — nothing is blocking it
- **#52** Release 11/11: design the data-correction intake (deferred; gates the API announcement) · `ws/corpus` — nothing is blocking it
- **#64** Partition-level caching: R/cache.R is still a stub, and the remote path pays for it every session · `ws/corpus` — nothing is blocking it
## Workstreams
| Stream | Commits since | Open | Debt | Owes docs |
|---|---|---|---|---|
| Query verbs and results | 77 | 2 | 0 | no |
| Corpus, mirror, provenance | 39 | 4 | 1 | no |
| Vignettes and guides | 34 | 0 | 0 | **yes** |
## CI
![R-CMD-check](https://gitea.civilytics.org/Civilytics/uscogdata/actions/workflows/ci.yml/badge.svg?branch=main)
![Mirror to GitHub](https://gitea.civilytics.org/Civilytics/uscogdata/actions/workflows/mirror-github.yml/badge.svg?branch=main)
<details>
<summary>Dependency graph and detail</summary>
```mermaid
graph TD
I34["#34 cog_revenue() offers expenditure recipes as sug…"]
I36["#36 n_units_reporting is category-conditional and c…"]
I2["#2 Extend population data to be households as an a…"]
I33["#33 Decompose .build_suggestions() (106 lines) into…"]
I52["#52 Release 11/11: design the data-correction intak…"]
I64["#64 Partition-level caching: R/cache.R is still a s…"]
class I34 ready;
class I36 ready;
class I2 ready;
class I33 ready;
class I52 ready;
class I64 ready;
classDef ready fill:#dafbe1,stroke:#2da44e;
```
- Marker: `none` (no journal entry yet)
- Commits since: 164
- Open issues: 6
</details>
<!-- compass:end -->
-44
View File
@@ -1,44 +0,0 @@
[project]
name = "uscogdata"
forge = "Civilytics/uscogdata"
# Three strands that go stale independently: what the verbs return, what the
# corpus is and how it is mounted, and how both are explained to a reader.
[[workstream]]
id = "api"
title = "Query verbs and results"
paths = [
"R/revenue.R", "R/spending.R", "R/balances.R", "R/peers.R", "R/search.R",
"R/categories.R", "R/recipes.R", "R/rollup.R", "R/explain.R", "R/basket.R",
"R/suggestions.R", "R/suppression.R", "R/complete.R", "R/cohort.R",
"R/basis.R", "R/adjust.R", "R/pagination.R",
]
docs = ["vignettes/*.Rmd", "README.md"]
[[workstream]]
id = "corpus"
title = "Corpus, mirror, provenance"
paths = [
"R/manifest.R", "R/mirror.R", "R/cache.R", "R/session.R", "R/provenance.R",
"R/coverage.R", "R/config.R", "R/views.R", "R/series_breaks.R",
"R/balance_caveats.R", "R/zzz.R", "data-raw/**", "inst/sql/**",
]
docs = ["vignettes/*.Rmd", "NEWS.md"]
[[workstream]]
id = "docs"
title = "Vignettes and guides"
paths = ["vignettes/**", "README.md", "_pkgdown.yml", "NEWS.md"]
docs = []
[roborev]
project_guidelines = [
"Every verb calls .ensure_session() first, then queries via DBI::dbGetQuery().",
"A verb's return value is always a tbl_df carrying a provenance attribute.",
"govid inputs always go through .coerce_govid_input(); it accepts a character vector or a data frame.",
"SQL has two layers: view definitions are numbered .sql files in inst/sql/ registered by .register_views(); query construction is inline sprintf() in R. Add a view as a file; build a query in R.",
"No arrow dependency -- DuckDB reads parquet natively.",
"withr is Suggests-only and must appear in tests alone.",
"Tests must pass offline against the bundled fixture; tests/testthat/setup.R sets USCOGDATA_URL for that.",
]
+1 -15
View File
@@ -13,9 +13,7 @@ cog_balances(
basis = c("harmonized", "raw"),
recipe = NULL,
state = NULL,
type = NULL,
limit = NULL,
offset = NULL
type = NULL
)
}
\arguments{
@@ -78,14 +76,6 @@ wide era to the modern one.}
and `govids_missing` are empty -- there is no id list to report against --
and `provenance$scope$cohort` carries `state`, `type` and
`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{
Tibble with columns `year`, `canonical_govid`, `gov_name`,
@@ -104,10 +94,6 @@ Tibble with columns `year`, `canonical_govid`, `gov_name`,
and `truncated` (the observed subtypes whose coverage falls short of the
requested years). `expenditure_concept`/`revenue_concept` are `NA` --
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{
Returns Census cash-and-security holdings (`category_type = "balance"`):
+4 -24
View File
@@ -4,13 +4,7 @@
\alias{cog_gov_search}
\title{Search for governments by name, state, and/or type}
\usage{
cog_gov_search(
name = NULL,
state = NULL,
type = NULL,
limit = NULL,
offset = NULL
)
cog_gov_search(name = NULL, state = NULL, type = NULL)
}
\arguments{
\item{name}{Character vector of place name(s). Length 1 = utility mode;
@@ -25,26 +19,12 @@ all entries; otherwise must match `length(name)`.}
in basket mode (recycles from length 1). Excluded types `4`/`5` (or
`"special_district"` / `"school_district"`) trigger an explanatory
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{
A tibble of `canonical_fips_xwalk` rows. In utility mode, all
matches sorted by `population_acs` desc, ties broken by
`canonical_govid`. In basket mode, resolved rows in input order, with
`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.
matches sorted by `population_acs` desc. In basket mode, resolved
rows in input order, with `attr(., "resolution")` set to the
sidecar tibble.
}
\description{
Resolves human-readable place names into rows of `canonical_fips_xwalk`,
-134
View File
@@ -1,134 +0,0 @@
# 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))
})
@@ -1,178 +0,0 @@
# 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))))
})