import sys sys.dont_write_bytecode = True if not sys.dont_write_bytecode: raise RuntimeError('JSONL projector could not disable bytecode writes') import argparse import hashlib import json import os import re import time from dataclasses import dataclass from lifecycle_authority import require_active_supervisor_child from paths import apply_path_config from process_identity import current_process_identity from runtime_security import ( PrivateFileLock, PrivatePathState, durable_publish, durable_unlink, harden_private_file, inspect_private_relative_path, private_file_ready, reject_reparse_components, require_private_directory, ) from scanner_db import ScannerDB STREAM_MASKS = {'scan_results': 1, 'found_secrets': 2, 'scan_errors': 4} MAX_SERIALIZED_EVENT_BYTES = 192 * 1024 * 1024 MAX_TAIL_QUARANTINE_BYTES = MAX_SERIALIZED_EVENT_BYTES @dataclass(frozen=True) class SerializedStream: stream_name: str path: str byte_length: int payload_sha256: str record_count: int artifact_id: int = 0 class _HashedWriter: def __init__(self, handle, max_bytes): self.handle = handle self.max_bytes = int(max_bytes) self.digest = hashlib.sha256() self.bytes_written = 0 def write(self, payload): payload = payload.encode('utf-8') if isinstance(payload, str) else bytes(payload) if self.bytes_written + len(payload) > self.max_bytes: raise ValueError('projection serialization exceeds its event byte bound') self.handle.write(payload) self.digest.update(payload) self.bytes_written += len(payload) def _write_json_line(writer, value): encoder = json.JSONEncoder( ensure_ascii=False, sort_keys=True, separators=(',', ':'), default=str, ) for chunk in encoder.iterencode(value): writer.write(chunk) writer.write(b'\n') def _write_json_value(writer, value): encoder = json.JSONEncoder( ensure_ascii=False, sort_keys=True, separators=(',', ':'), default=str, ) for chunk in encoder.iterencode(value): writer.write(chunk) def _write_json_array(writer, values): writer.write(b'[') first = True for value in values: if not first: writer.write(b',') _write_json_value(writer, value) first = False writer.write(b']') def _write_scan_result(writer, header, findings, errors): values = dict(header) keys = sorted(set(values) | {'findings', 'errors'}) writer.write(b'{') for index, key in enumerate(keys): if index: writer.write(b',') _write_json_value(writer, key) writer.write(b':') if key == 'findings': _write_json_array(writer, findings) elif key == 'errors': _write_json_array(writer, errors) else: _write_json_value(writer, values[key]) writer.write(b'}\n') class JsonlProjector: def __init__( self, db, results_dir, supervisor_instance_id, lease_seconds=300, fault=None, keycheck_dir=None, quarantine_max_items=10000, quarantine_max_bytes=1024 * 1024 * 1024, artifact_tracking=True, projection_max_bytes=2 * 1024 * 1024 * 1024, ): self.db = db self.results_dir = require_private_directory(results_dir, create=False) self.keycheck_dir = require_private_directory( keycheck_dir or os.path.join(os.path.dirname(self.results_dir), 'keychecks'), 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.file_lock = None self.temp_dir = require_private_directory( os.path.join(self.results_dir, '.projection-tmp'), create=True, ) self.quarantine_dir = require_private_directory( os.path.join(self.results_dir, '.projection-quarantine'), create=True, ) self.quarantine_max_items = max(0, int(quarantine_max_items)) self.quarantine_max_bytes = max(0, int(quarantine_max_bytes)) self.projection_max_bytes = max(0, int(projection_max_bytes)) self.artifact_tracking = bool(artifact_tracking) 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() self.file_lock = PrivateFileLock( os.path.join(self.results_dir, '.jsonl-projector.lock') ).acquire() self.lease = self.db.acquire_pipeline_lease( 'jsonl_projector', self.supervisor_instance_id, current_process_identity(), lease_seconds=self.lease_seconds, initial_state='starting', ) if not self.lease: self.file_lock.release() self.file_lock = None raise RuntimeError('another JSONL projector owns the singleton advisory lock') if self.artifact_tracking: self.reconcile_terminal_temps() self.recover_rotations() if not self.heartbeat(): raise RuntimeError('JSONL projector ready lease publication failed') return self def heartbeat(self, error=''): return self.db.heartbeat_pipeline_lease( 'jsonl_projector', self.lease['generation'], self.lease['lease_token'], lease_seconds=self.lease_seconds, state='ready', error=error, ) def _rollback_database(self): connection = getattr(self.db, 'conn', None) if connection is not None: try: connection.rollback() except BaseException: pass def stop(self, error=''): try: if self.lease: self._rollback_database() try: self.db.release_pipeline_lease( 'jsonl_projector', self.lease['generation'], self.lease['lease_token'], state='failed' if error else 'released', error=error, ) except BaseException: self._rollback_database() raise self.lease = None finally: if self.file_lock: self.file_lock.release() self.file_lock = None def _results_path(self, relative, stream_name=''): root = self.keycheck_dir if str(stream_name).startswith('keycheck:') else self.results_dir path = os.path.abspath(os.path.join(root, str(relative).replace('/', os.sep))) if os.path.commonpath((root, path)) != root or path == root: raise ValueError('projection path escapes results_dir') reject_reparse_components(os.path.dirname(path)) return path @staticmethod def _segment_relative(base_relative, generation): base, extension = os.path.splitext(base_relative) return f'{base}.g{int(generation):06d}{extension}' def _create_private_empty(self, path): if os.path.exists(path): if not private_file_ready(path): raise OSError(f'projection active file is not private: {path}') return descriptor = os.open( path, os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, 'O_BINARY', 0), 0o600, ) os.close(descriptor) harden_private_file(path) def _prepared_path(self, job, stream_name): safe_stream_name = re.sub(r'[^A-Za-z0-9_.-]+', '_', stream_name) event_hash = str(job['event_hash'] or '').lower() if not re.fullmatch(r'[a-f0-9]{64}', event_hash): raise ValueError('projection job event hash is not a canonical SHA-256 identity') path = os.path.join( self.temp_dir, f'job-{int(job["id"])}-{safe_stream_name}-{event_hash}.prepared', ) relative = os.path.relpath(path, self.results_dir).replace(os.sep, '/') artifact_id = 0 if self.artifact_tracking: artifact_id = self.db.register_pipeline_artifact( 'jsonl_projector', 'prepared_stream', job['id'], stream_name, relative, state='expected', byte_count=int(job['capacity_bytes']), ) if os.path.lexists(path): if not private_file_ready(path): raise OSError(f'projection prepared file is not private: {path}') else: descriptor = os.open( path, os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, 'O_BINARY', 0), 0o600, ) os.close(descriptor) harden_private_file(path) return path, artifact_id def _delete_registered_artifact(self, path, artifact_id): relative = os.path.relpath(path, self.results_dir).replace(os.sep, '/') inspection = inspect_private_relative_path(self.results_dir, relative) if inspection.state == PrivatePathState.UNKNOWN: raise OSError('projection artifact state is unknown during cleanup') if inspection.state == PrivatePathState.PRESENT: durable_unlink(inspection.path) inspection = inspect_private_relative_path(self.results_dir, relative) if inspection.state != PrivatePathState.ABSENT: raise OSError('projection artifact unlink was not confirmed') if artifact_id: self.db.mark_pipeline_artifact_deleted(artifact_id) def reconcile_terminal_temps(self, max_pages=100): if not self.artifact_tracking or not hasattr(self.db, 'projection_terminal_temp_artifacts'): return for _ in range(max(1, int(max_pages))): rows = self.db.projection_terminal_temp_artifacts(100) if not rows: return for row in rows: path = self._results_path(row['relative_path']) self._delete_registered_artifact(path, row['id']) if len(rows) < 100: return return def recover_rotations(self): for rotation in self.db.pending_projection_rotations(100): active = self._results_path(rotation['base_relative_path'], rotation['stream_name']) segment = self._results_path(rotation['segment_relative_path'], rotation['stream_name']) require_private_directory(os.path.dirname(active), create=True) active_exists = os.path.isfile(active) segment_exists = os.path.isfile(segment) if active_exists and segment_exists: if ( os.path.getsize(active) != 0 or os.path.getsize(segment) != int(rotation['source_bytes']) ): raise RuntimeError('projection rotation has conflicting active and immutable names') elif active_exists: if os.path.getsize(active) != int(rotation['source_bytes']): raise RuntimeError('projection rotation active size changed') durable_publish(active, segment) elif not segment_exists: raise RuntimeError('projection rotation lost both exact source names') self._create_private_empty(active) if not self.db.complete_projection_rotation(rotation['id']): raise RuntimeError('projection rotation completion fence failed') def _serialize(self, job): if job['job_kind'] == 'scan_event': scan = self.db.projection_scan_header( job['target_scan_id'], max_bytes=self.projection_max_bytes, ) if not isinstance(scan, dict) or not isinstance(scan.get('result'), dict): raise ValueError('authoritative scan result cannot be reconstructed') result = scan['result'] normalized = scan['storage'] == 'normalized_v2' stream_specs = [ (name, mask) for name, mask in STREAM_MASKS.items() if int(job['required_stream_mask']) & mask ] elif job['job_kind'] == 'keycheck_event': result = self.db.keycheck_result_for_projection(job['keycheck_result_id']) if not isinstance(result, dict): raise ValueError('authoritative keycheck result cannot be reconstructed') service = str(result['service']) stream_specs = [ (name, mask) for name, mask in ( (f'keycheck:{service}:results', 8), (f'keycheck:{service}:status', 16), ) if int(job['required_stream_mask']) & mask ] else: raise ValueError(f'unsupported projection job kind: {job["job_kind"]}') streams = [] prepared_paths = [] try: for stream_name, mask in stream_specs: path, artifact_id = self._prepared_path(job, stream_name) prepared_paths.append((path, artifact_id)) count = 0 with open(path, 'wb', buffering=0) as handle: writer = _HashedWriter(handle, self.projection_max_bytes) if stream_name == 'scan_results': findings = ( self.db.iter_projection_findings(job['target_scan_id']) if normalized else iter(result.get('findings') or ()) ) errors = ( self.db.iter_projection_errors(job['target_scan_id']) if normalized else iter(result.get('errors') or ()) ) _write_scan_result(writer, result, findings, errors) count = 1 elif stream_name == 'found_secrets': findings = ( self.db.iter_projection_findings(job['target_scan_id']) if normalized else iter(result.get('findings') or ()) ) for finding in findings: _write_json_line(writer, finding) count += 1 elif stream_name == 'scan_errors': timestamp = result.get('timestamp') or '' scan_type = result.get('scan_type') or '' target = result.get('target') or '' event_id = result.get('scan_event_id') or job['event_id'] errors = ( self.db.iter_projection_errors(job['target_scan_id']) if normalized else iter(result.get('errors') or ()) ) for index, error in enumerate(errors, 1): row_id = hashlib.sha256( f'{event_id}|{index}'.encode('utf-8') ).hexdigest() writer.write( f'{row_id}\t{timestamp}\t{scan_type}\t{target}\t{error}\n' ) count += 1 elif stream_name.endswith(':results'): payload = { 'event_id': result['event_id'], 'service': result['service'], 'status': result['status'], 'status_group': result['status_group'], 'checked_at': result['checked_at'], 'key_hash': result['key_hash'], 'secret_hash': result['secret_hash'], 'key_masked': result['key_masked'], 'finding_uid': result['finding_uid'], 'detector': result['detector_name'], 'source': result['source'], 'message': result['message'], 'metadata': json.loads(result['metadata_json'] or '{}'), 'result_source': result['result_source'], } _write_json_line(writer, payload) count = 1 else: secret = result.get('credential_secret_text') or result.get('credential_secret_json') or '' message = str(result.get('message') or '').replace('\r', ' ').replace('\n', ' ')[:1000] writer.write( f'{secret}\t{result["status"]}\t{result["checked_at"]}\t{message}\n' ) count = 1 handle.flush() os.fsync(handle.fileno()) streams.append(SerializedStream( stream_name, path, writer.bytes_written, writer.digest.hexdigest(), count, artifact_id, )) if self.artifact_tracking: self.db.register_pipeline_artifact( 'jsonl_projector', 'prepared_stream', job['id'], stream_name, os.path.relpath(path, self.results_dir).replace(os.sep, '/'), state='present', payload_sha256=writer.digest.hexdigest(), byte_count=writer.bytes_written, ) return streams except BaseException: self._rollback_database() for path, artifact_id in prepared_paths: try: self._delete_registered_artifact(path, artifact_id) except BaseException: pass raise def _hash_region(self, path, offset, length): digest = hashlib.sha256() remaining = int(length) with open(path, 'rb', buffering=0) as handle: handle.seek(int(offset)) while remaining: block = handle.read(min(1024 * 1024, remaining)) if not block: raise OSError('projection append proof is truncated') digest.update(block) remaining -= len(block) return digest.hexdigest() def _quarantine_tail(self, active, offset, job, stream_name, append): size = os.path.getsize(active) length = max(0, size - int(offset)) safe_stream_name = re.sub(r'[^A-Za-z0-9_.-]+', '_', stream_name) tail_hash = self._hash_region(active, offset, length) if length else hashlib.sha256(b'').hexdigest() path = os.path.join( self.quarantine_dir, f'append-{append["id"]}-o{int(offset)}-l{length}-{tail_hash}.tail', ) evidence_error = None evidence_registered = False try: if length: if length > MAX_TAIL_QUARANTINE_BYTES: raise ValueError('projection partial tail exceeds quarantine byte bound') relative = os.path.relpath(path, self.results_dir).replace(os.sep, '/') registration = self.db.register_projection_tail_quarantine( job['id'], append['id'], stream_name, relative, tail_hash, length, self.quarantine_max_items, self.quarantine_max_bytes, ) evidence_registered = True if not os.path.lexists(path): temporary = path + '.partial' temp_relative = os.path.relpath(temporary, self.results_dir).replace(os.sep, '/') temp_artifact = self.db.register_pipeline_artifact( 'jsonl_projector', 'projection_tail_temp', job['id'], f'{append["id"]}:{int(offset)}:{length}:{tail_hash}', temp_relative, state='expected', byte_count=length, ) if os.path.lexists(temporary): if not private_file_ready(temporary): raise OSError('existing projection tail temporary is not private') self._delete_registered_artifact(temporary, temp_artifact) temp_artifact = self.db.register_pipeline_artifact( 'jsonl_projector', 'projection_tail_temp', job['id'], f'{append["id"]}:{int(offset)}:{length}:{tail_hash}', temp_relative, state='expected', byte_count=length, ) descriptor = os.open( temporary, os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, 'O_BINARY', 0), 0o600, ) try: os.close(descriptor) descriptor = None harden_private_file(temporary) with open(active, 'rb', buffering=0) as source, open(temporary, 'wb', buffering=0) as target: source.seek(int(offset)) remaining = length while remaining: block = source.read(min(1024 * 1024, remaining)) if not block: raise OSError('projection partial tail changed while quarantining') target.write(block) remaining -= len(block) target.flush() os.fsync(target.fileno()) if ( os.path.getsize(temporary) != length or self._hash_region(temporary, 0, length) != tail_hash ): raise OSError('projection tail quarantine proof failed before publication') self.db.register_pipeline_artifact( 'jsonl_projector', 'projection_tail_temp', job['id'], f'{append["id"]}:{int(offset)}:{length}:{tail_hash}', temp_relative, state='present', payload_sha256=tail_hash, byte_count=length, ) durable_publish(temporary, path) temporary_state = inspect_private_relative_path( self.results_dir, temp_relative, ) if temporary_state.state != PrivatePathState.ABSENT: raise OSError('projection tail temporary retirement was not confirmed') self.db.mark_pipeline_artifact_deleted(temp_artifact) finally: if descriptor is not None: os.close(descriptor) elif ( not private_file_ready(path) or os.path.getsize(path) != length or self._hash_region(path, 0, length) != tail_hash ): raise OSError('immutable projection tail evidence is invalid') if not self.db.confirm_projection_tail_artifact( registration['artifact_id'], tail_hash, length, ): raise RuntimeError('projection tail artifact confirmation lost its fence') except BaseException as exc: evidence_error = exc finally: with open(active, 'r+b', buffering=0) as handle: handle.truncate(int(offset)) handle.flush() os.fsync(handle.fileno()) if os.path.getsize(active) != int(offset): raise OSError('projection partial tail corrective truncation was not confirmed') if evidence_error is not None: raise evidence_error def _rotate_if_needed(self, state, payload_bytes): active = self._results_path(state['base_relative_path'], state['stream_name']) require_private_directory(os.path.dirname(active), create=True) self._create_private_empty(active) current_size = os.path.getsize(active) if current_size != int(state['committed_offset']): state = self.db.initialize_projection_stream_offset( state['stream_name'], current_size, ) if current_size != int(state['committed_offset']): raise RuntimeError('projection active size does not match its committed cursor') if not current_size or current_size + payload_bytes <= int(state['rotation_bytes']): return state segment_relative = self._segment_relative( state['base_relative_path'], state['current_generation'], ) rotation = self.db.prepare_projection_rotation( state['stream_name'], current_size, segment_relative, ) segment = self._results_path(segment_relative, state['stream_name']) self._inject('before_rotation_rename', rotation) durable_publish(active, segment) self._inject('after_rotation_rename', rotation) self._create_private_empty(active) if not self.db.complete_projection_rotation(rotation['id']): raise RuntimeError('projection rotation completion fence failed') updated = self.db.projection_stream_state(state['stream_name']) oldest = int(updated['current_generation']) - int(updated['max_generations']) if oldest >= 0: old_relative = self._segment_relative(updated['base_relative_path'], oldest) old_path = self._results_path(old_relative, updated['stream_name']) if os.path.isfile(old_path): durable_unlink(old_path) return updated def _append_stream(self, job, serialized): state = self.db.projection_stream_state(serialized.stream_name) if not state: raise ValueError(f'projection stream is absent: {serialized.stream_name}') existing = self.db.projection_append_for_job(job['id'], serialized.stream_name) if existing is None: state = self._rotate_if_needed(state, serialized.byte_length) append = self.db.prepare_projection_append( job['id'], job['lease_token'], serialized.stream_name, state['generation'], serialized.byte_length, serialized.payload_sha256, serialized.record_count, ) else: append = existing if ( int(append['byte_length']) != serialized.byte_length or str(append['payload_sha256']) != serialized.payload_sha256 or int(append['record_count']) != serialized.record_count ): raise ValueError('prepared projection append conflicts with deterministic serialization') active = self._results_path(state['base_relative_path'], serialized.stream_name) require_private_directory(os.path.dirname(active), create=True) self._create_private_empty(active) offset = int(append['byte_offset']) end = offset + int(append['byte_length']) size = os.path.getsize(active) if append['state'] == 'appended': if size < end or self._hash_region(active, offset, append['byte_length']) != append['payload_sha256']: raise ValueError('acknowledged projection append proof is invalid') return if size >= end: if size == end and self._hash_region(active, offset, append['byte_length']) == append['payload_sha256']: if not self.db.complete_projection_append(append['id'], job['id'], job['lease_token']): raise RuntimeError('projection append replay acknowledgement failed') return self._quarantine_tail(active, offset, job, serialized.stream_name, append) self._inject('after_tail_recovery', append) elif size > offset: self._quarantine_tail(active, offset, job, serialized.stream_name, append) self._inject('after_tail_recovery', append) elif size < offset: raise ValueError('projection active file is shorter than its prepared offset') self._inject('before_append', append) with open(active, 'ab', buffering=0) as target, open(serialized.path, 'rb', buffering=0) as source: while True: block = source.read(1024 * 1024) if not block: break target.write(block) target.flush() os.fsync(target.fileno()) self._inject('after_append_fsync', append) if os.path.getsize(active) != end or self._hash_region(active, offset, append['byte_length']) != append['payload_sha256']: raise RuntimeError('projection append proof failed after fsync') if not self.db.complete_projection_append(append['id'], job['id'], job['lease_token']): raise RuntimeError('projection append completion fence failed') def process_one(self): self.reconcile_terminal_temps(max_pages=1) job = self.db.claim_projection_job( self.lease['generation'], self.lease['lease_token'], self.lease_seconds, ) if not job: return False streams = [] primary_failure = False try: streams = self._serialize(job) actual_bytes = sum(item.byte_length for item in streams) expanded = self.db.expand_projection_job_capacity( job['id'], job['lease_token'], actual_bytes, self.projection_max_bytes, ) if expanded is False: return True if expanded is None: raise RuntimeError('projection capacity expansion lost its lease fence') job = expanded for serialized in streams: self._append_stream(job, serialized) if not self.db.complete_projection_job(job['id'], job['lease_token']): raise RuntimeError('projection job completion fence failed') return True except (ValueError, TypeError, UnicodeError, json.JSONDecodeError) as exc: self._rollback_database() try: self.db.quarantine_projection_job( job['id'], job['lease_token'], 'deterministic_projection_error', str(exc), quarantine_max_items=self.quarantine_max_items, quarantine_max_bytes=self.quarantine_max_bytes, ) except BaseException: primary_failure = True self._rollback_database() raise return True except BaseException: primary_failure = True self._rollback_database() raise finally: for serialized in streams: try: self._delete_registered_artifact( serialized.path, serialized.artifact_id, ) except OSError: pass except BaseException: self._rollback_database() if not primary_failure: raise def parse_args(): parser = argparse.ArgumentParser(description='Singleton PostgreSQL-backed JSONL projector') parser.add_argument('--config', required=True) return parser.parse_args() def main(): metadata = require_active_supervisor_child(child_kind='jsonl-projector', 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 {} settings = ((config.get('supervisor') or {}).get('jsonl_projector') or {}) db = ScannerDB(db_url=global_config['database_url'], initialize=False) if not db.enabled: raise SystemExit('JSONL projector PostgreSQL connection is unavailable') db.set_application_name('truf-jsonl-projector') worker = JsonlProjector( db, global_config['results_dir'], metadata['instance_id'], lease_seconds=int(settings.get('lease_seconds', 300)), keycheck_dir=global_config['keycheck_dir'], 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)), projection_max_bytes=int(global_config.get( 'projection_backlog_max_bytes', 2 * 1024 * 1024 * 1024, )), ) error = '' try: worker.start() idle = max(0.05, float(settings.get('poll_sec', 0.2))) next_heartbeat = time.monotonic() + worker.lease_seconds / 3 while True: processed = worker.process_one() if time.monotonic() >= next_heartbeat: if not worker.heartbeat(): raise RuntimeError('JSONL projector 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()