functions
Shared helper functions for the file metadata runtime.
This module is the single source of truth for pure helpers and internal utilities used across the other modules in this package:
flow.py— Prefect flow entry-pointrefresh.py— cron-driven refresh flowtasks.py— Prefect tasks
Module
Functions
file_metadata_run_tags
def file_metadata_run_tags( source_path: str | None, datasource_type: str, *, scoped: bool,) ‑> dict[str, typing.Any]:Return the tags a file_metadata run is recorded with.
A scoped (only_paths) run sees only its own paths, so neither the
unchanged-tree skip nor refresh staleness may treat it as covering the
tree. The root and type identify which inputs the run indexed, so a pass
with different inputs is not measured against it.
Arguments
source_path: The datasource root or file path, as configured.datasource_type: The datasource type string.scoped: Whether the run was limited to explicit paths.
first_index_incomplete
def first_index_incomplete(cache: CacheProtocol, task_hash: str) ‑> bool:Whether the only inventory is what a cut-short first index reached.
True when some full file_metadata pass ended partial and none has ever
completed. Not "the latest pass is partial": a follow-up that is running,
or that failed, leaves the inventory just as incomplete. Only a complete
pass clears it, and until one exists the refresh keeps resubmitting the
datasource, since it has never been refreshed.
Arguments
cache: An open cache for the pod's background cache DB.task_hash: The datasource-level task hash.
generate_prefect_task_hash
def generate_prefect_task_hash(pod_name: str, datasource_name: str) ‑> str:Return a stable, fixed-length task hash for a (pod, datasource) pair.
The hash is derived from "{pod_name}/{datasource_name}" so that it
is deterministic across restarts and serves as the composite primary-key
prefix in FileMetadata and Run.
This is the datasource leaf hash, and it plays two roles:
- It keys the datasource/run-level tables directly —
file_metadata,scan_metadata,runs,schema_versions— none of which a DAG step writes. - It is the leaf mixed into every per-step Merkle
task_hash(bitfount.flows.dag.hashing), which is what keeps two datasources on the same pod in separate partitions of a step's table.
Arguments
pod_name: The name of the pod.datasource_name: The name of the datasource within the pod.
Returns A 32-character hex string (the first 128 bits of the SHA-256 digest).
inventory_is_partial
def inventory_is_partial(cache: CacheProtocol, task_hash: str) ‑> bool:Whether the inventory for task_hash is known to be incomplete.
Either the latest full file_metadata pass stopped part-way because its
source became unreachable, or a first index did and no full pass has
completed since (first_index_incomplete), which a later follow-up that
failed does not change. Scoped passes are ignored: they never claim to
cover the tree.
Arguments
cache: An open cache for the pod's background cache DB.task_hash: The datasource-level task hash.
is_deconfigured
def is_deconfigured(cache: CacheProtocol, task_hash: str) ‑> bool:Whether the pod has recorded task_hash's datasource as removed.
Checked by a metadata run before it touches the datasource. The pod records
a removal in configured_datasources before it cancels the removal's runs,
so a run submitted before the removal but started after it (a late start-up
notify, a queued or parked run) stands down here rather than walking a
datasource nobody has any more.
Only an explicit tombstone stands a run down. A datasource the snapshot has
no row for is not removed: the pod may know it under another name (identifier
enforcement can shorten one), or not have written its snapshot yet. A
snapshot that cannot be read reads as not removed. Off with the same switch
as the nightly refresh's gate, metadata_refresh_requires_configured_datasource.
Arguments
cache: The pod's open background cache.task_hash: The datasource's task hash.
Returns
True only when the pod has recorded the datasource as removed.
make_dataset_identifier
def make_dataset_identifier(pod_name: str, datasource_name: str) ‑> str:Return the "{pod_name}/{datasource_name}" dataset identifier.
This is the same slash-joined key generate_prefect_task_hash hashes, and the value
stored in FileMetadata.dataset_identifier.
open_background_cache
def open_background_cache( pod_name: str, *, pod_key_path: Path | None = None, create_dir: bool = False,) ‑> CacheProtocol:Open the per-pod background cache DB.
Single source of truth for resolving the per-pod background cache path
(PODS_CACHE_ROOT/<pod_name>/BACKGROUND_CACHE_FILENAME) and opening it,
shared by the DATASET_PROJECT_LINKED gRPC pod handler and the
run_background_task CLI so the two entry points cannot drift.
Arguments
pod_name: Name of the pod whose cache DB to open.pod_key_path: Path to the pod'spod_rsa.pem(typicallypod.pod_key_path). When provided, it is bound into the cache-encryption layer so encrypted columns can derive their key from this pod's RSA private key. Optional so test/CLI flows that don't need encryption can still call with justpod_name.create_dir: WhenTruethe parent directory is created if missing — the live pod handler path, where the DB is created on first write. WhenFalsethe cache file must already exist orFileNotFoundErroris raised — the CLI path, which requires the pod to have been started at least once to populate file metadata.
Returns
A CacheProtocol instance for the pod's background cache.
Raises
FileNotFoundError: When create_dir isFalseand the cache file does not exist.
resolve_datasource_inputs
def resolve_datasource_inputs(datasource: Any) ‑> tuple[str | None, str | None]:Return the (path, connection_string) a datasource should be indexed by.
file_metadata_runtime takes both as separate parameters and routes on
whichever is set, so every caller holding a live datasource has to make the
same decision: filesystem sources index by their directory path, SQL sources
by their connection string, anything else by neither.
Written once here because it was previously repeated at each trigger site,
and a site that forgot a branch would not fail — it would submit a run with
both inputs None, which the flow rejects only once it is already running.
Note this resolves from a live datasource instance. The flow itself routes
from a resolved class (via issubclass), because it receives a type string
rather than an object; the two are deliberately not merged.
Arguments
datasource: A constructed datasource instance.
Returns
(path, connection_string) with at most one set. Both are None for a
datasource of neither kind, which callers may treat as "nothing to index".
split_dataset_identifier
def split_dataset_identifier(dataset_identifier: str) ‑> tuple[str, str]:Split a dataset identifier back into (pod_name, datasource_name).
Splits on the first "/" so re-joining reproduces the original string
(and therefore the same task hash) even when the datasource name itself
contains a "/".
stands_down_as_deconfigured
def stands_down_as_deconfigured( cache_db_path: str, pod_name: str, datasource_name: str,) ‑> bool:Whether a metadata run for this datasource must not start at all.
Checked before the run writes any bookkeeping: a stand-down recorded as a
completed run would read as a fresh full pass, holding the datasource off
the nightly refresh and ranking it last. See is_deconfigured.
Never raises; a cache that cannot be opened reads as not removed.
Arguments
cache_db_path: The pod's background cache.pod_name: The pod.datasource_name: The datasource.
stands_down_for_partial_inventory
def stands_down_for_partial_inventory(cache_db_path: str, task_hash: str) ‑> bool:Whether a reader of the inventory must wait for a later full pass.
Never raises; a cache that cannot be opened reads as not partial.
Arguments
cache_db_path: The pod's background cache.task_hash: The datasource-level task hash.