tasks
Prefect tasks for the scan_metadata runtime.
collect_scan_metadata
Read the datasource's file inventory from file_metadata (never walks the
directory — file_metadata_runtime already paid the walk), keep only files
whose extension the datasource would load, skip files whose content is
unchanged since the last extraction (incremental, via the persisted
file_metadata hash), and for the rest call
process_file(skip_non_tabular_data=True) and normalize each returned
series into a ScanMetadataRecord. Files that cannot be processed get a
single reason-string row (series_index=0, null metadata), consistent with
the BIT-8339 gcc/cst reason-string standard. Rows are persisted in bounded
chunks as the inventory is processed, so the task returns a count rather
than a list.
store_scan_metadata
Persist one chunk, replacing the chunk's files' stored rows in a single
transaction.
Module
Functions
collect_scan_metadata
def collect_scan_metadata( datasource: Any, pod_name: str, datasource_name: str, cache: CacheProtocol, run_id: str | None = None, chunk_size: int = 500, only_paths: Collection[str] | None = None, stats: IndexingStats | None = None, datasource_factory: Callable[[], Any] | None = None, probe: RootProbe | None = None,) ‑> int:Extract and persist per-scan metadata for every changed file.
Chunked, crash-resumable: rather than accumulating every record and
writing once at the end, records are persisted in bounded chunks as the
inventory is processed, so a kill mid-parse keeps the scans of every file
completed so far. Chunks are flushed only at file boundaries — a
file's series rows are always persisted together — because a partially
written file behind a "stored" source_file_hash would be skipped (as
unchanged) on resume, stranding the missing series. Each flush replaces
its files' old rows in the same transaction as it writes the new ones, so
a kill before a flush leaves those files' old rows, which no longer match
the inventory and are re-extracted on resume.
A missing file or an OSError while parsing is checked against the
datasource root. If the root is unreachable the pass stops there: that
file and every one not yet parsed get no row, so a later pass extracts
them, and the failure is counted in stats.transient_failures. If the
root is reachable the failure is the file's own: a missing file gets no
row (the next index prunes it) and an OSError gets a
processing_failed row, as any other parse error does.
Arguments
datasource: A constructed ophthalmologyFileSystemIterableSource, used only toprocess_file(never to walk).pod_name: Name of the pod (for the file_metadata lookup and provenance).datasource_name: Name of the datasource (same).cache: An open cache instance for the pod's background cache DB.run_id: Optional UUID of theRunthat produced these rows; stamped on each persisted row for provenance.chunk_size: Persist once the pending buffer reaches this many rows (checked at file boundaries).only_paths: When set, restrict extraction to these inventory paths, read in scope (get_file_metadata_for_pathsplusget_file_metadata_under_dirs) rather than by loading the whole inventory. May name directories as well as files: a directory contributes every inventory row beneath it, so the file watcher can hand over a new folder without enumerating it. Any spelling is accepted, since each path is normalized before lookup. A path with no inventory row is skipped — extraction never walks, so a path the index has not reached yet has nothing to parse.stats: When set, per-file parse time, reuse counts and per-chunk flush figures are recorded on it, and a progress event is emitted per chunk.datasource_factory: Builds another datasource like datasource. When set, files are parsedscan_metadata_io_concurrencyat a time, each pool thread on an instance of its own, since a datasource records per-file state that is not safe to share. Without it every file is parsed serially by datasource.probe: Answers whether the datasource root is reachable. Built for the datasource's root when not given.
Returns
The total number of rows persisted. 0 when nothing changed or when
file_metadata has not populated yet (scan_metadata never walks; it is
chained after file_metadata).
store_scan_metadata
def store_scan_metadata( file_paths: list[str], records: list[ScanMetadataRecord], cache: CacheProtocol,) ‑> int:Replace the stored rows of file_paths with records.
Arguments
file_paths: The files this chunk re-extracted. Their existing rows are deleted in the same transaction that writes records.records: The chunk's rows, already stamped with theirrun_id.cache: An open cache instance for the pod's background cache DB.
Returns The number of rows written.