Run the staged PIP data pipeline incrementally
pd_run_pipeline.RdRuns the durable clean, metadata, and deflate stages in topological
waves under one dependency-manifest writer lease. Current nodes are returned
as cached units. Downstream waves are accepted only after committed upstream
receipts have been reloaded into the authoritative dependency facts.
Usage
pd_run_pipeline(
inv = NULL,
force = FALSE,
verbose = getOption("pipdata.verbose", default = TRUE),
force_surveys = NULL,
bootstrap = FALSE,
bootstrap_entities = NULL,
checkpoint_size = 25L,
checkpoint_seconds = Inf
)Arguments
- inv
The complete completed-validation inventory. If
NULL, the current durable validation inventory is loaded withpipload::load_gmd_valid_inv().- force
Logical scalar. Rebuild all selected nodes and temporarily use Stamp timestamp versioning. Mutually exclusive with
force_surveys.- verbose
Logical scalar passed to pipeline storage operations.
- force_surveys
Optional character vector of exact
survey_idorpip_idselectors. Selected chains are added to ordinary invalidation.- bootstrap
Logical scalar. Explicitly permit unknown C2 provenance.
- bootstrap_entities
Optional character vector of bootstrap survey or PIP selectors. A PIP selector includes its owning clean survey and complete atomic output chain. Requires
bootstrap = TRUE.- checkpoint_size
Positive whole-number metadata and deflate checkpoint batch size.
- checkpoint_seconds
Infor a positive numeric checkpoint interval in seconds.
Value
A visible pipdata_pipeline_result. This differs intentionally from
the legacy stage wrappers, which continue to return master inventories.
Details
The C2 dependency manifest and exact Stamp receipts are the only currentness authority. The function does not persist a run cursor. A restart creates a new run and replans from the latest valid manifest. Recoverable entity errors block only their descendants. Unknown storage, lease, fence, receipt, and checkpoint failures stop later writes and are captured only after a complete stage context exists. Interrupts and explicit cancellation conditions always propagate.
The only durable nodes are clean:<survey_id>, metadata:<pip_id>, and
deflate:<pip_id>. Load, PFW merge, recode, auxiliary attachment, and save
helpers are code-fingerprint components, not separately cached nodes. Each
selected node is reported as current or stale/forced and then as cached,
runnable, successful, failed, skipped, or blocked. Cached clean nodes do not
load household artifacts.
force = TRUE rebuilds the complete selected graph. force_surveys adds
the selected survey or PIP chain to ordinary invalidation without suppressing
unrelated stale work. An absent manifest or unknown pre-C2 provenance
requires bootstrap = TRUE; use bootstrap_entities for a restrictive
canary before a complete baseline rebuild.
Auxiliary invalidation is keyed. For example, a CPI change for Colombia 2018 refreshes only matching Colombia 2018 metadata and deflate nodes; another Colombia year and unrelated country/year nodes stay cached. Worker completion is not success until exact receipts, inventories, and a manifest checkpoint are finalized. Recoverable entity failures block only their descendants. A later call resumes by authoritative replan from Stamp and the last valid manifest, without a persisted run cursor or an exactly-once guarantee.
The top-level API always uses the canonical auxiliary measures pfw, cpi,
ppp, pop, gdp, and pce. Production activation remains blocked until
signed target Windows/SMB fencing and immutable unique-rename evidence are
complete.
See also
Other pd_process_data pipeline:
add_attr(),
aux_hash_candidates(),
build_pip_inventory(),
create_attr(),
data_to_dt(),
filter_aux_data(),
filter_aux_inv(),
fix_year_var(),
get_aux_hashes(),
inv_dlw_load(),
inv_to_process(),
log_report(),
pd_aux_attr(),
pd_deflation(),
resolve_force_surveys(),
save_pip_data(),
survey_id_to_attr(),
valid_dlw_load()
Examples
if (FALSE) { # \dontrun{
pipfun::setup_working_release("20260831", "TEST")
result <- pd_run_pipeline(verbose = FALSE)
print(result)
} # }