diff --git a/DESCRIPTION b/DESCRIPTION index 498ff41..7e248c8 100644 --- a/DESCRIPTION +++ b/DESCRIPTION @@ -32,4 +32,4 @@ Config/testthat/edition: 3 VignetteBuilder: knitr RoxygenNote: 7.3.3 MinCorpusSchema: 4 -MaxCorpusSchema: 4 +MaxCorpusSchema: 5 diff --git a/R/basis.R b/R/basis.R new file mode 100644 index 0000000..8fa1284 --- /dev/null +++ b/R/basis.R @@ -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 + ) +} diff --git a/R/manifest.R b/R/manifest.R index c248ec6..454eb63 100644 --- a/R/manifest.R +++ b/R/manifest.R @@ -128,11 +128,11 @@ } #' @noRd -.validate_schema <- function(manifest, expected_version) { - if (manifest$schema_version != expected_version) { +.validate_schema <- function(manifest, supported = c(4L, 5L)) { + if (!manifest$schema_version %in% supported) { cli::cli_abort(c( "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." )) } diff --git a/R/provenance.R b/R/provenance.R index 72dd42e..fbd3cb0 100644 --- a/R/provenance.R +++ b/R/provenance.R @@ -4,7 +4,9 @@ #' @noRd .build_provenance <- function(verb, call, govid, years, category, per_capita, adjust_to_year, result, sql, - subtype_col) { + subtype_col, basis = NA_character_, + basis_note = NA_character_, + harmonization = NULL) { manifest <- .uscogdata_env$manifest codes <- result[["codes_included"]] @@ -38,6 +40,12 @@ ), years = as.integer(years), 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( gov_types_included = as.integer(unlist(manifest$scope$gov_types_included)), gov_types_excluded = as.integer(unlist(manifest$scope$gov_types_excluded)), diff --git a/R/revenue.R b/R/revenue.R index c577556..5e94d9a 100644 --- a/R/revenue.R +++ b/R/revenue.R @@ -14,16 +14,19 @@ #' optional `pop_source`, `codes_included`, `aggregate_fallback`, `notes`. #' @export 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 = "cog_revenue", - view = "revenue_annotated", + view_base = "revenue_annotated", subtype_col = "revenue_subtype", + flow_prefixes = c("T", "A", "U", "B", "C", "D"), call = match.call(), govid = govid, years = years, category = category, per_capita = per_capita, - adjust_to_year = adjust_to_year + adjust_to_year = adjust_to_year, + basis = basis ) } diff --git a/R/session.R b/R/session.R index 1e9f313..66b6f4e 100644 --- a/R/session.R +++ b/R/session.R @@ -12,7 +12,7 @@ cog_open <- function(url = .resolve_url(), DBI::dbExecute(con, "INSTALL httpfs; LOAD httpfs;") manifest <- .fetch_or_cache_manifest(url, cache_dir) - .validate_schema(manifest, expected_version = 4L) + .validate_schema(manifest, supported = c(4L, 5L)) .validate_scope(manifest) .register_views(con, url, manifest) diff --git a/R/spending.R b/R/spending.R index 78ccd5b..8526103 100644 --- a/R/spending.R +++ b/R/spending.R @@ -20,6 +20,15 @@ #' in that year). #' @param adjust_to_year Integer base year for CPI-U real-dollar conversion, #' 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`, #' `spend_subtype`, `category`, `amt_nominal`, optional `amt_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`. #' @export 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 = "cog_spending", - view = "spending_annotated", + view_base = "spending_annotated", subtype_col = "spend_subtype", + flow_prefixes = c("E", "F", "G", "K"), call = match.call(), govid = govid, years = years, category = category, per_capita = per_capita, - adjust_to_year = adjust_to_year + adjust_to_year = adjust_to_year, + basis = basis ) } #' @noRd -.verb_spendrev <- function(verb, view, subtype_col, call, +.verb_spendrev <- function(verb, view_base, subtype_col, flow_prefixes, call, 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") .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) con <- .ensure_session() + manifest <- .uscogdata_env$manifest 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) result <- tibble::as_tibble(DBI::dbGetQuery(con, sql)) @@ -64,6 +84,10 @@ cog_spending <- function(govid, years, category = NULL, result$notes <- .notes_column(result) + harmonization <- .build_harmonization_block( + con, govid, years, resolved, flow_prefixes + ) + prov <- .build_provenance( verb = verb, call = call, @@ -74,7 +98,10 @@ cog_spending <- function(govid, years, category = NULL, adjust_to_year = adjust_to_year, result = result, 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_missing <- scope$missing @@ -107,6 +134,11 @@ cog_spending <- function(govid, years, category = NULL, invisible(TRUE) } +#' @noRd +.select_view <- function(view_base, basis) { + if (identical(basis, "harmonized")) paste0(view_base, "_harmonized") else view_base +} + #' @noRd .sql_lit_chr <- function(x) { safe <- gsub("'", "''", x, fixed = TRUE) diff --git a/R/views.R b/R/views.R index adfae8d..a5139aa 100644 --- a/R/views.R +++ b/R/views.R @@ -1,11 +1,33 @@ # 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 #' @noRd .register_views <- function(con, url, manifest) { 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) { + if (basename(f) %in% .harmonization_view_files && schema_version < 5L) next sql <- paste(readLines(f, warn = FALSE), collapse = "\n") sql <- gsub("\\{url\\}", url, sql, fixed = FALSE) DBI::dbExecute(con, sql) diff --git a/data-raw/regenerate_fixture_corpus.R b/data-raw/regenerate_fixture_corpus.R index 9e54756..086b633 100644 --- a/data-raw/regenerate_fixture_corpus.R +++ b/data-raw/regenerate_fixture_corpus.R @@ -3,12 +3,20 @@ # Regenerate inst/extdata/fixture_corpus/ from a cog_pipeline publish tree. # # What this does: -# 1. Copies the year=2019 and year=2020 long partitions as-is (byte-for- -# byte) from /data/long/ into the fixture. +# 1. Copies each requested year's long partition as-is (byte-for-byte) +# from /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, -# and summary_categories.parquet metadata 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). +# summary_categories.parquet, harmonization_map.parquet, +# harmonization_recipes.parquet, and series_breaks.parquet metadata +# 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, # reader-specification.md, README.md, series_breaks.md) from the # publish tree's docs/. @@ -35,7 +43,7 @@ regenerate_fixture_corpus <- function( "..", "cog_pipeline", "_targets", "publish_cache" ), fixture_dir = file.path("inst", "extdata", "fixture_corpus"), - fixture_years = c(2019L, 2020L)) { + fixture_years = c(2011L, 2012L, 2019L, 2020L)) { stopifnot( requireNamespace("digest", quietly = TRUE), requireNamespace("jsonlite", quietly = TRUE), @@ -92,14 +100,18 @@ regenerate_fixture_corpus <- function( invisible(NULL) } -# Copy the full (not year-scoped) canonical_fips_xwalk, canonical_alias, and -# summary_categories parquet tables. +# Copy the full (not year-scoped) canonical_fips_xwalk, canonical_alias, +# summary_categories, and (schema v5+) harmonization_map/ +# harmonization_recipes/series_breaks parquet tables. #' @noRd .copy_metadata_parquets <- function(publish_cache_dir, fixture_dir) { files <- c( "canonical_fips_xwalk.parquet", "canonical_alias.parquet", - "summary_categories.parquet" + "summary_categories.parquet", + "harmonization_map.parquet", + "harmonization_recipes.parquet", + "series_breaks.parquet" ) for (f in files) { src <- file.path(publish_cache_dir, "data", f) @@ -170,7 +182,10 @@ regenerate_fixture_corpus <- function( metadata_files <- c( "canonical_alias.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) { 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"), pipeline_commit = source_manifest$pipeline_commit, fixture_note = paste( - "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." + "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 = source_manifest$data_vintage, scope = source_manifest$scope, diff --git a/inst/extdata/fixture_corpus/data/canonical_alias.parquet b/inst/extdata/fixture_corpus/data/canonical_alias.parquet index 2b32b49..0ede2ed 100644 Binary files a/inst/extdata/fixture_corpus/data/canonical_alias.parquet and b/inst/extdata/fixture_corpus/data/canonical_alias.parquet differ diff --git a/inst/extdata/fixture_corpus/data/canonical_fips_xwalk.parquet b/inst/extdata/fixture_corpus/data/canonical_fips_xwalk.parquet index 6eed48d..72183b7 100644 Binary files a/inst/extdata/fixture_corpus/data/canonical_fips_xwalk.parquet and b/inst/extdata/fixture_corpus/data/canonical_fips_xwalk.parquet differ diff --git a/inst/extdata/fixture_corpus/data/harmonization_map.parquet b/inst/extdata/fixture_corpus/data/harmonization_map.parquet new file mode 100644 index 0000000..c41c577 Binary files /dev/null and b/inst/extdata/fixture_corpus/data/harmonization_map.parquet differ diff --git a/inst/extdata/fixture_corpus/data/harmonization_recipes.parquet b/inst/extdata/fixture_corpus/data/harmonization_recipes.parquet new file mode 100644 index 0000000..b4ac765 Binary files /dev/null and b/inst/extdata/fixture_corpus/data/harmonization_recipes.parquet differ diff --git a/inst/extdata/fixture_corpus/data/long/year=2011/part-0.parquet b/inst/extdata/fixture_corpus/data/long/year=2011/part-0.parquet new file mode 100644 index 0000000..2725698 Binary files /dev/null and b/inst/extdata/fixture_corpus/data/long/year=2011/part-0.parquet differ diff --git a/inst/extdata/fixture_corpus/data/long/year=2012/part-0.parquet b/inst/extdata/fixture_corpus/data/long/year=2012/part-0.parquet new file mode 100644 index 0000000..74879c8 Binary files /dev/null and b/inst/extdata/fixture_corpus/data/long/year=2012/part-0.parquet differ diff --git a/inst/extdata/fixture_corpus/data/long/year=2019/part-0.parquet b/inst/extdata/fixture_corpus/data/long/year=2019/part-0.parquet index 127794c..a1df4ca 100644 Binary files a/inst/extdata/fixture_corpus/data/long/year=2019/part-0.parquet and b/inst/extdata/fixture_corpus/data/long/year=2019/part-0.parquet differ diff --git a/inst/extdata/fixture_corpus/data/long/year=2020/part-0.parquet b/inst/extdata/fixture_corpus/data/long/year=2020/part-0.parquet index 4129531..5434811 100644 Binary files a/inst/extdata/fixture_corpus/data/long/year=2020/part-0.parquet and b/inst/extdata/fixture_corpus/data/long/year=2020/part-0.parquet differ diff --git a/inst/extdata/fixture_corpus/data/series_breaks.parquet b/inst/extdata/fixture_corpus/data/series_breaks.parquet new file mode 100644 index 0000000..c7f5fb9 Binary files /dev/null and b/inst/extdata/fixture_corpus/data/series_breaks.parquet differ diff --git a/inst/extdata/fixture_corpus/data/summary_categories.parquet b/inst/extdata/fixture_corpus/data/summary_categories.parquet index b58d39d..054efa2 100644 Binary files a/inst/extdata/fixture_corpus/data/summary_categories.parquet and b/inst/extdata/fixture_corpus/data/summary_categories.parquet differ diff --git a/inst/extdata/fixture_corpus/manifest.json b/inst/extdata/fixture_corpus/manifest.json index 32a78cb..8e2f99f 100644 --- a/inst/extdata/fixture_corpus/manifest.json +++ b/inst/extdata/fixture_corpus/manifest.json @@ -1,8 +1,8 @@ { - "schema_version": 4, - "built_at": "2026-07-13T23:35:09Z", - "pipeline_commit": "a082b26", - "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.", + "schema_version": 5, + "built_at": "2026-07-19T03:08:41Z", + "pipeline_commit": "ece9b32", + "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": { "census_source_downloaded": "unknown", "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." }, "schema": { - "long_column_count": 24, - "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_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", "harmonized_code", "survey_weight"], "data_dictionary": "docs/data_dictionary.md" }, "files": { "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, "path": "data/long/year=2019/part-0.parquet", - "sha256": "c0a2bf0758af129d5dfddb6ff6665cc435ddee87fd6879e788fb56ed53ab22b8", + "sha256": "475481fcfbb030f96cae457cbda824512a2da6f97f4c6d12a75130979e4331f9", "row_count": 318139, - "size_bytes": 1441404 + "size_bytes": 1717345 }, { "year": 2020, "path": "data/long/year=2020/part-0.parquet", - "sha256": "92570b9d55ec3425d034db37838f91c3b8359d0454d3d98730a6016b62e4bb48", + "sha256": "7cdf3eb1a73c34e6a3d22befdf7ccdd0b44b1e6b1cf4491159a38a5a24a1b005", "row_count": 317500, - "size_bytes": 1444011 + "size_bytes": 1720794 } ], "metadata": [ { "path": "data/canonical_alias.parquet", - "sha256": "feb8d01a640fb16c9a4b4ad66726b50b8fe8a1ce2a771bed8c1190fec51d5c8d", + "sha256": "fa27ce4b59e286ace05fbcb9f59a910b8e483fd6c58ddf01875a210b966809cb", "description": "canonical_alias.parquet" }, { "path": "data/canonical_fips_xwalk.parquet", - "sha256": "1ae47981531c7389f69eff3f7656045428564bfbe8032200eb9c039d32739a7e", + "sha256": "b6e2c4cb748141f6f2e324d6704c13830d2a68ad57bbb73d04ac54b881d1285b", "description": "canonical_fips_xwalk.parquet" }, { "path": "data/summary_categories.parquet", - "sha256": "dd59e7f58a022679ad43511c8c8e938b8dd4be81196bbeeee21e67bdcca2295b", + "sha256": "8e6fcd4dd9bb4723841a67233b19388c9762dfc23b4479501183cebf7ea3c1b5", "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" } ] }, diff --git a/inst/schemas/provenance-v1.json b/inst/schemas/provenance-v1.json index e970a3f..4c040b7 100644 --- a/inst/schemas/provenance-v1.json +++ b/inst/schemas/provenance-v1.json @@ -10,6 +10,11 @@ "target": { "type": "object" }, "years": { "type": "array", "items": { "type": "integer" } }, "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" }, "codes_summed": { "type": "object" }, "aggregate_fallback": { "type": ["object", "null"] }, diff --git a/inst/sql/22-spending_long_harmonized.sql b/inst/sql/22-spending_long_harmonized.sql new file mode 100644 index 0000000..ef70ac8 --- /dev/null +++ b/inst/sql/22-spending_long_harmonized.sql @@ -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'); diff --git a/inst/sql/23-revenue_long_harmonized.sql b/inst/sql/23-revenue_long_harmonized.sql new file mode 100644 index 0000000..69a4aa1 --- /dev/null +++ b/inst/sql/23-revenue_long_harmonized.sql @@ -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'); diff --git a/inst/sql/33-harmonization_map.sql b/inst/sql/33-harmonization_map.sql new file mode 100644 index 0000000..efa07fa --- /dev/null +++ b/inst/sql/33-harmonization_map.sql @@ -0,0 +1,3 @@ +CREATE OR REPLACE VIEW harmonization_map AS +SELECT * +FROM read_parquet('{url}data/harmonization_map.parquet'); diff --git a/inst/sql/34-harmonization_recipes.sql b/inst/sql/34-harmonization_recipes.sql new file mode 100644 index 0000000..ecaec4a --- /dev/null +++ b/inst/sql/34-harmonization_recipes.sql @@ -0,0 +1,3 @@ +CREATE OR REPLACE VIEW harmonization_recipes AS +SELECT * +FROM read_parquet('{url}data/harmonization_recipes.parquet'); diff --git a/inst/sql/35-series_breaks_pq.sql b/inst/sql/35-series_breaks_pq.sql new file mode 100644 index 0000000..2eb2207 --- /dev/null +++ b/inst/sql/35-series_breaks_pq.sql @@ -0,0 +1,3 @@ +CREATE OR REPLACE VIEW series_breaks_pq AS +SELECT * +FROM read_parquet('{url}data/series_breaks.parquet'); diff --git a/inst/sql/42-spending_annotated_harmonized.sql b/inst/sql/42-spending_annotated_harmonized.sql new file mode 100644 index 0000000..da77e55 --- /dev/null +++ b/inst/sql/42-spending_annotated_harmonized.sql @@ -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); diff --git a/inst/sql/43-revenue_annotated_harmonized.sql b/inst/sql/43-revenue_annotated_harmonized.sql new file mode 100644 index 0000000..d78b726 --- /dev/null +++ b/inst/sql/43-revenue_annotated_harmonized.sql @@ -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); diff --git a/man/cog_revenue.Rd b/man/cog_revenue.Rd index d05aac8..dd1a5e9 100644 --- a/man/cog_revenue.Rd +++ b/man/cog_revenue.Rd @@ -9,7 +9,8 @@ cog_revenue( years, category = NULL, per_capita = FALSE, - adjust_to_year = NULL + adjust_to_year = NULL, + basis = c("harmonized", "raw") ) } \arguments{ @@ -29,6 +30,16 @@ in that year).} \item{adjust_to_year}{Integer base year for CPI-U real-dollar conversion, 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{ Tibble with columns `year`, `canonical_govid`, `gov_name`, diff --git a/man/cog_spending.Rd b/man/cog_spending.Rd index a7b9bbb..876a48f 100644 --- a/man/cog_spending.Rd +++ b/man/cog_spending.Rd @@ -9,7 +9,8 @@ cog_spending( years, category = NULL, per_capita = FALSE, - adjust_to_year = NULL + adjust_to_year = NULL, + basis = c("harmonized", "raw") ) } \arguments{ @@ -29,6 +30,16 @@ in that year).} \item{adjust_to_year}{Integer base year for CPI-U real-dollar conversion, 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{ Tibble with columns `year`, `canonical_govid`, `gov_name`, diff --git a/tests/testthat/helper-fixture.R b/tests/testthat/helper-fixture.R index 12b79d3..0e650a6 100644 --- a/tests/testthat/helper-fixture.R +++ b/tests/testthat/helper-fixture.R @@ -26,3 +26,34 @@ with_fixture_corpus <- function(code) { }, add = TRUE) 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) +} diff --git a/tests/testthat/test-manifest.R b/tests/testthat/test-manifest.R index c8e5f54..c0c3ef4 100644 --- a/tests/testthat/test-manifest.R +++ b/tests/testthat/test-manifest.R @@ -107,6 +107,44 @@ test_that("cog_manifest returns the active session's parsed manifest", { expect_true(m$schema_version >= 4L) yrs <- vapply(m$files$long_partitions, function(p) as.integer(p$year), 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)) }) }) diff --git a/tests/testthat/test-revenue.R b/tests/testthat/test-revenue.R index a979996..c5b1e88 100644 --- a/tests/testthat/test-revenue.R +++ b/tests/testthat/test-revenue.R @@ -36,3 +36,20 @@ test_that("cog_revenue result has provenance attribute", { test_that("cog_revenue rejects invalid inputs", { 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) +}) diff --git a/tests/testthat/test-spending.R b/tests/testthat/test-spending.R index c767e6b..d6a6c8b 100644 --- a/tests/testthat/test-spending.R +++ b/tests/testthat/test-spending.R @@ -189,3 +189,114 @@ test_that("provenance records per-year denominator metadata", { 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) + }) +}) diff --git a/tests/testthat/test-views.R b/tests/testthat/test-views.R index 2d1255a..8cdd86d 100644 --- a/tests/testthat/test-views.R +++ b/tests/testthat/test-views.R @@ -14,6 +14,62 @@ test_that("all expected views register on session open", { 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", { skip_if_no_corpus() con <- cog_open() @@ -69,13 +125,15 @@ test_that("gov_population_yearly exposes one row per (year, canonical_govid)", { WHERE canonical_govid = '121011212191' ORDER BY year" ) - expect_setequal(df$year, c(2019L, 2020L)) - expect_equal(nrow(df), 2L) + expect_setequal(df$year, c(2011L, 2012L, 2019L, 2020L)) + expect_equal(nrow(df), 4L) expect_true(all(!is.na(df$population))) - # Hardcoded values are from the bundled fixture (regenerated 2026-07-11 - # against cog_pipeline publish tree, pipeline_commit 1a00925, Phase P - # schema_version 4). Update if the fixture is rebuilt against a - # different source vintage. + # Hardcoded values are from the bundled fixture (regenerated 2026-07-18 + # against cog_pipeline publish tree, pipeline_commit ece9b32, Phase R2 + # schema_version 5, years 2011/2012/2019/2020). Update if the fixture is + # 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 == 2020L], 1952778L) # Uniqueness on (year, canonical_govid) across the whole view.