feat: basis= harmonized/raw with v4/v5 dual-accept

Adds schema_version 5 support alongside the existing v4 corpus:
.validate_schema() now accepts a supported set (4, 5) instead of a single
expected version, and cog_spending()/cog_revenue() gain basis =
c("harmonized", "raw"). Harmonized basis routes to new
spending_annotated_harmonized / revenue_annotated_harmonized views built on
spending_long_harmonized / revenue_long_harmonized (REPLACE(harmonized_code
AS item_code), excluding aggregate and NA-harmonized rows); raw basis is
byte-identical to the pre-Phase-R2 behavior. On a v4 corpus, an unspecified
basis silently resolves to "raw" with a provenance note; an explicit
basis = "harmonized" aborts with an actionable message.

Provenance gains basis, basis_note, and a harmonization block
(applied/na_rows_excluded/na_amount_excluded). The five new schema-v5-only
SQL views (harmonized long/annotated views, harmonization_map,
harmonization_recipes, series_breaks_pq) are registered conditionally on
manifest$schema_version >= 5, since DuckDB's read_parquet() errors eagerly
at CREATE VIEW time when the backing file doesn't exist on a v4 corpus.

Fixture corpus regenerated to schema_version 5 / years 2011, 2012, 2019,
2020 (2011->2012 spans the wide-aggregate -> modern-leaf format boundary
needed for the harmonization/recipe work), with the harmonization_map /
harmonization_recipes / series_breaks parquet tables bundled alongside the
existing metadata registries.
This commit is contained in:
2026-07-18 23:19:17 -04:00
parent 3b725770d2
commit 7818cd2b1a
35 changed files with 579 additions and 53 deletions
+1 -1
View File
@@ -32,4 +32,4 @@ Config/testthat/edition: 3
VignetteBuilder: knitr VignetteBuilder: knitr
RoxygenNote: 7.3.3 RoxygenNote: 7.3.3
MinCorpusSchema: 4 MinCorpusSchema: 4
MaxCorpusSchema: 4 MaxCorpusSchema: 5
+78
View File
@@ -0,0 +1,78 @@
# R/basis.R
# basis= resolution (harmonized/raw, with v4/v5 dual-accept) and the
# harmonization exclusion-count block attached to provenance.
#' Resolve the requested `basis` against the active corpus's schema_version.
#'
#' On a `schema_version >= 5` corpus, the requested basis is used as-is. On
#' an older (`schema_version == 4`) corpus, which has no harmonization
#' tables: a caller who left `basis` at its default (`"harmonized"`, so
#' `explicit` is `FALSE`) silently gets `"raw"` back, with a note recorded
#' for provenance; a caller who explicitly asked for
#' `basis = "harmonized"` gets a hard abort instead of a silent downgrade.
#'
#' @param basis `"harmonized"` or `"raw"` (already resolved via `match.arg`).
#' @param explicit `TRUE` if the caller passed `basis` explicitly (as
#' opposed to relying on the default `c("harmonized", "raw")`).
#' @param manifest The active session's parsed manifest list.
#' @return List with `basis` (the resolved value) and `note` (character or
#' `NA_character_`).
#' @noRd
.resolve_basis <- function(basis, explicit, manifest) {
schema_version <- suppressWarnings(as.integer(manifest$schema_version %||% 0L))
if (schema_version >= 5L) {
return(list(basis = basis, note = NA_character_))
}
if (identical(basis, "harmonized") && explicit) {
cli::cli_abort(c(
"basis = \"harmonized\" requires corpus schema_version >= 5.",
x = "Active corpus has schema_version {schema_version}.",
i = "Use basis = \"raw\" (the default on this corpus), or point USCOGDATA_URL at a schema_version >= 5 corpus."
), class = "uscogdata_basis_unsupported")
}
list(
basis = "raw",
note = sprintf(
"basis resolved to \"raw\": corpus schema_version %d < 5 (harmonization tables unavailable)",
schema_version
)
)
}
#' Count + sum item-level rows that basis="harmonized" excludes because they
#' carry no harmonized_code (discontinued / not-yet-ruled codes) within the
#' requested flow type (spending or revenue), govids, and years. Only
#' meaningful when the resolved basis is "harmonized"; returns an
#' applied = FALSE stub otherwise (raw basis never excludes rows this way).
#' @noRd
.build_harmonization_block <- function(con, govid, years, resolved, flow_prefixes) {
if (!identical(resolved$basis, "harmonized")) {
return(list(
applied = FALSE,
na_rows_excluded = 0L,
na_amount_excluded = 0,
note = resolved$note
))
}
sql <- sprintf(
"SELECT COUNT(*) AS n, COALESCE(SUM(amt), 0) * 1000.0 AS amt
FROM long
WHERE canonical_govid IN (%s) AND year IN (%s)
AND NOT is_aggregate AND harmonized_code IS NULL
AND LEFT(item_code, 1) IN (%s)",
.sql_lit_chr(govid), paste(as.integer(years), collapse = ","),
.sql_lit_chr(flow_prefixes)
)
na <- DBI::dbGetQuery(con, sql)
list(
applied = TRUE,
na_rows_excluded = as.integer(na$n),
na_amount_excluded = as.numeric(na$amt),
note = resolved$note
)
}
+3 -3
View File
@@ -128,11 +128,11 @@
} }
#' @noRd #' @noRd
.validate_schema <- function(manifest, expected_version) { .validate_schema <- function(manifest, supported = c(4L, 5L)) {
if (manifest$schema_version != expected_version) { if (!manifest$schema_version %in% supported) {
cli::cli_abort(c( cli::cli_abort(c(
"Corpus schema version mismatch.", "Corpus schema version mismatch.",
x = "Package expects schema_version = {expected_version}; corpus has {manifest$schema_version}.", x = "Package supports schema_version in {paste(supported, collapse = ', ')}; corpus has {manifest$schema_version}.",
i = "Update uscogdata (install.packages or pak::pkg_install) or re-publish corpus." i = "Update uscogdata (install.packages or pak::pkg_install) or re-publish corpus."
)) ))
} }
+9 -1
View File
@@ -4,7 +4,9 @@
#' @noRd #' @noRd
.build_provenance <- function(verb, call, govid, years, category, .build_provenance <- function(verb, call, govid, years, category,
per_capita, adjust_to_year, result, sql, per_capita, adjust_to_year, result, sql,
subtype_col) { subtype_col, basis = NA_character_,
basis_note = NA_character_,
harmonization = NULL) {
manifest <- .uscogdata_env$manifest manifest <- .uscogdata_env$manifest
codes <- result[["codes_included"]] codes <- result[["codes_included"]]
@@ -38,6 +40,12 @@
), ),
years = as.integer(years), years = as.integer(years),
category = category, category = category,
basis = basis,
basis_note = basis_note,
harmonization = harmonization %||% list(
applied = FALSE, na_rows_excluded = 0L, na_amount_excluded = 0,
note = NA_character_
),
scope = list( scope = list(
gov_types_included = as.integer(unlist(manifest$scope$gov_types_included)), gov_types_included = as.integer(unlist(manifest$scope$gov_types_included)),
gov_types_excluded = as.integer(unlist(manifest$scope$gov_types_excluded)), gov_types_excluded = as.integer(unlist(manifest$scope$gov_types_excluded)),
+6 -3
View File
@@ -14,16 +14,19 @@
#' optional `pop_source`, `codes_included`, `aggregate_fallback`, `notes`. #' optional `pop_source`, `codes_included`, `aggregate_fallback`, `notes`.
#' @export #' @export
cog_revenue <- function(govid, years, category = NULL, cog_revenue <- function(govid, years, category = NULL,
per_capita = FALSE, adjust_to_year = NULL) { per_capita = FALSE, adjust_to_year = NULL,
basis = c("harmonized", "raw")) {
.verb_spendrev( .verb_spendrev(
verb = "cog_revenue", verb = "cog_revenue",
view = "revenue_annotated", view_base = "revenue_annotated",
subtype_col = "revenue_subtype", subtype_col = "revenue_subtype",
flow_prefixes = c("T", "A", "U", "B", "C", "D"),
call = match.call(), call = match.call(),
govid = govid, govid = govid,
years = years, years = years,
category = category, category = category,
per_capita = per_capita, per_capita = per_capita,
adjust_to_year = adjust_to_year adjust_to_year = adjust_to_year,
basis = basis
) )
} }
+1 -1
View File
@@ -12,7 +12,7 @@ cog_open <- function(url = .resolve_url(),
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)
.validate_schema(manifest, expected_version = 4L) .validate_schema(manifest, supported = c(4L, 5L))
.validate_scope(manifest) .validate_scope(manifest)
.register_views(con, url, manifest) .register_views(con, url, manifest)
+38 -6
View File
@@ -20,6 +20,15 @@
#' in that year). #' in that year).
#' @param adjust_to_year Integer base year for CPI-U real-dollar conversion, #' @param adjust_to_year Integer base year for CPI-U real-dollar conversion,
#' or `NULL` for nominal only. #' or `NULL` for nominal only.
#' @param basis `"harmonized"` (default) sums item codes through the
#' cross-vintage harmonization mapping (folding series-break-affected
#' codes onto a comparable target and excluding aggregate / discontinued
#' rows -- see the `harmonization` block in `cog_explain()`); `"raw"`
#' reproduces the pre-Phase-R2 behavior (published item codes, no
#' folding). On a corpus with `schema_version < 5` (no harmonization
#' tables), `basis` silently resolves to `"raw"` when left at its default
#' and the resolution is recorded in the provenance; explicitly passing
#' `basis = "harmonized"` on such a corpus aborts.
#' @return Tibble with columns `year`, `canonical_govid`, `gov_name`, #' @return Tibble with columns `year`, `canonical_govid`, `gov_name`,
#' `spend_subtype`, `category`, `amt_nominal`, optional `amt_real`, #' `spend_subtype`, `category`, `amt_nominal`, optional `amt_real`,
#' optional `amt_per_capita_nominal`, optional `amt_per_capita_real`, #' optional `amt_per_capita_nominal`, optional `amt_per_capita_real`,
@@ -27,24 +36,31 @@
#' Carries a `provenance` attribute matching `inst/schemas/provenance-v1.json`. #' Carries a `provenance` attribute matching `inst/schemas/provenance-v1.json`.
#' @export #' @export
cog_spending <- function(govid, years, category = NULL, cog_spending <- function(govid, years, category = NULL,
per_capita = FALSE, adjust_to_year = NULL) { per_capita = FALSE, adjust_to_year = NULL,
basis = c("harmonized", "raw")) {
.verb_spendrev( .verb_spendrev(
verb = "cog_spending", verb = "cog_spending",
view = "spending_annotated", view_base = "spending_annotated",
subtype_col = "spend_subtype", subtype_col = "spend_subtype",
flow_prefixes = c("E", "F", "G", "K"),
call = match.call(), call = match.call(),
govid = govid, govid = govid,
years = years, years = years,
category = category, category = category,
per_capita = per_capita, per_capita = per_capita,
adjust_to_year = adjust_to_year adjust_to_year = adjust_to_year,
basis = basis
) )
} }
#' @noRd #' @noRd
.verb_spendrev <- function(verb, view, subtype_col, call, .verb_spendrev <- function(verb, view_base, subtype_col, flow_prefixes, call,
govid, years, category, govid, years, category,
per_capita, adjust_to_year) { per_capita, adjust_to_year,
basis = c("harmonized", "raw")) {
basis_explicit <- length(basis) == 1L
basis <- match.arg(basis, c("harmonized", "raw"))
govid <- .coerce_govid_input(govid, arg = "govid") govid <- .coerce_govid_input(govid, arg = "govid")
.validate_verb_inputs(govid, years, category, per_capita, adjust_to_year) .validate_verb_inputs(govid, years, category, per_capita, adjust_to_year)
@@ -52,8 +68,12 @@ cog_spending <- function(govid, years, category = NULL,
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)
con <- .ensure_session() con <- .ensure_session()
manifest <- .uscogdata_env$manifest
scope <- .check_govids_in_scope(govid) scope <- .check_govids_in_scope(govid)
resolved <- .resolve_basis(basis, basis_explicit, manifest)
view <- .select_view(view_base, resolved$basis)
sql <- .build_verb_sql(view, subtype_col, govid, years, category) sql <- .build_verb_sql(view, subtype_col, govid, years, category)
result <- tibble::as_tibble(DBI::dbGetQuery(con, sql)) result <- tibble::as_tibble(DBI::dbGetQuery(con, sql))
@@ -64,6 +84,10 @@ cog_spending <- function(govid, years, category = NULL,
result$notes <- .notes_column(result) result$notes <- .notes_column(result)
harmonization <- .build_harmonization_block(
con, govid, years, resolved, flow_prefixes
)
prov <- .build_provenance( prov <- .build_provenance(
verb = verb, verb = verb,
call = call, call = call,
@@ -74,7 +98,10 @@ cog_spending <- function(govid, years, category = NULL,
adjust_to_year = adjust_to_year, adjust_to_year = adjust_to_year,
result = result, result = result,
sql = sql, sql = sql,
subtype_col = subtype_col subtype_col = subtype_col,
basis = resolved$basis,
basis_note = resolved$note,
harmonization = harmonization
) )
prov$scope$govids_found <- scope$found prov$scope$govids_found <- scope$found
prov$scope$govids_missing <- scope$missing prov$scope$govids_missing <- scope$missing
@@ -107,6 +134,11 @@ cog_spending <- function(govid, years, category = NULL,
invisible(TRUE) invisible(TRUE)
} }
#' @noRd
.select_view <- function(view_base, basis) {
if (identical(basis, "harmonized")) paste0(view_base, "_harmonized") else view_base
}
#' @noRd #' @noRd
.sql_lit_chr <- function(x) { .sql_lit_chr <- function(x) {
safe <- gsub("'", "''", x, fixed = TRUE) safe <- gsub("'", "''", x, fixed = TRUE)
+23 -1
View File
@@ -1,11 +1,33 @@
# R/views.R # R/views.R
# SQL files whose view definitions read schema-v5-only parquet tables
# (harmonization_map.parquet, harmonization_recipes.parquet,
# series_breaks.parquet) or select from views built on top of them. DuckDB's
# read_parquet() resolves the file at CREATE VIEW time (even for a view, it
# still needs the source schema) and errors immediately -- "IO Error: No
# files found" -- if the path doesn't exist, so these cannot be registered
# unconditionally against a v4 corpus the way the rest of inst/sql/ is.
# Registration is therefore gated on manifest$schema_version >= 5; verb-level
# *usage* of the resulting views is separately gated by .resolve_basis() /
# .require_schema_v5().
.harmonization_view_files <- c(
"22-spending_long_harmonized.sql",
"23-revenue_long_harmonized.sql",
"33-harmonization_map.sql",
"34-harmonization_recipes.sql",
"35-series_breaks_pq.sql",
"42-spending_annotated_harmonized.sql",
"43-revenue_annotated_harmonized.sql"
)
#' Register DuckDB views from inst/sql/ SQL files #' Register DuckDB views from inst/sql/ SQL files
#' @noRd #' @noRd
.register_views <- function(con, url, manifest) { .register_views <- function(con, url, manifest) {
sql_dir <- system.file("sql", package = "uscogdata") sql_dir <- system.file("sql", package = "uscogdata")
files <- list.files(sql_dir, pattern = "\\.sql$", full.names = TRUE) files <- sort(list.files(sql_dir, pattern = "\\.sql$", full.names = TRUE))
schema_version <- suppressWarnings(as.integer(manifest$schema_version %||% 0L))
for (f in files) { for (f in files) {
if (basename(f) %in% .harmonization_view_files && schema_version < 5L) next
sql <- paste(readLines(f, warn = FALSE), collapse = "\n") sql <- paste(readLines(f, warn = FALSE), collapse = "\n")
sql <- gsub("\\{url\\}", url, sql, fixed = FALSE) sql <- gsub("\\{url\\}", url, sql, fixed = FALSE)
DBI::dbExecute(con, sql) DBI::dbExecute(con, sql)
+34 -15
View File
@@ -3,12 +3,20 @@
# Regenerate inst/extdata/fixture_corpus/ from a cog_pipeline publish tree. # Regenerate inst/extdata/fixture_corpus/ from a cog_pipeline publish tree.
# #
# What this does: # What this does:
# 1. Copies the year=2019 and year=2020 long partitions as-is (byte-for- # 1. Copies each requested year's long partition as-is (byte-for-byte)
# byte) from <publish_cache>/data/long/ into the fixture. # from <publish_cache>/data/long/ into the fixture. Default years are
# c(2011L, 2012L, 2019L, 2020L): 2011/2012 straddle the wide-aggregate
# -> modern-leaf format boundary (the harmonization/recipe seam), and
# 2019/2020 are the pre-existing per-capita/CPI regression anchors.
# Each partition is a full year (all states/govs) as published, so
# Broward County FL and every other previously-pinned government stay
# covered without any per-gov slicing logic.
# 2. Copies the full canonical_fips_xwalk.parquet, canonical_alias.parquet, # 2. Copies the full canonical_fips_xwalk.parquet, canonical_alias.parquet,
# and summary_categories.parquet metadata tables as-is (these are small # summary_categories.parquet, harmonization_map.parquet,
# cross-vintage registries, not partitioned by year, so the fixture # harmonization_recipes.parquet, and series_breaks.parquet metadata
# ships the complete tables rather than a year-scoped subset). # tables as-is (these are small cross-vintage registries, not
# partitioned by year, so the fixture ships the complete tables rather
# than a year-scoped subset).
# 3. Resyncs the four reference docs (data_dictionary.md, # 3. Resyncs the four reference docs (data_dictionary.md,
# reader-specification.md, README.md, series_breaks.md) from the # reader-specification.md, README.md, series_breaks.md) from the
# publish tree's docs/. # publish tree's docs/.
@@ -35,7 +43,7 @@ regenerate_fixture_corpus <- function(
"..", "cog_pipeline", "_targets", "publish_cache" "..", "cog_pipeline", "_targets", "publish_cache"
), ),
fixture_dir = file.path("inst", "extdata", "fixture_corpus"), fixture_dir = file.path("inst", "extdata", "fixture_corpus"),
fixture_years = c(2019L, 2020L)) { fixture_years = c(2011L, 2012L, 2019L, 2020L)) {
stopifnot( stopifnot(
requireNamespace("digest", quietly = TRUE), requireNamespace("digest", quietly = TRUE),
requireNamespace("jsonlite", quietly = TRUE), requireNamespace("jsonlite", quietly = TRUE),
@@ -92,14 +100,18 @@ regenerate_fixture_corpus <- function(
invisible(NULL) invisible(NULL)
} }
# Copy the full (not year-scoped) canonical_fips_xwalk, canonical_alias, and # Copy the full (not year-scoped) canonical_fips_xwalk, canonical_alias,
# summary_categories parquet tables. # summary_categories, and (schema v5+) harmonization_map/
# harmonization_recipes/series_breaks parquet tables.
#' @noRd #' @noRd
.copy_metadata_parquets <- function(publish_cache_dir, fixture_dir) { .copy_metadata_parquets <- function(publish_cache_dir, fixture_dir) {
files <- c( files <- c(
"canonical_fips_xwalk.parquet", "canonical_fips_xwalk.parquet",
"canonical_alias.parquet", "canonical_alias.parquet",
"summary_categories.parquet" "summary_categories.parquet",
"harmonization_map.parquet",
"harmonization_recipes.parquet",
"series_breaks.parquet"
) )
for (f in files) { for (f in files) {
src <- file.path(publish_cache_dir, "data", f) src <- file.path(publish_cache_dir, "data", f)
@@ -170,7 +182,10 @@ regenerate_fixture_corpus <- function(
metadata_files <- c( metadata_files <- c(
"canonical_alias.parquet", "canonical_alias.parquet",
"canonical_fips_xwalk.parquet", "canonical_fips_xwalk.parquet",
"summary_categories.parquet" "summary_categories.parquet",
"harmonization_map.parquet",
"harmonization_recipes.parquet",
"series_breaks.parquet"
) )
metadata <- lapply(metadata_files, function(f) { metadata <- lapply(metadata_files, function(f) {
rel <- file.path("data", f) rel <- file.path("data", f)
@@ -187,11 +202,15 @@ regenerate_fixture_corpus <- function(
built_at = format(Sys.time(), "%Y-%m-%dT%H:%M:%SZ", tz = "UTC"), built_at = format(Sys.time(), "%Y-%m-%dT%H:%M:%SZ", tz = "UTC"),
pipeline_commit = source_manifest$pipeline_commit, pipeline_commit = source_manifest$pipeline_commit,
fixture_note = paste( fixture_note = paste(
"Two-year (2019-2020) fixture for uscogdata tests. Full corpus", "Four-year (2011, 2012, 2019, 2020) fixture for uscogdata tests. Full",
"available via USCOGDATA_URL. Regenerated for Phase P", "corpus available via USCOGDATA_URL. Regenerated for Phase R2",
"(schema_version 4, uniformly 12-char canonical_govid) with the full", "(schema_version 5, harmonization_map/harmonization_recipes/",
"canonical_fips_xwalk master and the new canonical_alias lookup", "series_breaks parquet tables added). 2011/2012 straddle the",
"table via data-raw/regenerate_fixture_corpus.R." "wide-aggregate -> modern-leaf format boundary exercised by basis=",
"\"harmonized\" and recipe= queries; 2019/2020 retain the prior",
"per-capita/CPI regression anchors. Full canonical_fips_xwalk master",
"and canonical_alias lookup table included via",
"data-raw/regenerate_fixture_corpus.R."
), ),
data_vintage = source_manifest$data_vintage, data_vintage = source_manifest$data_vintage,
scope = source_manifest$scope, scope = source_manifest$scope,
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
+42 -13
View File
@@ -1,8 +1,8 @@
{ {
"schema_version": 4, "schema_version": 5,
"built_at": "2026-07-13T23:35:09Z", "built_at": "2026-07-19T03:08:41Z",
"pipeline_commit": "a082b26", "pipeline_commit": "ece9b32",
"fixture_note": "Two-year (2019-2020) fixture for uscogdata tests. Full corpus available via USCOGDATA_URL. Regenerated for Phase P (schema_version 4, uniformly 12-char canonical_govid) with the full canonical_fips_xwalk master and the new canonical_alias lookup table via data-raw/regenerate_fixture_corpus.R.", "fixture_note": "Four-year (2011, 2012, 2019, 2020) fixture for uscogdata tests. Full corpus available via USCOGDATA_URL. Regenerated for Phase R2 (schema_version 5, harmonization_map/harmonization_recipes/ series_breaks parquet tables added). 2011/2012 straddle the wide-aggregate -> modern-leaf format boundary exercised by basis= \"harmonized\" and recipe= queries; 2019/2020 retain the prior per-capita/CPI regression anchors. Full canonical_fips_xwalk master and canonical_alias lookup table included via data-raw/regenerate_fixture_corpus.R.",
"data_vintage": { "data_vintage": {
"census_source_downloaded": "unknown", "census_source_downloaded": "unknown",
"cpi_vintage": "FRED CPIAUCSL", "cpi_vintage": "FRED CPIAUCSL",
@@ -14,42 +14,71 @@
"scope_note": "v0.1 covers state, county, city/municipality, and township governments. Special districts (type 4) and school districts (type 5) are excluded pending validation in a future cycle." "scope_note": "v0.1 covers state, county, city/municipality, and township governments. Special districts (type 4) and school districts (type 5) are excluded pending validation in a future cycle."
}, },
"schema": { "schema": {
"long_column_count": 24, "long_column_count": 26,
"long_columns": ["fips_state", "type", "fips_county", "govid", "gov_blank", "gov_name", "county_name", "fips_state_code", "fips_county_code", "fips_place_code", "population", "popyear", "enrollment", "enrollyear", "function_code", "sch_level_code", "fiscal_year_end", "srvy_year", "item_code", "amt", "srv_data", "impute_flag", "is_aggregate", "canonical_govid"], "long_columns": ["fips_state", "type", "fips_county", "govid", "gov_blank", "gov_name", "county_name", "fips_state_code", "fips_county_code", "fips_place_code", "population", "popyear", "enrollment", "enrollyear", "function_code", "sch_level_code", "fiscal_year_end", "srvy_year", "item_code", "amt", "srv_data", "impute_flag", "is_aggregate", "canonical_govid", "harmonized_code", "survey_weight"],
"data_dictionary": "docs/data_dictionary.md" "data_dictionary": "docs/data_dictionary.md"
}, },
"files": { "files": {
"long_partitions": [ "long_partitions": [
{
"year": 2011,
"path": "data/long/year=2011/part-0.parquet",
"sha256": "76c2153ef0a94c3551a24751aa08225fc937fafb579f10fd0750848d781dc471",
"row_count": 3037606,
"size_bytes": 4985667
},
{
"year": 2012,
"path": "data/long/year=2012/part-0.parquet",
"sha256": "a9bb10d04887490376a4b030a03a4ca1e0f1f9fbca5f7d99ea7ca7434c42b2bc",
"row_count": 1163338,
"size_bytes": 5900191
},
{ {
"year": 2019, "year": 2019,
"path": "data/long/year=2019/part-0.parquet", "path": "data/long/year=2019/part-0.parquet",
"sha256": "c0a2bf0758af129d5dfddb6ff6665cc435ddee87fd6879e788fb56ed53ab22b8", "sha256": "475481fcfbb030f96cae457cbda824512a2da6f97f4c6d12a75130979e4331f9",
"row_count": 318139, "row_count": 318139,
"size_bytes": 1441404 "size_bytes": 1717345
}, },
{ {
"year": 2020, "year": 2020,
"path": "data/long/year=2020/part-0.parquet", "path": "data/long/year=2020/part-0.parquet",
"sha256": "92570b9d55ec3425d034db37838f91c3b8359d0454d3d98730a6016b62e4bb48", "sha256": "7cdf3eb1a73c34e6a3d22befdf7ccdd0b44b1e6b1cf4491159a38a5a24a1b005",
"row_count": 317500, "row_count": 317500,
"size_bytes": 1444011 "size_bytes": 1720794
} }
], ],
"metadata": [ "metadata": [
{ {
"path": "data/canonical_alias.parquet", "path": "data/canonical_alias.parquet",
"sha256": "feb8d01a640fb16c9a4b4ad66726b50b8fe8a1ce2a771bed8c1190fec51d5c8d", "sha256": "fa27ce4b59e286ace05fbcb9f59a910b8e483fd6c58ddf01875a210b966809cb",
"description": "canonical_alias.parquet" "description": "canonical_alias.parquet"
}, },
{ {
"path": "data/canonical_fips_xwalk.parquet", "path": "data/canonical_fips_xwalk.parquet",
"sha256": "1ae47981531c7389f69eff3f7656045428564bfbe8032200eb9c039d32739a7e", "sha256": "b6e2c4cb748141f6f2e324d6704c13830d2a68ad57bbb73d04ac54b881d1285b",
"description": "canonical_fips_xwalk.parquet" "description": "canonical_fips_xwalk.parquet"
}, },
{ {
"path": "data/summary_categories.parquet", "path": "data/summary_categories.parquet",
"sha256": "dd59e7f58a022679ad43511c8c8e938b8dd4be81196bbeeee21e67bdcca2295b", "sha256": "8e6fcd4dd9bb4723841a67233b19388c9762dfc23b4479501183cebf7ea3c1b5",
"description": "summary_categories.parquet" "description": "summary_categories.parquet"
},
{
"path": "data/harmonization_map.parquet",
"sha256": "52b4f3f94e65231aa869445946df7a0cfebf8fef1940f800e08287b53d09ab42",
"description": "harmonization_map.parquet"
},
{
"path": "data/harmonization_recipes.parquet",
"sha256": "1133e9a0b02f8f34f5f936e55c5ecd596bb8a55d8425dcce76767f0f3203581c",
"description": "harmonization_recipes.parquet"
},
{
"path": "data/series_breaks.parquet",
"sha256": "049bd7365dae14e357f6e57765dc5a069440bab46facce237d7ba2ec94b0b113",
"description": "series_breaks.parquet"
} }
] ]
}, },
+5
View File
@@ -10,6 +10,11 @@
"target": { "type": "object" }, "target": { "type": "object" },
"years": { "type": "array", "items": { "type": "integer" } }, "years": { "type": "array", "items": { "type": "integer" } },
"category": { "type": ["string", "array", "null"] }, "category": { "type": ["string", "array", "null"] },
"basis": { "type": ["string", "null"] },
"basis_note": { "type": ["string", "null"] },
"harmonization": { "type": "object" },
"recipe": { "type": ["object", "null"] },
"suggestions": { "type": "array" },
"scope": { "type": "object" }, "scope": { "type": "object" },
"codes_summed": { "type": "object" }, "codes_summed": { "type": "object" },
"aggregate_fallback": { "type": ["object", "null"] }, "aggregate_fallback": { "type": ["object", "null"] },
+6
View File
@@ -0,0 +1,6 @@
CREATE OR REPLACE VIEW spending_long_harmonized AS
SELECT * REPLACE (harmonized_code AS item_code)
FROM long
WHERE NOT is_aggregate
AND harmonized_code IS NOT NULL
AND LEFT(harmonized_code, 1) IN ('E', 'F', 'G', 'K');
+6
View File
@@ -0,0 +1,6 @@
CREATE OR REPLACE VIEW revenue_long_harmonized AS
SELECT * REPLACE (harmonized_code AS item_code)
FROM long
WHERE NOT is_aggregate
AND harmonized_code IS NOT NULL
AND LEFT(harmonized_code, 1) IN ('T', 'A', 'U', 'B', 'C', 'D');
+3
View File
@@ -0,0 +1,3 @@
CREATE OR REPLACE VIEW harmonization_map AS
SELECT *
FROM read_parquet('{url}data/harmonization_map.parquet');
+3
View File
@@ -0,0 +1,3 @@
CREATE OR REPLACE VIEW harmonization_recipes AS
SELECT *
FROM read_parquet('{url}data/harmonization_recipes.parquet');
+3
View File
@@ -0,0 +1,3 @@
CREATE OR REPLACE VIEW series_breaks_pq AS
SELECT *
FROM read_parquet('{url}data/series_breaks.parquet');
@@ -0,0 +1,16 @@
CREATE OR REPLACE VIEW spending_annotated_harmonized AS
SELECT
s.*,
x.gov_name AS xwalk_gov_name,
x.govs_type,
x.type_label,
x.fips_state AS xwalk_fips_state,
x.fips_county AS xwalk_fips_county,
x.fips_place,
x.population_acs,
c.category,
c.category_type,
c.spend_subtype
FROM spending_long_harmonized s
LEFT JOIN canonical_fips_xwalk x USING (canonical_govid)
LEFT JOIN summary_categories c USING (item_code);
@@ -0,0 +1,16 @@
CREATE OR REPLACE VIEW revenue_annotated_harmonized AS
SELECT
s.*,
x.gov_name AS xwalk_gov_name,
x.govs_type,
x.type_label,
x.fips_state AS xwalk_fips_state,
x.fips_county AS xwalk_fips_county,
x.fips_place,
x.population_acs,
c.category,
c.category_type,
c.revenue_subtype
FROM revenue_long_harmonized s
LEFT JOIN canonical_fips_xwalk x USING (canonical_govid)
LEFT JOIN summary_categories c USING (item_code);
+12 -1
View File
@@ -9,7 +9,8 @@ cog_revenue(
years, years,
category = NULL, category = NULL,
per_capita = FALSE, per_capita = FALSE,
adjust_to_year = NULL adjust_to_year = NULL,
basis = c("harmonized", "raw")
) )
} }
\arguments{ \arguments{
@@ -29,6 +30,16 @@ in that year).}
\item{adjust_to_year}{Integer base year for CPI-U real-dollar conversion, \item{adjust_to_year}{Integer base year for CPI-U real-dollar conversion,
or `NULL` for nominal only.} or `NULL` for nominal only.}
\item{basis}{`"harmonized"` (default) sums item codes through the
cross-vintage harmonization mapping (folding series-break-affected
codes onto a comparable target and excluding aggregate / discontinued
rows -- see the `harmonization` block in `cog_explain()`); `"raw"`
reproduces the pre-Phase-R2 behavior (published item codes, no
folding). On a corpus with `schema_version < 5` (no harmonization
tables), `basis` silently resolves to `"raw"` when left at its default
and the resolution is recorded in the provenance; explicitly passing
`basis = "harmonized"` on such a corpus aborts.}
} }
\value{ \value{
Tibble with columns `year`, `canonical_govid`, `gov_name`, Tibble with columns `year`, `canonical_govid`, `gov_name`,
+12 -1
View File
@@ -9,7 +9,8 @@ cog_spending(
years, years,
category = NULL, category = NULL,
per_capita = FALSE, per_capita = FALSE,
adjust_to_year = NULL adjust_to_year = NULL,
basis = c("harmonized", "raw")
) )
} }
\arguments{ \arguments{
@@ -29,6 +30,16 @@ in that year).}
\item{adjust_to_year}{Integer base year for CPI-U real-dollar conversion, \item{adjust_to_year}{Integer base year for CPI-U real-dollar conversion,
or `NULL` for nominal only.} or `NULL` for nominal only.}
\item{basis}{`"harmonized"` (default) sums item codes through the
cross-vintage harmonization mapping (folding series-break-affected
codes onto a comparable target and excluding aggregate / discontinued
rows -- see the `harmonization` block in `cog_explain()`); `"raw"`
reproduces the pre-Phase-R2 behavior (published item codes, no
folding). On a corpus with `schema_version < 5` (no harmonization
tables), `basis` silently resolves to `"raw"` when left at its default
and the resolution is recorded in the provenance; explicitly passing
`basis = "harmonized"` on such a corpus aborts.}
} }
\value{ \value{
Tibble with columns `year`, `canonical_govid`, `gov_name`, Tibble with columns `year`, `canonical_govid`, `gov_name`,
+31
View File
@@ -26,3 +26,34 @@ with_fixture_corpus <- function(code) {
}, add = TRUE) }, add = TRUE)
force(code) force(code)
} }
# Copy the bundled fixture to a temp dir with manifest.json's schema_version
# patched to `version`, then run `code` against it with a clean session
# (mirrors with_fixture_corpus()). Used to exercise the v4/v5 dual-accept
# path without a second physical fixture tree: a real v4 corpus has no
# harmonization_map/harmonization_recipes/series_breaks parquet files, but
# .register_views() only *reads* those when schema_version >= 5 (see
# R/views.R), so a doctored copy of the (v5) bundled fixture with the
# manifest's schema_version knocked down to 4 is a faithful stand-in.
with_doctored_schema_version <- function(version, code) {
src <- fixture_corpus_path()
tmp <- withr::local_tempdir(.local_envir = parent.frame())
file.copy(list.files(src, full.names = TRUE), tmp, recursive = TRUE)
manifest_path <- file.path(tmp, "manifest.json")
m <- jsonlite::fromJSON(manifest_path, simplifyVector = FALSE)
m$schema_version <- as.integer(version)
writeLines(
jsonlite::toJSON(m, auto_unbox = TRUE, pretty = TRUE, null = "null"),
manifest_path
)
old_url <- Sys.getenv("USCOGDATA_URL", unset = NA)
uscogdata:::cog_close()
Sys.setenv(USCOGDATA_URL = paste0(tmp, "/"))
on.exit({
uscogdata:::cog_close()
if (is.na(old_url)) Sys.unsetenv("USCOGDATA_URL") else Sys.setenv(USCOGDATA_URL = old_url)
}, add = TRUE)
force(code)
}
+39 -1
View File
@@ -107,6 +107,44 @@ test_that("cog_manifest returns the active session's parsed manifest", {
expect_true(m$schema_version >= 4L) expect_true(m$schema_version >= 4L)
yrs <- vapply(m$files$long_partitions, function(p) as.integer(p$year), yrs <- vapply(m$files$long_partitions, function(p) as.integer(p$year),
integer(1)) integer(1))
expect_setequal(yrs, c(2019L, 2020L)) expect_setequal(yrs, c(2011L, 2012L, 2019L, 2020L))
})
})
test_that(".validate_schema accepts schema_version 4 and 5, rejects others", {
expect_silent(uscogdata:::.validate_schema(list(schema_version = 4L)))
expect_silent(uscogdata:::.validate_schema(list(schema_version = 5L)))
expect_error(
uscogdata:::.validate_schema(list(schema_version = 3L)),
"schema_version"
)
expect_error(
uscogdata:::.validate_schema(list(schema_version = 6L)),
"schema_version"
)
})
test_that("cog_open succeeds against a doctored schema_version 4 corpus (dual-accept)", {
skip_if_no_corpus()
with_doctored_schema_version(4L, {
con <- cog_open()
expect_true(DBI::dbIsValid(con))
expect_equal(as.integer(cog_manifest()$schema_version), 4L)
# Core (pre-Phase-R2) views must still register on a v4 corpus.
views <- DBI::dbGetQuery(con,
"SELECT table_name FROM information_schema.tables
WHERE table_schema = 'main' AND table_type = 'VIEW'"
)$table_name
expect_true(all(c("spending_annotated", "revenue_annotated") %in% views))
# Schema-v5-only harmonization views must NOT register on a v4 corpus:
# their parquet sources don't exist there and DuckDB's read_parquet()
# errors eagerly at CREATE VIEW time for a missing file/glob, so
# .register_views() gates these on manifest$schema_version >= 5.
expect_false(any(c(
"spending_long_harmonized", "spending_annotated_harmonized",
"harmonization_recipes", "harmonization_map", "series_breaks_pq"
) %in% views))
}) })
}) })
+17
View File
@@ -36,3 +36,20 @@ test_that("cog_revenue result has provenance attribute", {
test_that("cog_revenue rejects invalid inputs", { test_that("cog_revenue rejects invalid inputs", {
expect_error(cog_revenue(list(), 2020L), "character|data frame") expect_error(cog_revenue(list(), 2020L), "character|data frame")
}) })
test_that("cog_revenue basis = 'harmonized' (default) matches 'raw' in this fixture window", {
skip_if_no_corpus()
r_raw <- cog_revenue("121011212191", 2019:2020, basis = "raw")
r_harm <- cog_revenue("121011212191", 2019:2020, basis = "harmonized")
expect_equal(attr(r_raw, "provenance")$basis, "raw")
expect_equal(attr(r_harm, "provenance")$basis, "harmonized")
expect_equal(sum(r_raw$amt_nominal), sum(r_harm$amt_nominal))
})
test_that("cog_revenue provenance carries the harmonization block", {
skip_if_no_corpus()
r <- cog_revenue("121011212191", 2020L)
h <- attr(r, "provenance")$harmonization
expect_true(h$applied)
expect_true(h$na_rows_excluded >= 0L)
})
+111
View File
@@ -189,3 +189,114 @@ test_that("provenance records per-year denominator metadata", {
expect_equal(length(pc$popyear_range), 2L) expect_equal(length(pc$popyear_range), 2L)
}) })
}) })
# --- basis = "harmonized" / "raw" (Phase R2, schema v5) --------------------
test_that("basis = 'raw' reproduces the pre-harmonization Broward Police totals", {
skip_if_no_corpus()
with_fixture_corpus({
r <- cog_spending("121011212191", years = 2019:2020, category = "Police",
basis = "raw")
# Regression pin captured against the schema v5 fixture (2026-07-18,
# pipeline_commit ece9b32) before basis = "harmonized" existed as a
# concept; these are the same totals the pre-Phase-R2 default query
# returned (spending_annotated is untouched by the harmonized views).
ops <- r$amt_nominal[r$year == 2019L & r$spend_subtype == "operations"]
cap <- r$amt_nominal[r$year == 2020L & r$spend_subtype == "capital"]
expect_equal(ops, 483560000)
expect_equal(cap, 26693000)
expect_equal(attr(r, "provenance")$basis, "raw")
})
})
test_that("basis = 'harmonized' (default) matches 'raw' when no harmonization rule applies", {
skip_if_no_corpus()
with_fixture_corpus({
# Every `method = "collapse"` mapping in the curated harmonization_map
# ends by FY2004 for codes inside the spending/revenue flow-type
# prefixes (E/F/G/K, T/A/U/B/C/D); the one collapse extending to FY2011
# (L38/M38 -> L36/M36) is intergovernmental-transfer (L/M prefix) codes
# that were never part of spending_long/revenue_long to begin with. So
# for the fixture's 2011-2020 window, basis = "harmonized" is a
# data-verified no-op vs "raw" for in-scope codes -- this is the
# positive-control counterpart to the synthetic REPLACE-mechanism test
# in test-views.R, which proves the fold itself works when data exists.
r_raw <- cog_spending("121011212191", c(2011L, 2012L, 2019L, 2020L),
"Police", basis = "raw")
r_harm <- cog_spending("121011212191", c(2011L, 2012L, 2019L, 2020L),
"Police", basis = "harmonized")
expect_equal(attr(r_harm, "provenance")$basis, "harmonized")
expect_equal(
r_harm$amt_nominal[order(r_harm$year, r_harm$spend_subtype)],
r_raw$amt_nominal[order(r_raw$year, r_raw$spend_subtype)]
)
})
})
test_that("basis defaults to 'harmonized' when not passed", {
skip_if_no_corpus()
with_fixture_corpus({
r <- cog_spending("121011212191", 2020L, "Police")
expect_equal(attr(r, "provenance")$basis, "harmonized")
})
})
test_that("provenance carries basis + harmonization block with na_rows_excluded", {
skip_if_no_corpus()
with_fixture_corpus({
r <- cog_spending("121011212191", 2011:2012, "Corrections")
prov <- attr(r, "provenance")
expect_equal(prov$basis, "harmonized")
expect_true(prov$harmonization$applied)
expect_true(prov$harmonization$na_rows_excluded >= 0L)
expect_true(prov$harmonization$na_amount_excluded >= 0)
# Data-verified for this fixture: none of the discontinued_na rulings
# (S74, Z61, X04, X06, the debt-detail family, L24) fall inside the
# E/F/G/K spending prefixes, so the exclusion count is exactly zero for
# every year in the bundled window -- see
# docs/phase_r_harmonization_review.md § 1.3/1.4.
expect_equal(prov$harmonization$na_rows_excluded, 0L)
expect_equal(prov$harmonization$na_amount_excluded, 0)
})
})
test_that("basis = 'raw' never populates the harmonization exclusion block", {
skip_if_no_corpus()
with_fixture_corpus({
r <- cog_spending("121011212191", 2020L, "Corrections", basis = "raw")
h <- attr(r, "provenance")$harmonization
expect_false(h$applied)
expect_equal(h$na_rows_excluded, 0L)
})
})
test_that("v4 corpus: basis silently resolves to raw (default) with a provenance note", {
skip_if_no_corpus()
with_doctored_schema_version(4L, {
r <- cog_spending("121011212191", 2019L, "Police")
prov <- attr(r, "provenance")
expect_equal(prov$basis, "raw")
expect_match(prov$basis_note, "raw", fixed = TRUE)
expect_match(prov$basis_note, "schema_version", fixed = TRUE)
expect_false(prov$harmonization$applied)
})
})
test_that("v4 corpus: explicit basis = 'harmonized' aborts", {
skip_if_no_corpus()
with_doctored_schema_version(4L, {
expect_error(
cog_spending("121011212191", 2019L, "Police", basis = "harmonized"),
class = "uscogdata_basis_unsupported"
)
})
})
test_that("v4 corpus: explicit basis = 'raw' still works", {
skip_if_no_corpus()
with_doctored_schema_version(4L, {
r <- cog_spending("121011212191", 2019L, "Police", basis = "raw")
expect_equal(attr(r, "provenance")$basis, "raw")
expect_gt(nrow(r), 0L)
})
})
+64 -6
View File
@@ -14,6 +14,62 @@ test_that("all expected views register on session open", {
expect_true(all(expected %in% views$table_name)) expect_true(all(expected %in% views$table_name))
}) })
test_that("the REPLACE(harmonized_code AS item_code) pattern folds a collapsed code", {
# spending_long_harmonized / revenue_long_harmonized (inst/sql/22-, 23-)
# are defined as:
# SELECT * REPLACE (harmonized_code AS item_code) FROM long WHERE ...
# None of the curated harmonization_map's `collapse` rulings land inside
# the bundled fixture's 2011-2020 window for spending/revenue-prefixed
# codes (see the "basis = 'harmonized' (default) matches 'raw'" test in
# test-spending.R and docs/phase_r_harmonization_review.md § 0.2/§ 2), so
# there is no real fixture row that exercises a nonzero fold. This test
# proves the REPLACE mechanism itself is correct against a synthetic
# long-shaped table with a deliberate E38 -> E36 collapse, independent of
# whether the bundled data happens to contain one right now.
skip_if_no_corpus()
con <- cog_open()
on.exit(cog_close())
DBI::dbExecute(con, "
CREATE OR REPLACE TEMP TABLE synthetic_long AS
SELECT * FROM (VALUES
('121011212191', 2004, 'E36', 70942668, false, 'E36'),
('121011212191', 2004, 'E38', 837485, false, 'E36'),
('121011212191', 2004, 'E62', 100000, false, 'E62')
) AS t(canonical_govid, year, item_code, amt, is_aggregate, harmonized_code)
")
folded <- DBI::dbGetQuery(con, "
SELECT item_code, SUM(amt) AS amt
FROM (SELECT * REPLACE (harmonized_code AS item_code) FROM synthetic_long)
GROUP BY item_code ORDER BY item_code
")
expect_setequal(folded$item_code, c("E36", "E62"))
expect_equal(folded$amt[folded$item_code == "E36"], 70942668 + 837485)
expect_equal(folded$amt[folded$item_code == "E62"], 100000)
DBI::dbExecute(con, "DROP TABLE synthetic_long")
})
test_that("schema v5 harmonization views register when the corpus supports them", {
skip_if_no_corpus()
con <- cog_open()
on.exit(cog_close())
manifest <- uscogdata:::.uscogdata_env$manifest
skip_if(as.integer(manifest$schema_version) < 5L, "fixture is schema_version < 5")
views <- DBI::dbGetQuery(con,
"SELECT table_name FROM information_schema.tables
WHERE table_schema = 'main' AND table_type = 'VIEW'"
)
expected_v5 <- c(
"spending_long_harmonized", "revenue_long_harmonized",
"spending_annotated_harmonized", "revenue_annotated_harmonized",
"harmonization_map", "harmonization_recipes", "series_breaks_pq"
)
expect_true(all(expected_v5 %in% views$table_name))
})
test_that("spending_long filters to E/F/G/K prefixes and excludes aggregates", { test_that("spending_long filters to E/F/G/K prefixes and excludes aggregates", {
skip_if_no_corpus() skip_if_no_corpus()
con <- cog_open() con <- cog_open()
@@ -69,13 +125,15 @@ test_that("gov_population_yearly exposes one row per (year, canonical_govid)", {
WHERE canonical_govid = '121011212191' WHERE canonical_govid = '121011212191'
ORDER BY year" ORDER BY year"
) )
expect_setequal(df$year, c(2019L, 2020L)) expect_setequal(df$year, c(2011L, 2012L, 2019L, 2020L))
expect_equal(nrow(df), 2L) expect_equal(nrow(df), 4L)
expect_true(all(!is.na(df$population))) expect_true(all(!is.na(df$population)))
# Hardcoded values are from the bundled fixture (regenerated 2026-07-11 # Hardcoded values are from the bundled fixture (regenerated 2026-07-18
# against cog_pipeline publish tree, pipeline_commit 1a00925, Phase P # against cog_pipeline publish tree, pipeline_commit ece9b32, Phase R2
# schema_version 4). Update if the fixture is rebuilt against a # schema_version 5, years 2011/2012/2019/2020). Update if the fixture is
# different source vintage. # rebuilt against a different source vintage.
expect_equal(df$population[df$year == 2011L], 1759591L)
expect_equal(df$population[df$year == 2012L], 1819773L)
expect_equal(df$population[df$year == 2019L], 1935878L) expect_equal(df$population[df$year == 2019L], 1935878L)
expect_equal(df$population[df$year == 2020L], 1952778L) expect_equal(df$population[df$year == 2020L], 1952778L)
# Uniqueness on (year, canonical_govid) across the whole view. # Uniqueness on (year, canonical_govid) across the whole view.