Compare commits
18
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
62741343ee | ||
|
|
0c7c7eb299
|
||
|
|
7274ce3bfe
|
||
|
|
24e86ed598
|
||
|
|
1ec20174b7
|
||
|
|
2e317d3a0f
|
||
|
|
f7c606984f | ||
|
|
9617b86a26
|
||
|
|
224e5d0530 | ||
|
|
5f81ae386b
|
||
|
|
d2caa6de97 | ||
|
|
303aa59b07
|
||
|
|
81f72321ee
|
||
|
|
5cb83d8f1d | ||
|
|
2bb9d41d72
|
||
|
|
700ae93c9c | ||
|
|
f4ab9b6d90
|
||
|
|
9508b98676 |
@@ -3,6 +3,7 @@
|
|||||||
^\.Rproj\.user$
|
^\.Rproj\.user$
|
||||||
^_pkgdown\.yml$
|
^_pkgdown\.yml$
|
||||||
^docs$
|
^docs$
|
||||||
|
^pm$
|
||||||
^Meta$
|
^Meta$
|
||||||
^doc$
|
^doc$
|
||||||
^pkgdown$
|
^pkgdown$
|
||||||
|
|||||||
@@ -39,5 +39,17 @@ jobs:
|
|||||||
echo "PAT_GH is unset -- add it under Settings > Actions > Secrets." >&2
|
echo "PAT_GH is unset -- add it under Settings > Actions > Secrets." >&2
|
||||||
exit 1
|
exit 1
|
||||||
fi
|
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" \
|
git push "https://x-access-token:${PAT_GH}@github.com/civilytics/uscogdata.git" \
|
||||||
HEAD:refs/heads/main --tags
|
HEAD:refs/heads/main --tags
|
||||||
|
|||||||
@@ -0,0 +1,64 @@
|
|||||||
|
# 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,
|
||||||
|
});
|
||||||
@@ -4,6 +4,10 @@
|
|||||||
.Ruserdata
|
.Ruserdata
|
||||||
*.Rproj
|
*.Rproj
|
||||||
inst/doc
|
inst/doc
|
||||||
|
# pkgdown output. Compass used to keep its files in docs/pm/ and
|
||||||
|
# docs/decisions/, which forced this to be written as children with two
|
||||||
|
# re-includes -- git cannot re-include anything beneath an excluded directory.
|
||||||
|
# Compass lives in pm/ now, so the whole directory can be excluded again.
|
||||||
docs/
|
docs/
|
||||||
/doc/
|
/doc/
|
||||||
/Meta/
|
/Meta/
|
||||||
@@ -12,3 +16,6 @@ docs/
|
|||||||
|
|
||||||
# SDD working artifacts (ledger, briefs, review packages) — plans/ stays tracked
|
# SDD working artifacts (ledger, briefs, review packages) — plans/ stays tracked
|
||||||
.superpowers/sdd/
|
.superpowers/sdd/
|
||||||
|
.compass-cache/
|
||||||
|
# roborev snapshots
|
||||||
|
/.roborev/
|
||||||
|
|||||||
@@ -0,0 +1,60 @@
|
|||||||
|
# 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:', '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.
|
||||||
|
- Prose a person reads -- an issue title or body, a journal entry, a decision record,
|
||||||
|
the narrative on the status board -- names the action or the thing, not the shape of
|
||||||
|
the machinery. Flag "gate", "seam", "surface area", "load-bearing", "first-class",
|
||||||
|
"primitive", "blast radius". A project's own defined vocabulary is not the target.
|
||||||
|
- 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 ---
|
||||||
|
'''
|
||||||
+10
-2
@@ -101,5 +101,13 @@ but "usually" is not a release gate.
|
|||||||
variables set. This is the only check that catches a
|
variables set. This is the only check that catches a
|
||||||
corpus-unreachable defect, and its absence is why 0.3.0 needed fixing.
|
corpus-unreachable defect, and its absence is why 0.3.0 needed fixing.
|
||||||
8. Bump `Version` and add a `NEWS.md` section.
|
8. Bump `Version` and add a `NEWS.md` section.
|
||||||
9. Tag, then update the r-universe registry pin at
|
9. Tag on **Gitea** (`git tag -a vX.Y.Z && git push origin vX.Y.Z`). The mirror
|
||||||
`github.com/civilytics/civilytics.r-universe.dev`.
|
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.
|
||||||
|
|||||||
@@ -46,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
@@ -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
@@ -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)
|
||||||
|
}
|
||||||
|
|||||||
@@ -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
@@ -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
@@ -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
@@ -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.
|
||||||
|
|||||||
+136
-54
@@ -84,6 +84,21 @@
|
|||||||
#' @return List of `list(recipe_id, label, available_years, hint,
|
#' @return List of `list(recipe_id, label, available_years, hint,
|
||||||
#' ig_recipe_id, trigger, suppressed_amount, suppressed_years,
|
#' ig_recipe_id, trigger, suppressed_amount, suppressed_years,
|
||||||
#' suppressed_codes)`, possibly empty.
|
#' suppressed_codes)`, possibly empty.
|
||||||
|
#'
|
||||||
|
#' Decomposed (Issue #33) into three extracted helpers to stay within the
|
||||||
|
#' project's "functions under 50 lines" convention:
|
||||||
|
#' \itemize{
|
||||||
|
#' \item `.query_candidate_recipes()` -- candidate recipe lookup by
|
||||||
|
#' category/subtype scope + M/L exclusion.
|
||||||
|
#' \item `.query_recipe_meta()` -- metadata (label, year spans).
|
||||||
|
#' \item `.query_covered_years()` -- Path 1 gap-year coverage via the
|
||||||
|
#' recipe's own generic join.
|
||||||
|
#' }
|
||||||
|
#' The for-loop that merges covered-years + suppressed-components into
|
||||||
|
#' suggestion objects stays inline here because it interleaves
|
||||||
|
#' empty_hit/supp_hit precedence with field assembly. Likewise kept inline:
|
||||||
|
#' the M/L-exclusion design-comment block and the final
|
||||||
|
#' `.attach_ig_counterparts()` call.
|
||||||
#' @noRd
|
#' @noRd
|
||||||
.build_suggestions <- function(con, cohort, years, category, result, basis,
|
.build_suggestions <- function(con, cohort, years, category, result, basis,
|
||||||
flow_prefixes, long_view,
|
flow_prefixes, long_view,
|
||||||
@@ -111,28 +126,8 @@
|
|||||||
# by `category` (`.ALL_CATEGORIES` is never a row in
|
# by `category` (`.ALL_CATEGORIES` is never a row in
|
||||||
# `summary_categories.category`, so a category-keyed sub-select always
|
# `summary_categories.category`, so a category-keyed sub-select always
|
||||||
# came back empty here). The M/L exclusion below is unchanged either way.
|
# came back empty here). The M/L exclusion below is unchanged either way.
|
||||||
candidate_scope_sql <- if (isTRUE(all_categories)) {
|
candidates <- .query_candidate_recipes(con, category, all_categories,
|
||||||
sprintf(
|
subtype_col, subtype_scope)
|
||||||
"SELECT DISTINCT item_code FROM summary_categories WHERE %s IN (%s)",
|
|
||||||
subtype_col, .sql_lit_chr(subtype_scope)
|
|
||||||
)
|
|
||||||
} else {
|
|
||||||
sprintf(
|
|
||||||
"SELECT DISTINCT item_code FROM summary_categories WHERE category IN (%s)",
|
|
||||||
.sql_lit_chr(category)
|
|
||||||
)
|
|
||||||
}
|
|
||||||
candidates <- DBI::dbGetQuery(con, sprintf(
|
|
||||||
"SELECT DISTINCT recipe_id FROM harmonization_recipes
|
|
||||||
WHERE component_code IN (
|
|
||||||
%s
|
|
||||||
)
|
|
||||||
AND recipe_id NOT IN (
|
|
||||||
SELECT DISTINCT recipe_id FROM harmonization_recipes
|
|
||||||
WHERE LEFT(component_code, 1) IN ('M', 'L')
|
|
||||||
)",
|
|
||||||
candidate_scope_sql
|
|
||||||
))$recipe_id
|
|
||||||
if (length(candidates) == 0L) return(list())
|
if (length(candidates) == 0L) return(list())
|
||||||
|
|
||||||
result_years <- if (is.null(result) || nrow(result) == 0L) {
|
result_years <- if (is.null(result) || nrow(result) == 0L) {
|
||||||
@@ -164,36 +159,11 @@
|
|||||||
|
|
||||||
if (length(gap_years) == 0L && nrow(supp) == 0L) return(list())
|
if (length(gap_years) == 0L && nrow(supp) == 0L) return(list())
|
||||||
|
|
||||||
meta <- tibble::as_tibble(DBI::dbGetQuery(con, sprintf(
|
meta <- .query_recipe_meta(con, candidates)
|
||||||
"SELECT recipe_id, any_value(label) AS label,
|
|
||||||
MIN(year_min) AS year_min, MAX(year_max) AS year_max
|
|
||||||
FROM harmonization_recipes
|
|
||||||
WHERE recipe_id IN (%s)
|
|
||||||
GROUP BY recipe_id",
|
|
||||||
.sql_lit_chr(candidates)
|
|
||||||
)))
|
|
||||||
|
|
||||||
# Path 1 (unchanged): (recipe, year) pairs the recipe's own generic join
|
# Path 1 (unchanged): (recipe, year) pairs the recipe's own generic join
|
||||||
# covers for this government, restricted to the gap years.
|
# covers for this government, restricted to the gap years.
|
||||||
covered <- if (length(gap_years) == 0L) {
|
covered <- .query_covered_years(con, candidates, cohort, gap_years)
|
||||||
data.frame(recipe_id = character(0), year = integer(0))
|
|
||||||
} else {
|
|
||||||
DBI::dbGetQuery(con, sprintf(
|
|
||||||
"SELECT DISTINCT r.recipe_id, l.year
|
|
||||||
FROM long l
|
|
||||||
JOIN harmonization_recipes r
|
|
||||||
ON l.item_code = r.component_code
|
|
||||||
AND l.year BETWEEN r.year_min AND r.year_max
|
|
||||||
AND (r.gov_type_scope = 'all'
|
|
||||||
OR (r.gov_type_scope = 'state' AND l.type = 0)
|
|
||||||
OR (r.gov_type_scope = 'local' AND l.type BETWEEN 1 AND 3))
|
|
||||||
WHERE r.recipe_id IN (%s)
|
|
||||||
AND %s
|
|
||||||
AND l.year IN (%s)",
|
|
||||||
.sql_lit_chr(candidates), .cohort_sql(cohort, "l.canonical_govid"),
|
|
||||||
paste(gap_years, collapse = ",")
|
|
||||||
))
|
|
||||||
}
|
|
||||||
|
|
||||||
suggestions <- list()
|
suggestions <- list()
|
||||||
for (rid in candidates) {
|
for (rid in candidates) {
|
||||||
@@ -227,6 +197,122 @@
|
|||||||
.attach_ig_counterparts(con, suggestions, flow_prefixes)
|
.attach_ig_counterparts(con, suggestions, flow_prefixes)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#' Query candidate harmonization recipe IDs for a coverage-gap suggestion.
|
||||||
|
#'
|
||||||
|
#' Selects recipes whose component codes fall within the requested scope
|
||||||
|
#' (category or subtype allowlist), excluding any recipe that is ITSELF an
|
||||||
|
#' intergovernmental (M/L) recipe -- i.e. every one of its own component
|
||||||
|
#' codes is M/L-prefixed. Without this exclusion, a category whose
|
||||||
|
#' summary_categories rows span both a Direct family (e.g. E04/E05,
|
||||||
|
#' "Corrections") and its M/L counterpart (M04/M05) makes the M/L recipe
|
||||||
|
#' itself a raw top-level candidate for a plain `cog_spending()` call --
|
||||||
|
#' following that hint would silently return intergovernmental dollars
|
||||||
|
#' under `expenditure_concept = "direct"` provenance.
|
||||||
|
#'
|
||||||
|
#' In all-categories mode (`all_categories = TRUE`) the inner sub-select is
|
||||||
|
#' scoped by `subtype_col`/`subtype_scope` -- the same allowlist
|
||||||
|
#' `.build_verb_sql()` applies as a WHERE predicate to make the summed
|
||||||
|
#' result a *concept* (see R/spending.R), not by `category`.
|
||||||
|
#' `.ALL_CATEGORIES` ("All Categories") is never itself a row in
|
||||||
|
#' `summary_categories.category`, so a category-keyed sub-select always
|
||||||
|
#' returns zero candidates and silently disables signposting.
|
||||||
|
#'
|
||||||
|
#' @param con Active DuckDB connection.
|
||||||
|
#' @param category Category name, or `NULL`.
|
||||||
|
#' @param all_categories `TRUE` when the caller used `.ALL_CATEGORIES`.
|
||||||
|
#' @param subtype_col Name of the summary_categories subtype column to
|
||||||
|
#' scope by when `all_categories = TRUE`; ignored otherwise.
|
||||||
|
#' @param subtype_scope Character vector of subtype values to scope by
|
||||||
|
#' when `all_categories = TRUE`; ignored otherwise.
|
||||||
|
#' @return Character vector of recipe IDs (possibly empty).
|
||||||
|
#' @noRd
|
||||||
|
.query_candidate_recipes <- function(con, category, all_categories = FALSE,
|
||||||
|
subtype_col = NULL,
|
||||||
|
subtype_scope = NULL) {
|
||||||
|
candidate_scope_sql <- if (isTRUE(all_categories)) {
|
||||||
|
sprintf(
|
||||||
|
"SELECT DISTINCT item_code FROM summary_categories WHERE %s IN (%s)",
|
||||||
|
subtype_col, .sql_lit_chr(subtype_scope)
|
||||||
|
)
|
||||||
|
} else {
|
||||||
|
sprintf(
|
||||||
|
"SELECT DISTINCT item_code FROM summary_categories WHERE category IN (%s)",
|
||||||
|
.sql_lit_chr(category)
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
DBI::dbGetQuery(con, sprintf(
|
||||||
|
"SELECT DISTINCT recipe_id FROM harmonization_recipes
|
||||||
|
WHERE component_code IN (
|
||||||
|
%s
|
||||||
|
)
|
||||||
|
AND recipe_id NOT IN (
|
||||||
|
SELECT DISTINCT recipe_id FROM harmonization_recipes
|
||||||
|
WHERE LEFT(component_code, 1) IN ('M', 'L')
|
||||||
|
)",
|
||||||
|
candidate_scope_sql
|
||||||
|
))$recipe_id
|
||||||
|
}
|
||||||
|
|
||||||
|
#' Query gap-year coverage: which (recipe, year) pairs the recipe's own
|
||||||
|
#' generic join covers for this government, restricted to `gap_years`.
|
||||||
|
#'
|
||||||
|
#' This is Path 1 of a suggestion (unchanged): it finds recipes whose
|
||||||
|
#' component codes' generic join produces at least one row for this
|
||||||
|
#' government in each gap year -- i.e. the category returned nothing in
|
||||||
|
#' that year but a recipe would fill it.
|
||||||
|
#'
|
||||||
|
#' @param con Active DuckDB connection.
|
||||||
|
#' @param candidates Character vector of recipe IDs to check coverage for.
|
||||||
|
#' @param cohort The verb's cohort object (see `.make_cohort()`), rendered
|
||||||
|
#' into the govid predicate on the joined `long` scan via `.cohort_sql()`.
|
||||||
|
#' @param gap_years Integer vector of requested years absent from the
|
||||||
|
#' result.
|
||||||
|
#' @return Data frame with columns `recipe_id` (character) and `year`
|
||||||
|
#' (integer). Returns an empty data frame (`recipe_id = character(0)`,
|
||||||
|
#' `year = integer(0)`) when `gap_years` is empty, so callers can safely
|
||||||
|
#' reference `$recipe_id`.
|
||||||
|
#' @noRd
|
||||||
|
.query_covered_years <- function(con, candidates, cohort, gap_years) {
|
||||||
|
if (length(gap_years) == 0L) {
|
||||||
|
return(data.frame(recipe_id = character(0), year = integer(0)))
|
||||||
|
}
|
||||||
|
DBI::dbGetQuery(con, sprintf(
|
||||||
|
"SELECT DISTINCT r.recipe_id, l.year
|
||||||
|
FROM long l
|
||||||
|
JOIN harmonization_recipes r
|
||||||
|
ON l.item_code = r.component_code
|
||||||
|
AND l.year BETWEEN r.year_min AND r.year_max
|
||||||
|
AND (r.gov_type_scope = 'all'
|
||||||
|
OR (r.gov_type_scope = 'state' AND l.type = 0)
|
||||||
|
OR (r.gov_type_scope = 'local' AND l.type BETWEEN 1 AND 3))
|
||||||
|
WHERE r.recipe_id IN (%s)
|
||||||
|
AND %s
|
||||||
|
AND l.year IN (%s)",
|
||||||
|
.sql_lit_chr(candidates), .cohort_sql(cohort, "l.canonical_govid"),
|
||||||
|
paste(gap_years, collapse = ",")
|
||||||
|
))
|
||||||
|
}
|
||||||
|
|
||||||
|
#' Query recipe metadata: labels and year spans for a set of candidate
|
||||||
|
#' recipes.
|
||||||
|
#'
|
||||||
|
#' @param con Active DuckDB connection.
|
||||||
|
#' @param candidates Character vector of recipe IDs to look up.
|
||||||
|
#' @return Tibble with columns `recipe_id`, `label`, `year_min` (int), and
|
||||||
|
#' `year_max` (int).
|
||||||
|
#' @noRd
|
||||||
|
.query_recipe_meta <- function(con, candidates) {
|
||||||
|
tibble::as_tibble(DBI::dbGetQuery(con, sprintf(
|
||||||
|
"SELECT recipe_id, any_value(label) AS label,
|
||||||
|
MIN(year_min) AS year_min, MAX(year_max) AS year_max
|
||||||
|
FROM harmonization_recipes
|
||||||
|
WHERE recipe_id IN (%s)
|
||||||
|
GROUP BY recipe_id",
|
||||||
|
.sql_lit_chr(candidates)
|
||||||
|
)))
|
||||||
|
}
|
||||||
|
|
||||||
#' Attach `ig_recipe_id` to each suggestion: the intergovernmental-expenditure
|
#' Attach `ig_recipe_id` to each suggestion: the intergovernmental-expenditure
|
||||||
#' recipe (an M-to-local or L-to-state recipe) whose component codes cover
|
#' recipe (an M-to-local or L-to-state recipe) whose component codes cover
|
||||||
#' exactly the same set of function suffixes as the firing recipe's own
|
#' exactly the same set of function suffixes as the firing recipe's own
|
||||||
@@ -272,7 +358,7 @@
|
|||||||
#' `R/basis.R`). This blocks a recipe surfaced through a mis-scoped
|
#' `R/basis.R`). This blocks a recipe surfaced through a mis-scoped
|
||||||
#' category from ever reaching the M/L search, e.g. `cog_spending()`'s
|
#' category from ever reaching the M/L search, e.g. `cog_spending()`'s
|
||||||
#' flow_prefixes are `c("E","F","G")`, which `ig_federal_b47_wide`'s own
|
#' flow_prefixes are `c("E","F","G")`, which `ig_federal_b47_wide`'s own
|
||||||
#' `"B"` is not part of.
|
#' "B" is not part of.
|
||||||
#' 2. `own_prefix %in% c("E","F","G")`: M/L only ever pairs with the
|
#' 2. `own_prefix %in% c("E","F","G")`: M/L only ever pairs with the
|
||||||
#' DIRECT-expenditure family, never with revenue (`cog_revenue()`'s
|
#' DIRECT-expenditure family, never with revenue (`cog_revenue()`'s
|
||||||
#' flow_prefixes already fold B/C/D in as ordinary revenue -- there is
|
#' flow_prefixes already fold B/C/D in as ordinary revenue -- there is
|
||||||
@@ -280,10 +366,6 @@
|
|||||||
#' adds one for spending) and never with ANOTHER M/L recipe (without
|
#' adds one for spending) and never with ANOTHER M/L recipe (without
|
||||||
#' this check, `ige_local_m47_wide` would wrongly match sibling
|
#' this check, `ige_local_m47_wide` would wrongly match sibling
|
||||||
#' `ige_state_l47_wide` on their shared {"47","94"} suffix set).
|
#' `ige_state_l47_wide` on their shared {"47","94"} suffix set).
|
||||||
#' Condition 1 alone does not catch this: under `cog_revenue()`,
|
|
||||||
#' `ig_federal_b47_wide`'s own `"B"` IS inside revenue's own
|
|
||||||
#' `flow_prefixes`, so only this second, family-specific check blocks
|
|
||||||
#' the search.
|
|
||||||
#' @noRd
|
#' @noRd
|
||||||
.attach_ig_counterparts <- function(con, suggestions, flow_prefixes) {
|
.attach_ig_counterparts <- function(con, suggestions, flow_prefixes) {
|
||||||
if (length(suggestions) == 0L) return(suggestions)
|
if (length(suggestions) == 0L) return(suggestions)
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
# uscogdata
|
# uscogdata
|
||||||
|
|
||||||
<!-- badges: start -->
|
<!-- badges: start -->
|
||||||
|
[](https://github.com/civilytics/uscogdata/actions/workflows/R-CMD-check.yaml)
|
||||||
[](https://civilytics.r-universe.dev/uscogdata)
|
[](https://civilytics.r-universe.dev/uscogdata)
|
||||||
[](LICENSE.md)
|
[](LICENSE.md)
|
||||||
<!-- badges: end -->
|
<!-- badges: end -->
|
||||||
@@ -18,7 +19,7 @@ carries provenance describing what was converted, what was aggregated, and
|
|||||||
which known series breaks intersect your query.
|
which known series breaks intersect your query.
|
||||||
|
|
||||||
**Scope:** government types 0–3 (state, county, municipality, township).
|
**Scope:** government types 0–3 (state, county, municipality, township).
|
||||||
56 fiscal years, 46,148,034 rows, 190.6 MB. There is no source data for FY1968
|
56 fiscal years, 46,148,034 rows, ~201 MB. There is no source data for FY1968
|
||||||
or FY1969. Special districts (type 4) and school districts (type 5) are
|
or FY1969. Special districts (type 4) and school districts (type 5) are
|
||||||
excluded pending validation.
|
excluded pending validation.
|
||||||
|
|
||||||
@@ -85,14 +86,38 @@ cog_explain(spend)
|
|||||||
|
|
||||||
| | Remote (default) | Mirrored |
|
| | Remote (default) | Mirrored |
|
||||||
|---|---|---|
|
|---|---|---|
|
||||||
| Setup | none | `cog_mirror(dest)`, 190.6 MB once |
|
| Setup | none | `cog_mirror(dest)`, ~201 MB once |
|
||||||
| Disk used | **0 MB** — HTTP range requests only | 190.6 MB |
|
| Disk used | **0 MB** — HTTP range requests only | ~201 MB |
|
||||||
| Per query | ~4 s (one government, one year)<br>~6 s (one government, 23 years) | local speed |
|
| Opening a session | ~7.5 s | ~0.1 s |
|
||||||
|
| One government, one year | ~4 s | ~0.05 s |
|
||||||
|
| One government, full history | ~7 s | ~0.1 s |
|
||||||
|
| Later queries, same session | ~1.5 s | ~0.05 s |
|
||||||
| Good for | trying it out, teaching, one-off questions | repeated analysis, offline work, reproducibility |
|
| Good for | trying it out, teaching, one-off questions | repeated analysis, offline work, reproducibility |
|
||||||
|
|
||||||
|
**A local mirror is roughly 60–80x faster, and it is one function call.** That is
|
||||||
|
by far the largest difference any of these settings makes. If you are going to
|
||||||
|
ask more than a handful of questions, mirror first.
|
||||||
|
|
||||||
|
Measured 2026-08-10 on a 16-core Linux workstation against the published corpus
|
||||||
|
(schema v7, `pipeline_commit 3d28ddd`), fresh R session per arm. A one-off
|
||||||
|
question costs about **12 seconds end to end remotely and 0.15 seconds
|
||||||
|
mirrored**, session setup included.
|
||||||
|
|
||||||
|
Two things the per-query rows hide:
|
||||||
|
|
||||||
|
- **Opening the session is the single largest remote cost** — larger than any
|
||||||
|
one query. It fetches the manifest and registers 23 SQL views over HTTPS, and
|
||||||
|
it lands on your first query, not on `library(uscogdata)`.
|
||||||
|
- **The cost is network round-trips, not scanning.** A repeat query against
|
||||||
|
partitions this session has already touched is ~1.5 s rather than ~4 s, and a
|
||||||
|
full-history query costs ~7 s whether it runs first or last. What you are
|
||||||
|
paying for is reaching each of the 56 yearly files over HTTPS the first time.
|
||||||
|
|
||||||
Nothing is written to disk in remote mode: DuckDB fetches the parquet footer,
|
Nothing is written to disk in remote mode: DuckDB fetches the parquet footer,
|
||||||
works out which row groups it needs, and reads only those. Nothing is cached
|
works out which row groups it needs, and reads only those. Nothing is cached
|
||||||
between sessions either, so every query goes back to the network.
|
between sessions either, so every query goes back to the network — and a session
|
||||||
|
that issues many remote queries in quick succession can be rate-limited by the
|
||||||
|
host (`HTTP Error: ... 429`). Both are further reasons to mirror for real work.
|
||||||
|
|
||||||
The default points at a public HuggingFace mirror of the corpus. If you would
|
The default points at a public HuggingFace mirror of the corpus. If you would
|
||||||
rather not depend on a third party — for reproducibility, for an air-gapped
|
rather not depend on a third party — for reproducibility, for an air-gapped
|
||||||
@@ -110,6 +135,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
@@ -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
@@ -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`,
|
||||||
|
|||||||
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,14 @@
|
|||||||
|
# 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.
|
||||||
|
|
||||||
|
---
|
||||||
@@ -0,0 +1,67 @@
|
|||||||
|
# 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 (#52), is a decision rather than a task: it was parked during
|
||||||
|
the 0.3.0 design, and the API announcement waits on it, 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.
|
||||||
|
|
||||||
|
Compass's own files moved out of `docs/` this session. They were sitting inside
|
||||||
|
pkgdown's output directory, and `pkgdown::clean_site()` deletes every top-level entry
|
||||||
|
there except `CNAME` and `dev` — asked directly, it listed `docs/pm` and
|
||||||
|
`docs/decisions` among the 28 it would remove, with the guard that would have stopped
|
||||||
|
it satisfied by `docs/pkgdown.yml`. They are in `pm/` now. Nothing was lost: the
|
||||||
|
journal had no entries and there were no decision records yet, which made this the
|
||||||
|
cheapest moment to move. The `.gitignore` workaround that re-included two children of
|
||||||
|
an excluded `docs/` is gone with it.
|
||||||
|
|
||||||
|
## 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
|
||||||
|
|
||||||
|

|
||||||
|

|
||||||
|
|
||||||
|
<details>
|
||||||
|
<summary>Dependency graph and detail</summary>
|
||||||
|
|
||||||
|
_Nothing blocks anything else, so there is no graph to draw._
|
||||||
|
|
||||||
|
- Marker: `none` (no journal entry yet)
|
||||||
|
- Commits since: 165
|
||||||
|
- Open issues: 6
|
||||||
|
|
||||||
|
</details>
|
||||||
|
|
||||||
|
<!-- compass:end -->
|
||||||
@@ -0,0 +1,44 @@
|
|||||||
|
[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.",
|
||||||
|
]
|
||||||
@@ -0,0 +1,12 @@
|
|||||||
|
# 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 -->
|
||||||
@@ -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))))
|
||||||
|
})
|
||||||
Reference in New Issue
Block a user