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/NAMESPACE b/NAMESPACE index 19baa13..1c67f0b 100644 --- a/NAMESPACE +++ b/NAMESPACE @@ -10,5 +10,6 @@ export(cog_gov_search) export(cog_manifest) export(cog_mirror) export(cog_peer_compare) +export(cog_recipes) export(cog_revenue) export(cog_spending) 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/explain.R b/R/explain.R index 0cec305..2f2d44b 100644 --- a/R/explain.R +++ b/R/explain.R @@ -51,6 +51,15 @@ cog_explain <- function(result, format = c("print", "list")) { cli::cli_text("Category: (all)") } + if (!is.null(prov$basis)) { + note <- if (!is.null(prov$basis_note) && !is.na(prov$basis_note)) { + sprintf(" (%s)", prov$basis_note) + } else { + "" + } + cli::cli_text("Basis: {prov$basis}{note}") + } + cli::cli_h2("Codes observed") codes <- prov$codes_summed$observed if (length(codes) == 0L) { @@ -66,6 +75,39 @@ cog_explain <- function(result, format = c("print", "list")) { ) } + h <- prov$harmonization + if (!is.null(h) && isTRUE(h$applied)) { + cli::cli_h2("Harmonization") + cli::cli_text( + "Excluded {h$na_rows_excluded} row(s) with no harmonized_code (${format(h$na_amount_excluded, big.mark = ',')})" + ) + } + + rc <- prov$recipe + if (!is.null(rc)) { + cli::cli_h2("Recipe") + cli::cli_text("{rc$recipe_id}: {rc$label}") + comp_lines <- vapply(rc$components, function(x) { + sprintf("%s (%s, %s-%s, weight=%s)", x$component_code, x$gov_type_scope, + x$year_min, x$year_max, x$weight) + }, character(1)) + cli::cli_ul(comp_lines) + } + + if (length(prov$suggestions) > 0L) { + cli::cli_h2("Suggestions") + sugg_lines <- vapply(prov$suggestions, function(s) { + sprintf("%s -- %s (years %s-%s): %s", s$recipe_id, s$label, + s$available_years[1], s$available_years[2], s$hint) + }, character(1)) + cli::cli_ul(sugg_lines) + } + + if (length(prov$series_break_refs) > 0L) { + cli::cli_h2("Series breaks") + cli::cli_ul(.series_break_story_lines(prov$series_break_refs)) + } + cli::cli_h2("Transformations") uc <- prov$transformations$units_conversion if (isTRUE(uc$applied)) { @@ -109,6 +151,28 @@ cog_explain <- function(result, format = c("print", "list")) { invisible(NULL) } +# One "break-story" line per referenced break_id: "SB109 (2005): ". +# Re-queries series_breaks_pq for the detail (break_year, join_advice) that +# provenance$series_break_refs deliberately doesn't carry (the schema keeps +# that field to a plain id array). Falls back to bare ids if no session is +# available (e.g. explaining a result after cog_close()) rather than +# erroring cog_explain() over a cosmetic detail. +#' @noRd +.series_break_story_lines <- function(break_ids) { + con <- tryCatch(.ensure_session(), error = function(e) NULL) + if (is.null(con) || !DBI::dbIsValid(con)) return(break_ids) + detail <- tryCatch( + DBI::dbGetQuery(con, sprintf( + "SELECT break_id, break_year, join_advice FROM series_breaks_pq + WHERE break_id IN (%s) ORDER BY break_id", + .sql_lit_chr(break_ids) + )), + error = function(e) NULL + ) + if (is.null(detail) || nrow(detail) == 0L) return(break_ids) + sprintf("%s (%s): %s", detail$break_id, detail$break_year, detail$join_advice) +} + # Expand a 2-digit Census popyear (e.g. 19) to a 4-digit calendar year (2019). # F-33 metadata stores popyear as 2 digits; pivot at 70 to handle a future # corpus that ever spans pre-1970 vintages, though current scope is 2000+. 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..bfd7a03 100644 --- a/R/provenance.R +++ b/R/provenance.R @@ -4,7 +4,10 @@ #' @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, recipe = NULL, + suggestions = list()) { manifest <- .uscogdata_env$manifest codes <- result[["codes_included"]] @@ -29,6 +32,14 @@ unique(result$gov_name) } + schema_version <- suppressWarnings(as.integer(manifest$schema_version %||% 0L)) + con <- .uscogdata_env$con + break_refs <- if (!is.null(con) && DBI::dbIsValid(con)) { + .build_series_break_refs(con, codes_observed, years, schema_version) + } else { + character(0) + } + list( verb = verb, call = paste(deparse(call), collapse = " "), @@ -38,6 +49,14 @@ ), 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_ + ), + recipe = recipe, + suggestions = suggestions, scope = list( gov_types_included = as.integer(unlist(manifest$scope$gov_types_included)), gov_types_excluded = as.integer(unlist(manifest$scope$gov_types_excluded)), @@ -90,7 +109,7 @@ index = if (is.null(adjust_to_year)) NA_character_ else "CPI-U (BLS CPIAUCSL annual average, bundled)" ) ), - series_break_refs = character(0), + series_break_refs = break_refs, manifest = list( schema_version = as.integer(manifest$schema_version), pipeline_commit = manifest$pipeline_commit %||% NA_character_, diff --git a/R/recipes.R b/R/recipes.R new file mode 100644 index 0000000..11a55be --- /dev/null +++ b/R/recipes.R @@ -0,0 +1,167 @@ +# R/recipes.R +# Harmonization recipes: multi-code, cross-vintage series built by summing a +# fixed set of component item codes with per-component weights and +# year/gov-type scoping (see the `harmonization_recipes` view, registered +# from data/harmonization_recipes.parquet, schema_version >= 5 only). +# +# Recipes exist because some cross-vintage series can't be expressed as a +# 1:1 harmonized_code mapping (basis = "harmonized"): the wide era (pre-2012) +# publishes only a combined aggregate row for these families (e.g. +# corrections functions 04+05), while the modern era splits them into leaf +# codes. A recipe's generic join sums whichever of its component codes are +# present for a given year, so the resulting series is continuous across +# that format boundary. + +#' List available harmonization recipes +#' +#' Recipes are multi-code cross-vintage series (see [cog_spending()]'s +#' `recipe` argument) catalogued in the corpus's `harmonization_recipes` +#' table. Use this to discover valid `recipe` ids. +#' +#' @param pattern Optional regex matched case-insensitively against +#' `recipe_id` or `label`. +#' @return Tibble with columns `recipe_id`, `label`, `n_components`, +#' `year_min`, `year_max` (the min/max component year coverage), sorted by +#' `recipe_id`. +#' @export +cog_recipes <- function(pattern = NULL) { + if (!is.null(pattern) && + (!is.character(pattern) || length(pattern) != 1L)) { + cli::cli_abort("`pattern` must be a length-1 character string or NULL.") + } + con <- .ensure_session() + .require_schema_v5(con, .uscogdata_env$manifest, "cog_recipes()") + + where <- if (is.null(pattern)) { + "" + } else { + sprintf( + "WHERE regexp_matches(recipe_id, %1$s, 'i') OR regexp_matches(label, %1$s, 'i')", + .sql_lit_chr(pattern) + ) + } + sql <- paste( + "SELECT recipe_id, any_value(label) AS label, + COUNT(*) AS n_components, + MIN(year_min) AS year_min, MAX(year_max) AS year_max + FROM harmonization_recipes", + where, + "GROUP BY recipe_id + ORDER BY recipe_id" + ) + out <- tibble::as_tibble(DBI::dbGetQuery(con, sql)) + out$year_min <- as.integer(out$year_min) + out$year_max <- as.integer(out$year_max) + out$n_components <- as.integer(out$n_components) + out +} + +#' Abort unless the active corpus has schema_version >= 5. +#' @noRd +.require_schema_v5 <- function(con, manifest, what) { + sv <- suppressWarnings(as.integer(manifest$schema_version %||% 0L)) + if (sv < 5L) { + cli::cli_abort(c( + sprintf("%s requires corpus schema_version >= 5.", what), + x = "Active corpus has schema_version {sv}.", + i = "Point USCOGDATA_URL at a schema_version >= 5 corpus to use harmonization recipes." + ), class = "uscogdata_schema_unsupported") + } + invisible(sv) +} + +#' Abort with the valid id list unless `recipe_id` exists in the catalog. +#' @noRd +.validate_recipe_id <- function(con, recipe_id) { + ids <- DBI::dbGetQuery( + con, "SELECT DISTINCT recipe_id FROM harmonization_recipes" + )$recipe_id + if (!recipe_id %in% ids) { + cli::cli_abort(c( + "Unknown recipe = {.val {recipe_id}}.", + i = "Valid ids: {paste(sort(ids), collapse = ', ')}", + i = "See cog_recipes() for labels and year coverage." + ), class = "uscogdata_unknown_recipe") + } + invisible(TRUE) +} + +#' Fetch the component rows for one recipe (label, component codes, scope, +#' year ranges, weights) -- both for running the recipe and for the +#' `recipe` provenance block. +#' @noRd +.recipe_components <- function(con, recipe_id) { + sql <- sprintf( + "SELECT recipe_id, label, component_code, gov_type_scope, + year_min, year_max, weight, source_break_ids, notes + FROM harmonization_recipes + WHERE recipe_id = %s + ORDER BY component_code", + .sql_lit_chr(recipe_id) + ) + tibble::as_tibble(DBI::dbGetQuery(con, sql)) +} + +#' Run a recipe's generic join: sum `amt * weight` across whichever +#' component codes are present for each (year, canonical_govid), scoped by +#' gov_type_scope. Deliberately does NOT filter `NOT is_aggregate`: in the +#' wide era (<= 2011) these families' component codes exist ONLY as +#' aggregate rows (leaves first appear 2012), so excluding aggregates would +#' zero out the wide-era half of every recipe. This is safe by corpus +#' construction -- wide-era rows for these codes are aggregate-only, modern +#' rows are leaf-only, and every component row is year-scoped via +#' `year_min`/`year_max` -- so there is no double-counting. (Checkpoint +#' review docs/phase_r_harmonization_review.md § 0.2.) +#' @noRd +.run_recipe <- function(con, recipe_id, govid, years) { + sql <- sprintf( + "SELECT l.year, l.canonical_govid, + COALESCE(x.gov_name, l.gov_name) AS gov_name, + SUM(l.amt * r.weight) * 1000.0 AS amt_nominal, + string_agg(DISTINCT l.item_code, ',' ORDER BY l.item_code) AS codes_included + 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)) + LEFT JOIN canonical_fips_xwalk x USING (canonical_govid) + WHERE r.recipe_id = %1$s + AND l.canonical_govid IN (%2$s) + AND l.year IN (%3$s) + GROUP BY 1, 2, 3 + ORDER BY 1, 2", + .sql_lit_chr(recipe_id), .sql_lit_chr(govid), + paste(as.integer(years), collapse = ",") + ) + result <- tibble::as_tibble(DBI::dbGetQuery(con, sql)) + attr(result, "sql_query") <- sql + result +} + +#' Shape a raw .run_recipe() result into the standard cog_spending()/ +#' cog_revenue() column layout: subtype = "recipe", category = the recipe's +#' label, aggregate_fallback = FALSE (recipes resolve coverage gaps by +#' construction, not by falling back to an aggregate row). +#' @noRd +.shape_recipe_result <- function(result, subtype_col, label) { + sql_query <- attr(result, "sql_query") + n <- nrow(result) + result[[subtype_col]] <- rep("recipe", n) + result$category <- rep(label, n) + result$aggregate_fallback <- rep(FALSE, n) + result <- result[, c( + "year", "canonical_govid", "gov_name", subtype_col, "category", + "amt_nominal", "codes_included", "aggregate_fallback" + ), drop = FALSE] + attr(result, "sql_query") <- sql_query + result +} + +#' Turn a small data.frame into a list-of-lists (one list per row), the +#' shape used for the `recipe$components` provenance block. +#' @noRd +.df_to_row_list <- function(df) { + lapply(seq_len(nrow(df)), function(i) as.list(df[i, , drop = FALSE])) +} diff --git a/R/revenue.R b/R/revenue.R index c577556..ed71d08 100644 --- a/R/revenue.R +++ b/R/revenue.R @@ -14,16 +14,20 @@ #' 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"), recipe = NULL) { .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, + recipe = recipe ) } diff --git a/R/series_breaks.R b/R/series_breaks.R new file mode 100644 index 0000000..d4f0c1c --- /dev/null +++ b/R/series_breaks.R @@ -0,0 +1,22 @@ +# R/series_breaks.R +# Populates prov$series_break_refs (schema in inst/schemas/provenance-v1.json +# defines the field; it was always present but always empty pre-Phase-R2) +# with the ids of any catalogued series break whose fin_code appears among +# the result's observed item codes and whose break_year falls inside the +# requested year span -- the "break warnings in the provenance envelope" +# spec § 5 promises downstream consumers (cog-api passes provenance through +# verbatim). schema_version >= 5 only: series_breaks_pq isn't registered on +# an older corpus. + +#' @noRd +.build_series_break_refs <- function(con, codes_observed, years, schema_version) { + if (schema_version < 5L || length(codes_observed) == 0L) return(character(0)) + sql <- sprintf( + "SELECT DISTINCT break_id + FROM series_breaks_pq + WHERE fin_code IN (%s) AND break_year BETWEEN %d AND %d + ORDER BY break_id", + .sql_lit_chr(codes_observed), min(as.integer(years)), max(as.integer(years)) + ) + DBI::dbGetQuery(con, sql)$break_id +} 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..bb4f641 100644 --- a/R/spending.R +++ b/R/spending.R @@ -20,6 +20,28 @@ #' 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. Ignored when `recipe` +#' is set (see below). +#' @param recipe Optional harmonization recipe id (see [cog_recipes()]) for +#' multi-code cross-vintage series that a 1:1 harmonized_code mapping +#' can't express (e.g. a wide-era aggregate that only splits into leaf +#' codes in the modern era). Mutually exclusive with `category`. The +#' result's subtype column reads `"recipe"` and `category` reads the +#' recipe's label. Requires `schema_version >= 5`. A recipe query bypasses +#' `basis` entirely (it joins `long` directly rather than going through +#' the `*_annotated`/`*_annotated_harmonized` views), so the `basis` +#' argument is ignored and the result's provenance reports +#' `basis = "recipe"` with an inert `harmonization` block (`applied = +#' FALSE`, pointing at the `recipe` block instead) rather than a +#' possibly-misleading `"harmonized"`/`"raw"` value. #' @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,35 +49,65 @@ #' 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"), recipe = NULL) { .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, + recipe = recipe ) } #' @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"), recipe = NULL) { + 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) + .validate_verb_inputs(govid, years, category, per_capita, adjust_to_year, + recipe) years <- as.integer(years) 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) - sql <- .build_verb_sql(view, subtype_col, govid, years, category) - result <- tibble::as_tibble(DBI::dbGetQuery(con, sql)) + resolved <- .resolve_basis(basis, basis_explicit, manifest) + + recipe_block <- NULL + category_for_prov <- category + if (!is.null(recipe)) { + .require_schema_v5(con, manifest, "recipe =") + .validate_recipe_id(con, recipe) + comps <- .recipe_components(con, recipe) + recipe_label <- comps$label[[1]] + result <- .run_recipe(con, recipe, govid, years) + sql <- attr(result, "sql_query") + result <- .shape_recipe_result(result, subtype_col, recipe_label) + recipe_block <- list( + recipe_id = recipe, label = recipe_label, + components = .df_to_row_list(comps) + ) + category_for_prov <- recipe_label + } else { + 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)) + } if (per_capita) result <- .attach_per_capita(result, con, govid) if (!is.null(adjust_to_year)) { @@ -64,28 +116,64 @@ cog_spending <- function(govid, years, category = NULL, result$notes <- .notes_column(result) + # A recipe result doesn't go through spending_annotated(_harmonized) / + # revenue_annotated(_harmonized) at all -- .run_recipe()'s generic join + # reads `long` directly -- so `basis` and the `harmonization` exclusion + # count (which is itself computed from `long`, independent of which view + # a non-recipe query used) would describe a code path this result never + # took. Rather than report a technically-still-computed but misleading + # basis = "harmonized"/"raw" + harmonization$applied combo, recipe + # results report basis = "recipe" and an explicit, inert harmonization + # block pointing at the `recipe` block instead. Task 12 (cog-api) passes + # provenance through verbatim, so this needs to be unambiguous rather + # than technically-defensible-but-confusing. + if (!is.null(recipe)) { + basis_for_prov <- "recipe" + basis_note_for_prov <- NA_character_ + harmonization <- list( + applied = FALSE, na_rows_excluded = 0L, na_amount_excluded = 0, + note = "basis/harmonization not applicable to recipe results; see the recipe block instead" + ) + suggestions <- list() + } else { + basis_for_prov <- resolved$basis + basis_note_for_prov <- resolved$note + harmonization <- .build_harmonization_block( + con, govid, years, resolved, flow_prefixes + ) + suggestions <- .build_suggestions(con, govid, years, category, result, resolved$basis) + } + prov <- .build_provenance( verb = verb, call = call, govid = govid, years = years, - category = category, + category = category_for_prov, per_capita = per_capita, adjust_to_year = adjust_to_year, result = result, sql = sql, - subtype_col = subtype_col + subtype_col = subtype_col, + basis = basis_for_prov, + basis_note = basis_note_for_prov, + harmonization = harmonization, + recipe = recipe_block, + suggestions = suggestions ) prov$scope$govids_found <- scope$found prov$scope$govids_missing <- scope$missing attr(result, "provenance") <- prov attr(result, ".popyear_range") <- NULL + + if (length(suggestions) > 0L) .inform_suggestions(suggestions) + result } #' @noRd .validate_verb_inputs <- function(govid, years, category, - per_capita, adjust_to_year) { + per_capita, adjust_to_year, recipe = NULL) { if (!is.character(govid) || length(govid) == 0L) { cli::cli_abort("`govid` must be a non-empty character vector.") } @@ -104,9 +192,25 @@ cog_spending <- function(govid, years, category = NULL, cli::cli_abort("`adjust_to_year` must be NULL or a length-1 integer.") } } + if (!is.null(recipe)) { + if (!is.character(recipe) || length(recipe) != 1L) { + cli::cli_abort("`recipe` must be NULL or a length-1 character string.") + } + if (!is.null(category)) { + cli::cli_abort(c( + "`recipe` and `category` are mutually exclusive.", + i = "Pass one or the other, not both." + ), class = "uscogdata_recipe_category_conflict") + } + } 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/suggestions.R b/R/suggestions.R new file mode 100644 index 0000000..fc59d21 --- /dev/null +++ b/R/suggestions.R @@ -0,0 +1,119 @@ +# R/suggestions.R +# Recipe-component-driven signposting: when a basis = "harmonized" query for +# a category comes back with a coverage gap in some requested years (the +# result has no rows at all in that year) that a harmonization recipe would +# actually fill for this government, surface that recipe as a suggestion. +# +# This is deliberately keyed off the recipe catalog's component codes, not +# off harmonization_map rows: no live map row carries a non-blank +# suggested_recipe_id (the corpus's wide era exposes split families like +# corrections functions 04+05 ONLY as aggregate rows, which basis = +# "harmonized" excludes by construction -- there's no NA ruling to hang a +# suggestion off of, just a leaf-code absence a recipe happens to fill). +# See docs/phase_r_harmonization_review.md § 0.3. +# +# Scope is deliberately narrow: signposting only runs when the caller +# supplied a `category` (an un-scoped, all-categories query has no single +# coverage question to answer) and only flags a recipe when the ACTUAL +# result has zero rows in a requested year AND the candidate recipe's own +# generic join (same join .run_recipe() uses, including its wide-era +# aggregate rows) produces at least one row for this government in that +# year. Checking presence per-government (not corpus-wide) avoids false +# positives from ordinary reporting variance -- most governments don't use +# every sibling code in a multi-code category every year, and that is not +# a format-boundary gap worth signposting. + +#' Build the `prov$suggestions` list for a (non-recipe) basis = "harmonized" +#' verb call: recipes whose generic join would fill a real gap in `result`. +#' +#' @param con Active DuckDB connection. +#' @param govid Character vector of canonical_govid values (the verb's raw +#' `govid`). +#' @param years Integer vector of requested years. +#' @param category `category` argument as passed to the verb (character +#' vector or `NULL`; suggestions are only computed when non-NULL). +#' @param result The verb's already-computed result tibble (post basis +#' query, pre per_capita/adjust_to_year). +#' @param basis The *resolved* basis (`"harmonized"` or `"raw"`). +#' @return List of `list(recipe_id, label, available_years, hint)`, possibly +#' empty. +#' @noRd +.build_suggestions <- function(con, govid, years, category, result, basis) { + if (!identical(basis, "harmonized") || is.null(category)) return(list()) + + candidates <- DBI::dbGetQuery(con, sprintf( + "SELECT DISTINCT recipe_id FROM harmonization_recipes + WHERE component_code IN ( + SELECT DISTINCT item_code FROM summary_categories WHERE category IN (%s) + )", + .sql_lit_chr(category) + ))$recipe_id + if (length(candidates) == 0L) return(list()) + + result_years <- if (is.null(result) || nrow(result) == 0L) { + integer(0) + } else { + unique(as.integer(result$year)) + } + gap_years <- setdiff(as.integer(years), result_years) + if (length(gap_years) == 0L) return(list()) + + meta <- 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) + ))) + + # Which (recipe_id, year) pairs the recipe's own generic join actually + # covers for this government, restricted to the gap years -- the same + # join .run_recipe() uses (component year_min/year_max + gov_type_scope, + # no is_aggregate filter), just checking existence instead of summing. + covered <- 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 l.canonical_govid IN (%s) + AND l.year IN (%s)", + .sql_lit_chr(candidates), .sql_lit_chr(govid), + paste(gap_years, collapse = ",") + )) + + suggestions <- list() + for (rid in candidates) { + if (!rid %in% covered$recipe_id) next + m <- meta[meta$recipe_id == rid, ] + suggestions[[length(suggestions) + 1L]] <- list( + recipe_id = rid, + label = m$label[[1]], + available_years = c(as.integer(m$year_min), as.integer(m$year_max)), + hint = sprintf("re-run with recipe = '%s'", rid) + ) + } + suggestions +} + +#' Emit the single cli::cli_inform() message summarizing all suggestions +#' for a verb call (the brief's "one message", not one per suggestion). +#' Bullet text is pre-formatted plain text (no cli/glue `{}` markup) since +#' recipe ids/labels are untrusted-ish data values, not literal call-site +#' expressions. +#' @noRd +.inform_suggestions <- function(suggestions) { + bullets <- vapply(suggestions, function(s) { + sprintf("%s (%d-%d): %s", s$recipe_id, + s$available_years[1], s$available_years[2], s$hint) + }, character(1)) + cli::cli_inform(c( + i = "Coverage gap detected for the requested years; a harmonization recipe may fill it:", + stats::setNames(bullets, rep("*", length(bullets))) + )) +} 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_recipes.Rd b/man/cog_recipes.Rd new file mode 100644 index 0000000..149fb76 --- /dev/null +++ b/man/cog_recipes.Rd @@ -0,0 +1,22 @@ +% Generated by roxygen2: do not edit by hand +% Please edit documentation in R/recipes.R +\name{cog_recipes} +\alias{cog_recipes} +\title{List available harmonization recipes} +\usage{ +cog_recipes(pattern = NULL) +} +\arguments{ +\item{pattern}{Optional regex matched case-insensitively against +`recipe_id` or `label`.} +} +\value{ +Tibble with columns `recipe_id`, `label`, `n_components`, + `year_min`, `year_max` (the min/max component year coverage), sorted by + `recipe_id`. +} +\description{ +Recipes are multi-code cross-vintage series (see [cog_spending()]'s +`recipe` argument) catalogued in the corpus's `harmonization_recipes` +table. Use this to discover valid `recipe` ids. +} diff --git a/man/cog_revenue.Rd b/man/cog_revenue.Rd index d05aac8..843d48e 100644 --- a/man/cog_revenue.Rd +++ b/man/cog_revenue.Rd @@ -9,7 +9,9 @@ cog_revenue( years, category = NULL, per_capita = FALSE, - adjust_to_year = NULL + adjust_to_year = NULL, + basis = c("harmonized", "raw"), + recipe = NULL ) } \arguments{ @@ -29,6 +31,30 @@ 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. Ignored when `recipe` +is set (see below).} + +\item{recipe}{Optional harmonization recipe id (see [cog_recipes()]) for +multi-code cross-vintage series that a 1:1 harmonized_code mapping +can't express (e.g. a wide-era aggregate that only splits into leaf +codes in the modern era). Mutually exclusive with `category`. The +result's subtype column reads `"recipe"` and `category` reads the +recipe's label. Requires `schema_version >= 5`. A recipe query bypasses +`basis` entirely (it joins `long` directly rather than going through +the `*_annotated`/`*_annotated_harmonized` views), so the `basis` +argument is ignored and the result's provenance reports +`basis = "recipe"` with an inert `harmonization` block (`applied = +FALSE`, pointing at the `recipe` block instead) rather than a +possibly-misleading `"harmonized"`/`"raw"` value.} } \value{ Tibble with columns `year`, `canonical_govid`, `gov_name`, diff --git a/man/cog_spending.Rd b/man/cog_spending.Rd index a7b9bbb..f1016ec 100644 --- a/man/cog_spending.Rd +++ b/man/cog_spending.Rd @@ -9,7 +9,9 @@ cog_spending( years, category = NULL, per_capita = FALSE, - adjust_to_year = NULL + adjust_to_year = NULL, + basis = c("harmonized", "raw"), + recipe = NULL ) } \arguments{ @@ -29,6 +31,30 @@ 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. Ignored when `recipe` +is set (see below).} + +\item{recipe}{Optional harmonization recipe id (see [cog_recipes()]) for +multi-code cross-vintage series that a 1:1 harmonized_code mapping +can't express (e.g. a wide-era aggregate that only splits into leaf +codes in the modern era). Mutually exclusive with `category`. The +result's subtype column reads `"recipe"` and `category` reads the +recipe's label. Requires `schema_version >= 5`. A recipe query bypasses +`basis` entirely (it joins `long` directly rather than going through +the `*_annotated`/`*_annotated_harmonized` views), so the `basis` +argument is ignored and the result's provenance reports +`basis = "recipe"` with an inert `harmonization` block (`applied = +FALSE`, pointing at the `recipe` block instead) rather than a +possibly-misleading `"harmonized"`/`"raw"` value.} } \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-explain.R b/tests/testthat/test-explain.R index 83f18ab..e824807 100644 --- a/tests/testthat/test-explain.R +++ b/tests/testthat/test-explain.R @@ -31,6 +31,45 @@ test_that("cog_explain errors on non-verb input", { expect_error(cog_explain(df), "provenance") }) +test_that("cog_explain prints basis + harmonization block", { + skip_if_no_corpus() + r <- cog_spending("121011212191", 2020L, "Corrections") + txt <- paste(c( + capture.output(cog_explain(r)), + capture.output(cog_explain(r), type = "message") + ), collapse = "\n") + expect_true(grepl("Basis: harmonized", txt)) + expect_true(grepl("Harmonization", txt)) + expect_true(grepl("Excluded 0 row", txt)) +}) + +test_that("cog_explain prints a Recipe section for recipe = results", { + skip_if_no_corpus() + r <- cog_spending("121011212191", c(2011L, 2012L), recipe = "corrections_combined") + txt <- paste(c( + capture.output(cog_explain(r)), + capture.output(cog_explain(r), type = "message") + ), collapse = "\n") + expect_true(grepl("Recipe", txt)) + expect_true(grepl("corrections_combined", txt)) + expect_true(grepl("E04", txt)) + expect_true(grepl("E05", txt)) +}) + +test_that("cog_explain prints a Suggestions section when the provenance has one", { + skip_if_no_corpus() + r <- suppressMessages( + cog_spending("121011212191", c(2011L, 2012L), category = "Corrections") + ) + txt <- paste(c( + capture.output(cog_explain(r)), + capture.output(cog_explain(r), type = "message") + ), collapse = "\n") + expect_true(grepl("Suggestions", txt)) + expect_true(grepl("corrections_combined", txt)) + expect_true(grepl("re-run with recipe", txt)) +}) + test_that("cog_explain prints denominator + popyear_range + counts", { skip_if_no_corpus() with_fixture_corpus({ 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-recipes.R b/tests/testthat/test-recipes.R new file mode 100644 index 0000000..0bf9518 --- /dev/null +++ b/tests/testthat/test-recipes.R @@ -0,0 +1,216 @@ +# tests/testthat/test-recipes.R +# +# cog_recipes(), recipe = in cog_spending()/cog_revenue(), and the +# recipe-component-driven signposting in prov$suggestions (Phase R2 / +# Task 11, schema_version 5). + +test_that("cog_recipes lists the curated catalog including corrections_combined", { + skip_if_no_corpus() + r <- cog_recipes() + expect_s3_class(r, "tbl_df") + expect_equal(names(r), c("recipe_id", "label", "n_components", "year_min", "year_max")) + expect_equal(nrow(r), 24L) + expect_true("corrections_combined" %in% r$recipe_id) + expect_true("t19_selective_sales_wide" %in% r$recipe_id) + expect_true("ig_federal_b89_wide" %in% r$recipe_id) + expect_true("rents_royalties_u4_wide" %in% r$recipe_id) + expect_true("higher_ed_e18_wide" %in% r$recipe_id) + expect_true("cash_securities_z77_wide" %in% r$recipe_id) + # Superseded id from the pre-curation brief text must NOT be present. + expect_false("corrections_judicial_combined" %in% r$recipe_id) +}) + +test_that("cog_recipes(pattern=) filters by recipe_id or label", { + skip_if_no_corpus() + r <- cog_recipes("corrections") + expect_true(nrow(r) >= 1L) + expect_true(all(grepl("corrections", r$recipe_id, ignore.case = TRUE) | + grepl("corrections", r$label, ignore.case = TRUE))) +}) + +test_that("cog_recipes requires schema_version >= 5", { + skip_if_no_corpus() + with_doctored_schema_version(4L, { + expect_error(cog_recipes(), class = "uscogdata_schema_unsupported") + }) +}) + +# --- recipe = : generic join, no is_aggregate filter ----------------------- + +test_that("recipe = 'corrections_combined' is continuous across the 2011->2012 seam", { + skip_if_no_corpus() + r <- cog_spending("121011212191", years = c(2011L, 2012L), + recipe = "corrections_combined") + expect_equal(nrow(r), 2L) + expect_true(all(c("year", "canonical_govid", "gov_name", "spend_subtype", + "category", "amt_nominal", "codes_included", + "aggregate_fallback", "notes") %in% names(r))) + expect_equal(unique(r$spend_subtype), "recipe") + expect_equal(unique(r$category), "Corrections (functions 04+05 combined)") + expect_false(any(r$aggregate_fallback)) + + r2011 <- r$amt_nominal[r$year == 2011L] + r2012 <- r$amt_nominal[r$year == 2012L] + # 2011: E05 only exists as a wide-era AGGREGATE row (is_aggregate = TRUE) + # for Broward -- data-verified $216,088,000. Since .run_recipe() does NOT + # filter is_aggregate (amendment: the recipe join must not, because these + # families exist ONLY as aggregate rows in the wide era), the recipe + # correctly picks this up. + expect_equal(r2011, 216088000) + # 2012: modern E04 leaf ($213,056,000); Broward reports no E05 leaf that + # year, so the recipe total equals E04 alone -- still continuous with the + # 2011 aggregate, proving the wide-aggregate -> modern-leaf handoff. + expect_equal(r2012, 213056000) + expect_true(all(grepl("E04|E05", r$codes_included))) +}) + +test_that("recipe result carries a recipe provenance block with component rows", { + skip_if_no_corpus() + r <- cog_spending("121011212191", years = c(2011L, 2012L), + recipe = "corrections_combined") + prov <- attr(r, "provenance") + expect_equal(prov$basis, "recipe") + expect_equal(prov$category, "Corrections (functions 04+05 combined)") + expect_type(prov$recipe, "list") + expect_equal(prov$recipe$recipe_id, "corrections_combined") + expect_equal(prov$recipe$label, "Corrections (functions 04+05 combined)") + expect_length(prov$recipe$components, 2L) + comp_codes <- vapply(prov$recipe$components, function(x) x$component_code, character(1)) + expect_setequal(comp_codes, c("E04", "E05")) + # A recipe query resolves its own coverage; it should never also carry + # suggestions for itself. + expect_length(prov$suggestions, 0L) +}) + +test_that("recipe results report an unambiguous basis/harmonization, ignoring basis=", { + skip_if_no_corpus() + # A recipe query bypasses spending_annotated(_harmonized) entirely -- + # .run_recipe() joins `long` directly -- so `basis` must never read + # "harmonized"/"raw" (which would describe a code path this query never + # took) regardless of what the caller passed for `basis`. Task 12 + # consumes provenance verbatim, so this needs to be unambiguous. + r_default <- cog_spending("121011212191", years = c(2011L, 2012L), + recipe = "corrections_combined") + r_raw <- cog_spending("121011212191", years = c(2011L, 2012L), + recipe = "corrections_combined", basis = "raw") + r_harm <- cog_spending("121011212191", years = c(2011L, 2012L), + recipe = "corrections_combined", basis = "harmonized") + + for (r in list(r_default, r_raw, r_harm)) { + prov <- attr(r, "provenance") + expect_equal(prov$basis, "recipe") + expect_true(is.na(prov$basis_note)) + expect_false(prov$harmonization$applied) + expect_equal(prov$harmonization$na_rows_excluded, 0L) + expect_match(prov$harmonization$note, "recipe", ignore.case = TRUE) + } + + # basis= truly has zero effect on a recipe query's actual numbers. + expect_equal(r_raw$amt_nominal, r_harm$amt_nominal) + expect_equal(r_default$amt_nominal, r_raw$amt_nominal) +}) + +test_that("recipe = 't19_selective_sales_wide' sums the local T11/T14 legs when present", { + skip_if_no_corpus() + # Westminster City, CA (canonical_govid 082001211654): T11 = 0 in 2011, + # T11 = 568 (T14 = 0/absent) in 2012 -- a real, data-verified equality/ + # inequality pair inside the amended fixture window (2011-2012), standing + # in for the brief's original 2004/2005 example (out of scope per the + # amended fixture years; the underlying local-tax-split boundary is + # nationally FY2005, but this government's own T11 reporting activates + # within our 2011-2012 window). + r <- cog_revenue("082001211654", years = c(2011L, 2012L), + recipe = "t19_selective_sales_wide") + + # Raw, single-code T19 total (not the "Other Taxes" category total, which + # would also sum in T11/T14/T21/T23/T27/T29/T53/T99 -- queried directly to + # isolate exactly the code the brief's equality/inequality check is about). + con <- uscogdata:::.ensure_session() + raw_t19 <- DBI::dbGetQuery(con, " + SELECT year, SUM(amt) * 1000.0 AS amt + FROM revenue_long + WHERE canonical_govid = '082001211654' AND item_code = 'T19' + AND year IN (2011, 2012) + GROUP BY year ORDER BY year + ") + raw_t19_2011 <- raw_t19$amt[raw_t19$year == 2011L] + raw_t19_2012 <- raw_t19$amt[raw_t19$year == 2012L] + expect_equal(raw_t19_2011, 2231000) + expect_equal(raw_t19_2012, 2365000) + + recipe_2011 <- r$amt_nominal[r$year == 2011L] + recipe_2012 <- r$amt_nominal[r$year == 2012L] + + expect_equal(recipe_2011, raw_t19_2011) # equality: no local T11/T14 yet + expect_gt(recipe_2012, raw_t19_2012) # inequality: local T11 joins in + expect_equal(recipe_2012, raw_t19_2012 + 568000) +}) + +test_that("recipe = and category = together aborts", { + skip_if_no_corpus() + expect_error( + cog_spending("121011212191", 2020L, category = "Corrections", + recipe = "corrections_combined"), + class = "uscogdata_recipe_category_conflict" + ) +}) + +test_that("unknown recipe id aborts and lists valid ids", { + skip_if_no_corpus() + err <- tryCatch( + cog_spending("121011212191", 2020L, recipe = "does_not_exist"), + error = identity + ) + expect_s3_class(err, "uscogdata_unknown_recipe") + expect_match(conditionMessage(err), "corrections_combined") +}) + +test_that("recipe = requires schema_version >= 5", { + skip_if_no_corpus() + with_doctored_schema_version(4L, { + expect_error( + cog_spending("121011212191", 2020L, recipe = "corrections_combined"), + class = "uscogdata_schema_unsupported" + ) + }) +}) + +# --- signposting ------------------------------------------------------- + +test_that("signposting suggests corrections_combined across the 2011->2012 gap", { + skip_if_no_corpus() + expect_message( + r <- cog_spending("121011212191", years = c(2011L, 2012L), + category = "Corrections"), + "recipe" + ) + prov <- attr(r, "provenance") + expect_true(length(prov$suggestions) >= 1L) + ids <- vapply(prov$suggestions, function(s) s$recipe_id, character(1)) + expect_true("corrections_combined" %in% ids) + hit <- prov$suggestions[[which(ids == "corrections_combined")]] + expect_equal(hit$hint, "re-run with recipe = 'corrections_combined'") + expect_equal(hit$available_years, c(1967L, 2023L)) +}) + +test_that("no signposting when the result already has full year coverage", { + skip_if_no_corpus() + r <- cog_spending("121011212191", years = 2019:2020, category = "Corrections") + prov <- attr(r, "provenance") + expect_length(prov$suggestions, 0L) +}) + +test_that("no signposting when category is NULL (unscoped query)", { + skip_if_no_corpus() + r <- cog_spending("121011212191", years = c(2011L, 2012L)) + prov <- attr(r, "provenance") + expect_length(prov$suggestions, 0L) +}) + +test_that("no signposting under basis = 'raw'", { + skip_if_no_corpus() + r <- cog_spending("121011212191", years = c(2011L, 2012L), + category = "Corrections", basis = "raw") + prov <- attr(r, "provenance") + expect_length(prov$suggestions, 0L) +}) 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..b8e7c08 100644 --- a/tests/testthat/test-spending.R +++ b/tests/testthat/test-spending.R @@ -189,3 +189,129 @@ 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("provenance$series_break_refs is a populated-when-applicable character vector", { + skip_if_no_corpus() + with_fixture_corpus({ + r <- cog_spending("121011212191", 2020L, "Corrections") + refs <- attr(r, "provenance")$series_break_refs + expect_type(refs, "character") + # No catalogued series_breaks_pq row falls inside this fixture's + # 2011/2012/2019/2020 window for the codes this query touches (E04/G04) + # -- data-verified; the mechanism itself is what's under test here, via + # a query-shaped unit test in test-views.R since the fixture has no + # positive case to pin against. + expect_equal(refs, character(0)) + }) +}) + +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..75eaa5b 100644 --- a/tests/testthat/test-views.R +++ b/tests/testthat/test-views.R @@ -14,6 +14,144 @@ test_that("all expected views register on session open", { expect_true(all(expected %in% views$table_name)) }) +test_that("inst/sql/22- and 23- harmonized views enforce every WHERE predicate (real SQL text, synthetic parquet)", { + # spending_long_harmonized / revenue_long_harmonized are three-predicate + # views: + # 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 () + # 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 or a + # predicate-excluded row. Rather than re-implement the WHERE clause by + # hand against an in-memory VALUES table (which would only prove the SQL + # *pattern* works, not that the deployed inst/sql/22-/23- text actually + # applies it), this test reads the real SQL files off disk, substitutes + # {url} exactly as .register_views() does, and executes them -- plus + # their 10-long.sql dependency -- against a synthetic hive-partitioned + # parquet tree written to a temp dir. A regression in any predicate (e.g. + # `NOT is_aggregate` dropped, the prefix list changed, the NULL guard + # removed) would change which of the rows below survive. + # + # The synthetic parquet is written with DuckDB's own COPY ... TO (FORMAT + # PARQUET) rather than the arrow package: this package has no arrow + # dependency (CLAUDE.md "No arrow dependency -- DuckDB reads parquet + # natively"), and DuckDB can round-trip its own parquet writer/reader + # without adding one for tests either. + skip_if_no_corpus() + + tmp <- withr::local_tempdir() + part_dir <- file.path(tmp, "data", "long", "year=2004") + dir.create(part_dir, recursive = TRUE) + part_path <- file.path(part_dir, "part-0.parquet") + + write_con <- DBI::dbConnect(duckdb::duckdb()) + on.exit(DBI::dbDisconnect(write_con, shutdown = TRUE), add = TRUE) + DBI::dbExecute(write_con, sprintf(" + COPY ( + SELECT * FROM (VALUES + -- Spending (E/F/G/K) rows, exercised against spending_long_harmonized: + ('spend-A', 'E36', 100, false, 'E36'), -- control: passes every predicate as-is + ('spend-B', 'E38', 50, false, 'E36'), -- collapse-fold: passes every predicate, renamed to E36 + ('spend-C', 'E05', 999999, true, 'E05'), -- excluded ONLY by `NOT is_aggregate` + ('spend-D', 'E99', 888888, false, NULL), -- excluded by `harmonized_code IS NOT NULL` + -- 'S74' is outside BOTH flow families (E/F/G/K spending and + -- T/A/U/B/C/D revenue -- it mirrors the real corpus's own + -- non-flow-type codes like S74/Z61), so it can only leak into + -- EITHER view via the E/F/G/K or T/A/U/B/C/D prefix filter, never + -- both at once -- a prefix drawn from the other view's own family + -- (e.g. a real T-code for the spending row) would incorrectly + -- leak into the other view's assertion below and not discriminate + -- the predicate under test. + ('spend-E', 'S74', 777777, false, 'S74'), -- excluded ONLY by the E/F/G/K prefix filter + -- Revenue (T/A/U/B/C/D) rows, exercised against revenue_long_harmonized: + ('rev-A', 'U11', 200, false, 'U11'), -- control: passes every predicate as-is + ('rev-B', 'U10', 25, false, 'U11'), -- collapse-fold: passes every predicate, renamed to U11 + ('rev-C', 'T29', 555555, true, 'T29'), -- excluded ONLY by `NOT is_aggregate` + ('rev-D', 'T88', 444444, false, NULL), -- excluded by `harmonized_code IS NOT NULL` + ('rev-E', 'Z61', 333333, false, 'Z61') -- excluded ONLY by the T/A/U/B/C/D prefix filter + ) AS t(canonical_govid, item_code, amt, is_aggregate, harmonized_code) + ) TO %s (FORMAT PARQUET) + ", uscogdata:::.sql_lit_chr(part_path))) + + sql_dir <- system.file("sql", package = "uscogdata") + .read_view_sql <- function(filename) { + txt <- paste(readLines(file.path(sql_dir, filename), warn = FALSE), collapse = "\n") + gsub("\\{url\\}", paste0(tmp, "/"), txt, fixed = FALSE) + } + + con <- DBI::dbConnect(duckdb::duckdb()) + on.exit(DBI::dbDisconnect(con, shutdown = TRUE), add = TRUE) + DBI::dbExecute(con, .read_view_sql("10-long.sql")) + DBI::dbExecute(con, .read_view_sql("22-spending_long_harmonized.sql")) + DBI::dbExecute(con, .read_view_sql("23-revenue_long_harmonized.sql")) + + spend <- DBI::dbGetQuery(con, + "SELECT item_code, SUM(amt) AS amt FROM spending_long_harmonized + GROUP BY item_code ORDER BY item_code" + ) + # Exactly one surviving row: spend-C (aggregate), spend-D (NULL + # harmonized_code), and spend-E (wrong prefix family) must all be gone, + # and spend-A + spend-B must be folded together under E36. + expect_equal(nrow(spend), 1L) + expect_equal(spend$item_code, "E36") + expect_equal(spend$amt, 150) + + rev <- DBI::dbGetQuery(con, + "SELECT item_code, SUM(amt) AS amt FROM revenue_long_harmonized + GROUP BY item_code ORDER BY item_code" + ) + expect_equal(nrow(rev), 1L) + expect_equal(rev$item_code, "U11") + expect_equal(rev$amt, 225) +}) + +test_that(".build_series_break_refs matches fin_code + break_year window", { + # No series_breaks_pq row falls inside the bundled fixture's 2011-2020 + # window (data-verified; see the "series_break_refs" test in + # test-spending.R), so this proves the matching logic itself against the + # live view + a synthetic year window that DOES hit a cataloged break + # (SB075, fin_code E62, break_year 2005). + skip_if_no_corpus() + con <- cog_open() + on.exit(cog_close()) + refs <- uscogdata:::.build_series_break_refs( + con, codes_observed = c("E62", "E04"), years = c(2003L, 2006L), + schema_version = 5L + ) + expect_true("SB075" %in% refs) + expect_true("SB071" %in% refs) + + # Gated on schema_version >= 5 even when the codes/years would otherwise match. + refs_v4 <- uscogdata:::.build_series_break_refs( + con, codes_observed = c("E62"), years = c(2003L, 2006L), schema_version = 4L + ) + expect_equal(refs_v4, character(0)) +}) + +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 +207,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.