606 lines
29 KiB
Python
606 lines
29 KiB
Python
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()
|