import sys sys.dont_write_bytecode = True if not sys.dont_write_bytecode: raise RuntimeError('result ingester could not disable bytecode writes') import argparse import json import logging import os import time from lifecycle_authority import require_active_supervisor_child from paths import apply_path_config from process_identity import current_process_identity, exact_process_identity_state from result_bundle import ( ResultBundleError, ResultBundleReader, bundle_partial_relative_path, ) from runtime_security import ( PrivatePathState, durable_publish, durable_unlink, ensure_private_directory, private_file_ready, inspect_private_relative_path, require_private_directory, sha256_file, ) from scanner_db import ( DockerCoverageDispositionConflictError, DockerFindingAttributionLimitError, ScanEventConflictError, ScannerDB, ) logger = logging.getLogger(__name__) class ResultIngester: def __init__( self, db, bundle_root, supervisor_instance_id, lease_seconds=300, fault=None, quarantine_max_items=10000, quarantine_max_bytes=1024 * 1024 * 1024, metadata_retention_days=30, metadata_retirement_batch=100, recover_expired_ready=False, ): self.db = db self.bundle_root = require_private_directory(bundle_root, create=False) self.supervisor_instance_id = str(supervisor_instance_id) self.lease_seconds = max(30, int(lease_seconds)) self.fault = fault self.lease = None self.recovery_after_id = 0 self.quarantine_max_items = max(0, int(quarantine_max_items)) self.quarantine_max_bytes = max(0, int(quarantine_max_bytes)) self.metadata_retention_seconds = max(1, int(metadata_retention_days)) * 86400 self.metadata_retirement_batch = min(500, max(1, int(metadata_retirement_batch))) self.next_metadata_retirement = 0.0 self.recover_expired_ready = recover_expired_ready is True def _inject(self, stage, value=None): if self.fault is not None: self.fault(stage, value) def start(self): self.db.require_runtime_safety_schema() self.db.require_final_cutover() identity = current_process_identity() self.lease = self.db.acquire_pipeline_lease( 'result_ingester', self.supervisor_instance_id, identity, lease_seconds=self.lease_seconds, initial_state='recovering', ) if not self.lease: raise RuntimeError('another result ingester owns the singleton advisory lock') self.reconcile_terminal_artifacts() self.retire_terminal_metadata() self.recover() if not self.heartbeat('ready'): raise RuntimeError('result ingester ready lease publication failed') return self def heartbeat(self, state='ready', error=''): return self.db.heartbeat_pipeline_lease( 'result_ingester', self.lease['generation'], self.lease['lease_token'], lease_seconds=self.lease_seconds, state=state, error=error, ) def stop(self, error=''): if self.lease: released = self.db.release_pipeline_lease( 'result_ingester', self.lease['generation'], self.lease['lease_token'], state='failed' if error else 'released', error=error, ) self.lease = None return released return True def _path(self, relative): normalized = str(relative or '').replace('/', os.sep) path = os.path.abspath(os.path.join(self.bundle_root, normalized)) if os.path.commonpath((self.bundle_root, path)) != self.bundle_root or path == self.bundle_root: raise ResultBundleError('bundle database path escapes its configured root') return path def _quarantine_path(self, reservation): bundle_id = str(reservation['bundle_id']) return os.path.join( self.bundle_root, 'quarantine', bundle_id[:2], f'{bundle_id}.trb', ) def _quarantine_relative_path(self, reservation): bundle_id = str(reservation['bundle_id']) return f'quarantine/{bundle_id[:2]}/{bundle_id}.trb' def _ensure_quarantine_shard(self, reservation): bundle_id = str(reservation['bundle_id']) return ensure_private_directory( os.path.join(self.bundle_root, 'quarantine', bundle_id[:2]), reject_reparse=True, ) def _inspect(self, relative_path): return inspect_private_relative_path(self.bundle_root, relative_path) def _defer_reservation_cleanup(self, reservation, error): try: self.db.defer_result_reservation_cleanup(reservation['id'], str(error)) except Exception: logger.warning( 'reservation cleanup backoff could not be recorded: %s', reservation['id'] ) def reconcile_terminal_artifacts(self, max_pages=100): if not hasattr(self.db, 'bundle_terminal_temp_artifacts'): return for _ in range(max(1, int(max_pages))): rows = self.db.bundle_terminal_temp_artifacts(100) if not rows: return progressed = False for row in rows: try: inspection = self._inspect(row['relative_path']) if inspection.state == PrivatePathState.UNKNOWN: self.db.defer_pipeline_artifact_cleanup( row['id'], inspection.detail or 'artifact storage state is unknown', ) continue if inspection.state == PrivatePathState.PRESENT: if not private_file_ready(inspection.path): self.db.defer_pipeline_artifact_cleanup( row['id'], 'artifact is not an exact private file', ) continue durable_unlink(inspection.path) inspection = self._inspect(row['relative_path']) if inspection.state == PrivatePathState.ABSENT: self.db.mark_pipeline_artifact_deleted(row['id']) progressed = True else: self.db.defer_pipeline_artifact_cleanup( row['id'], 'artifact unlink was not confirmed', ) except OSError as exc: try: self.db.defer_pipeline_artifact_cleanup(row['id'], str(exc)) except Exception: logger.warning( 'artifact cleanup backoff could not be recorded: %s', row['id'] ) continue if not progressed: return if len(rows) < 100: return return def retire_terminal_metadata(self): if time.monotonic() < self.next_metadata_retirement: return self.next_metadata_retirement = time.monotonic() + 60 try: self.db.retire_admission_intents( self.metadata_retention_seconds, self.metadata_retirement_batch, ) self.db.retire_deleted_pipeline_artifacts( self.metadata_retention_seconds, self.metadata_retirement_batch, ) except Exception as exc: logger.warning('bounded terminal metadata retirement deferred: %s', type(exc).__name__) def quarantine(self, reservation, ready_path, reason_code, detail): ready_relative = str(reservation['ready_relative_path']).replace('\\', '/') quarantine_relative = self._quarantine_relative_path(reservation) self._ensure_quarantine_shard(reservation) ready = self._inspect(ready_relative) quarantine = self._inspect(quarantine_relative) if ready.state == PrivatePathState.UNKNOWN or quarantine.state == PrivatePathState.UNKNOWN: raise ResultBundleError('bundle quarantine path state is unknown') if ready.state == PrivatePathState.PRESENT and quarantine.state == PrivatePathState.PRESENT: raise ResultBundleError('both ready and quarantine paths exist for one reservation') if ready.state == PrivatePathState.ABSENT and quarantine.state == PrivatePathState.ABSENT: raise ResultBundleError('bundle disappeared before quarantine') quarantine_path = quarantine.path ensure_private_directory(os.path.dirname(quarantine_path), reject_reparse=True) source_path = ready.path if ready.state == PrivatePathState.PRESENT else quarantine.path if not private_file_ready(source_path): raise ResultBundleError('bundle quarantine source is not an exact private file') byte_count = os.path.getsize(source_path) payload_hash = sha256_file(source_path) if byte_count else '' relative = quarantine_relative quarantine_id = self.db.quarantine_result_bundle( reservation['id'], reason_code, detail, relative, payload_sha256=payload_hash, byte_count=byte_count, quarantine_max_items=self.quarantine_max_items, quarantine_max_bytes=self.quarantine_max_bytes, physical_confirmed=ready.state != PrivatePathState.PRESENT, ) if ready.state == PrivatePathState.PRESENT: durable_publish(ready.path, quarantine_path) confirmed = self._inspect(quarantine_relative) if confirmed.state != PrivatePathState.PRESENT: raise ResultBundleError('bundle quarantine publication was not confirmed') quarantine_id = self.db.quarantine_result_bundle( reservation['id'], reason_code, detail, relative, payload_sha256=payload_hash, byte_count=byte_count, quarantine_max_items=self.quarantine_max_items, quarantine_max_bytes=self.quarantine_max_bytes, physical_confirmed=True, ) return quarantine_id def recover(self, page_size=100, max_pages=100): pages = 0 while pages < max(1, int(max_pages)): rows = self.db.active_result_reservations(self.recovery_after_id, page_size) if not rows: self.recovery_after_id = 0 return pages pages += 1 for reservation in rows: self.recovery_after_id = int(reservation['id']) if ( reservation['state'] == 'scanning' and str(reservation.get('assignment_kind') or 'local') == 'remote' ): continue ready_relative = str(reservation['ready_relative_path']).replace('\\', '/') ready = self._inspect(ready_relative) identity_state = None if reservation['state'] == 'scanning': identity_state = exact_process_identity_state( reservation['producer_pid'], reservation['producer_creation_time'], reservation['producer_executable'], ) if ( ready.state == PrivatePathState.UNKNOWN and identity_state in ('dead', 'reused') ): try: ensure_private_directory( os.path.join( self.bundle_root, 'ready', str(reservation['bundle_id'])[:2], ), reject_reparse=True, ) ensure_private_directory( os.path.join( self.bundle_root, 'tmp', str(reservation['bundle_id'])[:2], ), reject_reparse=True, ) ready = self._inspect(ready_relative) except OSError as exc: self._defer_reservation_cleanup(reservation, exc) continue quarantine_relative = self._quarantine_relative_path(reservation) try: self._ensure_quarantine_shard(reservation) except OSError as exc: self._defer_reservation_cleanup(reservation, exc) continue quarantine = self._inspect(quarantine_relative) if quarantine.state == PrivatePathState.PRESENT: try: self.db.quarantine_result_bundle( reservation['id'], reservation.get('last_error_code') or 'recovered_quarantine', reservation.get('last_error_detail') or 'recovered deterministic quarantine file', quarantine_relative, payload_sha256=sha256_file(quarantine.path), byte_count=os.path.getsize(quarantine.path), quarantine_max_items=self.quarantine_max_items, quarantine_max_bytes=self.quarantine_max_bytes, physical_confirmed=True, ) except OSError as exc: self._defer_reservation_cleanup(reservation, exc) continue if quarantine.state == PrivatePathState.UNKNOWN: self._defer_reservation_cleanup( reservation, quarantine.detail or 'quarantine state unknown', ) continue if ready.state == PrivatePathState.UNKNOWN: self._defer_reservation_cleanup( reservation, ready.detail or 'ready state unknown', ) continue prepared_quarantine = self.db.pending_result_bundle_quarantine( reservation['id'] ) if prepared_quarantine and ready.state == PrivatePathState.PRESENT: try: self.quarantine( reservation, ready.path, prepared_quarantine['reason_code'], prepared_quarantine.get('reason_detail') or 'recovered prepared quarantine', ) except OSError as exc: self._defer_reservation_cleanup(reservation, exc) continue if reservation['state'] == 'scanning': if ready.state == PrivatePathState.PRESENT: try: metadata = ResultBundleReader(ready.path).validate() recovered = metadata.as_dict() recovered['relative_path'] = str( reservation['ready_relative_path'] ).replace('\\', '/') marked_ready = self.db.mark_result_bundle_ready( reservation['id'], recovered, ) if ( not marked_ready and self.recover_expired_ready and identity_state in ('dead', 'reused') ): self.db.recover_expired_result_bundle_ready( reservation['id'], recovered, ) except OSError as exc: self._defer_reservation_cleanup(reservation, exc) except (ValueError, ResultBundleError) as exc: try: self.quarantine( reservation, ready.path, 'bundle_validation_failed', str(exc), ) except OSError as cleanup_exc: self._defer_reservation_cleanup(reservation, cleanup_exc) continue if identity_state in ('dead', 'reused'): try: ensure_private_directory( os.path.join(self.bundle_root, 'ready', str(reservation['bundle_id'])[:2]), reject_reparse=True, ) ensure_private_directory( os.path.join(self.bundle_root, 'tmp', str(reservation['bundle_id'])[:2]), reject_reparse=True, ) except OSError as exc: self._defer_reservation_cleanup(reservation, exc) continue ready = self._inspect(ready_relative) if ready.state != PrivatePathState.ABSENT: continue if identity_state in ('dead', 'reused'): partial_relative = bundle_partial_relative_path( reservation['bundle_id'], reservation['reservation_token'], ).replace(os.sep, '/') partial = self._inspect(partial_relative) if partial.state == PrivatePathState.UNKNOWN: self._defer_reservation_cleanup( reservation, partial.detail or 'partial state unknown', ) continue if partial.state == PrivatePathState.PRESENT: if not private_file_ready(partial.path): continue try: durable_unlink(partial.path) except OSError as exc: self._defer_reservation_cleanup(reservation, exc) continue partial = self._inspect(partial_relative) if partial.state != PrivatePathState.ABSENT: self._defer_reservation_cleanup( reservation, 'partial unlink was not confirmed', ) continue self.db.refund_uncommitted_reservation( reservation['id'], { 'pid': reservation['producer_pid'], 'creation_time': reservation['producer_creation_time'], 'executable': reservation['producer_executable'], }, f'producer identity is {identity_state} and exact ready path is absent', partial_absence_confirmed=True, ) elif reservation['state'] in ('ready', 'ingesting'): if ready.state == PrivatePathState.ABSENT: self.db.quarantine_result_bundle( reservation['id'], 'ready_bundle_missing', 'database ready row has a definitively absent exact ready path', '', byte_count=0, quarantine_max_items=self.quarantine_max_items, quarantine_max_bytes=self.quarantine_max_bytes, physical_confirmed=True, ) elif reservation['state'] == 'db_committed': bundle = self.db.result_bundle_for_reservation(reservation['id']) if not bundle: continue if ready.state == PrivatePathState.PRESENT: if not private_file_ready(ready.path): continue try: durable_unlink(ready.path) except OSError as exc: self._defer_reservation_cleanup(reservation, exc) continue ready = self._inspect(ready_relative) if ready.state == PrivatePathState.ABSENT: event = self.db.confirm_scan_event( reservation['scan_event_id'], bundle['scan_event_hash'], ) if event: self.db.acknowledge_removed_bundle( reservation['id'], reservation['scan_event_id'], bundle['scan_event_hash'], ) return pages def process_one(self): claimed = self.db.claim_ready_result_bundle( self.lease['generation'], self.lease['lease_token'], self.lease_seconds, ) if not claimed: return False reservation = claimed['reservation'] bundle = claimed['bundle'] ready_relative = str(bundle['relative_path']).replace('\\', '/') ready = self._inspect(ready_relative) try: if bundle['state'] == 'db_committed': event = self.db.confirm_scan_event(bundle['scan_event_id'], bundle['scan_event_hash']) if not event: raise ResultBundleError('db_committed bundle has no exact authoritative event') else: if ready.state == PrivatePathState.UNKNOWN: return False if ready.state == PrivatePathState.ABSENT: self.db.quarantine_result_bundle( reservation['id'], 'ready_bundle_missing', 'claimed ready bundle is definitively absent', '', quarantine_max_items=self.quarantine_max_items, quarantine_max_bytes=self.quarantine_max_bytes, physical_confirmed=True, ) return True self._inject('before_validation', reservation) reader = ResultBundleReader(ready.path) validated = reader.validate() self._inject('after_validation', validated) if ( validated.bundle_id != str(bundle['bundle_id']) or validated.scan_event_id != str(bundle['scan_event_id']) or validated.scan_event_hash != str(bundle['scan_event_hash']) or validated.actual_bytes != int(bundle['actual_bytes']) ): raise ResultBundleError('validated bundle totals conflict with its claimed database row') self._inject('before_db_commit', reservation) self.db.ingest_result_bundle(reader, reservation, bundle) self._inject('after_db_commit', reservation) event = self.db.confirm_scan_event(bundle['scan_event_id'], bundle['scan_event_hash']) self._inject('after_confirmation', event) if not event: raise ResultBundleError('database commit was not confirmed by exact event ID and hash') ready = self._inspect(ready_relative) if ready.state == PrivatePathState.UNKNOWN: return False if ready.state == PrivatePathState.PRESENT: self._inject('before_unlink', reservation) durable_unlink(ready.path) self._inject('after_unlink', reservation) ready = self._inspect(ready_relative) if ready.state != PrivatePathState.ABSENT: return False self._inject('before_capacity_release', reservation) if not self.db.acknowledge_removed_bundle( reservation['id'], bundle['scan_event_id'], bundle['scan_event_hash'], ): raise RuntimeError('bundle capacity acknowledgement was not fenced') self._inject('after_capacity_release', reservation) return True except DockerCoverageDispositionConflictError as exc: ready = self._inspect(ready_relative) if ready.state == PrivatePathState.PRESENT: try: self.quarantine( reservation, ready.path, 'docker_coverage_disposition_conflict', str(exc), ) except OSError as cleanup_exc: self._defer_reservation_cleanup(reservation, cleanup_exc) return False return True except DockerFindingAttributionLimitError as exc: ready = self._inspect(ready_relative) if ready.state == PrivatePathState.PRESENT: try: self.quarantine( reservation, ready.path, 'docker_attribution_limit_exceeded', str(exc), ) except OSError as cleanup_exc: self._defer_reservation_cleanup(reservation, cleanup_exc) return False return True except ScanEventConflictError as exc: ready = self._inspect(ready_relative) if ready.state == PrivatePathState.PRESENT: try: self.quarantine(reservation, ready.path, 'scan_event_hash_conflict', str(exc)) except OSError as cleanup_exc: self._defer_reservation_cleanup(reservation, cleanup_exc) return False return True except (ResultBundleError, ValueError) as exc: ready = self._inspect(ready_relative) if ready.state == PrivatePathState.PRESENT: try: self.quarantine(reservation, ready.path, 'bundle_validation_failed', str(exc)) except OSError as cleanup_exc: self._defer_reservation_cleanup(reservation, cleanup_exc) return False return True raise except OSError as exc: self._defer_reservation_cleanup(reservation, exc) return False def parse_args(): parser = argparse.ArgumentParser(description='Singleton durable result bundle ingester') parser.add_argument('--config', required=True) return parser.parse_args() def main(): metadata = require_active_supervisor_child(child_kind='result-ingester', require_dsn=True) args = parse_args() import yaml with open(args.config, 'r', encoding='utf-8') as handle: config = apply_path_config(yaml.safe_load(handle) or {}, args.config) global_config = config.get('global') or {} db = ScannerDB(db_url=global_config['database_url'], initialize=False) if not db.enabled: raise SystemExit('result ingester PostgreSQL connection is unavailable') db.set_application_name('truf-result-ingester') worker = ResultIngester( db, global_config['result_bundle_dir'], metadata['instance_id'], lease_seconds=int(((config.get('supervisor') or {}).get('result_ingester') or {}).get('lease_seconds', 300)), quarantine_max_items=int(global_config.get('pipeline_quarantine_max_items', 10000)), quarantine_max_bytes=int(global_config.get('pipeline_quarantine_max_bytes', 1024 * 1024 * 1024)), metadata_retention_days=int(global_config.get('pipeline_metadata_retention_days', 30)), metadata_retirement_batch=int(global_config.get('pipeline_metadata_retirement_batch', 100)), ) error = '' try: worker.start() idle = max(0.05, float(((config.get('supervisor') or {}).get('result_ingester') or {}).get('poll_sec', 0.2))) next_heartbeat = time.monotonic() + worker.lease_seconds / 3 while True: worker.retire_terminal_metadata() worker.reconcile_terminal_artifacts(max_pages=1) worker.recover(max_pages=1) processed = worker.process_one() if time.monotonic() >= next_heartbeat: if not worker.heartbeat('ready'): raise RuntimeError('result ingester heartbeat fence was lost') next_heartbeat = time.monotonic() + worker.lease_seconds / 3 if not processed: time.sleep(idle) except KeyboardInterrupt: pass except BaseException as exc: error = f'{type(exc).__name__}: {exc}' raise finally: try: worker.stop(error) finally: db.close() if __name__ == '__main__': main()