Files
truf-server/tests/container_e2e.py
T
2026-09-30 20:30:56 +03:00

3516 lines
168 KiB
Python

"""Offline, disposable-volume container E2E driver. Never run against user data.
Install this file in the TEST image at /opt/truf/tests/container_e2e.py. All
commands use `python -I -S -B /opt/truf/tests/container_e2e.py MODE`. The image
must run as 10001:10001, with a fresh provisioned /data named volume, read-only
application code, network_mode: none, and /run/truf as a 0700 UID-10001 tmpfs.
Integration order (the provisioner creates directories, not a database):
1. prepare: copy the full Linux profile, commit the known-fake Git fixture,
and check the native TruffleHog Git command. No database is opened.
2. Start the REAL entrypoint: python -I -S -B app/container_runtime.py run
--config /data/config/e2e.yaml. `run` performs real first initialization.
Explicit `initialize --config /data/config/e2e.yaml` before run is also OK.
3. assert-pipeline --timeout 180: wait for the real scanner/ingester/projector
and save private artifact identities in /data/fixture/result.json.
4. keycheck-fixture --timeout 180 (optional): run the actual OpenAI provider
through child_bootstrap, replacing ONLY requests.Session.request.
5. Recreate the container with the SAME named volume, config, and isolation.
Do not prepare again. assert-persisted waits for the new one-shot GitHub
pass and requires unchanged output counts, identities, offsets and bytes.
6. If step 4 ran, keycheck-fixture again must make zero HTTP requests and
produce no duplicate keycheck results or projections.
7. assert-remote-recovery proves cross-device quota serialization, fixed
expiry, recovery, and stale fencing without contacting a source.
8. prepare-remote-transport streams one canonical synthetic bundle through
the Worker API, reconciles publication/ready failures, and waits for the
real ingester and projector to durably consume and remove it.
9. prepare-dockerhub-canary proves a bounded synthetic DockerHub discovery,
protocol-2 result, expiry, replay, and drain flow using test-only provider
transport substitution.
10. Recreate the container again. assert-remote-transport-replay and
finish-dockerhub-canary prove receipts survive restart, former expiry,
and spool cleanup.
Only /data/fixture and the new /data/config/e2e.yaml are written by this driver;
the real application owns its database and output writes. stdout is one JSON
object containing counts and hashes, including on failure. No child output,
credentials, findings, control tokens, or exception messages are printed.
"""
import argparse
from contextlib import redirect_stderr, redirect_stdout
import copy
import hashlib
import json
import logging
import math
import os
from pathlib import Path
import re
import runpy
import socket
import stat
import subprocess
import sys
import tempfile
import time
from urllib.parse import quote
APP = Path(__file__).resolve().parents[1] / 'app'
DATA = Path('/data')
FIXTURE = DATA / 'fixture'
CONFIG = DATA / 'config/e2e.yaml'
BASE_CONFIG = APP / 'config.linux.yaml'
RUNTIME = DATA / 'runtime-linux'
PG_DATA = DATA / 'postgres-linux'
BUNDLES = DATA / 'scanner-result-bundles'
CONTROL = Path('/run/truf/control')
REPO = FIXTURE / 'repo'
TARGET_FILE = RUNTIME / 'queues/e2e-targets.txt'
TARGET = 'https://gitlab.com/truf-e2e/container-pipeline-fixture.git'
CORE_SOURCES = ('gitlab', 'dockerhub', 'huggingface')
REMOTE_TARGET = 'https://gitlab.com/truf-e2e/container-fixture.git'
FIXTURE_REPOSITORY = 'file:///data/fixture/repo'
PREPARED = FIXTURE / 'prepared.json'
RESULT = FIXTURE / 'result.json'
REMOTE_TRANSPORT = FIXTURE / 'remote-transport.json'
REMOTE_FULL = FIXTURE / 'remote-full.json'
DOCKERHUB_CANARY = FIXTURE / 'dockerhub-canary.json'
REMOTE_CLIENT_BUNDLES = FIXTURE / 'remote-client-bundles'
REMOTE_LOCAL_BUNDLES = FIXTURE / 'remote-local-bundles'
MAX_BYTES = 4 * 1024 * 1024
WORKERS = ('result-ingester', 'jsonl-projector', 'janitor')
DOCKERHUB_CANARY_QUERY = 'truf-e2e-dockerhub-canary'
DOCKERHUB_CANARY_DISCOVERY_TOKEN = 'test-only-dockerhub-discovery-token'
DOCKERHUB_CANARY_REPOSITORIES = tuple(
'truf-e2e/canary-' + name for name in ('permanent', 'retryable', 'success', 'pending')
)
DOCKERHUB_CANARY_TARGETS = tuple(
repository + '@sha256:' + marker * 64
for repository, marker in zip(DOCKERHUB_CANARY_REPOSITORIES, 'abcd')
)
DOCKERHUB_CANARY_PER_PAGE = len(DOCKERHUB_CANARY_REPOSITORIES)
BASE_ENVIRONMENT = {
'PATH': '/usr/local/bin:/usr/bin:/bin',
'LANG': 'C.UTF-8', 'LC_ALL': 'C.UTF-8',
'PYTHONDONTWRITEBYTECODE': '1', 'PYTHONNOUSERSITE': '1', 'PYTHONPATH': '',
'HTTP_PROXY': '', 'http_proxy': '', 'HTTPS_PROXY': '', 'https_proxy': '',
'ALL_PROXY': '', 'all_proxy': '', 'NO_PROXY': '*', 'no_proxy': '*',
}
KEYCHECK_CHILD_ENVIRONMENT = frozenset({
'SCANNER_SUPERVISED', 'TRUF_SUPERVISOR_INSTANCE_FILE',
'TRUF_SUPERVISOR_INSTANCE_ID', 'TRUF_SUPERVISOR_TOKEN',
'TRUF_SUPERVISOR_CONFIG_SHA256', 'TRUF_SUPERVISOR_SHA256',
'TRUF_SUPERVISOR_CODE_MANIFEST_SHA256', 'TRUF_SUPERVISOR_DSN_SHA256',
'TRUF_SUPERVISOR_CHILD_KIND', 'TRUF_MANAGED_POSTGRES_DSN',
'SCANNER_DB_URL', 'DATABASE_URL', 'KEYCHECK_DB_URL', 'KEYCHECK_SERVICE',
'KEYCHECK_INPUT_MODE', 'KEYCHECK_OUTPUT_DIR', 'KEYCHECK_STATE_DIR',
'KEYCHECK_PROVIDER_SLICE_KEYS', 'KEYCHECK_DB_INLINE', 'KEYCHECK_PROXY_FILE',
'KEYCHECK_RESULT_PROJECTION_RESERVE_BYTES', 'KEYCHECK_PROJECTION_MAX_ITEMS',
'KEYCHECK_PROJECTION_MAX_BYTES',
})
class E2EFailure(AssertionError):
"""Only fixed, credential-free check names may be used as messages."""
class NotReady(E2EFailure):
pass
def require(condition, check):
if not condition:
raise E2EFailure(check)
def neutralize_environment(mode):
preserved = {
name: os.environ[name]
for name in KEYCHECK_CHILD_ENVIRONMENT
if mode == '_keycheck-child' and name in os.environ
}
os.environ.clear()
os.environ.update(BASE_ENVIRONMENT)
os.environ.update(preserved)
tempfile.tempdir = None
def json_bytes(value):
return json.dumps(
value, ensure_ascii=False, sort_keys=True, separators=(',', ':'),
).encode('utf-8')
def digest(value):
return hashlib.sha256(value).hexdigest()
def synthetic_openai_token():
# Identical known-fake canary to tests/test_synthetic_llm_pipeline.py.
alphabet = 'abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789'
def material(label):
seed = hashlib.sha512(label.encode('ascii')).hexdigest()
return ''.join(
alphabet[int(seed[index:index + 2], 16) % len(alphabet)]
for index in range(0, len(seed), 2)
)[:24]
return (
'sk-proj-' + material('local-openai-canary-prefix') + 'T3BlbkFJ'
+ material('local-openai-canary-suffix')
)
def fixture_config(base):
config = copy.deepcopy(base)
global_config = config['global']
global_config.update({
'root_dir': '/opt/truf', 'project_dir': '/opt/truf/app',
'runtime_dir': RUNTIME.as_posix(), 'postgres_data_dir': PG_DATA.as_posix(),
'postgres_bin_dir': '/usr/lib/postgresql/16/bin',
'result_bundle_dir': BUNDLES.as_posix(), 'control_dir': CONTROL.as_posix(),
'secrets_file': '/data/config/secrets.yaml',
'work_dir': '/data/scanner-work', 'database_url': '',
'dashboard_db_url': '', 'api_proxy_enabled': False,
'download_proxy_enabled': False, 'sync_file_queues': False,
'max_active_scans': 1, 'opportunistic_scan_slots': 0,
'min_free_gb': 0,
'result_bundle_min_free_bytes': 0, 'loop': False,
'detectors': 'OpenAI', 'exclude_detectors': '', 'drop_detectors': '',
'no_verification': True, 'trufflehog_config': '',
})
for name, source in config['sources'].items():
source['enabled'] = name in CORE_SOURCES
if 'auth_pool' in source:
source['auth_pool'] = ''
config['sources']['gitlab'].update({
'enabled': True, 'mode': 'custom', 'target_file': TARGET_FILE.as_posix(),
'queries': ['e2e'], 'query_overrides': {}, 'workers': 1,
'timeout': 60, 'max_commit_age_days': 0,
'exact_git_planning_enabled': False, 'updated_target_rescan_enabled': False,
'external_trufflehog_lifecycle': True, 'auth_pool': '',
})
supervisor = config['supervisor']
supervisor.update({
'enabled_sources': list(CORE_SOURCES), 'interactive': False, 'autostart': True,
'postgres_stable_ready_sec': 2, 'postgres_health_interval_sec': 1,
'interval': 3600,
})
supervisor.setdefault('defaults', {}).update({
'enabled': False, 'once': True, 'repeat': False, 'restart': False,
})
supervisor['sources']['gitlab'].update({
'enabled': True, 'once': True,
'env': {
'GOMAXPROCS': '1', 'GIT_CONFIG_NOSYSTEM': '1',
'GIT_CONFIG_GLOBAL': '/dev/null', 'GIT_ALLOW_PROTOCOL': 'file',
'GIT_CONFIG_COUNT': '2',
'GIT_CONFIG_KEY_0': 'url.' + FIXTURE_REPOSITORY + '.insteadOf',
'GIT_CONFIG_VALUE_0': TARGET,
'GIT_CONFIG_KEY_1': 'url.' + FIXTURE_REPOSITORY + '.insteadOf',
'GIT_CONFIG_VALUE_1': REMOTE_TARGET,
},
})
for name in ('result_ingester', 'jsonl_projector', 'janitor'):
supervisor[name]['enabled'] = True
supervisor['docker_shadow']['enabled'] = False
supervisor['dashboard']['enabled'] = False
config['keychecks'].update({'enabled': False, 'autostart': False})
return config
def require_container(config_path):
require(sys.platform == 'linux', 'linux_test_container_required')
require(os.getuid() == os.geteuid() == 10001 and os.getgid() == 10001,
'uid_10001_required')
require(APP == Path('/opt/truf/app'), 'test_image_path_required')
require(Path(config_path) == CONFIG, 'fixture_config_path_required')
require(sys.flags.isolated and sys.flags.no_site and sys.flags.dont_write_bytecode,
'isolated_python_required')
require({name for _, name in socket.if_nameindex()} == {'lo'},
'network_none_required')
require(not os.access(APP, os.W_OK), 'readonly_application_required')
run = CONTROL.parent.stat(follow_symlinks=False)
require(stat.S_ISDIR(run.st_mode) and stat.S_IMODE(run.st_mode) == 0o700
and run.st_uid == 10001, 'private_run_directory_required')
with open('/proc/self/mountinfo', 'rb') as handle:
mounts = handle.read(1024 * 1024 + 1)
require(len(mounts) <= 1024 * 1024, 'mountinfo_byte_bound')
require(any(
len(fields) > 6 and fields[4] == b'/run/truf' and b'-' in fields
and fields[fields.index(b'-') + 1] == b'tmpfs'
for fields in (line.split() for line in mounts.splitlines())
), 'run_tmpfs_required')
os.umask(0o077)
def read_bytes(path, limit=MAX_BYTES, private=True):
from runtime_security import reject_reparse_components, require_private_file
reject_reparse_components(str(path))
if private:
require_private_file(str(path))
details = path.stat(follow_symlinks=False)
require(stat.S_ISREG(details.st_mode) and details.st_size <= limit,
'bounded_regular_file_required')
with path.open('rb') as handle:
value = handle.read(limit + 1)
require(len(value) <= limit, 'file_byte_bound')
return value
def write_new(path, payload):
from runtime_security import fsync_directory, require_private_directory
require_private_directory(str(path.parent), create=False)
descriptor = os.open(path, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600)
with os.fdopen(descriptor, 'wb') as handle:
handle.write(payload)
handle.flush()
os.fsync(handle.fileno())
fsync_directory(str(path.parent))
def fixture_environment():
# Git must not read user/global configuration, hooks, or credentials.
return {
'PATH': '/usr/local/bin:/usr/bin:/bin', 'HOME': str(FIXTURE),
'TMPDIR': str(FIXTURE / 'tmp'), 'LANG': 'C.UTF-8', 'GOMAXPROCS': '1',
'GIT_CONFIG_NOSYSTEM': '1', 'GIT_CONFIG_GLOBAL': '/dev/null',
'GIT_TERMINAL_PROMPT': '0', 'GIT_ALLOW_PROTOCOL': 'file',
'GIT_OPTIONAL_LOCKS': '0', 'NO_PROXY': '*',
'GIT_CONFIG_COUNT': '2',
'GIT_CONFIG_KEY_0': 'url.' + FIXTURE_REPOSITORY + '.insteadOf',
'GIT_CONFIG_VALUE_0': TARGET,
'GIT_CONFIG_KEY_1': 'url.' + FIXTURE_REPOSITORY + '.insteadOf',
'GIT_CONFIG_VALUE_1': REMOTE_TARGET,
}
def git(*arguments):
command = [
'/usr/bin/git', '-c', 'user.name=Container E2E Fixture',
'-c', 'user.email=fixture@example.invalid', '-c', 'commit.gpgsign=false',
'-c', 'core.hooksPath=/dev/null', *arguments,
]
env = fixture_environment()
env.update({
'GIT_AUTHOR_DATE': '2026-01-01T00:00:00+00:00',
'GIT_COMMITTER_DATE': '2026-01-01T00:00:00+00:00',
})
completed = subprocess.run(
command, cwd=REPO, env=env, stdin=subprocess.DEVNULL,
stdout=subprocess.PIPE, stderr=subprocess.PIPE, timeout=15, check=False,
)
require(completed.returncode == 0, 'fixture_git_failed')
require(len(completed.stdout) <= 65536, 'fixture_git_output_bound')
return completed.stdout.strip()
def native_fixture_scan(config):
from runtime_security import require_trusted_native_executable
command = [
require_trusted_native_executable(config['global']['trufflehog_path']),
'git', TARGET, '--json', '--no-update', '--local-dev',
'--include-detectors', 'OpenAI', '--no-verification', '--concurrency', '1',
]
# Capture even the native self-check in private storage, never the terminal.
with tempfile.TemporaryFile(dir=FIXTURE) as output, tempfile.TemporaryFile(dir=FIXTURE) as errors:
completed = subprocess.run(
command, cwd=FIXTURE, env=fixture_environment(), stdin=subprocess.DEVNULL,
stdout=output, stderr=errors, timeout=60, check=False,
)
require(completed.returncode == 0, 'native_scan_exit')
output.seek(0)
errors.seek(0)
raw = output.read(MAX_BYTES + 1)
diagnostics = errors.read(MAX_BYTES + 1)
require(max(len(raw), len(diagnostics)) <= MAX_BYTES, 'native_scan_output_bound')
findings = [json.loads(line) for line in raw.splitlines() if line.strip()]
require(len(findings) == 1, 'native_scan_finding_count')
require(findings[0].get('DetectorName') == 'OpenAI'
and findings[0].get('Raw') == synthetic_openai_token()
and findings[0].get('Verified') is False, 'native_scan_fixture_identity')
finished = 0
for line in diagnostics.splitlines():
try:
value = json.loads(line)
except ValueError:
continue
if isinstance(value, dict) and value.get('msg') == 'finished scanning':
finished += 1
require(finished == 1, 'native_scan_completion_marker')
return {'counts': {'native_findings': 1, 'native_finished': 1},
'hashes': {'native_command_sha256': digest(json_bytes(command))}}
def require_fresh_volume():
from runtime_security import require_private_directory
# Only the provisioner's exact empty proxy file may exist in the fresh tree.
proxy = RUNTIME / 'proxy.txt'
entries = 0
for root in (RUNTIME, PG_DATA, BUNDLES, DATA / 'scanner-work', FIXTURE):
if not root.exists():
continue
require_private_directory(str(root), create=False)
for current, directories, files in os.walk(root, followlinks=False):
entries += len(directories) + len(files)
require(entries <= 512, 'fresh_volume_required')
for name in files:
path = Path(current) / name
require(path == proxy, 'fresh_volume_required')
details = path.lstat()
require(stat.S_ISREG(details.st_mode)
and stat.S_IMODE(details.st_mode) == 0o600
and details.st_uid == details.st_gid == 10001
and details.st_size == 0 and details.st_nlink == 1,
'fresh_volume_required')
for name in directories:
require(Path(current) / name != proxy, 'fresh_volume_required')
require_private_directory(os.path.join(current, name), create=False)
def prepare():
import yaml
from runtime_security import require_private_directory, write_private_json_exclusive
require_private_directory(str(DATA), create=False)
require_private_directory(str(CONFIG.parent), create=False)
require(not os.path.lexists(CONFIG) and not os.path.lexists(PREPARED),
'fixture_already_prepared')
require(not os.path.lexists(PG_DATA / 'PG_VERSION'), 'fresh_volume_required')
require_fresh_volume()
require_private_directory(str(FIXTURE), create=True)
require_private_directory(str(FIXTURE / 'tmp'), create=True)
require_private_directory(str(REPO), create=True)
config = fixture_config(yaml.safe_load(read_bytes(BASE_CONFIG, private=False)))
payload = ('OPENAI_API_KEY=' + synthetic_openai_token() + '\n').encode('ascii')
git('init', '--quiet', '--initial-branch=main', '--template=')
write_new(REPO / 'synthetic.env', payload)
git('add', '--', 'synthetic.env')
git('commit', '--quiet', '-m', 'Known-fake offline OpenAI fixture')
commit = git('rev-parse', 'HEAD').decode('ascii')
require(re.fullmatch(r'[a-f0-9]{40}', commit), 'fixture_git_commit')
write_new(TARGET_FILE, (TARGET + '\n').encode('ascii'))
native = native_fixture_scan(config)
config_bytes = yaml.safe_dump(config, sort_keys=False).encode('utf-8')
write_new(CONFIG, config_bytes)
marker = {
'schema': 1, 'repo_commit': commit, 'config_sha256': digest(config_bytes),
'fixture_sha256': digest(payload), 'native': native,
}
write_private_json_exclusive(str(PREPARED), marker)
return {
'counts': {'ok': 1, 'prepared': 1, **native['counts']},
'hashes': {'config_sha256': marker['config_sha256'], **native['hashes']},
}
def load_fixture():
import yaml
marker = json.loads(read_bytes(PREPARED))
require(marker.get('schema') == 1, 'prepared_fixture_required')
config_bytes = read_bytes(CONFIG)
require(digest(config_bytes) == marker['config_sha256'], 'fixture_config_changed')
config = yaml.safe_load(config_bytes)
require(config == fixture_config(yaml.safe_load(read_bytes(BASE_CONFIG, private=False))),
'fixture_config_contract_changed')
payload = ('OPENAI_API_KEY=' + synthetic_openai_token() + '\n').encode('ascii')
require(read_bytes(REPO / 'synthetic.env') == payload
and marker['fixture_sha256'] == digest(payload), 'fixture_material_changed')
require(read_bytes(TARGET_FILE) == (TARGET + '\n').encode('ascii'),
'fixture_target_changed')
require(git('rev-parse', 'HEAD').decode('ascii') == marker['repo_commit']
and git('rev-list', '--all', '--count') == b'1'
and git('status', '--porcelain', '--untracked-files=all') == b'',
'fixture_repository_changed')
return config, marker
def database_environment():
import yaml
password = read_bytes(DATA / 'postgres-password', limit=4096).decode('utf-8').rstrip('\n')
require(len(password) >= 32 and not any(char.isspace() for char in password),
'generated_postgres_password_required')
require(yaml.safe_load(read_bytes(DATA / 'config/secrets.yaml')) == {},
'empty_provider_secrets_required')
url = 'postgresql://truf:' + quote(password, safe='') + '@127.0.0.1:5432/truf'
expected = {
'TRUF_POSTGRES_DB': 'truf', 'TRUF_POSTGRES_USER': 'truf',
'TRUF_POSTGRES_PORT': '5432', 'TRUF_POSTGRES_PASSWORD': password,
'SCANNER_DB_URL': url, 'DATABASE_URL': url, 'TRUF_MANAGED_POSTGRES_DSN': url,
}
# docker exec does not inherit the environment derived inside PID 1.
for name, value in expected.items():
require(name not in os.environ or os.environ[name] == value,
'conflicting_database_environment')
require(not any(name.upper().startswith('PG') for name in os.environ),
'libpq_overrides_forbidden')
os.environ.update(expected)
os.environ.update({
'TRUF_DB_CONNECT_TIMEOUT_SEC': '2', 'TRUF_DB_STATEMENT_TIMEOUT_MS': '3000',
'TRUF_DB_LOCK_TIMEOUT_MS': '1000', 'TRUF_DB_IDLE_TRANSACTION_TIMEOUT_MS': '10000',
'TRUF_DB_TCP_USER_TIMEOUT_MS': '3000',
})
return url
def control_snapshot(marker, url):
from supervisor import get_control_snapshot
from supervisor_instance import load_instance_metadata
from runtime_security import require_private_directory
require_private_directory(str(CONTROL.parent), create=False)
try:
metadata = load_instance_metadata(str(CONTROL / 'supervisor.instance.json'))
snapshot = get_control_snapshot(metadata)
except (OSError, ValueError, RuntimeError):
raise NotReady('authenticated_control_not_ready') from None
require(metadata['config_path'] == str(CONFIG)
and metadata['config_sha256'] == marker['config_sha256']
and metadata['canonical_dsn_sha256'] == digest(url.encode('utf-8')),
'supervisor_authority_mismatch')
require(snapshot.get('activation_state') not in ('STOPPING', 'FAILED_HOLD'),
'supervisor_failed_hold')
if snapshot.get('activation_state') != 'ACTIVE':
raise NotReady('supervisor_not_active')
postgres = snapshot.get('postgres') or {}
if postgres.get('state') != 'READY' or postgres.get('ready') is not True:
raise NotReady('postgres_not_ready')
signatures = snapshot.get('signature') or []
require(all(isinstance(row, (list, tuple)) and len(row) == 9 for row in signatures),
'worker_signature_shape')
workers = {row[0]: row for row in signatures}
require(len(workers) == len(signatures)
and set(workers) == {'gitlab', *WORKERS}, 'unexpected_worker_set')
for name, row in workers.items():
require(row[1] != 'failed' and not row[6] and not row[8],
'worker_failed_or_restarted')
if name in WORKERS:
require(not row[5], 'pipeline_worker_restarted')
if tuple(row[1:3]) != ('running', 'running') or not row[3] or row[7]:
raise NotReady('pipeline_workers_not_running')
gitlab = workers['gitlab']
if gitlab[1] != 'waiting':
raise NotReady('gitlab_discovery_cycle_not_finished')
require(gitlab[2] == 'running' and gitlab[3] is None
and gitlab[4] == 0 and gitlab[5] == 1 and not gitlab[7],
'gitlab_discovery_cycle_failed')
require((snapshot.get('dashboard') or {}).get('status') == 'disabled',
'dashboard_must_be_disabled')
return metadata
def rows(db, statement, parameters=()):
require(statement.lstrip().startswith('SELECT '), 'read_only_assertions_required')
result = [dict(row) for row in db.conn.execute(statement, parameters).fetchall()]
db.conn.commit()
return result
async def invoke_bundle_upload(
service, token, reservation_id, body, *, content_length=None,
payload_sha256=None, disconnect_after=None,
):
from types import SimpleNamespace
import worker_api
class Request:
def __init__(self):
self.headers = {
'authorization': 'Bearer ' + token,
'content-type': 'application/octet-stream',
'content-length': str(len(body) if content_length is None else content_length),
'x-truf-payload-sha256': payload_sha256 or digest(body),
}
self.path_params = {'reservation_id': str(reservation_id)}
self.app = SimpleNamespace(state=SimpleNamespace(worker_service=service))
async def stream(self):
sent = 0
while sent < len(body):
chunk = body[sent:sent + 257]
sent += len(chunk)
yield chunk
if disconnect_after is not None and sent >= disconnect_after:
raise ConnectionError('synthetic disconnect')
try:
response = await worker_api.upload_bundle(Request())
except worker_api.WorkerAPIError as exc:
return exc.status_code, exc.code, None
except worker_api.ResultBundleError:
return 400, 'invalid_bundle', None
except worker_api.ScanEventConflictError:
return 409, 'reservation_conflict', None
return response.status_code, None, json.loads(bytes(response.body))
async def invoke_worker_claim(service, token, request_id, build):
from types import SimpleNamespace
import worker_api
body = json_bytes({'request_id': request_id, 'build': build})
class Request:
headers = {
'authorization': 'Bearer ' + token,
'content-type': 'application/json',
}
app = SimpleNamespace(state=SimpleNamespace(worker_service=service))
async def stream(self):
for offset in range(0, len(body), 257):
yield body[offset:offset + 257]
response = await worker_api.claim(Request())
body = bytes(response.body)
return response.status_code, json.loads(body) if body else None
async def invoke_terminal_report(service, token, reservation_id, payload):
from types import SimpleNamespace
import worker_api
body = json_bytes(payload)
class Request:
headers = {
'authorization': 'Bearer ' + token,
'content-type': 'application/json',
}
path_params = {'reservation_id': str(reservation_id)}
app = SimpleNamespace(state=SimpleNamespace(worker_service=service))
async def stream(self):
yield body
try:
response = await worker_api.terminal_report(Request())
except worker_api.WorkerAPIError as exc:
return exc.status_code, exc.code, None
except worker_api.ScanEventConflictError:
return 409, 'reservation_conflict', None
return response.status_code, None, json.loads(bytes(response.body))
def remote_device_token(role):
require(role in ('expired', 'winner'), 'remote_fixture_role')
return 'container-e2e-' + role + '-device-token'
def remote_source_args(config):
from types import SimpleNamespace
source = config['sources']['gitlab']
return SimpleNamespace(
platform='gitlab', exact_git_planning_enabled=True,
workers=1, timeout=int(source['timeout']), save_dir=str(FIXTURE),
detectors='OpenAI', exclude_detectors='', drop_detectors=[],
no_verification=True,
trufflehog_config=str(APP / 'trufflehog-custom-detectors.yaml'), token='',
external_trufflehog_lifecycle=True,
scan_full_history=False, max_depth=0, git_baseline_depth=1,
max_commit_age_days=0, commit_lookup_pages=1,
skip_if_commit_lookup_fails=True,
result_bundle_max_event_bytes=MAX_BYTES,
result_bundle_max_items=4, result_bundle_max_total_bytes=4 * MAX_BYTES,
projection_backlog_max_items=8, projection_backlog_max_bytes=8 * MAX_BYTES,
projection_backlog_headroom_bytes=2 * MAX_BYTES,
keycheck_queue_max_items=8, keycheck_queue_max_bytes=4 * MAX_BYTES,
pipeline_quarantine_max_items=4, pipeline_quarantine_max_bytes=4 * MAX_BYTES,
keycheck_candidates_per_event=8, keycheck_candidate_bytes_per_event=MAX_BYTES,
target_retry_max_attempts=3, target_retry_base_delay_sec=60,
target_retry_max_delay_sec=600, target_timeout_retry_delay_sec=300,
max_active_scans=1, admission_resolution_attempts=2,
admission_resolution_seconds=2, admission_resolution_retry_delay_sec=0.01,
target_claim_order='oldest', git_ref_resolution_attempts=1,
git_ref_resolution_timeout_sec=1, git_ref_resolution_max_bytes=1 << 20,
strict_git_provider_token_filter=True, trufflehog_stdout_max_mb=2,
trufflehog_stderr_max_mb=1, trufflehog_max_findings_per_target=100,
trufflehog_job_memory_limit_bytes=0, trufflehog_windows_job_cpu_weight=0,
trufflehog_windows_memory_priority=0, trufflehog_diagnostic_max_lines=200,
trufflehog_diagnostic_max_line_chars=2048,
trufflehog_diagnostic_max_line_bytes=2048,
trufflehog_diagnostic_max_errors=20, trufflehog_diagnostic_max_warnings=20,
trufflehog_diagnostic_max_unclassified=10,
)
def remote_package_manifest(metadata, capabilities=None):
from lifecycle_authority import (
GIT_MANIFEST_NAME, REMOTE_WORKER_CODE_AUTHORITY_FILES,
TRUFFLEHOG_MANIFEST_NAME,
)
from result_bundle import FORMAT_VERSION
from scan_execution import PROTOCOL_VERSION, local_platform_tag
authority = metadata['code_manifest']
require(set(REMOTE_WORKER_CODE_AUTHORITY_FILES) <= set(authority['files']),
'remote_code_authority_fixture')
policy = APP / 'trufflehog-custom-detectors.yaml'
policy_sha256 = digest(read_bytes(policy, private=False))
files = {
name: {
'path': 'app/' + name,
'sha256': authority['files'][name]['sha256'],
}
for name in REMOTE_WORKER_CODE_AUTHORITY_FILES
}
executables = authority['executables']
capabilities = capabilities or [{
'source': 'gitlab', 'platform': 'gitlab',
'planning_kind': 'exact_git_v1',
}]
return {
'schema': 3, 'protocol_version': PROTOCOL_VERSION,
'bundle_format_version': FORMAT_VERSION,
'platform_tag': local_platform_tag(),
'capabilities': [dict(value) for value in capabilities],
'app_root': 'app', 'files': files,
'executables': {
TRUFFLEHOG_MANIFEST_NAME: {
'path': 'bin/trufflehog',
'sha256': executables[TRUFFLEHOG_MANIFEST_NAME]['sha256'],
},
GIT_MANIFEST_NAME: {
'path': 'runtime/git/bin/git',
'sha256': executables[GIT_MANIFEST_NAME]['sha256'],
},
},
'assets': {'detector_policy': {
'path': 'app/trufflehog-custom-detectors.yaml',
'sha256': policy_sha256,
}},
'runtime_trees': {'git': {
'path': 'runtime/git',
'sha256': digest(json_bytes({
'git': executables[GIT_MANIFEST_NAME]['sha256'],
})),
'file_count': 1,
}},
}
def remote_worker_service(config, marker, metadata, url):
import scanner_db
from worker_api import WorkerService
from worker_assignment import RemoteGitAssignmentBuilder
from worker_package import worker_package_build_compatibility
def planner(args, db_url, source, claim, scan_kwargs, remote_credential=None):
del args, scan_kwargs
require(source == 'gitlab' and claim['target'] == REMOTE_TARGET,
'remote_fixture_plan_target')
planning = scanner_db.ScannerDB(db_url=db_url, initialize=False)
try:
resolution = {
'provider': 'gitlab', 'repo_url': REMOTE_TARGET,
'repo_path': 'truf-e2e/container-fixture', 'branch': 'main',
'ref': 'refs/heads/main', 'head_sha': marker['repo_commit'],
'ref_source': 'provider_default',
}
return planning.bind_git_scan_plan(
claim['reservation_id'], claim['claim_lease_token'], resolution, 1,
remote_credential=remote_credential,
)
finally:
planning.close()
package = remote_package_manifest(metadata)
builder = RemoteGitAssignmentBuilder(
url, str(BUNDLES), {'gitlab': remote_source_args(config)},
{'container-e2e': {'package_manifest': package, 'sources': ['gitlab']}},
metadata['instance_id'], assignment_ttl_seconds=60, planner=planner,
)
return (
WorkerService(url, str(BUNDLES), builder, max_bundle_bytes=MAX_BYTES),
worker_package_build_compatibility(package),
)
def execute_fixture_assignment(config, metadata, assignment, *, compare_local):
from console_runner import apply_global_config
from lifecycle_authority import build_code_manifest, code_manifest_sha256
from paths import apply_path_config
from result_bundle import BundleReservation
from runtime_security import ensure_private_directory
from scan_execution import (
PACKAGE_DETECTOR_POLICY, execute_planned_claim, stage_scan_result_in_scope,
)
import scanner
stage = 'reservation'
reservation = BundleReservation.from_mapping(assignment['reservation'])
require(reservation.target == REMOTE_TARGET and assignment['scan_kwargs'].get(
'trufflehog_config') == PACKAGE_DETECTOR_POLICY, 'remote_execution_assignment')
stage = 'directories'
ensure_private_directory(str(REMOTE_CLIENT_BUNDLES), reject_reparse=True)
if compare_local:
ensure_private_directory(str(REMOTE_LOCAL_BUNDLES), reject_reparse=True)
stage = 'config'
apply_global_config(apply_path_config(config, str(CONFIG))['global'])
scanner.initialize_scanner_runtime(preflight_complete=True, register_cleanup=False)
scan_kwargs = dict(assignment['scan_kwargs'])
policy_path = str(APP / 'trufflehog-custom-detectors.yaml')
scan_kwargs['trufflehog_config'] = policy_path
client_manifest = build_code_manifest(
str(APP), scanner.scan_config.trufflehog_path, (policy_path,),
git_path=metadata['code_manifest']['executables']['git']['path'],
)
limits = assignment['limits']
environment = fixture_environment()
previous = {name: os.environ.get(name) for name in environment}
try:
os.environ.update(environment)
stage = 'authority'
with scanner.client_scan_launch_authority(
client_manifest, code_manifest_sha256(client_manifest),
):
local = None
if compare_local:
stage = 'local_scan'
with scanner.client_scan_execution_policy(assignment['scan_policy']):
result = scanner.scan_target_result(
reservation.target, reservation.platform,
reservation.scan_event_id, scan_kwargs,
)
stage = 'local_bundle'
local = stage_scan_result_in_scope(
result, reservation, str(REMOTE_LOCAL_BUNDLES),
assignment['event_scan_options'], assignment['queue_policy'],
attempts=int(assignment['reservation']['attempts']),
candidate_max_items=int(limits['candidate_max_items']),
candidate_max_bytes=int(limits['candidate_max_bytes']),
)
stage = 'remote_scan_bundle'
remote = execute_planned_claim(
reservation, str(REMOTE_CLIENT_BUNDLES), scan_kwargs,
assignment['event_scan_options'], assignment['queue_policy'],
assignment['scan_policy'],
attempts=int(assignment['reservation']['attempts']),
candidate_max_items=int(limits['candidate_max_items']),
candidate_max_bytes=int(limits['candidate_max_bytes']),
)
except E2EFailure:
raise
except Exception as exc:
trace = exc.__traceback__
origins = []
while trace:
origins.append(trace.tb_frame.f_code.co_name)
trace = trace.tb_next
origin = '_'.join([type(exc).__name__, *origins[-3:]])
origin = re.sub(r'[^a-z0-9_]', '_', origin.lower())[:80]
raise E2EFailure('remote_execution_' + stage + '_' + origin) from None
finally:
for name, value in previous.items():
if value is None:
os.environ.pop(name, None)
else:
os.environ[name] = value
stage = 'bundle_read'
try:
remote_path = REMOTE_CLIENT_BUNDLES / remote.relative_path
payload = read_bytes(remote_path)
except E2EFailure:
raise
except Exception:
raise E2EFailure('remote_execution_' + stage + '_exception') from None
parity_sha256 = ''
if compare_local:
stage = 'parity'
try:
helpers = runpy.run_path(str(Path(__file__).with_name('parity_helpers.py')))
local_evidence = helpers['normalized_bundle_evidence'](
str(REMOTE_LOCAL_BUNDLES / local.relative_path),
)
remote_evidence = helpers['normalized_bundle_evidence'](str(remote_path))
except Exception:
raise E2EFailure('remote_execution_' + stage + '_exception') from None
require(helpers['bundle_evidence_difference_paths'](
local_evidence, remote_evidence,
) == [], 'remote_local_bundle_parity')
require(local.queue_status == remote.queue_status == 'done'
and not local.first_error and not remote.first_error,
'remote_local_disposition_parity')
parity_sha256 = digest(json_bytes(remote_evidence))
return payload, remote.as_dict(), parity_sha256
def verify_projection_region(payload, expected, append, state, job, records):
offset = int(append['byte_offset'])
end = offset + len(expected)
require(0 <= offset <= end <= len(payload)
and payload[offset:end] == expected, 'projection_region_bytes_differ')
require(append['state'] == 'appended' and append['job_id'] == job['id']
and append['stream_name'] == state['stream_name']
and append['event_id'] == job['event_id']
and append['event_hash'] == job['event_hash'], 'projection_region_identity')
require(append['byte_length'] == len(expected)
and append['payload_sha256'] == digest(expected)
and append['record_count'] == records, 'projection_region_length')
require(append['generation'] == state['generation'] == state['current_generation'],
'projection_region_generation')
require(state['committed_offset'] == len(payload)
and state['committed_offset'] >= end,
'projection_region_committed_size')
require(state['last_append_id'] >= append['id'],
'projection_region_cursor_coverage')
def verify_projection(payload, expected, append, state, job, records):
require(payload == expected, 'projection_bytes_differ')
require(append['state'] == 'appended'
and append['job_id'] == job['id']
and append['stream_name'] == state['stream_name']
and append['event_id'] == job['event_id']
and append['event_hash'] == job['event_hash'], 'projection_append_identity')
require(append['byte_offset'] == 0 and append['byte_length'] == len(payload)
and append['payload_sha256'] == digest(payload)
and append['record_count'] == records, 'projection_append_bytes')
require(append['generation'] == state['generation'] == state['current_generation'] == 0
and state['committed_offset'] == len(payload)
and state['last_append_id'] == append['id']
and state['last_job_id'] == job['id']
and state['last_event_id'] == job['event_id']
and state['last_event_hash'] == job['event_hash'], 'projection_cursor_identity')
def pipeline_snapshot(config, marker, url, checked):
import scanner_db
from keycheck_candidates import candidate_uid, extract_candidates
metadata = control_snapshot(marker, url)
db = scanner_db.ScannerDB(db_url=url, initialize=False)
if not db.enabled:
raise NotReady('database_not_ready')
try:
require(db.conn.is_postgres, 'postgres_required')
db.set_application_name('truf-container-e2e:read-only')
db.conn.execute('SET default_transaction_read_only = on')
db.conn.commit()
identity = rows(db, '''SELECT current_database() AS database, current_user AS username,
current_setting('data_directory') AS data_directory,
current_setting('server_version_num') AS server_version''')[0]
require(identity['database'] == identity['username'] == 'truf'
and identity['data_directory'] == str(PG_DATA)
and int(identity['server_version']) // 10000 == 16, 'postgres_identity')
for worker in ('result_ingester', 'jsonl_projector'):
if not db.pipeline_worker_health(worker, metadata['instance_id'])['healthy']:
raise NotReady('pipeline_lease_not_ready')
db.require_runtime_safety_schema()
db.require_final_cutover()
migrations = rows(db, 'SELECT version, code_sha256 FROM runtime_schema_migrations ORDER BY version')
require(len(migrations) == len(scanner_db.PIPELINE_MIGRATION_VERSIONS) == 33
and {row['version'] for row in migrations} == set(scanner_db.PIPELINE_MIGRATION_VERSIONS),
'fresh_migration_versions')
sql_hash = digest(scanner_db.PIPELINE_SCHEMA_SQL.encode('utf-8'))
require(next(row['code_sha256'] for row in migrations
if row['version'] == scanner_db.PIPELINE_MIGRATION_VERSIONS[-1]) == sql_hash,
'latest_migration_sql_hash')
cutover = db.final_cutover_status()
require(cutover is not None, 'valid_final_cutover_required')
expected_counts = {
'target_queue': 1, 'result_reservations': 1, 'result_bundles': 1,
'target_scans': 1, 'scan_result_compat': 1, 'findings': 1,
'finding_compat_payloads': 1, 'finding_uid_map': 1,
'keycheck_candidates': 1, 'keycheck_credentials': 1,
'projection_jobs': 1 + checked, 'projection_appends': 2 + checked,
'keycheck_results': checked, 'keycheck_current_state': checked,
'keycheck_event_map': checked, 'errors': 0, 'pipeline_quarantine': 0,
'projection_append_audit': 0, 'projection_rotations': 0,
'scan_publication_outbox': 0,
}
counts = {
name: int(rows(db, f'SELECT COUNT(*) AS count FROM {name}')[0]['count'])
for name in expected_counts
}
for name, expected in expected_counts.items():
require(counts[name] <= expected, 'unexpected_or_duplicate_' + name)
if counts != expected_counts:
raise NotReady('pipeline_output_counts_not_ready')
queue = rows(db, 'SELECT * FROM target_queue LIMIT 2')[0]
reservation = rows(db, 'SELECT * FROM result_reservations LIMIT 2')[0]
scan = rows(db, 'SELECT * FROM target_scans LIMIT 2')[0]
finding = rows(db, 'SELECT * FROM findings LIMIT 2')[0]
candidate = rows(db, 'SELECT * FROM keycheck_candidates LIMIT 2')[0]
credential = rows(db, 'SELECT * FROM keycheck_credentials LIMIT 2')[0]
bundle = db.result_bundle_for_reservation(reservation['id'])
jobs = rows(db, 'SELECT * FROM projection_jobs ORDER BY id LIMIT 3')
if (queue['status'] != 'done' or reservation['state'] != 'acknowledged'
or bundle['state'] != 'acknowledged'
or any(job['status'] != 'completed' for job in jobs)):
raise NotReady('pipeline_durable_completion_not_ready')
for row in (queue, reservation, scan, finding, candidate):
require(row['source'] == 'gitlab' and row['query'] == 'e2e'
and row['target'] == TARGET, 'pipeline_fixture_attribution')
require(queue['target_scan_id'] == scan['id']
and queue['attempts'] == 1
and queue['lease_token'] is None
and queue['current_result_reservation_id'] is None
and queue['claim_event_id'] is None and not queue['last_error'], 'queue_completion')
require(scan['queue_id'] == reservation['queue_id'] == queue['id']
and scan['result_reservation_id'] == reservation['id']
and scan['scan_event_id'] == reservation['scan_event_id'] == bundle['scan_event_id']
and scan['scan_event_hash'] == bundle['scan_event_hash']
and scan['claim_lease_token'] == reservation['claim_lease_token']
and scan['queue_completion_applied'] == 1
and scan['queue_completion_disposition'] == 'applied', 'scan_reservation_identity')
require(scan['raw_result_storage'] == 'normalized_v2' and scan['raw_result_json'] is None
and scan['status'] == 'found' and scan['scan_type'] == 'gitlab'
and scan['findings_count'] == 1 and scan['verified_findings_count'] == 0
and scan['error_count'] == 0, 'normalized_scan_result')
confirmed = db.confirm_scan_event(scan['scan_event_id'], scan['scan_event_hash'])
require(confirmed and confirmed['id'] == scan['id'], 'confirmed_scan_event')
require(bundle['target_scan_id'] == scan['id'] and bundle['format_version'] == 2
and bundle['bundle_id'] == reservation['bundle_id']
and bundle['finding_count'] == bundle['candidate_count'] == 1
and bundle['error_count'] == 0 and bundle['actual_bytes'] > 0
and reservation['bundle_credit_released'] == 1
and reservation['candidate_credit_transferred'] == 1
and reservation['projection_credit_transferred'] == 1, 'bundle_completion')
relative = Path(bundle['relative_path'])
require(not relative.is_absolute() and '..' not in relative.parts, 'bundle_path')
require(not os.path.lexists(BUNDLES / relative), 'acknowledged_bundle_not_removed')
require(finding['target_scan_id'] == scan['id'] and finding['detector_name'] == 'OpenAI'
and not finding['verified'] and finding['file_path'] == 'synthetic.env'
and finding['commit_hash'] == marker['repo_commit'], 'native_git_finding_provenance')
rebuilt = db.reconstruct_scan_result(scan['id'], max_bytes=MAX_BYTES, max_findings=2, max_errors=0)
require(rebuilt and len(rebuilt['findings']) == 1 and rebuilt['errors'] == [],
'reconstructed_scan_result')
scan_meta = rebuilt.get('scan_meta') or {}
require(scan_meta.get('trufflehog_returncode') == 0
and scan_meta.get('trufflehog_finished') is True
and scan_meta.get('command_timed_out') is False, 'native_pipeline_completion')
options = rebuilt.get('scan_options') or {}
require(options.get('no_verification') is True and options.get('detectors') == 'OpenAI'
and not options.get('exclude_detectors') and not options.get('trufflehog_config')
and not options.get('max_commit_age_days'), 'native_pipeline_scan_options')
rebuilt_finding = rebuilt['findings'][0]
require(rebuilt_finding.get('Raw') == synthetic_openai_token()
and rebuilt_finding.get('finding_uid') == finding['finding_uid'], 'synthetic_finding_identity')
specs = list(extract_candidates(rebuilt_finding))
require(len(specs) == 1 and specs[0].service == 'openai', 'candidate_extraction')
spec = specs[0]
require(candidate['finding_id'] == finding['id']
and candidate['finding_uid'] == finding['finding_uid']
and candidate['target_scan_id'] == scan['id']
and candidate['scan_event_id'] == scan['scan_event_id']
and candidate['credential_id'] == credential['id']
and candidate['service'] == candidate['routed_service'] == credential['service'] == 'openai'
and candidate['secret_hash'] == finding['secret_hash'] == spec.secret_hash
and candidate['candidate_uid'] == candidate_uid(
scan['scan_event_id'], finding['finding_uid'], 'openai', spec.credential_hash),
'candidate_exact_attribution')
for name in ('credential_hash', 'provider_key_hash', 'candidate_kind', 'secret_text', 'secret_json'):
require(credential[name] == getattr(spec, name), 'candidate_credential_identity')
require(candidate['state'] == ('completed' if checked else 'pending')
and candidate['attempts'] == checked and candidate['lease_token'] is None
and candidate['capacity_released'] == checked, 'candidate_completion_state')
compat = rows(db, 'SELECT * FROM scan_result_compat LIMIT 2')[0]
compat_bytes = compat['metadata_json'].encode('utf-8')
require(compat['metadata_sha256'] == digest(compat_bytes)
and compat['metadata_bytes'] == len(compat_bytes), 'scan_metadata_hash')
payload_row = rows(db, 'SELECT * FROM finding_compat_payloads LIMIT 2')[0]
require(payload_row['finding_id'] == finding['id'] and not payload_row['payload_omitted']
and payload_row['payload_sha256'] == finding['raw_payload_sha256'], 'finding_payload_hash')
require(rows(db, 'SELECT finding_id FROM finding_uid_map WHERE finding_uid = ?',
(finding['finding_uid'],)) == [{'finding_id': finding['id']}], 'finding_uid_mapping')
scan_job = next(job for job in jobs if job['job_kind'] == 'scan_event')
require(scan_job['target_scan_id'] == scan['id']
and scan_job['event_id'] == scan['scan_event_id']
and scan_job['event_hash'] == scan['scan_event_hash']
and scan_job['required_stream_mask'] == 3, 'scan_projection_job')
projections = [
(scan_job, 'scan_results', 'scan_results.jsonl', json_bytes(rebuilt) + b'\n'),
(scan_job, 'found_secrets', 'found_secrets.jsonl', json_bytes(rebuilt_finding) + b'\n'),
]
ids = {
'queue_id': queue['id'], 'reservation_id': reservation['id'],
'bundle_id': bundle['bundle_id'], 'scan_id': scan['id'],
'scan_event_id': scan['scan_event_id'], 'finding_id': finding['id'],
'finding_uid': finding['finding_uid'], 'candidate_id': candidate['id'],
'candidate_uid': candidate['candidate_uid'], 'credential_id': credential['id'],
'scan_projection_id': scan_job['id'],
}
hashes = {
'migrations_sha256': digest(json_bytes(migrations)), 'pipeline_sql_sha256': sql_hash,
'cutover_sha256': digest(json_bytes(cutover)),
'scan_event_sha256': scan['scan_event_hash'], 'scan_metadata_sha256': digest(compat_bytes),
'finding_payload_sha256': payload_row['payload_sha256'],
'fixture_sha256': marker['fixture_sha256'], 'config_sha256': marker['config_sha256'],
}
if checked:
result = db.keycheck_result_for_projection(candidate['keycheck_result_id'])
require(result and result['status'] == 'INVALID_OR_REVOKED'
and result['status_group'] == 'dead' and result['result_source'] == 'api_check'
and result['link_status'] == 'linked' and result['service'] == 'openai'
and result['candidate_id'] == candidate['id']
and result['credential_id'] == credential['id']
and result['finding_id'] == finding['id']
and result['finding_uid'] == finding['finding_uid']
and result['target_scan_id'] == scan['id']
and result['key_hash'] == spec.provider_key_hash
and result['secret_hash'] == spec.secret_hash, 'real_provider_completion')
current = rows(db, 'SELECT * FROM keycheck_current_state LIMIT 2')[0]
require(current['last_result_id'] == result['id'] and current['state_version'] == 1
and current['credential_id'] == credential['id']
and current['status'] == 'INVALID_OR_REVOKED'
and current['status_group'] == 'dead', 'keycheck_current_state')
require(rows(db, 'SELECT keycheck_result_id FROM keycheck_event_map WHERE event_id = ?',
(result['event_id'],)) == [{'keycheck_result_id': result['id']}],
'keycheck_event_mapping')
key_job = next(job for job in jobs if job['job_kind'] == 'keycheck_event')
require(key_job['keycheck_result_id'] == result['id']
and key_job['event_id'] == result['event_id']
and key_job['required_stream_mask'] == 8
and key_job['event_hash'] == digest(
result['metadata_json'].encode('utf-8') + result['event_id'].encode('ascii')),
'keycheck_projection_job')
projected = {key: result[key] for key in (
'event_id', 'service', 'status', 'status_group', 'checked_at', 'key_hash',
'secret_hash', 'key_masked', 'finding_uid', 'source', 'message', 'result_source',
)}
projected.update({'detector': result['detector_name'],
'metadata': json.loads(result['metadata_json'])})
projections.append((key_job, 'keycheck:openai:results', 'openai/openaiResults.jsonl',
json_bytes(projected) + b'\n'))
ids.update({'keycheck_result_id': result['id'], 'keycheck_event_id': result['event_id'],
'keycheck_projection_id': key_job['id']})
for job, stream, relative_path, expected in projections:
require(job['capacity_released'] == 1, 'projection_capacity_not_released')
state = db.projection_stream_state(stream)
append = db.projection_append_for_job(job['id'], stream)
require(state and append and state['base_relative_path'] == relative_path,
'projection_stream_path')
root = RUNTIME / ('keychecks' if stream.startswith('keycheck:') else 'results')
payload = read_bytes(root / relative_path)
verify_projection(payload, expected, append, state, job, 1)
hashes[stream + '_sha256'] = digest(payload)
hashes[stream + '_ledger_sha256'] = digest(json_bytes({
'append': append, 'cursor': state,
}))
counts[stream + '_bytes'] = len(payload)
capacity = db.pipeline_capacity_snapshot()
for name in ('bundle_items', 'bundle_bytes', 'projection_items', 'projection_bytes',
'quarantine_items', 'quarantine_bytes'):
require(capacity[name] == 0, 'pipeline_capacity_not_drained')
require(capacity['keycheck_items'] == 1 - checked
and capacity['keycheck_bytes'] == (0 if checked else candidate['capacity_bytes']),
'keycheck_capacity_accounting')
hashes['artifact_ids_sha256'] = digest(json_bytes(ids))
counts.update({'queue_attempts': queue['attempts'],
'migration_versions': len(migrations), 'healthy_pipeline_leases': 2,
'pipeline_workers': len(WORKERS), 'native_finished': 1})
return {'schema': 1, 'checked': checked, 'counts': counts, 'hashes': hashes, 'ids': ids,
'origin_instance_sha256': digest(metadata['instance_id'].encode('utf-8'))}
finally:
db.close()
def wait_for_pipeline(config, marker, url, checked, timeout):
deadline = time.monotonic() + timeout
last_check = 'pipeline_timeout'
while time.monotonic() < deadline:
try:
return pipeline_snapshot(config, marker, url, checked)
except NotReady as exc:
last_check = str(exc)
time.sleep(min(0.5, max(0, deadline - time.monotonic())))
raise E2EFailure('timeout_' + last_check)
def run_local_pipeline_fixture(config, marker, url):
from types import SimpleNamespace
import console_runner
from lifecycle_authority import supervised_child_environment
import scanner_db
metadata = control_snapshot(marker, url)
db = scanner_db.ScannerDB(db_url=url, initialize=False)
try:
queue = rows(db, 'SELECT source, status FROM target_queue LIMIT 2')
require(queue == [{'source': 'gitlab', 'status': 'pending'}],
'pipeline_fixture_queue_required')
finally:
db.close()
authority_environment = supervised_child_environment(metadata, url, 'scanner')
metadata = None
os.environ.update(fixture_environment())
os.environ.update(authority_environment)
try:
console_runner.run_config_mode(SimpleNamespace(
config=str(CONFIG), source='gitlab', once=True, show_state=False,
cleanup_only=False, cooldown=0,
), runtime_role='scanner')
finally:
for name in authority_environment:
os.environ.pop(name, None)
authority_environment = None
return {'counts': {'ok': 1, 'local_pipeline_scans': 1}, 'hashes': {}}
def compare_persisted(before, after, require_restart=False):
require(before.get('schema') == after.get('schema') == 1, 'result_baseline_schema')
for name in ('checked', 'counts', 'hashes', 'ids'):
require(before[name] == after[name], 'persisted_' + name + '_changed')
if require_restart:
require(before['origin_instance_sha256'] != after['origin_instance_sha256'],
'container_recreation_required')
def assert_remote_recovery(config, marker, url):
del config
import threading
from process_identity import current_process_identity
from result_bundle import FORMAT_VERSION
from scan_execution import (
PACKAGE_DETECTOR_POLICY, PROTOCOL_VERSION, QueueDispositionPolicy,
remote_execution_identity,
)
import scanner_db
from unittest import mock
metadata = control_snapshot(marker, url)
db = scanner_db.ScannerDB(db_url=url, initialize=False)
connections = []
try:
require(db.enabled and db.conn.is_postgres, 'remote_postgres_required')
db.set_application_name('truf-container-e2e:remote-recovery')
db.require_runtime_safety_schema()
db.require_final_cutover()
require(db.pipeline_worker_health(
'result_ingester', metadata['instance_id'],
)['healthy'], 'remote_ingester_lease_required')
require(rows(db, 'SELECT COUNT(*) AS count FROM remote_worker_users')[0]['count'] == 0,
'remote_fixture_database_not_fresh')
capacity_names = (
'bundle_items', 'bundle_bytes', 'projection_items', 'projection_bytes',
'keycheck_items', 'keycheck_bytes', 'quarantine_items', 'quarantine_bytes',
)
baseline = db.pipeline_capacity_snapshot()
baseline = {name: int(baseline[name]) for name in capacity_names}
require(not any(baseline.values()), 'remote_capacity_not_drained')
limits = {
'bundle_items': 2, 'bundle_bytes': 16384,
'projection_items': 2, 'projection_bytes': 16384,
'keycheck_items': 2, 'keycheck_bytes': 16384,
'quarantine_items': 1, 'quarantine_bytes': 16384,
}
query = 'remote-recovery'
targets = (
'https://github.com/truf-e2e/remote-race.git',
'https://github.com/truf-e2e/remote-spare.git',
)
require(db.enqueue_targets('github', 'github', query, targets) == 2,
'remote_fixture_enqueue')
user_key = 'container-e2e-remote-user'
device_keys = ('container-e2e-device-a', 'container-e2e-device-b')
token_hashes = tuple(digest(name.encode('ascii')) for name in device_keys)
devices = [
db.provision_remote_worker_device(user_key, key, token_hash, 1)
for key, token_hash in zip(device_keys, token_hashes)
]
require(devices[0]['user_id'] == devices[1]['user_id']
and devices[0]['device_id'] != devices[1]['device_id']
and all(not device['revoked'] for device in devices),
'remote_shared_user_fixture')
scan_kwargs = {
'timeout_sec': 30.0,
'trufflehog_config': PACKAGE_DETECTOR_POLICY,
}
execution_limits = {'candidate_max_items': 10, 'candidate_max_bytes': 4096}
scan_policy = {
'drop_detectors': [],
'strict_git_provider_token_filter': True,
'trufflehog_stdout_max_mb': 32,
'trufflehog_stderr_max_mb': 8,
'result_bundle_max_event_bytes': 1024 * 1024,
'trufflehog_max_findings_per_target': 20000,
'trufflehog_job_memory_limit_bytes': 0,
'trufflehog_windows_job_cpu_weight': 0,
'trufflehog_windows_memory_priority': 0,
'trufflehog_diagnostic_max_lines': 2000,
'trufflehog_diagnostic_max_line_chars': 8192,
'trufflehog_diagnostic_max_line_bytes': 8192,
'trufflehog_diagnostic_max_errors': 200,
'trufflehog_diagnostic_max_warnings': 200,
'trufflehog_diagnostic_max_unclassified': 20,
}
effective, execution = remote_execution_identity(
'github', scan_kwargs, scan_kwargs, QueueDispositionPolicy(),
execution_limits, scan_policy,
)
snapshot = {
'schema': 1,
'compatibility': {
'protocol_version': PROTOCOL_VERSION,
'bundle_format_version': FORMAT_VERSION,
'platform_tag': 'linux-x86_64',
'code_manifest_sha256': digest(b'remote-recovery-code'),
'effective_config_sha256': effective,
'detector_policy_sha256': digest(b'remote-recovery-policy'),
},
'execution': execution,
'planning': {
'kind': 'exact_git_v1', 'git_baseline_depth': 100,
'git_ref_resolution_attempts': 2,
'git_ref_resolution_timeout_sec': 10.0,
'git_ref_resolution_max_bytes': 1 << 20,
},
'credential_ref': {'source': 'github', 'auth_entry': 'primary'},
}
producer = current_process_identity()
requests = [{
'reservation_token': digest(f'remote-request-{index}'.encode('ascii')),
'bundle_id': digest(f'remote-bundle-{index}'.encode('ascii')),
'scan_event_id': digest(f'remote-event-{index}'.encode('ascii')),
'device': devices[index], 'token_sha256': token_hashes[index],
'client_compat_sha256': digest(f'remote-client-{index}'.encode('ascii')),
} for index in range(2)]
def claim(connection, request):
return connection.reserve_and_claim_target(
'github', 'github', producer, metadata['instance_id'],
4096, 4096, 0, 0, lease_seconds=90, max_attempts=3,
capacity_limits=limits,
reservation_token=request['reservation_token'],
bundle_id=request['bundle_id'],
scan_event_id=request['scan_event_id'], claim_order='oldest',
final_cutover=True, remote_assignment={
'user_id': request['device']['user_id'],
'device_id': request['device']['device_id'],
'effective_config_sha256': effective,
'client_compat_sha256': request['client_compat_sha256'],
'token_sha256': request['token_sha256'],
'result_upload_body_timeout_seconds': 1800,
'execution_snapshot': snapshot,
},
)
barrier = threading.Barrier(2)
claimed = [None, None]
failures = []
def race(index):
connection = scanner_db.ScannerDB(db_url=url, initialize=False)
connections.append(connection)
try:
connection.set_application_name(
f'truf-container-e2e:remote-race-{index}'
)
barrier.wait(timeout=10)
claimed[index] = claim(connection, requests[index])
except Exception:
failures.append(index)
finally:
connection.close()
threads = [threading.Thread(target=race, args=(index,)) for index in range(2)]
for thread in threads:
thread.start()
for thread in threads:
thread.join(30)
require(not failures and all(not thread.is_alive() for thread in threads),
'remote_quota_race_failed')
winners = [index for index, value in enumerate(claimed) if value is not None]
require(len(winners) == 1, 'remote_quota_not_atomic')
winner_index = winners[0]
loser_index = 1 - winner_index
winner = requests[winner_index]
active_claim = claimed[winner_index]
intents = rows(db, '''SELECT reservation_token, state FROM admission_intents
WHERE reservation_token IN (?, ?) ORDER BY reservation_token''',
tuple(request['reservation_token'] for request in requests))
require({row['state'] for row in intents} == {'committed', 'aborted'},
'remote_intent_race_resolution')
require(rows(db, '''SELECT COUNT(*) AS count FROM result_reservations
WHERE assignment_kind = 'remote' AND remote_user_id = ?
AND remote_resolved_at IS NULL''', (winner['device']['user_id'],))[0]['count'] == 1,
'remote_shared_quota_count')
current_capacity = db.pipeline_capacity_snapshot()
expected_delta = {
'bundle_items': 1, 'bundle_bytes': 4096,
'projection_items': 1, 'projection_bytes': 4096,
'keycheck_items': 0, 'keycheck_bytes': 0,
'quarantine_items': 0, 'quarantine_bytes': 0,
}
require(all(int(current_capacity[name]) == baseline[name] + expected_delta[name]
for name in capacity_names), 'remote_capacity_single_charge')
db.close()
db = scanner_db.ScannerDB(db_url=url, initialize=False)
db.set_application_name('truf-container-e2e:remote-reconcile')
recovered = db.reconcile_remote_assignment_request(
winner['reservation_token'], winner['device']['device_id'],
winner['token_sha256'],
)
rejected = db.reconcile_remote_assignment_request(
requests[loser_index]['reservation_token'],
requests[loser_index]['device']['device_id'],
requests[loser_index]['token_sha256'],
)
require(recovered['state'] == 'committed'
and recovered['claim'] == active_claim
and recovered['execution_snapshot'] == snapshot
and recovered['git_plan'] is None and recovered['receipt'] is None
and rejected == {'state': 'aborted'}, 'remote_lost_reply_recovery')
before_lower = rows(db, '''SELECT remote_expires_at, producer_lease_expires_at
FROM result_reservations WHERE id = ?''',
(active_claim['reservation_id'],))[0]
lowered = db.provision_remote_worker_device(
user_key, winner['device']['device_key'], winner['token_sha256'], 0,
)
authenticated = db.authenticate_remote_worker(winner['token_sha256'])
require(lowered['active_assignment_cap'] == 0
and authenticated['active_assignment_cap'] == 0,
'remote_lower_cap_not_applied')
require(claim(db, winner) == active_claim,
'remote_lower_cap_cancelled_existing_claim')
require(not db.renew_result_claim(
active_claim['reservation_id'], active_claim['claim_lease_token'], 3600,
), 'remote_claim_was_renewed')
require(rows(db, '''SELECT remote_expires_at, producer_lease_expires_at
FROM result_reservations WHERE id = ?''',
(active_claim['reservation_id'],))[0] == before_lower,
'remote_contact_changed_deadline')
resolution = {
'provider': 'github', 'repo_url': active_claim['target'],
'repo_path': 'truf-e2e/remote-race', 'branch': 'main',
'ref': 'refs/heads/main', 'head_sha': 'a' * 40,
'ref_source': 'provider_default',
}
credential = {
'device_id': winner['device']['device_id'],
'token_sha256': winner['token_sha256'],
}
plan = db.bind_git_scan_plan(
active_claim['reservation_id'], active_claim['claim_lease_token'],
resolution, 100, remote_credential=credential,
)
require(plan['mode'] == 'baseline'
and db.remote_bound_git_scan_plan(
active_claim['reservation_id'], winner['device']['device_id'],
active_claim['claim_lease_token'], winner['token_sha256'],
) == plan, 'remote_plan_recovery')
lease = rows(db, '''SELECT r.remote_issued_at, r.remote_expires_at,
r.producer_lease_expires_at, q.leased_at, q.lease_expires_at
FROM result_reservations r JOIN target_queue q ON q.id = r.queue_id
WHERE r.id = ?''', (active_claim['reservation_id'],))[0]
require(lease['remote_issued_at'] == lease['leased_at']
and lease['remote_expires_at'] == lease['producer_lease_expires_at']
== lease['lease_expires_at'], 'remote_dependent_lease_alignment')
require((scanner_db.parse_time(lease['remote_expires_at'])
- scanner_db.parse_time(lease['remote_issued_at'])).total_seconds() == 90,
'remote_fixed_lease_window')
with mock.patch.object(
scanner_db, 'utc_now_iso', return_value=lease['remote_expires_at'],
):
expired = db.reap_expired_remote_assignments(limit=10)
require(len(expired) == 1 and expired[0]['resolution'] == 'expired'
and expired[0]['reservation_id'] == active_claim['reservation_id'],
'remote_expiry_reaper')
expired_state = rows(db, '''SELECT r.state, r.remote_resolution_kind,
q.status, q.attempts, q.lease_token, q.current_result_reservation_id
FROM result_reservations r JOIN target_queue q ON q.id = r.queue_id
WHERE r.id = ?''', (active_claim['reservation_id'],))[0]
require(expired_state == {
'state': 'refunded', 'remote_resolution_kind': 'expired',
'status': 'pending', 'attempts': 0, 'lease_token': None,
'current_result_reservation_id': None,
}, 'remote_expiry_refund_state')
require(all(int(db.pipeline_capacity_snapshot()[name]) == baseline[name]
for name in capacity_names), 'remote_expiry_capacity_release')
db.provision_remote_worker_device(
user_key, winner['device']['device_key'], winner['token_sha256'], 1,
)
reclaim_request = dict(winner)
reclaim_request.update({
'reservation_token': digest(b'remote-reclaim-request'),
'bundle_id': digest(b'remote-reclaim-bundle'),
'scan_event_id': digest(b'remote-reclaim-event'),
})
reclaimed = claim(db, reclaim_request)
require(reclaimed is not None and reclaimed['queue_id'] == active_claim['queue_id']
and reclaimed['reservation_id'] != active_claim['reservation_id']
and reclaimed['claim_lease_token'] != active_claim['claim_lease_token'],
'remote_same_worker_reclaim')
try:
db.bind_git_scan_plan(
active_claim['reservation_id'], active_claim['claim_lease_token'],
resolution, 100, remote_credential=credential,
)
except scanner_db.ScanEventConflictError:
pass
else:
raise E2EFailure('remote_stale_plan_not_fenced')
try:
db.remote_bound_git_scan_plan(
active_claim['reservation_id'], winner['device']['device_id'],
active_claim['claim_lease_token'], winner['token_sha256'],
)
except scanner_db.ScanEventConflictError:
pass
else:
raise E2EFailure('remote_stale_plan_recovery_not_fenced')
try:
db.report_remote_prebundle_failure(
active_claim['reservation_id'], winner['device']['device_id'],
winner['token_sha256'], {
'failure_code': 'client_process_failed', 'detail': 'stale fixture',
},
)
except scanner_db.ScanEventConflictError:
pass
else:
raise E2EFailure('remote_stale_report_not_fenced')
current_fence = rows(db, '''SELECT status, lease_token,
current_result_reservation_id, claim_event_id
FROM target_queue WHERE id = ?''', (reclaimed['queue_id'],))[0]
require(current_fence == {
'status': 'in_progress', 'lease_token': reclaimed['claim_lease_token'],
'current_result_reservation_id': reclaimed['reservation_id'],
'claim_event_id': reclaimed['scan_event_id'],
}, 'remote_stale_response_changed_reclaim')
report = {
'failure_code': 'client_process_failed',
'detail': 'synthetic remote recovery fixture',
}
receipt = db.report_remote_prebundle_failure(
reclaimed['reservation_id'], winner['device']['device_id'],
winner['token_sha256'], report,
)
require(receipt['resolution'] == 'prebundle_report'
and db.report_remote_prebundle_failure(
reclaimed['reservation_id'], winner['device']['device_id'],
winner['token_sha256'], report,
) == receipt, 'remote_terminal_report_replay')
try:
db.report_remote_prebundle_failure(
reclaimed['reservation_id'], winner['device']['device_id'],
winner['token_sha256'], {
'failure_code': 'client_storage_failed', 'detail': 'conflict fixture',
},
)
except scanner_db.ScanEventConflictError:
pass
else:
raise E2EFailure('remote_conflicting_report_not_fenced')
require(all(int(db.pipeline_capacity_snapshot()[name]) == baseline[name]
for name in capacity_names), 'remote_terminal_capacity_release')
require(rows(db, '''SELECT COUNT(*) AS count FROM result_reservations
WHERE assignment_kind = 'remote' AND remote_user_id = ?
AND remote_resolved_at IS NULL''', (winner['device']['user_id'],))[0]['count'] == 0,
'remote_quota_not_released_once')
final_queue = rows(db, '''SELECT status, attempts, lease_token,
current_result_reservation_id FROM target_queue WHERE id = ?''',
(reclaimed['queue_id'],))[0]
require(final_queue == {
'status': 'pending', 'attempts': 0, 'lease_token': None,
'current_result_reservation_id': None,
}, 'remote_terminal_queue_release')
evidence = {
'race_devices': 2, 'admitted': 1, 'expired': 1,
'same_worker_reclaims': 1, 'terminal_replays': 1,
}
return {
'counts': {'ok': 1, **evidence},
'hashes': {'remote_recovery_sha256': digest(json_bytes(evidence))},
}
finally:
for connection in connections:
connection.close()
db.close()
def prepare_remote_transport(config, marker, url, timeout):
del config
import asyncio
from process_identity import current_process_identity
from result_bundle import bundle_partial_path, bundle_ready_path
from runtime_security import ensure_private_directory, write_private_json_exclusive
from scanner import stage_result_bundle
import scanner_db
from unittest import mock
from worker_api import WorkerService
metadata = control_snapshot(marker, url)
db = scanner_db.ScannerDB(db_url=url, initialize=False)
stage = 'fixture'
try:
db.set_application_name('truf-container-e2e:remote-transport')
capacity_names = (
'bundle_items', 'bundle_bytes', 'projection_items', 'projection_bytes',
'keycheck_items', 'keycheck_bytes', 'quarantine_items', 'quarantine_bytes',
)
baseline = {name: int(db.pipeline_capacity_snapshot()[name]) for name in capacity_names}
require(not any(baseline.values()), 'remote_transport_capacity_not_drained')
before_scans = rows(db, 'SELECT COUNT(*) AS count FROM target_scans')[0]['count']
device_key = 'container-e2e-device-a'
token_sha256 = digest(device_key.encode('ascii'))
fixture = rows(db, '''SELECT u.id AS user_id, d.id AS device_id
FROM remote_worker_users u JOIN remote_worker_devices d ON d.user_id = u.id
WHERE u.user_key = ? AND d.device_key = ?''',
('container-e2e-remote-user', device_key))
require(len(fixture) == 1, 'remote_transport_device_fixture')
device = db.provision_remote_worker_device(
'container-e2e-remote-user', device_key, token_sha256, 1,
)
require(device['user_id'] == fixture[0]['user_id']
and device['device_id'] == fixture[0]['device_id'],
'remote_transport_device_identity')
prior = rows(db, '''SELECT remote_execution_snapshot_json,
remote_effective_config_sha256
FROM result_reservations
WHERE assignment_kind = 'remote' AND remote_user_id = ?
ORDER BY id LIMIT 1''', (device['user_id'],))
require(len(prior) == 1, 'remote_transport_snapshot_fixture')
snapshot = json.loads(prior[0]['remote_execution_snapshot_json'])
effective = str(prior[0]['remote_effective_config_sha256'])
stage = 'claim'
reservation_token = digest(b'remote-transport-request')
bundle_id = digest(b'remote-transport-bundle')
scan_event_id = digest(b'remote-transport-event')
declared_bytes = 16384
limits = {
'bundle_items': 1, 'bundle_bytes': declared_bytes,
'projection_items': 1, 'projection_bytes': declared_bytes,
'keycheck_items': 1, 'keycheck_bytes': declared_bytes,
'quarantine_items': 1, 'quarantine_bytes': declared_bytes,
}
claimed = db.reserve_and_claim_target(
'github', 'github', current_process_identity(), metadata['instance_id'],
declared_bytes, declared_bytes, 0, 0, lease_seconds=300,
max_attempts=3, capacity_limits=limits,
reservation_token=reservation_token, bundle_id=bundle_id,
scan_event_id=scan_event_id, claim_order='oldest', final_cutover=True,
remote_assignment={
'user_id': device['user_id'], 'device_id': device['device_id'],
'effective_config_sha256': effective,
'client_compat_sha256': digest(b'remote-transport-client'),
'token_sha256': token_sha256,
'result_upload_body_timeout_seconds': 1800,
'execution_snapshot': snapshot,
},
)
require(claimed is not None, 'remote_transport_claim')
repo_path = str(claimed['target']).split('github.com/', 1)[-1]
if repo_path.endswith('.git'):
repo_path = repo_path[:-4]
resolution = {
'provider': 'github', 'repo_url': claimed['target'],
'repo_path': repo_path, 'branch': 'main', 'ref': 'refs/heads/main',
'head_sha': 'b' * 40, 'ref_source': 'provider_default',
}
plan = db.bind_git_scan_plan(
claimed['reservation_id'], claimed['claim_lease_token'], resolution, 100,
remote_credential={
'device_id': device['device_id'], 'token_sha256': token_sha256,
},
)
stage = 'bundle_plan'
plan_sha256 = digest(scanner_db.canonical_git_scan_plan_bytes(plan))
exact_scope = {
key: plan.get(key) for key in (
'provider', 'ref', 'head_sha', 'base_sha', 'mode',
'baseline_depth', 'ref_source',
)
}
stage = 'bundle_directory'
ensure_private_directory(str(REMOTE_CLIENT_BUNDLES), reject_reparse=True)
stage = 'bundle_stage'
staged = stage_result_bundle(
{
'scan_event_id': claimed['scan_event_id'],
'target': claimed['target'], 'scan_type': 'github',
'timestamp': '2026-09-18T00:00:00+00:00',
'scan_started_at': '2026-09-18T00:00:00+00:00',
'duration_sec': 1.0, 'findings': [], 'errors': [],
'git_scan_plan': plan,
'git_scan_execution': {
'mode': plan['mode'], 'pinned': True, 'success': True,
'continuity_reset': False, 'plan_sha256': plan_sha256,
},
'scan_meta': {'exact_git_scope': exact_scope},
},
claimed, str(REMOTE_CLIENT_BUNDLES), {},
{'queue_status': 'done', 'queue_error': None, 'available_after': None},
candidate_max_items=0, candidate_max_bytes=0,
)
stage = 'bundle_read'
client_ready = Path(bundle_ready_path(str(REMOTE_CLIENT_BUNDLES), bundle_id))
payload = read_bytes(client_ready, limit=declared_bytes)
payload_sha256 = digest(payload)
require(staged.actual_bytes == len(payload) and len(payload) < declared_bytes,
'remote_transport_client_bundle_bound')
service = WorkerService(
url, str(BUNDLES), lambda *_args: None, max_bundle_bytes=declared_bytes,
)
server_ready = bundle_ready_path(str(BUNDLES), bundle_id)
server_partial = bundle_partial_path(str(BUNDLES), bundle_id, reservation_token)
def upload(body=payload, **kwargs):
return asyncio.run(invoke_bundle_upload(
service, device_key, claimed['reservation_id'], body, **kwargs,
))
stage = 'negative_upload'
status, code, _ = upload(
payload[:-1], content_length=len(payload),
payload_sha256=digest(payload[:-1]),
)
require((status, code) == (400, 'length_mismatch'),
'remote_transport_truncation_not_rejected')
corrupted = bytearray(payload)
corrupted[len(corrupted) // 2] ^= 1
status, code, _ = upload(bytes(corrupted))
require((status, code) == (400, 'invalid_bundle'),
'remote_transport_invalid_bundle_not_rejected')
try:
upload(disconnect_after=257)
except ConnectionError:
pass
else:
raise E2EFailure('remote_transport_disconnect_not_rejected')
require(not os.path.lexists(server_ready) and not os.path.lexists(server_partial),
'remote_transport_failed_upload_not_cleaned')
stage = 'deadline'
with mock.patch.object(
scanner_db, 'utc_now_iso', return_value=claimed['remote_expires_at'],
):
status, code, _ = upload()
require((status, code) == (410, 'assignment_expired')
and not os.path.lexists(server_ready),
'remote_transport_deadline_crossing_accepted')
stage = 'publication'
with mock.patch.object(
service, 'accept_ready', side_effect=RuntimeError('synthetic ready failure'),
):
try:
upload()
except RuntimeError:
pass
else:
raise E2EFailure('remote_transport_ready_crash_not_injected')
unresolved = rows(db, '''SELECT state, remote_resolution_kind
FROM result_reservations WHERE id = ?''', (claimed['reservation_id'],))[0]
require(unresolved == {'state': 'scanning', 'remote_resolution_kind': None}
and os.path.isfile(server_ready),
'remote_transport_publication_not_recoverable')
service = WorkerService(
url, str(BUNDLES), lambda *_args: None, max_bundle_bytes=declared_bytes,
)
stage = 'ready_commit'
committed_receipts = []
accept_ready = service.accept_ready
def accept_then_crash(*args):
committed_receipts.append(accept_ready(*args))
raise RuntimeError('synthetic lost ready commit acknowledgement')
with mock.patch.object(service, 'accept_ready', side_effect=accept_then_crash):
try:
upload()
except RuntimeError:
pass
else:
raise E2EFailure('remote_transport_commit_crash_not_injected')
require(len(committed_receipts) == 1
and committed_receipts[0]['resolution'] == 'bundle_accepted',
'remote_transport_ready_commit_not_durable')
stage = 'retry'
async def simultaneous_retry():
return await asyncio.gather(*(
invoke_bundle_upload(
service, device_key, claimed['reservation_id'], payload,
payload_sha256=payload_sha256,
) for _ in range(2)
))
retries = asyncio.run(simultaneous_retry())
require(all(status == 200 and code is None for status, code, _ in retries)
and retries[0][2] == retries[1][2] == committed_receipts[0],
'remote_transport_simultaneous_retry')
receipt = retries[0][2]
require(receipt['resolution'] == 'bundle_accepted'
and receipt['payload_sha256'] == payload_sha256,
'remote_transport_receipt_identity')
status, code, _ = asyncio.run(invoke_bundle_upload(
service, device_key, claimed['reservation_id'], payload + b'x',
))
require((status, code) == (409, 'resolution_conflict'),
'remote_transport_conflicting_replay')
stage = 'ingestion'
deadline = time.monotonic() + timeout
complete = None
while time.monotonic() < deadline:
state = rows(db, '''SELECT r.state, r.remote_resolution_kind,
r.remote_resolution_json, b.state AS bundle_state,
q.status AS queue_status, q.target_scan_id
FROM result_reservations r
JOIN result_bundles b ON b.reservation_id = r.id
JOIN target_queue q ON q.id = r.queue_id
WHERE r.id = ?''', (claimed['reservation_id'],))
capacity = db.pipeline_capacity_snapshot()
if (
len(state) == 1 and state[0]['state'] == 'acknowledged'
and state[0]['bundle_state'] == 'acknowledged'
and state[0]['queue_status'] == 'done'
and state[0]['target_scan_id'] is not None
and not os.path.lexists(server_ready)
and all(int(capacity[name]) == baseline[name] for name in capacity_names)
):
complete = state[0]
break
time.sleep(min(0.25, max(0, deadline - time.monotonic())))
require(complete is not None, 'remote_transport_ingestion_timeout')
require(json.loads(complete['remote_resolution_json']) == receipt,
'remote_transport_receipt_changed_during_ingestion')
require(rows(db, 'SELECT COUNT(*) AS count FROM target_scans WHERE scan_event_id = ?',
(scan_event_id,))[0]['count'] == 1
and rows(db, 'SELECT COUNT(*) AS count FROM scan_result_compat WHERE target_scan_id = ?',
(complete['target_scan_id'],))[0]['count'] == 1
and rows(db, 'SELECT COUNT(*) AS count FROM target_scans')[0]['count'] == before_scans + 1,
'remote_transport_not_ingested_once')
require(rows(db, '''SELECT COUNT(*) AS count FROM result_reservations
WHERE assignment_kind = 'remote' AND remote_user_id = ?
AND remote_resolved_at IS NULL''', (device['user_id'],))[0]['count'] == 0,
'remote_transport_quota_not_released')
status_receipt = service.status({
'device_id': device['device_id'], 'token_sha256': token_sha256,
}, claimed['reservation_id'])
require(status_receipt == receipt, 'remote_transport_status_receipt')
stage = 'evidence'
evidence = {
'schema': 1, 'reservation_id': claimed['reservation_id'],
'bundle_id': bundle_id, 'scan_event_id': scan_event_id,
'payload_sha256': payload_sha256, 'receipt': receipt,
'target_scan_id': complete['target_scan_id'],
'before_scans': before_scans,
}
write_private_json_exclusive(str(REMOTE_TRANSPORT), evidence)
public = {
'truncated_rejected': 1, 'invalid_rejected': 1,
'disconnect_rejected': 1, 'deadline_rejected': 1,
'publication_recovered': 1, 'simultaneous_replays': 2,
'commit_reply_recovered': 1, 'ingested': 1,
}
return {
'counts': {'ok': 1, **public},
'hashes': {'remote_transport_sha256': digest(json_bytes(public))},
}
except E2EFailure:
raise
except Exception:
raise E2EFailure('remote_transport_' + stage + '_exception') from None
finally:
db.close()
def assert_remote_transport_replay(config, marker, url, timeout):
del config
import asyncio
from result_bundle import bundle_ready_path
import worker_api
from unittest import mock
deadline = time.monotonic() + timeout
while True:
try:
control_snapshot(marker, url)
break
except NotReady:
if time.monotonic() >= deadline:
raise
time.sleep(min(0.25, max(0, deadline - time.monotonic())))
evidence = json.loads(read_bytes(REMOTE_TRANSPORT))
require(evidence.get('schema') == 1, 'remote_transport_evidence_required')
db = worker_api.ScannerDB(db_url=url, initialize=False)
try:
db.set_application_name('truf-container-e2e:remote-transport-replay')
owner = rows(db, '''SELECT d.device_key, d.token_sha256, r.remote_expires_at,
r.remote_resolution_json, r.state, r.queue_id
FROM result_reservations r
JOIN remote_worker_devices d ON d.id = r.remote_device_id
WHERE r.id = ?''', (evidence['reservation_id'],))
require(len(owner) == 1 and owner[0]['state'] == 'acknowledged',
'remote_transport_restart_state')
client_ready = Path(bundle_ready_path(
str(REMOTE_CLIENT_BUNDLES), evidence['bundle_id'],
))
payload = read_bytes(client_ready, limit=16384)
require(digest(payload) == evidence['payload_sha256'],
'remote_transport_client_bundle_changed')
service = worker_api.WorkerService(
url, str(BUNDLES), lambda *_args: None, max_bundle_bytes=16384,
)
with mock.patch.object(worker_api, 'utc_now_iso', return_value='9999-12-31T23:59:59+00:00'):
status, code, receipt = asyncio.run(invoke_bundle_upload(
service, owner[0]['device_key'], evidence['reservation_id'], payload,
))
require(status == 200 and code is None and receipt == evidence['receipt'],
'remote_transport_restart_replay')
status, code, _ = asyncio.run(invoke_bundle_upload(
service, owner[0]['device_key'], evidence['reservation_id'], payload + b'x',
))
require((status, code) == (409, 'resolution_conflict'),
'remote_transport_restart_conflict')
require(json.loads(owner[0]['remote_resolution_json']) == evidence['receipt']
and service.status({
'device_id': rows(db, '''SELECT remote_device_id FROM result_reservations
WHERE id = ?''', (evidence['reservation_id'],))[0]['remote_device_id'],
'token_sha256': owner[0]['token_sha256'],
}, evidence['reservation_id']) == evidence['receipt'],
'remote_transport_restart_receipt_changed')
require(not os.path.lexists(bundle_ready_path(str(BUNDLES), evidence['bundle_id']))
and rows(db, 'SELECT COUNT(*) AS count FROM target_scans WHERE scan_event_id = ?',
(evidence['scan_event_id'],))[0]['count'] == 1
and rows(db, '''SELECT COUNT(*) AS count FROM target_scans
WHERE id = ? AND scan_event_id = ?''',
(evidence['target_scan_id'], evidence['scan_event_id']))[0]['count'] == 1
and rows(db, 'SELECT COUNT(*) AS count FROM result_bundles WHERE reservation_id = ?',
(evidence['reservation_id'],))[0]['count'] == 1,
'remote_transport_restart_exactly_once')
capacity = db.pipeline_capacity_snapshot()
require(not any(int(capacity[name]) for name in (
'bundle_items', 'bundle_bytes', 'projection_items', 'projection_bytes',
'keycheck_items', 'keycheck_bytes', 'quarantine_items', 'quarantine_bytes',
)), 'remote_transport_restart_capacity')
public = {'receipt_replayed': 1, 'conflict_rejected': 1, 'exactly_once': 1}
return {
'counts': {'ok': 1, **public},
'hashes': {'remote_transport_replay_sha256': digest(json_bytes(public))},
}
finally:
db.close()
def prepare_remote_full_race(config, marker, url):
import asyncio
from result_bundle import bundle_partial_path, bundle_ready_path
from runtime_security import ensure_private_directory, write_private_json_exclusive
import scanner_db
import worker_api
from unittest import mock
metadata = control_snapshot(marker, url)
baseline = json.loads(read_bytes(RESULT))
require(baseline.get('schema') == 1 and baseline.get('checked') == 1,
'remote_full_checked_local_baseline')
require(not os.path.lexists(REMOTE_FULL), 'remote_full_evidence_exists')
ensure_private_directory(str(REMOTE_CLIENT_BUNDLES), reject_reparse=True)
ensure_private_directory(str(REMOTE_LOCAL_BUNDLES), reject_reparse=True)
db = scanner_db.ScannerDB(db_url=url, initialize=False)
stage = 'database'
try:
db.set_application_name('truf-container-e2e:remote-full-prepare')
capacity_names = (
'bundle_items', 'bundle_bytes', 'projection_items', 'projection_bytes',
'keycheck_items', 'keycheck_bytes', 'quarantine_items', 'quarantine_bytes',
)
require(not any(int(db.pipeline_capacity_snapshot()[name]) for name in capacity_names),
'remote_full_initial_capacity')
user_key = 'container-e2e-full-user'
require(not rows(db, 'SELECT id FROM remote_worker_users WHERE user_key = ?',
(user_key,)), 'remote_full_user_not_fresh')
require(db.enqueue_targets('gitlab', 'gitlab', 'e2e', [REMOTE_TARGET]) == 1,
'remote_full_requeue')
devices = {}
stage = 'devices'
for role in ('expired', 'winner'):
token = remote_device_token(role)
device = db.provision_remote_worker_device(
user_key, 'container-e2e-full-' + role,
digest(token.encode('ascii')), 1,
)
require(device and not device['revoked'], 'remote_full_device_provision')
devices[role] = device
stage = 'service'
service, build = remote_worker_service(config, marker, metadata, url)
def claim(role):
status, response = asyncio.run(invoke_worker_claim(
service, remote_device_token(role),
digest(('remote-full-' + role + '-request').encode('ascii')), build,
))
require(status == 201 and set(response) == {'assignment'},
'remote_full_claim_transport')
assignment = response['assignment']
require(assignment['reservation']['target'] == REMOTE_TARGET
and assignment['reservation']['remote_device_id'] == devices[role]['device_id'],
'remote_full_claim_identity')
return assignment
stage = 'expired_claim'
expired_assignment = claim('expired')
expired = expired_assignment['reservation']
stage = 'expired_scan'
expired_payload, expired_bundle, _ = execute_fixture_assignment(
config, metadata, expired_assignment, compare_local=False,
)
if (
expired_bundle['finding_count'], expired_bundle['candidate_count'],
expired_bundle['error_count'], expired_bundle['queue_status'],
) != (1, 1, 0, 'done'):
status = re.sub(r'[^a-z0-9_]', '_', expired_bundle['queue_status'])[:20]
category = re.sub(
r'[^a-z0-9_]', '_', expired_bundle['source_failure_category'],
)[:30]
error_text = str(expired_bundle['first_error'] or '').lower()
error_classes = '_'.join(
label for marker, label in (
('protocol', 'protocol'), ('error preparing repo', 'prepare_repo'),
('clone', 'clone'), ('branch', 'branch'), ('reference', 'reference'),
('commit', 'commit'), ('exit status', 'exit_status'),
('authentication', 'auth'), ('certificate', 'certificate'),
('config', 'config'), ('permission', 'permission'),
('no such file', 'missing_file'), ('deadline', 'deadline'),
('timed out', 'timeout'), ('scan-slot', 'scan_slot'),
('authority', 'authority'), ('supervisor', 'supervisor'),
('lifecycle', 'lifecycle'), ('executable', 'executable'),
('completion', 'completion'), ('error running scan', 'running_scan'),
('unable', 'unable'), ('failed', 'failed'), ('private', 'private'),
) if marker in error_text
)[:40] or 'unclassified'
raise E2EFailure(
f'remote_full_expired_counts_{expired_bundle["finding_count"]}_'
f'{expired_bundle["candidate_count"]}_{expired_bundle["error_count"]}_'
f'{status}_{int(expired_bundle["source_failure"])}_{category}_{error_classes}'
)
stage = 'partial_upload'
try:
asyncio.run(invoke_bundle_upload(
service, remote_device_token('expired'), expired['reservation_id'],
expired_payload, disconnect_after=257,
))
except ConnectionError:
pass
else:
raise E2EFailure('remote_full_disconnect_not_observed')
expired_partial = bundle_partial_path(
str(BUNDLES), expired['bundle_id'], expired['reservation_token'],
)
require(not os.path.lexists(expired_partial)
and not os.path.lexists(bundle_ready_path(str(BUNDLES), expired['bundle_id'])),
'remote_full_partial_not_cleaned')
stage = 'exact_deadline'
with mock.patch.object(
worker_api, 'utc_now_iso', return_value=expired['remote_expires_at'],
):
status, code, _ = asyncio.run(invoke_bundle_upload(
service, remote_device_token('expired'), expired['reservation_id'],
expired_payload,
))
require((status, code) == (410, 'assignment_expired'),
'remote_full_exact_expiry_upload')
stage = 'expiry_reaper'
with mock.patch.object(
scanner_db, 'utc_now_iso', return_value=expired['remote_expires_at'],
):
expiry_receipts = service.reap()
require(len(expiry_receipts) == 1
and expiry_receipts[0]['reservation_id'] == expired['reservation_id']
and expiry_receipts[0]['resolution'] == 'expired'
and expiry_receipts[0]['resolved_at'] == expired['remote_expires_at'],
'remote_full_exact_expiry_reap')
stage = 'winner_claim'
winner_assignment = claim('winner')
winner = winner_assignment['reservation']
require(winner['queue_id'] == expired['queue_id']
and winner['reservation_id'] != expired['reservation_id']
and winner['reservation_token'] != expired['reservation_token']
and winner['claim_lease_token'] != expired['claim_lease_token']
and winner['scan_event_id'] != expired['scan_event_id']
and winner['bundle_id'] != expired['bundle_id'],
'remote_full_reissue_identity')
stage = 'stale_report'
status, code, _ = asyncio.run(invoke_terminal_report(
service, remote_device_token('expired'), expired['reservation_id'], {
'failure_code': 'client_process_failed',
'detail': 'synthetic stale issuance fixture',
},
))
require((status, code) == (409, 'reservation_conflict'),
'remote_full_stale_report_fence')
stage = 'winner_scan'
winner_payload, winner_bundle, parity_sha256 = execute_fixture_assignment(
config, metadata, winner_assignment, compare_local=True,
)
if (
winner_bundle['finding_count'], winner_bundle['candidate_count'],
winner_bundle['error_count'], winner_bundle['queue_status'], bool(parity_sha256),
) != (1, 1, 0, 'done', True):
status = re.sub(r'[^a-z0-9_]', '_', winner_bundle['queue_status'])[:20]
raise E2EFailure(
f'remote_full_winner_counts_{winner_bundle["finding_count"]}_'
f'{winner_bundle["candidate_count"]}_{winner_bundle["error_count"]}_'
f'{status}_{int(bool(parity_sha256))}'
)
plan = winner_assignment['scan_kwargs']['git_plan']
expired_plan = expired_assignment['scan_kwargs']['git_plan']
require(plan == expired_plan and plan['ref'] == 'refs/heads/main'
and plan['head_sha'] == marker['repo_commit']
and plan['mode'] == 'baseline' and plan['baseline_depth'] == 1,
'remote_full_exact_plans')
stage = 'plan_bindings'
bound = rows(db, '''SELECT id, git_scan_plan_json, git_scan_plan_sha256
FROM result_reservations WHERE id IN (?, ?) ORDER BY id''',
(expired['reservation_id'], winner['reservation_id']))
plan_sha256 = digest(scanner_db.canonical_git_scan_plan_bytes(plan))
require(len(bound) == 2
and all(json.loads(row['git_scan_plan_json']) == plan
and row['git_scan_plan_sha256'] == plan_sha256 for row in bound),
'remote_full_separate_plan_bindings')
stage = 'capacity'
capacity = db.pipeline_capacity_snapshot()
expected_capacity = {
'bundle_items': 1, 'bundle_bytes': MAX_BYTES,
'projection_items': 1, 'projection_bytes': 2 * MAX_BYTES,
'keycheck_items': 8, 'keycheck_bytes': MAX_BYTES,
'quarantine_items': 0, 'quarantine_bytes': 0,
}
require(all(int(capacity[name]) == value
for name, value in expected_capacity.items()),
'remote_full_reissue_capacity')
interim = db.admin_remote_worker_snapshot(limit=200)
workers = {
row['device_key']: row for row in interim['workers']
if row['user_key'] == user_key
}
require(set(workers) == {
'container-e2e-full-expired', 'container-e2e-full-winner',
} and {
key: int(workers['container-e2e-full-expired'][key])
for key in ('unfinished_count', 'completed_count', 'failed_count', 'expired_count')
} == {
'unfinished_count': 0, 'completed_count': 0,
'failed_count': 0, 'expired_count': 1,
} and {
key: int(workers['container-e2e-full-winner'][key])
for key in ('unfinished_count', 'completed_count', 'failed_count', 'expired_count')
} == {
'unfinished_count': 1, 'completed_count': 0,
'failed_count': 0, 'expired_count': 0,
}, 'remote_full_interim_statistics')
stage = 'evidence'
evidence = {
'schema': 1, 'user_key': user_key,
'origin_instance_sha256': digest(metadata['instance_id'].encode('utf-8')),
'queue_id': winner['queue_id'],
'expired': {
'reservation_id': expired['reservation_id'],
'bundle_id': expired['bundle_id'], 'scan_event_id': expired['scan_event_id'],
'payload_sha256': digest(expired_payload),
'lease_sha256': digest(expired['claim_lease_token'].encode('utf-8')),
},
'winner': {
'reservation_id': winner['reservation_id'],
'bundle_id': winner['bundle_id'], 'scan_event_id': winner['scan_event_id'],
'payload_sha256': digest(winner_payload),
'lease_sha256': digest(winner['claim_lease_token'].encode('utf-8')),
},
'plan': plan, 'plan_sha256': plan_sha256,
'parity_sha256': parity_sha256,
'baseline_keycheck_result_id': baseline['ids']['keycheck_result_id'],
}
write_private_json_exclusive(str(REMOTE_FULL), evidence)
public = {
'real_claims': 2, 'native_scans': 3, 'partial_disconnects': 1,
'exact_expiries': 1, 'stale_reports': 1, 'reissues': 1,
'parity_comparisons': 1, 'pending_restart': 1,
}
return {
'counts': {'ok': 1, **public},
'hashes': {
'remote_full_prepare_sha256': digest(json_bytes(public)),
'remote_full_parity_sha256': parity_sha256,
},
}
except E2EFailure:
raise
except Exception as exc:
trace = exc.__traceback__
origins = []
while trace:
origins.append(trace.tb_frame.f_code.co_name)
trace = trace.tb_next
origin = '_'.join([type(exc).__name__, *origins[-3:]])
origin = re.sub(r'[^a-z0-9_]', '_', origin.lower())[:70]
raise E2EFailure('remote_full_prepare_' + stage + '_' + origin) from None
finally:
db.close()
def remote_full_snapshot(config, marker, url, evidence, checked):
from keycheck_candidates import candidate_uid, extract_candidates
from result_bundle import bundle_ready_path
import scanner_db
metadata = control_snapshot(marker, url)
require(digest(metadata['instance_id'].encode('utf-8'))
!= evidence['origin_instance_sha256'], 'remote_full_runtime_restart')
db = scanner_db.ScannerDB(db_url=url, initialize=False)
try:
db.set_application_name('truf-container-e2e:remote-full-snapshot')
for worker in ('result_ingester', 'jsonl_projector'):
if not db.pipeline_worker_health(worker, metadata['instance_id'])['healthy']:
raise NotReady('remote_full_pipeline_lease')
reservations = rows(db, '''SELECT r.*, b.reservation_id AS result_bundle_id,
b.state AS bundle_state, b.actual_bytes AS bundle_actual_bytes,
b.finding_count AS bundle_finding_count,
b.error_count AS bundle_error_count,
b.candidate_count AS bundle_candidate_count
FROM result_reservations r
JOIN remote_worker_users u ON u.id = r.remote_user_id
LEFT JOIN result_bundles b ON b.reservation_id = r.id
WHERE u.user_key = ? ORDER BY r.id''', (evidence['user_key'],))
require(len(reservations) == 2, 'remote_full_reservation_count')
by_id = {row['id']: row for row in reservations}
expired = by_id[evidence['expired']['reservation_id']]
winner = by_id[evidence['winner']['reservation_id']]
require(expired['state'] == 'refunded'
and expired['remote_resolution_kind'] == 'expired'
and expired['bundle_credit_released'] == 1
and expired['projection_credit_transferred'] == 0
and expired['candidate_credit_transferred'] == 0
and expired['result_bundle_id'] is None,
'remote_full_expired_credit_state')
if winner['state'] != 'acknowledged' or winner['bundle_state'] != 'acknowledged':
raise NotReady('remote_full_ingestion')
require(winner['remote_resolution_kind'] == 'bundle_accepted'
and winner['bundle_credit_released'] == 1
and winner['projection_credit_transferred'] == 1
and winner['candidate_credit_transferred'] == 1
and winner['bundle_finding_count'] == winner['bundle_candidate_count'] == 1
and winner['bundle_error_count'] == 0
and 0 < int(winner['bundle_actual_bytes']) < MAX_BYTES,
'remote_full_winner_credit_state')
receipt = json.loads(winner['remote_resolution_json'])
require(receipt['resolution'] == 'bundle_accepted'
and receipt['reservation_id'] == winner['id']
and receipt['bundle_id'] == evidence['winner']['bundle_id']
and receipt['scan_event_id'] == evidence['winner']['scan_event_id']
and receipt['payload_sha256'] == evidence['winner']['payload_sha256'],
'remote_full_receipt_identity')
require(not os.path.lexists(bundle_ready_path(str(BUNDLES), expired['bundle_id']))
and not os.path.lexists(bundle_ready_path(str(BUNDLES), winner['bundle_id'])),
'remote_full_server_spool_cleanup')
scans = rows(db, '''SELECT s.* FROM target_scans s
JOIN result_reservations r ON r.id = s.result_reservation_id
JOIN remote_worker_users u ON u.id = r.remote_user_id
WHERE u.user_key = ? ORDER BY s.id''', (evidence['user_key'],))
require(len(scans) == 1 and scans[0]['result_reservation_id'] == winner['id'],
'remote_full_authoritative_scan_count')
scan = scans[0]
queue = rows(db, 'SELECT * FROM target_queue WHERE id = ?',
(evidence['queue_id'],))[0]
require(queue['status'] == 'done' and queue['target_scan_id'] == scan['id']
and queue['attempts'] == 1 and queue['lease_token'] is None
and queue['current_result_reservation_id'] is None
and queue['claim_event_id'] is None and not queue['last_error'],
'remote_full_queue_disposition')
require(queue['covered_ref'] == evidence['plan']['ref']
and queue['covered_head'] == evidence['plan']['head_sha'],
'remote_full_winning_plan_coverage')
require(scan['scan_event_id'] == evidence['winner']['scan_event_id']
and scan['status'] == 'found' and scan['findings_count'] == 1
and scan['error_count'] == 0
and scan['queue_completion_disposition'] == 'applied',
'remote_full_scan_result')
rebuilt = db.reconstruct_scan_result(
scan['id'], max_bytes=MAX_BYTES, max_findings=2, max_errors=0,
)
require(rebuilt and len(rebuilt['findings']) == 1 and rebuilt['errors'] == [],
'remote_full_reconstructed_result')
finding = rows(db, 'SELECT * FROM findings WHERE target_scan_id = ?',
(scan['id'],))
candidate = rows(db, 'SELECT * FROM keycheck_candidates WHERE target_scan_id = ?',
(scan['id'],))
require(len(finding) == len(candidate) == 1, 'remote_full_normalized_counts')
finding, candidate = finding[0], candidate[0]
rebuilt_finding = rebuilt['findings'][0]
specs = list(extract_candidates(rebuilt_finding))
require(len(specs) == 1 and specs[0].service == 'openai'
and rebuilt_finding['Raw'] == synthetic_openai_token()
and finding['file_path'] == 'synthetic.env'
and finding['commit_hash'] == marker['repo_commit']
and candidate['finding_id'] == finding['id']
and candidate['finding_uid'] == finding['finding_uid']
and candidate['candidate_uid'] == candidate_uid(
scan['scan_event_id'], finding['finding_uid'], 'openai',
specs[0].credential_hash,
), 'remote_full_finding_candidate_identity')
require(rebuilt['git_scan_plan'] == evidence['plan']
and rebuilt['git_scan_execution']['success'] is True
and rebuilt['git_scan_execution']['pinned'] is True
and rebuilt['git_scan_execution']['plan_sha256'] == evidence['plan_sha256']
and rebuilt['scan_meta']['exact_git_scope']['ref'] == evidence['plan']['ref']
and rebuilt['scan_meta']['exact_git_scope']['head_sha'] == evidence['plan']['head_sha'],
'remote_full_persisted_exact_metadata')
scan_jobs = rows(db, 'SELECT * FROM projection_jobs WHERE target_scan_id = ?',
(scan['id'],))
if len(scan_jobs) != 1 or scan_jobs[0]['status'] != 'completed':
raise NotReady('remote_full_scan_projection')
scan_job = scan_jobs[0]
require(scan_job['event_id'] == scan['scan_event_id']
and scan_job['event_hash'] == scan['scan_event_hash']
and scan_job['required_stream_mask'] == 3
and scan_job['capacity_released'] == 1,
'remote_full_scan_projection_lineage')
projection_hashes = {}
for stream, relative, expected in (
('scan_results', 'scan_results.jsonl', json_bytes(rebuilt) + b'\n'),
('found_secrets', 'found_secrets.jsonl', json_bytes(rebuilt_finding) + b'\n'),
):
state = db.projection_stream_state(stream)
append = db.projection_append_for_job(scan_job['id'], stream)
if not state or not append:
raise NotReady('remote_full_scan_projection_append')
require(state['base_relative_path'] == relative,
'remote_full_scan_projection_path')
payload = read_bytes(RUNTIME / 'results' / relative)
verify_projection_region(payload, expected, append, state, scan_job, 1)
projection_hashes[stream + '_sha256'] = digest(expected)
if checked:
if candidate['state'] != 'completed':
raise NotReady('remote_full_keycheck_completion')
require(candidate['attempts'] == 1 and candidate['lease_token'] is None
and candidate['capacity_released'] == 1
and candidate['result_projection_credit_transferred'] == 1,
'remote_full_keycheck_candidate_credit')
result = db.keycheck_result_for_projection(candidate['keycheck_result_id'])
require(result and result['status'] == 'INVALID_OR_REVOKED'
and result['status_group'] == 'dead'
and result['result_source'] == 'cached_status'
and result['link_status'] == 'linked'
and result['candidate_id'] == candidate['id']
and result['finding_id'] == finding['id']
and result['finding_uid'] == finding['finding_uid']
and result['target_scan_id'] == scan['id'],
'remote_full_detailed_keycheck_result')
result_metadata = json.loads(result['metadata_json'])
require(result_metadata['cached_status'] is True
and result_metadata['cached_state_version'] == 1
and result_metadata['cached_result_id']
== evidence['baseline_keycheck_result_id']
and result_metadata['cached_result_source'] == 'api_check',
'remote_full_cached_keycheck_lineage')
current = rows(db, 'SELECT * FROM keycheck_current_state WHERE credential_id = ?',
(candidate['credential_id'],))[0]
require(current['last_result_id'] == evidence['baseline_keycheck_result_id']
and current['state_version'] == 1
and current['status'] == 'INVALID_OR_REVOKED'
and current['status_group'] == 'dead',
'remote_full_keycheck_current_state')
require(rows(db, 'SELECT keycheck_result_id FROM keycheck_event_map WHERE event_id = ?',
(result['event_id'],)) == [{'keycheck_result_id': result['id']}],
'remote_full_keycheck_event_map')
key_jobs = rows(db, 'SELECT * FROM projection_jobs WHERE keycheck_result_id = ?',
(result['id'],))
if len(key_jobs) != 1 or key_jobs[0]['status'] != 'completed':
raise NotReady('remote_full_keycheck_projection')
key_job = key_jobs[0]
require(key_job['event_id'] == result['event_id']
and key_job['required_stream_mask'] == 8
and key_job['capacity_released'] == 1,
'remote_full_keycheck_projection_lineage')
projected = {key: result[key] for key in (
'event_id', 'service', 'status', 'status_group', 'checked_at', 'key_hash',
'secret_hash', 'key_masked', 'finding_uid', 'source', 'message',
'result_source',
)}
projected.update({'detector': result['detector_name'],
'metadata': result_metadata})
expected = json_bytes(projected) + b'\n'
stream = 'keycheck:openai:results'
state = db.projection_stream_state(stream)
append = db.projection_append_for_job(key_job['id'], stream)
if not state or not append:
raise NotReady('remote_full_keycheck_projection_append')
payload = read_bytes(RUNTIME / 'keychecks/openai/openaiResults.jsonl')
verify_projection_region(payload, expected, append, state, key_job, 1)
projection_hashes[stream + '_sha256'] = digest(expected)
else:
require(candidate['state'] == 'pending' and candidate['attempts'] == 0
and candidate['capacity_released'] == 0
and candidate['result_projection_credit_transferred'] == 0,
'remote_full_keycheck_pending_state')
unreleased_jobs = rows(db, '''SELECT job_kind, status, capacity_items,
capacity_bytes FROM projection_jobs
WHERE capacity_released = 0 ORDER BY id''')
if unreleased_jobs:
raise NotReady('remote_full_projection_drain')
capacity = db.pipeline_capacity_snapshot()
for name in (
'bundle_items', 'bundle_bytes', 'projection_items', 'projection_bytes',
'quarantine_items', 'quarantine_bytes',
):
require(int(capacity[name]) == 0,
'remote_full_capacity_' + name + '_' + str(int(capacity[name])))
require(int(capacity['keycheck_items']) == (0 if checked else 1)
and int(capacity['keycheck_bytes'])
== (0 if checked else int(candidate['capacity_bytes'])),
'remote_full_keycheck_capacity')
admin = db.admin_remote_worker_snapshot(limit=200)
workers = {
row['device_key']: row for row in admin['workers']
if row['user_key'] == evidence['user_key']
}
require(set(workers) == {
'container-e2e-full-expired', 'container-e2e-full-winner',
}, 'remote_full_admin_worker_identity')
expected_statistics = {
'container-e2e-full-expired': (0, 0, 0, 1),
'container-e2e-full-winner': (0, 1, 0, 0),
}
for device_key, expected in expected_statistics.items():
actual = tuple(int(workers[device_key][name]) for name in (
'unfinished_count', 'completed_count', 'failed_count', 'expired_count',
))
require(actual == expected, 'remote_full_admin_statistics')
assignments = [row for row in admin['assignments']
if row['user_key'] == evidence['user_key']]
require(len(assignments) == 2
and {row['outcome'] for row in assignments} == {'completed', 'expired'}
and sum(int(row['accepted']) for row in assignments) == 1
and sum(int(row['ingested']) for row in assignments) == 1,
'remote_full_admin_assignment_counts')
return {
'scan_id': scan['id'], 'candidate_id': candidate['id'],
'receipt_id': receipt['receipt_id'], 'projection_hashes': projection_hashes,
}
finally:
db.close()
def wait_for_remote_full(config, marker, url, evidence, checked, timeout):
deadline = time.monotonic() + timeout
last_check = 'remote_full_timeout'
while time.monotonic() < deadline:
try:
return remote_full_snapshot(config, marker, url, evidence, checked)
except NotReady as exc:
last_check = str(exc)
time.sleep(min(0.25, max(0, deadline - time.monotonic())))
raise E2EFailure('timeout_' + last_check)
def dockerhub_canary_args(config, *, discovery):
from console_runner import build_args_from_source_config
source = copy.deepcopy(config['sources']['dockerhub'])
source.update({
'enabled': True, 'mode': 'search',
'queries': [DOCKERHUB_CANARY_QUERY], 'query_overrides': {},
'pages': 1, 'per_page': DOCKERHUB_CANARY_PER_PAGE,
'workers': 1, 'max_targets': len(DOCKERHUB_CANARY_TARGETS),
'timeout': 60, 'sync_file_queues': False,
'tag_fetch_workers': 1, 'tag_retry_count': 1, 'tag_retry_delay': 1,
'tag_resolve_limit': len(DOCKERHUB_CANARY_TARGETS),
'docker_images_per_repository': 1,
'docker_platform_filter_enabled': True,
'docker_platform_os': 'linux', 'docker_platform_arch': 'amd64',
'docker_platform_candidate_tags': len(DOCKERHUB_CANARY_TARGETS),
'docker_repository_refresh_interval_sec': 0,
'docker_repository_refresh_max_per_cycle': 0,
'docker_content_scan_mode': 'full',
'detectors': 'OpenAI', 'exclude_detectors': '', 'drop_detectors': [],
'no_verification': True, 'external_trufflehog_lifecycle': True,
'trufflehog_config': str(APP / 'trufflehog-custom-detectors.yaml'),
'token': DOCKERHUB_CANARY_DISCOVERY_TOKEN if discovery else '',
'docker_username': '', 'docker_token': '', 'auth_pool': '',
})
args = build_args_from_source_config(
'dockerhub', source, copy.deepcopy(config['global']),
DOCKERHUB_CANARY_QUERY, auth_entry=None,
)
args.max_active_scans = 1
args.result_bundle_max_event_bytes = MAX_BYTES
args.result_bundle_max_items = 1
args.result_bundle_max_total_bytes = MAX_BYTES
args.projection_backlog_max_items = 1
args.projection_backlog_max_bytes = 4 * MAX_BYTES
args.projection_backlog_headroom_bytes = 2 * MAX_BYTES
args.keycheck_queue_max_items = 1
args.keycheck_queue_max_bytes = MAX_BYTES
args.keycheck_candidates_per_event = 1
args.keycheck_candidate_bytes_per_event = 1024
args.pipeline_quarantine_max_items = 1
args.pipeline_quarantine_max_bytes = MAX_BYTES
return args
def dockerhub_canary_transport():
from types import SimpleNamespace
calls = {'search': [], 'tags': []}
targets = dict(zip(DOCKERHUB_CANARY_REPOSITORIES, DOCKERHUB_CANARY_TARGETS))
def search(query, page, **kwargs):
require(
query == DOCKERHUB_CANARY_QUERY and page == 1
and kwargs == {
'per_page': DOCKERHUB_CANARY_PER_PAGE,
'sort_by': 'updated_at', 'sort_order': 'desc',
'request_timeout': 15,
}
and not calls['search'],
'dockerhub_canary_search_bound',
)
calls['search'].append((query, page, dict(kwargs)))
return {
'page': 1, 'total_count': len(DOCKERHUB_CANARY_REPOSITORIES),
'repositories': [
{'repo_name': repository}
for repository in DOCKERHUB_CANARY_REPOSITORIES
],
}
def tags(repository, *args, **kwargs):
require(
repository in targets and repository not in calls['tags']
and args == (
None, 1, 1, 1, True, 'linux', 'amd64',
len(DOCKERHUB_CANARY_TARGETS),
)
and kwargs == {'return_outcome': True},
'dockerhub_canary_tag_bound',
)
calls['tags'].append(repository)
return SimpleNamespace(
tags=(targets[repository],), status='ok', remote_attempted=True,
retry_at=None, error='',
)
return search, tags, calls
def dockerhub_canary_service(config, metadata, url):
from worker_api import WorkerService
from worker_assignment import RemoteAssignmentBuilder
from worker_package import worker_package_build_compatibility
capability = {
'source': 'dockerhub', 'platform': 'docker',
'planning_kind': 'docker_direct_v1',
}
package = remote_package_manifest(metadata, [capability])
args = dockerhub_canary_args(config, discovery=False)
bundle_body_timeout = int(
config['supervisor']['worker_api']['bundle_body_timeout_seconds']
)
assignment_ttl = int(args.timeout) + bundle_body_timeout + 60
builder = RemoteAssignmentBuilder(
url, str(BUNDLES), {'dockerhub': args},
{'dockerhub-canary': {
'package_manifest': package, 'sources': ['dockerhub'],
}},
metadata['instance_id'], assignment_ttl_seconds=assignment_ttl,
)
return (
WorkerService(url, str(BUNDLES), builder, max_bundle_bytes=MAX_BYTES),
worker_package_build_compatibility(package), args,
)
def require_dockerhub_canary_assignment(assignment):
reservation = assignment['reservation']
snapshot = assignment['execution_snapshot']
plan = assignment['execution_plan']
require(
reservation['source'] == 'dockerhub'
and reservation['platform'] == 'docker'
and reservation['target'] in DOCKERHUB_CANARY_TARGETS
and '@sha256:' in reservation['target']
and snapshot['planning'] == {'kind': 'docker_direct_v1'}
and snapshot['credential_ref'] == {
'source': 'dockerhub', 'auth_entry': '',
}
and plan == {
'kind': 'docker_direct_v1',
'execution_target': reservation['target'], 'bound_plan': None,
}
and assignment['scan_kwargs'] == assignment['event_scan_options'],
'dockerhub_canary_assignment_identity',
)
encoded = json.dumps(assignment, ensure_ascii=True, sort_keys=True).lower()
require(
DOCKERHUB_CANARY_DISCOVERY_TOKEN not in encoded
and all(name not in encoded for name in (
'docker_registry_auth', 'access_probe', 'public_access_proof',
'proof_freshness', 'docker_layer_work', 'docker_layer_plan',
)),
'dockerhub_canary_assignment_minimal',
)
def stage_dockerhub_canary_result(metadata, assignment, result_kind):
from lifecycle_authority import build_code_manifest, code_manifest_sha256
from result_bundle import BundleReservation, bundle_ready_path
import scan_execution
import scanner
diagnostics = {
'permanent': json.dumps({
'level': 'error', 'msg': 'provider access failed',
'error': 'pull access denied',
}),
'retryable': json.dumps({
'level': 'error', 'msg': 'provider request failed',
'error': 'HTTP 503 service unavailable',
}),
'success': json.dumps({
'level': 'info-0', 'logger': 'trufflehog',
'msg': 'finished scanning',
}),
}
require(result_kind in diagnostics, 'dockerhub_canary_result_kind')
reservation = BundleReservation.from_mapping(assignment['reservation'])
result = {'findings': [], 'errors': []}
policy_path = str(APP / 'trufflehog-custom-detectors.yaml')
client_manifest = build_code_manifest(
str(APP), scanner.scan_config.trufflehog_path, (policy_path,),
git_path=metadata['code_manifest']['executables']['git']['path'],
)
with scanner.client_scan_launch_authority(
client_manifest, code_manifest_sha256(client_manifest),
), scanner.client_scan_execution_policy(
assignment['scan_policy'],
), scanner.client_remote_execution_binding('docker_direct_v1'):
scanner.apply_trufflehog_diagnostics(
result, diagnostics[result_kind], 0 if result_kind == 'success' else 1,
'docker',
)
result.update({
'target': reservation.target, 'scan_type': 'docker',
'scan_event_id': reservation.scan_event_id,
'scan_started_at': '2026-09-20T00:00:00+00:00',
'duration_sec': 0.0, 'timestamp': '2026-09-20T00:00:00+00:00',
})
classification = {
'error_class': str(result.get('error_class') or ''),
'retryable': bool(result.get('retryable', False)),
'skipped': bool(result.get('skipped')),
}
staged = scan_execution.stage_scan_result_in_scope(
result, reservation, str(REMOTE_CLIENT_BUNDLES),
assignment['event_scan_options'], assignment['queue_policy'],
attempts=int(assignment['reservation'].get('attempts') or 0),
candidate_max_items=int(assignment['limits']['candidate_max_items']),
candidate_max_bytes=int(assignment['limits']['candidate_max_bytes']),
)
payload = read_bytes(Path(bundle_ready_path(
str(REMOTE_CLIENT_BUNDLES), reservation.bundle_id,
)), limit=MAX_BYTES)
require(
staged.actual_bytes == len(payload) and 0 < len(payload) < MAX_BYTES,
'dockerhub_canary_bundle_bound',
)
expected = {
'permanent': ('docker_registry_access', False, True, 'done', 0),
'retryable': ('remote_transient', True, False, 'deferred', 1),
'success': ('', False, False, 'done', 0),
}[result_kind]
require(
(
classification['error_class'], classification['retryable'],
classification['skipped'], staged.queue_status, staged.error_count,
) == expected,
'dockerhub_canary_worker_disposition',
)
return payload, staged, classification
def dockerhub_canary_result_snapshot(db, reservation_id, result_kind, receipt):
state = rows(db, '''SELECT r.state, r.remote_resolution_json,
b.state AS bundle_state, q.status AS queue_status, q.target_scan_id
FROM result_reservations r
JOIN result_bundles b ON b.reservation_id = r.id
JOIN target_queue q ON q.id = r.queue_id
WHERE r.id = ?''', (reservation_id,))
if (
len(state) != 1 or state[0]['state'] != 'acknowledged'
or state[0]['bundle_state'] != 'acknowledged'
or state[0]['target_scan_id'] is None
):
raise NotReady('dockerhub_canary_ingestion')
expected_queue = 'deferred' if result_kind == 'retryable' else 'done'
require(state[0]['queue_status'] == expected_queue,
'dockerhub_canary_queue_disposition')
require(json.loads(state[0]['remote_resolution_json']) == receipt,
'dockerhub_canary_receipt_persistence')
scan = rows(db, 'SELECT * FROM target_scans WHERE id = ?',
(state[0]['target_scan_id'],))[0]
rebuilt = db.reconstruct_scan_result(
scan['id'], max_bytes=MAX_BYTES, max_findings=1, max_errors=2,
)
require(rebuilt and rebuilt['target'] in DOCKERHUB_CANARY_TARGETS
and rebuilt['scan_type'] == 'docker',
'dockerhub_canary_reconstructed_identity')
if result_kind == 'permanent':
require(rebuilt.get('error_class') == 'docker_registry_access'
and rebuilt.get('retryable') is False
and rebuilt.get('skipped') == 'Docker image is unavailable to the worker'
and rebuilt['errors'] == [],
'dockerhub_canary_permanent_classification')
elif result_kind == 'retryable':
require(rebuilt.get('error_class') == 'remote_transient'
and rebuilt.get('retryable') is True
and len(rebuilt['errors']) == 1,
'dockerhub_canary_retryable_classification')
else:
require(not rebuilt['errors'] and not rebuilt.get('skipped'),
'dockerhub_canary_success_classification')
jobs = rows(db, 'SELECT * FROM projection_jobs WHERE target_scan_id = ?',
(scan['id'],))
if len(jobs) != 1 or jobs[0]['status'] != 'completed':
raise NotReady('dockerhub_canary_projection')
job = jobs[0]
expected_stream_mask = 5 if result_kind == 'retryable' else 1
require(job['job_kind'] == 'scan_event'
and job['event_id'] == scan['scan_event_id']
and job['event_hash'] == scan['scan_event_hash']
and job['required_stream_mask'] == expected_stream_mask
and job['capacity_released'] == 1,
'dockerhub_canary_projection_lineage')
state_row = db.projection_stream_state('scan_results')
append = db.projection_append_for_job(job['id'], 'scan_results')
if not state_row or not append:
raise NotReady('dockerhub_canary_projection_append')
payload = read_bytes(RUNTIME / 'results' / state_row['base_relative_path'])
expected = json_bytes(rebuilt) + b'\n'
verify_projection_region(payload, expected, append, state_row, job, 1)
return {
'scan_id': scan['id'], 'projection_job_id': job['id'],
'projection_sha256': digest(expected),
'projection_completed_at': job['completed_at'],
}
def wait_for_dockerhub_canary_result(
db, reservation_id, result_kind, receipt, timeout,
):
deadline = time.monotonic() + timeout
last_check = 'dockerhub_canary_pipeline'
while time.monotonic() < deadline:
try:
return dockerhub_canary_result_snapshot(
db, reservation_id, result_kind, receipt,
)
except NotReady as exc:
last_check = str(exc)
time.sleep(min(0.1, max(0, deadline - time.monotonic())))
raise E2EFailure('timeout_' + last_check)
def dockerhub_canary_operation_id(label):
value = digest(('dockerhub-canary:' + label).encode('ascii'))[:32]
return '-'.join((value[:8], value[8:12], value[12:16], value[16:20], value[20:]))
def dockerhub_canary_cycle(db, args, label):
run_id = db.start_run(
'container-e2e-dockerhub-canary', selected_source='dockerhub',
selected_platform='docker', config_path=str(CONFIG),
enabled_sources=['dockerhub'],
)
cycle_id = db.start_source_cycle(
run_id, 'dockerhub', 'docker', 'search', DOCKERHUB_CANARY_QUERY, 1, 1,
)
return run_id, cycle_id
def prepare_dockerhub_canary(config, marker, url, timeout):
import asyncio
from datetime import datetime, timedelta, timezone
from runtime_security import ensure_private_directory, write_private_json_exclusive
import console_runner
import scanner_db
from unittest import mock
metadata = control_snapshot(marker, url)
require(not os.path.lexists(DOCKERHUB_CANARY),
'dockerhub_canary_evidence_exists')
ensure_private_directory(str(REMOTE_CLIENT_BUNDLES), reject_reparse=True)
db = scanner_db.ScannerDB(db_url=url, initialize=False)
stage = 'database'
try:
db.set_application_name('truf-container-e2e:dockerhub-canary-prepare')
require(db.enabled and db.conn.is_postgres, 'dockerhub_canary_postgres')
db.require_runtime_safety_schema()
db.require_final_cutover()
control = db.runtime_control_state()
require(control['drain_state'] == 'normal'
and not control['effective_discovery_paused']
and not control['effective_dispatch_paused'],
'dockerhub_canary_initial_control')
require(not rows(db, 'SELECT id FROM target_queue WHERE source = ? AND query = ?',
('dockerhub', DOCKERHUB_CANARY_QUERY)),
'dockerhub_canary_queue_not_fresh')
authority_columns = rows(db, '''SELECT table_name, column_name
FROM information_schema.columns
WHERE table_schema = current_schema()
ORDER BY table_name, ordinal_position''')
require(not any(
forbidden in (
str(column['table_name']) + '.' + str(column['column_name'])
).lower()
for column in authority_columns
for forbidden in ('access_probe', 'public_access_proof', 'proof_freshness')
), 'dockerhub_canary_access_state_forbidden')
stage = 'discovery'
discovery_args = dockerhub_canary_args(config, discovery=True)
search, tags, calls = dockerhub_canary_transport()
run_id, cycle_id = dockerhub_canary_cycle(db, discovery_args, 'initial')
with mock.patch.object(
console_runner, 'fetch_dockerhub_search_page', new=search,
), mock.patch.object(
console_runner, 'fetch_dockerhub_tags', new=tags,
):
metrics = console_runner.run_discovery_cycle(
discovery_args, db, run_id, cycle_id, 'dockerhub',
)
db.finish_run(run_id)
require(metrics['cycle_status'] == 'completed'
and metrics['fetched_count'] == len(DOCKERHUB_CANARY_REPOSITORIES)
and metrics['queued_new_count'] == len(DOCKERHUB_CANARY_REPOSITORIES)
and metrics['discovery_pages_fetched'] == 1
and metrics['scan_requested_count'] == 0
and len(calls['search']) == 1 and not calls['tags'],
'dockerhub_canary_discovery_result')
stage = 'resolution'
resolver_now = (
datetime.now(timezone.utc) + timedelta(seconds=7200)
).isoformat(timespec='seconds')
with mock.patch.object(
console_runner, 'fetch_dockerhub_tags', new=tags,
), mock.patch.object(
scanner_db, 'utc_now_iso', return_value=resolver_now,
):
resolved = console_runner.resolve_due_docker_queue_targets(
db, 'dockerhub', discovery_args,
)
if (
resolved != len(DOCKERHUB_CANARY_TARGETS)
or calls['tags'] != list(DOCKERHUB_CANARY_REPOSITORIES)
):
raise E2EFailure(
f'dockerhub_canary_resolution_result_{resolved}_{len(calls["tags"])}'
)
queue = rows(db, '''SELECT target, status, resolver_state FROM target_queue
WHERE source = ? AND query = ? ORDER BY id''',
('dockerhub', DOCKERHUB_CANARY_QUERY))
repositories = [row for row in queue if '@sha256:' not in row['target']]
targets = [row for row in queue if '@sha256:' in row['target']]
require([row['target'] for row in repositories]
== list(DOCKERHUB_CANARY_REPOSITORIES)
and all(row['status'] == 'done' and row['resolver_state'] == 'resolved'
for row in repositories)
and [row['target'] for row in targets] == list(DOCKERHUB_CANARY_TARGETS)
and all(row['status'] == 'pending' for row in targets),
'dockerhub_canary_immutable_queue')
stage = 'worker'
user_key = 'container-e2e-dockerhub-canary-user'
device_key = 'container-e2e-dockerhub-canary-device'
device_token = 'container-e2e-dockerhub-canary-token'
device = db.provision_remote_worker_device(
user_key, device_key, digest(device_token.encode('ascii')), 1,
)
require(device['active_assignment_cap'] == 1 and not device['revoked'],
'dockerhub_canary_device')
service, build, assignment_args = dockerhub_canary_service(
config, metadata, url,
)
require(assignment_args.workers == assignment_args.max_active_scans == 1
and not assignment_args.token and not assignment_args.docker_username
and not assignment_args.docker_token and not assignment_args.auth_name,
'dockerhub_canary_assignment_configuration')
def claim(label):
request_id = digest(('dockerhub-canary:' + label).encode('ascii'))
status, response = asyncio.run(invoke_worker_claim(
service, device_token, request_id, build,
))
require(status == 201 and set(response or {}) == {'assignment'},
'dockerhub_canary_claim_transport')
assignment = response['assignment']
require_dockerhub_canary_assignment(assignment)
require(assignment['reservation']['remote_device_id'] == device['device_id'],
'dockerhub_canary_claim_device')
return request_id, assignment
def upload(assignment, result_kind):
payload, staged, classification = stage_dockerhub_canary_result(
metadata, assignment, result_kind,
)
status, code, receipt = asyncio.run(invoke_bundle_upload(
service, device_token,
assignment['reservation']['reservation_id'], payload,
))
require(status == 201 and code is None
and receipt['resolution'] == 'bundle_accepted'
and receipt['payload_sha256'] == digest(payload),
'dockerhub_canary_upload')
snapshot = wait_for_dockerhub_canary_result(
db, assignment['reservation']['reservation_id'], result_kind,
receipt, timeout,
)
return payload, staged, classification, receipt, snapshot
stage = 'permanent'
lost_request_id, lost_assignment = claim('permanent-lost-response')
status, replayed = asyncio.run(invoke_worker_claim(
service, device_token, lost_request_id, build,
))
require(status == 201 and replayed == {'assignment': lost_assignment},
'dockerhub_canary_lost_claim_replay')
permanent = upload(lost_assignment, 'permanent')
stage = 'retryable'
_, retryable_assignment = claim('retryable')
retryable = upload(retryable_assignment, 'retryable')
stage = 'expiry_claim'
_, expired_assignment = claim('unfinished-expiry')
expired_payload, _, _, = stage_dockerhub_canary_result(
metadata, expired_assignment, 'success',
)
expired = expired_assignment['reservation']
state = db.runtime_control_state()
db.start_runtime_drain(
expected_revision=state['revision'], actor='test:container-e2e',
operation_id=dockerhub_canary_operation_id('expiry-drain-start'),
)
status_value = service.status({
'device_id': device['device_id'],
'token_sha256': digest(device_token.encode('ascii')),
}, expired['reservation_id'])
require(status_value['state'] == 'scanning',
'dockerhub_canary_drain_status')
stage = 'paused_discovery'
paused_run, paused_cycle = dockerhub_canary_cycle(
db, discovery_args, 'paused',
)
with mock.patch.object(
console_runner, 'fetch_dockerhub_search_page', new=search,
), mock.patch.object(
console_runner, 'fetch_dockerhub_tags', new=tags,
):
paused = console_runner.run_discovery_cycle(
discovery_args, db, paused_run, paused_cycle, 'dockerhub',
)
db.finish_run(paused_run)
require(paused['cycle_status'] == 'paused'
and len(calls['search']) == 1
and len(calls['tags']) == len(DOCKERHUB_CANARY_TARGETS),
'dockerhub_canary_drain_discovery_gate')
stage = 'expiry'
with mock.patch.object(
scanner_db, 'utc_now_iso', return_value=expired['remote_expires_at'],
):
expiry_receipts = service.reap()
require(len(expiry_receipts) == 1
and expiry_receipts[0]['reservation_id'] == expired['reservation_id']
and expiry_receipts[0]['resolution'] == 'expired',
'dockerhub_canary_expiry_receipt')
status, code, _ = asyncio.run(invoke_bundle_upload(
service, device_token, expired['reservation_id'], expired_payload,
))
require((status, code) == (409, 'resolution_conflict'),
'dockerhub_canary_stale_upload')
first_drained = db.runtime_control_state()
require(first_drained['drain_state'] == 'drained'
and first_drained['effective_discovery_paused']
and first_drained['effective_dispatch_paused'],
'dockerhub_canary_expiry_drain')
pending_after_expiry = rows(db, '''SELECT target FROM target_queue
WHERE source = ? AND query = ? AND status = 'pending'
AND target LIKE '%@sha256:%' ORDER BY id''',
('dockerhub', DOCKERHUB_CANARY_QUERY))
require(len(pending_after_expiry) == 2,
'dockerhub_canary_expiry_pending_targets')
db.cancel_runtime_drain(
expected_revision=first_drained['revision'], actor='test:container-e2e',
operation_id=dockerhub_canary_operation_id('expiry-drain-cancel'),
)
stage = 'success_claim'
success_request_id, success_assignment = claim('success')
require(success_assignment['reservation']['queue_id'] == expired['queue_id'],
'dockerhub_canary_expired_target_reclaim')
success_payload, success_staged, success_classification = (
stage_dockerhub_canary_result(metadata, success_assignment, 'success')
)
success = success_assignment['reservation']
state = db.runtime_control_state()
drain_start = db.start_runtime_drain(
expected_revision=state['revision'], actor='test:container-e2e',
operation_id=dockerhub_canary_operation_id('upload-drain-start'),
)
require(service.status({
'device_id': device['device_id'],
'token_sha256': digest(device_token.encode('ascii')),
}, success['reservation_id'])['state'] == 'scanning',
'dockerhub_canary_upload_drain_status')
stage = 'drain_upload'
status, code, success_receipt = asyncio.run(invoke_bundle_upload(
service, device_token, success['reservation_id'], success_payload,
))
require(status == 201 and code is None
and success_receipt['resolution'] == 'bundle_accepted',
'dockerhub_canary_drain_upload')
status, code, replay_receipt = asyncio.run(invoke_bundle_upload(
service, device_token, success['reservation_id'], success_payload,
))
require(status == 200 and code is None and replay_receipt == success_receipt,
'dockerhub_canary_drain_replay')
status, code, _ = asyncio.run(invoke_bundle_upload(
service, device_token, success['reservation_id'], success_payload + b'x',
))
require((status, code) == (409, 'resolution_conflict'),
'dockerhub_canary_drain_conflict')
reservations_before = rows(db, '''SELECT COUNT(*) AS count
FROM result_reservations WHERE remote_user_id = ?''',
(device['user_id'],))[0]['count']
blocked_status, blocked = asyncio.run(invoke_worker_claim(
service, device_token,
digest(b'dockerhub-canary-drain-blocked-claim'), build,
))
reservations_after = rows(db, '''SELECT COUNT(*) AS count
FROM result_reservations WHERE remote_user_id = ?''',
(device['user_id'],))[0]['count']
require(blocked_status == 204 and blocked is None
and reservations_before == reservations_after,
'dockerhub_canary_drain_dispatch_gate')
stage = 'drain_pipeline'
success_snapshot = wait_for_dockerhub_canary_result(
db, success['reservation_id'], 'success', success_receipt, timeout,
)
state = db.runtime_control_state()
if state['drain_state'] == 'draining':
completed = db.reconcile_runtime_drain()
require(completed is not None, 'dockerhub_canary_drain_reconcile')
state = db.runtime_control_state()
require(state['drain_state'] == 'drained',
'dockerhub_canary_final_drained')
pending = rows(db, '''SELECT target FROM target_queue
WHERE source = ? AND query = ? AND status = 'pending'
AND target LIKE '%@sha256:%' ORDER BY id''',
('dockerhub', DOCKERHUB_CANARY_QUERY))
require(pending == [{'target': DOCKERHUB_CANARY_TARGETS[-1]}],
'dockerhub_canary_final_pending_target')
drain_operation = rows(db, '''SELECT completed_at FROM runtime_operations
WHERE operation_id = ?''', (drain_start['operation_id'],))[0]
require(success_snapshot['projection_completed_at'] >= drain_operation['completed_at'],
'dockerhub_canary_projection_during_drain')
db.cancel_runtime_drain(
expected_revision=state['revision'], actor='test:container-e2e',
operation_id=dockerhub_canary_operation_id('upload-drain-cancel'),
)
capacity = db.pipeline_capacity_snapshot()
require(not any(int(capacity[name]) for name in (
'bundle_items', 'bundle_bytes', 'projection_items', 'projection_bytes',
'keycheck_items', 'keycheck_bytes', 'quarantine_items',
'quarantine_bytes',
)), 'dockerhub_canary_capacity_release')
stage = 'evidence'
evidence = {
'schema': 1,
'origin_instance_sha256': digest(metadata['instance_id'].encode('utf-8')),
'user_key': user_key, 'device_key': device_key,
'success_request_id': success_request_id,
'build': build,
'success': {
'reservation_id': success['reservation_id'],
'bundle_id': success['bundle_id'],
'scan_event_id': success['scan_event_id'],
'payload_sha256': digest(success_payload),
'remote_expires_at': success['remote_expires_at'],
'receipt': success_receipt,
'scan_id': success_snapshot['scan_id'],
'projection_job_id': success_snapshot['projection_job_id'],
},
'expired': {
'reservation_id': expired['reservation_id'],
'receipt': expiry_receipts[0],
},
'classifications': {
'permanent': permanent[2], 'retryable': retryable[2],
'success': success_classification,
},
}
write_private_json_exclusive(str(DOCKERHUB_CANARY), evidence)
public = {
'search_pages': 1,
'cohort_targets': len(DOCKERHUB_CANARY_TARGETS),
'digest_resolutions': len(calls['tags']),
'protocol2_assignments': 4, 'lost_claim_replays': 1,
'worker_provider_failures': 2, 'accepted_uploads': 3,
'expired_assignments': 1, 'stale_uploads_rejected': 1,
'drain_cycles': 2, 'drain_replays': 1,
'drain_conflicts_rejected': 1, 'pending_targets': len(pending),
'projected_results': 3,
}
return {
'counts': {'ok': 1, **public},
'hashes': {
'dockerhub_canary_prepare_sha256': digest(json_bytes(public)),
'dockerhub_canary_projection_sha256': success_snapshot['projection_sha256'],
},
}
except E2EFailure:
raise
except Exception as exc:
trace = exc.__traceback__
origins = []
while trace:
origins.append(trace.tb_frame.f_code.co_name)
trace = trace.tb_next
origin = '_'.join([type(exc).__name__, *origins[-3:]])
origin = re.sub(r'[^a-z0-9_]', '_', origin.lower())[:80]
raise E2EFailure(
'dockerhub_canary_prepare_' + stage + '_' + origin
) from None
finally:
db.close()
def finish_dockerhub_canary(config, marker, url, timeout):
del timeout
import asyncio
from result_bundle import bundle_ready_path
import scanner_db
import worker_api
from unittest import mock
metadata = control_snapshot(marker, url)
evidence = json.loads(read_bytes(DOCKERHUB_CANARY))
require(evidence.get('schema') == 1
and digest(metadata['instance_id'].encode('utf-8'))
!= evidence['origin_instance_sha256'],
'dockerhub_canary_restart_evidence')
service, build, _args = dockerhub_canary_service(config, metadata, url)
require(build == evidence['build'], 'dockerhub_canary_restart_build')
success = evidence['success']
path = Path(bundle_ready_path(str(REMOTE_CLIENT_BUNDLES), success['bundle_id']))
payload = read_bytes(path, limit=MAX_BYTES)
require(digest(payload) == success['payload_sha256'],
'dockerhub_canary_restart_payload')
status, response = asyncio.run(invoke_worker_claim(
service, 'container-e2e-dockerhub-canary-token',
evidence['success_request_id'], build,
))
require(status == 200 and response == {'resolution': success['receipt']},
'dockerhub_canary_restart_claim_receipt')
with mock.patch.object(
worker_api, 'utc_now_iso', return_value='9999-12-31T23:59:59+00:00',
):
status, code, receipt = asyncio.run(invoke_bundle_upload(
service, 'container-e2e-dockerhub-canary-token',
success['reservation_id'], payload,
))
require(status == 200 and code is None and receipt == success['receipt'],
'dockerhub_canary_restart_receipt_replay')
status, code, _ = asyncio.run(invoke_bundle_upload(
service, 'container-e2e-dockerhub-canary-token',
success['reservation_id'], payload + b'x',
))
require((status, code) == (409, 'resolution_conflict'),
'dockerhub_canary_restart_conflict')
db = scanner_db.ScannerDB(db_url=url, initialize=False)
try:
db.set_application_name('truf-container-e2e:dockerhub-canary-finish')
assignments = rows(db, '''SELECT r.id, r.state, r.remote_resolution_kind,
r.remote_resolution_json
FROM result_reservations r
JOIN remote_worker_users u ON u.id = r.remote_user_id
WHERE u.user_key = ? ORDER BY r.id''', (evidence['user_key'],))
require(len(assignments) == 4
and sum(row['remote_resolution_kind'] == 'bundle_accepted'
for row in assignments) == 3
and sum(row['remote_resolution_kind'] == 'expired'
for row in assignments) == 1,
'dockerhub_canary_restart_assignments')
scans = rows(db, '''SELECT id, scan_event_id, scan_event_hash
FROM target_scans WHERE scan_event_id = ?''', (success['scan_event_id'],))
bundles = rows(db, '''SELECT reservation_id FROM result_bundles
WHERE reservation_id = ?''', (success['reservation_id'],))
jobs = rows(db, 'SELECT * FROM projection_jobs WHERE target_scan_id = ?',
(success['scan_id'],))
appends = rows(db, '''SELECT * FROM projection_appends
WHERE job_id = ? ORDER BY id''', (success['projection_job_id'],))
require(len(scans) == len(bundles) == len(jobs) == len(appends) == 1
and scans[0]['id'] == success['scan_id']
and bundles[0]['reservation_id'] == success['reservation_id']
and jobs[0]['id'] == success['projection_job_id']
and jobs[0]['job_kind'] == 'scan_event'
and jobs[0]['status'] == 'completed'
and jobs[0]['event_id'] == scans[0]['scan_event_id']
and jobs[0]['event_hash'] == scans[0]['scan_event_hash']
and jobs[0]['required_stream_mask'] == 1
and jobs[0]['capacity_released'] == 1
and appends[0]['stream_name'] == 'scan_results'
and appends[0]['event_id'] == jobs[0]['event_id']
and appends[0]['event_hash'] == jobs[0]['event_hash'],
'dockerhub_canary_restart_exactly_once')
require(not os.path.lexists(bundle_ready_path(
str(BUNDLES), success['bundle_id'],
)), 'dockerhub_canary_restart_spool_cleanup')
capacity = db.pipeline_capacity_snapshot()
require(not any(int(capacity[name]) for name in (
'bundle_items', 'bundle_bytes', 'projection_items', 'projection_bytes',
'keycheck_items', 'keycheck_bytes', 'quarantine_items',
'quarantine_bytes',
)), 'dockerhub_canary_restart_capacity')
finally:
db.close()
public = {
'runtime_restarts': 1, 'lost_claim_receipts': 1,
'receipt_replays': 1, 'conflicts_rejected': 1,
'exactly_once': 1, 'former_expiry_replays': 1,
}
return {
'counts': {'ok': 1, **public},
'hashes': {
'dockerhub_canary_finish_sha256': digest(json_bytes(public)),
},
}
def fixture_transport(expected_id, calls):
import requests
def request(session, method, url, **kwargs):
del session
common = sys.modules.get('keycheck_common')
active = getattr(common, '_ACTIVE_DB_CANDIDATE', None)
require(active and active['id'] == expected_id and active.get('lease_token')
and not active.get('_completed'), 'provider_active_candidate_fence')
require(not calls and str(method).upper() == 'GET'
and url == 'https://api.openai.com/v1/models'
and kwargs.get('headers') == {'Authorization': 'Bearer ' + synthetic_openai_token()}
and set(kwargs) <= {'headers', 'proxies', 'params', 'timeout', 'allow_redirects'}
and not kwargs.get('proxies') and not kwargs.get('params'),
'unexpected_provider_network')
calls.append(1)
response = requests.Response()
response.status_code = 401
response.url = url
response.headers['Content-Type'] = 'application/json'
response._content = b'{"error":{"message":"Known-fake offline fixture","type":"invalid_request_error","code":"invalid_api_key"}}'
return response
return request
def keycheck_child(bootstrap):
# This is a test transport wrapper, not an alternate production bootstrap.
# main() below still authenticates the provider, verifies its immutable
# manifest, and invokes the actual leaf entrypoint and database candidate loop.
from unittest import mock
import requests
baseline = json.loads(read_bytes(RESULT))
calls = []
sys.argv = [str(APP / 'child_bootstrap.py'), 'keycheck-provider',
'keycheckers/openai/Keycheck.py', '--', '--proxy-file', str(FIXTURE / 'no-proxy')]
with mock.patch.object(requests.Session, 'request', new=fixture_transport(baseline['ids']['candidate_id'], calls)):
try:
bootstrap['main']()
except SystemExit as exc:
require(exc.code in (None, 0), 'provider_exit')
expected_calls = 1 - baseline['checked']
require(len(calls) == expected_calls, 'provider_request_count')
return {'counts': {'ok': 1, 'http_requests': len(calls)}, 'hashes': {}}
def run_keycheck_process(config, marker, url, expected_http_requests, timeout):
from lifecycle_authority import strip_supervisor_credentials, supervised_child_environment
metadata = control_snapshot(marker, url)
env = os.environ.copy()
strip_supervisor_credentials(env)
for name in list(env):
if name.startswith('KEYCHECK_') or name.lower().endswith('_proxy'):
env.pop(name)
env.update(supervised_child_environment(metadata, url, 'keycheck-provider'))
output = str(RUNTIME / 'keychecks/openai')
env.update({
'SCANNER_DB_URL': url, 'DATABASE_URL': url, 'KEYCHECK_DB_URL': url,
'KEYCHECK_SERVICE': 'openai', 'KEYCHECK_INPUT_MODE': 'postgres',
'KEYCHECK_OUTPUT_DIR': output, 'KEYCHECK_STATE_DIR': output,
'KEYCHECK_PROVIDER_SLICE_KEYS': '1', 'KEYCHECK_DB_INLINE': '0',
'KEYCHECK_PROXY_FILE': str(FIXTURE / 'no-proxy'),
'KEYCHECK_RESULT_PROJECTION_RESERVE_BYTES': str(config['global']['keycheck_result_projection_reserve_bytes']),
'KEYCHECK_PROJECTION_MAX_ITEMS': str(config['global']['projection_backlog_max_items']),
'KEYCHECK_PROJECTION_MAX_BYTES': str(config['global']['projection_backlog_max_bytes']),
})
require(not os.path.lexists(FIXTURE / 'no-proxy'), 'fixture_proxy_forbidden')
completed = subprocess.run(
[sys.executable, '-I', '-S', '-B', str(Path(__file__).resolve()), '_keycheck-child'],
cwd=FIXTURE, env=env, stdin=subprocess.DEVNULL, stdout=subprocess.PIPE,
stderr=subprocess.PIPE, timeout=min(90, timeout), check=False,
)
require(completed.returncode == 0 and len(completed.stdout) <= 4096, 'provider_subprocess_failed')
summary = json.loads(completed.stdout)
require(summary == {'counts': {
'ok': 1, 'http_requests': expected_http_requests,
}, 'hashes': {}},
'provider_subprocess_summary')
return summary
def run_keycheck(config, marker, url, baseline, timeout):
summary = run_keycheck_process(
config, marker, url, 1 - baseline['checked'], timeout,
)
after = wait_for_pipeline(config, marker, url, 1, timeout)
if baseline['checked']:
compare_persisted(baseline, after)
else:
for name, value in baseline['ids'].items():
require(after['ids'][name] == value, 'keycheck_changed_scanner_identity')
for name, value in baseline['hashes'].items():
if name != 'artifact_ids_sha256':
require(after['hashes'][name] == value, 'keycheck_changed_scanner_projection')
after['origin_instance_sha256'] = baseline['origin_instance_sha256']
from runtime_security import atomic_write_private_json
atomic_write_private_json(str(RESULT), after)
return {'counts': {'ok': 1, 'http_requests': summary['counts']['http_requests'], **after['counts']},
'hashes': after['hashes']}
def _finish_remote_full_race(config, marker, url, timeout):
import asyncio
from worker_api import WorkerService
evidence = json.loads(read_bytes(REMOTE_FULL))
require(evidence.get('schema') == 1, 'remote_full_evidence_required')
service = WorkerService(
url, str(BUNDLES), lambda *_args: None, max_bundle_bytes=MAX_BYTES,
)
payloads = {}
for role in ('expired', 'winner'):
item = evidence[role]
path = REMOTE_CLIENT_BUNDLES / 'ready' / item['bundle_id'][:2] / (
item['bundle_id'] + '.trb'
)
payloads[role] = read_bytes(path)
require(digest(payloads[role]) == item['payload_sha256'],
'remote_full_client_bundle_changed')
async def race_uploads():
barrier = asyncio.Barrier(2)
async def upload(role):
await barrier.wait()
return await invoke_bundle_upload(
service, remote_device_token(role),
evidence[role]['reservation_id'], payloads[role],
)
return await asyncio.gather(upload('expired'), upload('winner'))
expired_result, winner_result = asyncio.run(race_uploads())
require(expired_result[:2] == (409, 'resolution_conflict')
and expired_result[2] is None,
'remote_full_concurrent_stale_upload')
require(winner_result[0] == 201 and winner_result[1] is None
and winner_result[2]['resolution'] == 'bundle_accepted'
and winner_result[2]['payload_sha256'] == evidence['winner']['payload_sha256'],
'remote_full_concurrent_winner_upload')
before_keycheck = wait_for_remote_full(
config, marker, url, evidence, False, timeout,
)
summary = run_keycheck_process(config, marker, url, 0, timeout)
complete = wait_for_remote_full(
config, marker, url, evidence, True, timeout,
)
require(before_keycheck['scan_id'] == complete['scan_id']
and before_keycheck['candidate_id'] == complete['candidate_id']
and complete['receipt_id'] == winner_result[2]['receipt_id'],
'remote_full_processing_identity')
public = {
'concurrent_uploads': 2, 'authoritative_receipts': 1,
'authoritative_scans': 1, 'expired_losers': 1,
'completed_winners': 1, 'cached_keychecks': 1,
'provider_http_requests': summary['counts']['http_requests'],
'projection_streams': len(complete['projection_hashes']),
}
return {
'counts': {'ok': 1, **public},
'hashes': {
'remote_full_finish_sha256': digest(json_bytes(public)),
'remote_full_parity_sha256': evidence['parity_sha256'],
**complete['projection_hashes'],
},
}
def finish_remote_full_race(config, marker, url, timeout):
try:
return _finish_remote_full_race(config, marker, url, timeout)
except E2EFailure:
raise
except Exception as exc:
trace = exc.__traceback__
origins = []
while trace:
origins.append(trace.tb_frame.f_code.co_name)
trace = trace.tb_next
origin = '_'.join([type(exc).__name__, *origins[-3:]])
origin = re.sub(r'[^a-z0-9_]', '_', origin.lower())[:90]
raise E2EFailure('remote_full_finish_' + origin) from None
def main(argv=None):
class Parser(argparse.ArgumentParser):
def error(self, message):
raise E2EFailure('invalid_arguments')
try:
parser = Parser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter)
parser.add_argument('mode', choices=(
'prepare', 'run-local-pipeline', 'assert-pipeline',
'keycheck-fixture', 'assert-persisted',
'assert-remote-recovery', 'prepare-remote-transport',
'assert-remote-transport-replay', 'prepare-remote-full-race',
'finish-remote-full-race', 'prepare-dockerhub-canary',
'finish-dockerhub-canary', '_keycheck-child',
))
parser.add_argument('--config', default=str(CONFIG))
parser.add_argument('--timeout', type=float, default=180)
args = parser.parse_args(argv)
require(math.isfinite(args.timeout) and 1 <= args.timeout <= 600, 'bounded_timeout_required')
require_container(args.config)
neutralize_environment(args.mode)
# Application/provider logging can contain masked secrets or source data.
# The only terminal output from this driver is the explicit safe summary.
logging.disable(logging.CRITICAL)
with open(os.devnull, 'w', encoding='utf-8') as quiet, redirect_stdout(quiet), redirect_stderr(quiet):
bootstrap = runpy.run_path(str(APP / 'child_bootstrap.py'))
if args.mode == '_keycheck-child':
bootstrap['_authenticate']('keycheck-provider')
bootstrap['_enable_dependency_paths']('keycheck-provider' if args.mode == '_keycheck-child' else 'scanner')
sys.path.insert(0, str(APP))
if args.mode == 'prepare':
summary = prepare()
else:
config, marker = load_fixture()
url = database_environment()
if args.mode == '_keycheck-child':
summary = keycheck_child(bootstrap)
elif args.mode == 'run-local-pipeline':
summary = run_local_pipeline_fixture(config, marker, url)
elif args.mode == 'assert-pipeline':
snapshot = wait_for_pipeline(config, marker, url, 0, args.timeout)
if RESULT.exists():
compare_persisted(json.loads(read_bytes(RESULT)), snapshot)
else:
from runtime_security import write_private_json_exclusive
write_private_json_exclusive(str(RESULT), snapshot)
summary = {'counts': {'ok': 1, **snapshot['counts']}, 'hashes': snapshot['hashes']}
elif args.mode == 'assert-remote-recovery':
summary = assert_remote_recovery(config, marker, url)
elif args.mode == 'prepare-remote-transport':
summary = prepare_remote_transport(
config, marker, url, args.timeout,
)
elif args.mode == 'assert-remote-transport-replay':
summary = assert_remote_transport_replay(
config, marker, url, args.timeout,
)
elif args.mode == 'prepare-remote-full-race':
summary = prepare_remote_full_race(config, marker, url)
elif args.mode == 'finish-remote-full-race':
summary = finish_remote_full_race(
config, marker, url, args.timeout,
)
elif args.mode == 'prepare-dockerhub-canary':
summary = prepare_dockerhub_canary(
config, marker, url, args.timeout,
)
elif args.mode == 'finish-dockerhub-canary':
summary = finish_dockerhub_canary(
config, marker, url, args.timeout,
)
else:
baseline = json.loads(read_bytes(RESULT))
require(baseline.get('checked') in (0, 1), 'pipeline_baseline_required')
snapshot = wait_for_pipeline(config, marker, url, baseline['checked'], args.timeout)
compare_persisted(baseline, snapshot, require_restart=args.mode == 'assert-persisted')
if args.mode == 'keycheck-fixture':
summary = run_keycheck(config, marker, url, baseline, args.timeout)
else:
summary = {'counts': {'ok': 1, **snapshot['counts']}, 'hashes': snapshot['hashes']}
require(set(summary) == {'counts', 'hashes'}
and all(type(value) is int and value >= 0 for value in summary['counts'].values())
and all(isinstance(value, str) and re.fullmatch(r'[a-f0-9]{64}', value)
for value in summary['hashes'].values()), 'invalid_public_summary')
except Exception as exc:
check = str(exc) if isinstance(exc, E2EFailure) else 'helper_exception'
if not re.fullmatch(r'[a-z_][a-z0-9_]{0,120}', check):
check = 'helper_exception'
summary = {'counts': {'ok': 0, check: 1}, 'hashes': {}}
print(json.dumps(summary, sort_keys=True, separators=(',', ':')))
return 0 if summary['counts']['ok'] else 1
if __name__ == '__main__':
raise SystemExit(main())