3516 lines
168 KiB
Python
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())
|