Files
truf-server/app/jsonl_projector.py
2026-09-30 20:30:56 +03:00

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()