737 lines
33 KiB
Python
737 lines
33 KiB
Python
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()
|