import hashlib import json import os from pathlib import Path import re import sqlite3 import sys import tempfile from types import SimpleNamespace import unittest import uuid from unittest import mock ROOT = Path(__file__).resolve().parents[1] APP_DIR = ROOT / 'app' sys.path.insert(0, str(APP_DIR)) import console_runner import docker_depth_experiment as depth import scanner import scanner_db from scanner_db import ScannerDB from process_identity import current_process_identity from result_bundle import ResultBundleReader from runtime_security import ensure_private_directory QUERIES = tuple(f'resolver-query-{index:02d}' for index in range(61)) POLICY_SHA256 = 'b' * 64 PLAN_SHA256 = 'c' * 64 HOLD_SHA256 = 'd' * 64 def digest(value): return 'sha256:' + f'{value:064x}' def experiment_config(enabled=True, selector_version=depth.DOCKER_DEPTH_SELECTOR_VERSION): profile = depth.reviewed_docker_experiment_profile(selector_version) config = { 'global': { 'database_url': 'postgresql://example.invalid/truf', 'sync_file_queues': False, }, 'sources': {'dockerhub': { 'mode': 'search', 'require_digest': True, 'queries': list(QUERIES), 'pages': 30, 'per_page': 100, 'docker_platform_filter_enabled': True, 'docker_platform_os': 'linux', 'docker_platform_arch': 'amd64', 'docker_platform_candidate_tags': 20, 'docker_images_per_repository': 3, 'docker_depth_experiment': { 'experiment_key': ( 'resolver-breadth-test-v1' if selector_version == depth.DOCKER_RANK1_BREADTH_SELECTOR_VERSION else 'resolver-test-v1' ), 'enabled': enabled, 'queries': list(QUERIES), 'repositories_per_query': profile['repositories_per_query'], 'shallow_images_per_repository': profile['shallow_images_per_repository'], 'deep_repositories_per_query': profile['deep_repositories_per_query'], 'deep_images_per_repository': profile['images_per_repository'], 'target_limit': profile['target_limit'], 'selector_version': selector_version, }, }}, } return depth.validate_docker_depth_config(config).experiment def resolver_authority( enabled=True, selector_version=depth.DOCKER_DEPTH_SELECTOR_VERSION, ): return depth.docker_depth_resolver_authority( experiment_config(enabled, selector_version), POLICY_SHA256, ) class SQLitePostgresResolverShape: is_postgres = True is_sqlite = False def __init__(self, connection): self.connection = connection def execute(self, sql, params=None): sql = re.sub( r'\s+FOR (?:UPDATE(?:\s+OF\s+[a-z0-9_,.\s]+)?' r'(?:\s+SKIP LOCKED)?|SHARE)\s*$', '', str(sql).strip(), flags=re.IGNORECASE, ) return self.connection.execute(sql, params) def commit(self): return self.connection.commit() def rollback(self): return self.connection.rollback() def __getattr__(self, name): return getattr(self.connection, name) def rich_outcome( repository, graph_layers=((101, 102),), candidate_count=None, selector_version=depth.DOCKER_DEPTH_SELECTOR_VERSION, ): candidates = [] for index, values in enumerate(graph_layers, 1): manifest_digest = digest(1000 + index) layers = tuple(digest(value) for value in values) candidates.append({ 'target': f'{repository}@{manifest_digest}', 'repository': repository, 'manifest_digest': manifest_digest, 'manifest_media_type': 'application/vnd.oci.image.manifest.v1+json', 'manifest_size_bytes': 4000 + index, 'config_digest': digest(2000 + index), 'layers': layers, 'layer_descriptors': tuple({ 'digest': layer_digest, 'media_type': 'application/vnd.oci.image.layer.v1.tar+gzip', 'size': 100 + position, } for position, layer_digest in enumerate(layers, 1)), 'updated_at': 1000 - index, 'source_index': index - 1, }) selected = scanner.select_docker_layer_graphs(candidates, len(candidates)) distinct_count = candidate_count if candidate_count is not None else len({ tuple(candidate['layers']) for candidate in candidates }) selected = tuple( scanner._docker_depth_selection_evidence( record, distinct_count, selector_version, ) for record in selected ) return scanner.DockerTagResolutionOutcome( tags=tuple(record['target'] for record in selected), status='ok', remote_attempted=True, selection_records=selected, selector_version=selector_version, selector_hash=depth.canonical_selector_hash(selector_version), candidate_distinct_graph_count=distinct_count, fresh_graph_evidence=True, cache_bypassed=True, ) class DockerDepthResolverSQLiteShapeTests(unittest.TestCase): def setUp(self): self.environment = mock.patch.dict( os.environ, {'SCANNER_DB_URL': '', 'DATABASE_URL': ''}, ) self.environment.start() self.temp = tempfile.TemporaryDirectory() self.db = ScannerDB(db_path=os.path.join(self.temp.name, 'scanner.db')) def tearDown(self): self.db.close() self.temp.cleanup() self.environment.stop() def seed_cohort(self, shared_first=False, experiment=None): experiment = experiment or experiment_config() now = scanner_db.utc_now_iso() experiment_row = self.db.conn.execute( '''INSERT INTO docker_depth_experiments( experiment_key, source, state, config_sha256, ordered_queries_sha256, selector_version, selector_sha256, provenance_policy_sha256, query_count, repositories_per_query, images_per_repository, target_limit, plan_sha256, hold_manifest_sha256, created_at, updated_at ) VALUES (?, 'dockerhub', 'resolving', ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) RETURNING id''', ( experiment.experiment_key, experiment.config_hash, experiment.ordered_query_hash, experiment.selector_version, experiment.selector_hash, POLICY_SHA256, len(experiment.queries), experiment.repositories_per_query, experiment.deep_images_per_repository, experiment.target_limit, PLAN_SHA256, HOLD_SHA256, now, now, ), ).fetchone() discovery_pass = self.db.conn.execute( '''INSERT INTO docker_discovery_passes( experiment_id, pass_token, source, pass_kind, policy_sha256, ordered_queries_sha256, expected_query_count, completed_query_count, state, started_at, completed_at, created_at, updated_at ) VALUES (?, 'resolver-pass-token-0001', 'dockerhub', 'deep', ?, ?, 61, 61, 'complete', ?, ?, ?, ?) RETURNING id''', ( experiment_row['id'], POLICY_SHA256, experiment.ordered_query_hash, now, now, now, now, ), ).fetchone() member_ids = {} queue_by_repository = {} for query_ordinal, query in enumerate(QUERIES): self.db.conn.execute( '''INSERT INTO docker_depth_experiment_queries( experiment_id, source, query_ordinal, query, query_sha256, required_repository_count, selected_repository_count, created_at ) VALUES (?, 'dockerhub', ?, ?, ?, ?, 10, ?)''', ( experiment_row['id'], query_ordinal, query, depth._canonical_sha256(query), experiment.repositories_per_query, now, ), ) page = self.db.conn.execute( '''INSERT INTO docker_discovery_pages( pass_id, query, query_ordinal, page_number, result_count, admitted_count, query_complete, admission_kind, page_sha256, observed_at, created_at ) VALUES (?, ?, ?, 1, 10, 10, 1, 'main', ?, ?, ?) RETURNING id''', ( discovery_pass['id'], query, query_ordinal, f'{query_ordinal + 1:064x}', now, now, ), ).fetchone() for repository_rank in range(1, 11): repository = ( 'owner/shared-repository' if shared_first and repository_rank == 1 and query_ordinal in (0, 1) else f'owner/q{query_ordinal:02d}-repository-{repository_rank:02d}' ) queue_id = queue_by_repository.get(repository) if queue_id is None: queue_id = self.db.conn.execute( '''INSERT INTO target_queue( source, platform, query, target, normalized_target, status, resolver_state, resolver_due_at, created_at, updated_at ) VALUES ('dockerhub', 'docker', ?, ?, ?, 'deferred', 'pending', ?, ?, ?) RETURNING id''', (query, repository, repository, now, now, now), ).fetchone()['id'] queue_by_repository[repository] = queue_id self.db.conn.execute( '''INSERT INTO docker_repository_query_provenance( source, query, repository_queue_id, provenance_kind, first_observed_at, last_observed_at, first_search_rank, best_search_rank, last_search_rank, first_page_id, last_page_id, first_policy_sha256, last_policy_sha256, observation_count, fresh_observation_count, fresh_complete_observation_count, fresh_coverage_eligible, created_at, updated_at ) VALUES ('dockerhub', ?, ?, 'fresh_page', ?, ?, ?, ?, ?, ?, ?, ?, ?, 1, 1, 1, 1, ?, ?)''', ( query, queue_id, now, now, repository_rank, repository_rank, repository_rank, page['id'], page['id'], POLICY_SHA256, POLICY_SHA256, now, now, ), ) self.db.conn.execute( '''INSERT INTO docker_repository_query_observations( page_id, repository_queue_id, source, query, search_rank, observed_at ) VALUES (?, ?, 'dockerhub', ?, ?, ?)''', (page['id'], queue_id, query, repository_rank, now), ) member = self.db.conn.execute( '''INSERT INTO docker_depth_experiment_repositories( experiment_id, query_ordinal, source, query, repository_queue_id, eligibility_page_id, repository_rank, planned_is_deep_probe, is_deep_probe, work_state, created_at, updated_at ) VALUES (?, ?, 'dockerhub', ?, ?, ?, ?, ?, ?, 'pending', ?, ?) RETURNING id''', ( experiment_row['id'], query_ordinal, query, queue_id, page['id'], repository_rank, int(repository_rank == 1), int(repository_rank == 1), now, now, ), ).fetchone() member_ids[(query_ordinal, repository_rank)] = int(member['id']) authority = depth._docker_depth_authority(experiment, POLICY_SHA256) stored_experiment = self.db.conn.execute( 'SELECT * FROM docker_depth_experiments WHERE id = ?', (experiment_row['id'],), ).fetchone() plan_sha256 = depth.canonical_docker_depth_plan_hash( depth._stored_cohort_plan(self.db.conn, stored_experiment, authority) ) self.db.conn.execute( 'UPDATE docker_depth_experiments SET plan_sha256 = ? WHERE id = ?', (plan_sha256, experiment_row['id']), ) self.db.conn.execute( '''INSERT OR REPLACE INTO runtime_final_cutover( id, marker, checked_at, evidence_sha256 ) VALUES (1, ?, ?, ?)''', (scanner_db.FINAL_CUTOVER_MARKER, now, 'e' * 64), ) self.db.conn.commit() return int(experiment_row['id']), member_ids def add_repository_candidate(self, experiment_id, query_ordinal=0, rank=11): query = QUERIES[query_ordinal] repository = f'owner/q{query_ordinal:02d}-repository-{rank:02d}' now = scanner_db.utc_now_iso() page = self.db.conn.execute( '''SELECT page.id FROM docker_discovery_pages page JOIN docker_discovery_passes discovery_pass ON discovery_pass.id = page.pass_id WHERE discovery_pass.experiment_id = ? AND page.query_ordinal = ?''', (experiment_id, query_ordinal), ).fetchone() queue = self.db.conn.execute( '''INSERT INTO target_queue( source, platform, query, target, normalized_target, status, resolver_state, resolver_due_at, created_at, updated_at ) VALUES ('dockerhub', 'docker', ?, ?, ?, 'deferred', 'pending', ?, ?, ?) RETURNING id''', (query, repository, repository, now, now, now), ).fetchone() self.db.conn.execute( '''INSERT INTO docker_repository_query_provenance( source, query, repository_queue_id, provenance_kind, first_observed_at, last_observed_at, first_search_rank, best_search_rank, last_search_rank, first_page_id, last_page_id, first_policy_sha256, last_policy_sha256, observation_count, fresh_observation_count, fresh_complete_observation_count, fresh_coverage_eligible, created_at, updated_at ) VALUES ('dockerhub', ?, ?, 'fresh_page', ?, ?, ?, ?, ?, ?, ?, ?, ?, 1, 1, 1, 1, ?, ?)''', ( query, queue['id'], now, now, rank, rank, rank, page['id'], page['id'], POLICY_SHA256, POLICY_SHA256, now, now, ), ) self.db.conn.execute( '''INSERT INTO docker_repository_query_observations( page_id, repository_queue_id, source, query, search_rank, observed_at ) VALUES (?, ?, 'dockerhub', ?, ?, ?)''', (page['id'], queue['id'], query, rank, now), ) self.db.conn.commit() return int(queue['id']), repository def hold_repository_candidate(self, experiment_id, queue_id): experiment = self.db.conn.execute( 'SELECT * FROM docker_depth_experiments WHERE id = ?', (experiment_id,), ).fetchone() queue = self.db.conn.execute( 'SELECT * FROM target_queue WHERE id = ?', (queue_id,), ).fetchone() entry = { 'queue_id': int(queue_id), 'source': 'dockerhub', 'platform': 'docker', 'query': str(queue['query']), 'prior_status': str(queue['status']), 'prior_updated_at': str(queue['updated_at']), } audit = self.db._target_queue_policy_audit_sha256( 'cold', HOLD_SHA256, entry, experiment_id, ) now = scanner_db.utc_now_iso() self.db.conn.execute( 'UPDATE target_queue SET status = \'cold\', updated_at = ? WHERE id = ?', (now, queue_id), ) self.db.conn.execute( '''INSERT INTO target_queue_policy_events( queue_id, action, prior_status, next_status, source, platform, query, reason_code, config_sha256, policy_sha256, manifest_sha256, review_audit_sha256, experiment_id, prior_updated_at, created_at ) VALUES (?, 'cold', ?, 'cold', 'dockerhub', 'docker', ?, ?, ?, ?, ?, ?, ?, ?, ?)''', ( queue_id, entry['prior_status'], entry['query'], depth.DOCKER_DEPTH_HOLD_REASON, experiment['config_sha256'], POLICY_SHA256, HOLD_SHA256, audit, experiment_id, entry['prior_updated_at'], now, ), ) self.db.conn.commit() def hold_member_at_resolver_attempt_limit(self, experiment_id, member_id): now = scanner_db.utc_now_iso() self.db.conn.execute( '''UPDATE docker_depth_experiments SET state = 'held', hold_reason_code = 'resolver_attempt_limit', held_at = ?, updated_at = ? WHERE id = ?''', (now, now, experiment_id), ) self.db.conn.execute( '''UPDATE docker_depth_experiment_repositories SET work_state = 'held', resolver_attempts = 3, resolver_owner = NULL, resolver_token = NULL, resolver_expires_at = NULL, resolver_due_at = NULL, last_error_code = 'resolver_attempt_limit', updated_at = ? WHERE id = ?''', (now, member_id), ) self.db.conn.commit() def shrink_query_cohort(self, experiment_id, query_ordinal, selected_count): self.db.conn.execute( '''DELETE FROM docker_depth_experiment_repositories WHERE experiment_id = ? AND query_ordinal = ? AND repository_rank > ?''', (experiment_id, query_ordinal, selected_count), ) self.db.conn.execute( '''UPDATE docker_depth_experiment_queries SET selected_repository_count = ? WHERE experiment_id = ? AND query_ordinal = ?''', (selected_count, experiment_id, query_ordinal), ) experiment = self.db.conn.execute( 'SELECT * FROM docker_depth_experiments WHERE id = ?', (experiment_id,), ).fetchone() authority = resolver_authority() plan_sha256 = depth.canonical_docker_depth_plan_hash( depth._stored_cohort_plan(self.db.conn, experiment, authority) ) self.db.conn.execute( 'UPDATE docker_depth_experiments SET plan_sha256 = ? WHERE id = ?', (plan_sha256, experiment_id), ) self.db.conn.commit() def use_postgres_shape(self): self.db.conn = SQLitePostgresResolverShape(self.db.conn) def claim_one(self): return self.db.claim_docker_depth_experiment_resolutions( 'dockerhub', 1, 'resolver-owner', authority=resolver_authority(), final_cutover=True, )[0] def finish(self, claim, outcome, complete=True): return self.db.finish_docker_depth_experiment_resolution( 'dockerhub', claim['id'], claim['resolver_generation'], claim['resolver_token'], outcome, complete=complete, resolver_owner=claim['resolver_owner'], authority=resolver_authority(), final_cutover=True, ) def set_discovery_paused(self, paused): state = self.db.runtime_control_state() return self.db.set_runtime_discovery_paused( paused, expected_revision=state['revision'], actor='test:docker-depth-gate', operation_id=str(uuid.uuid4()), ) def test_discovery_pause_blocks_depth_claim_and_completion_without_mutation(self): _experiment_id, member_ids = self.seed_cohort() first_member_id = member_ids[(0, 1)] self.set_discovery_paused(True) self.use_postgres_shape() self.assertEqual(self.db.claim_docker_depth_experiment_resolutions( 'dockerhub', 1, 'paused-owner', authority=resolver_authority(), final_cutover=True, ), []) pending = self.db.conn.execute( '''SELECT work_state, resolver_attempts, resolver_token FROM docker_depth_experiment_repositories WHERE id = ?''', (first_member_id,), ).fetchone() self.assertEqual((pending['work_state'], pending['resolver_attempts'], pending['resolver_token']), ( 'pending', 0, None, )) self.set_discovery_paused(False) claim = self.claim_one() self.set_discovery_paused(True) with self.assertRaises(scanner_db.DiscoveryPausedError): self.finish(claim, rich_outcome(claim['target'])) resolving = self.db.conn.execute( '''SELECT work_state, resolver_token, selected_image_count FROM docker_depth_experiment_repositories WHERE id = ?''', (claim['id'],), ).fetchone() self.assertEqual(( resolving['work_state'], resolving['resolver_token'], resolving['selected_image_count'], ), ('resolving', claim['resolver_token'], 0)) def mark_members_skipped(self, experiment_id, excluded_member_ids=()): now = scanner_db.utc_now_iso() excluded_member_ids = {int(member_id) for member_id in excluded_member_ids} experiment = self.db.conn.execute( 'SELECT * FROM docker_depth_experiments WHERE id = ?', (experiment_id,), ).fetchone() for member in self.db.conn.execute( '''SELECT * FROM docker_depth_experiment_repositories WHERE experiment_id = ? ORDER BY id''', (experiment_id,), ).fetchall(): if int(member['id']) in excluded_member_ids: continue self.db.conn.execute( '''UPDATE docker_depth_experiment_repositories SET work_state = 'skipped', resolver_owner = NULL, resolver_token = NULL, resolver_expires_at = NULL, resolver_due_at = NULL, is_deep_probe = 0, candidate_distinct_graph_count = 0, selected_image_count = 0, last_error_code = ?, resolved_at = ?, updated_at = ? WHERE id = ?''', ( depth.DOCKER_DEPTH_REPOSITORY_SKIP_REASON, now, now, member['id'], ), ) repository_queue_id = int( member['replacement_repository_queue_id'] or member['repository_queue_id'] ) self.db._record_docker_depth_candidate_skip_locked( experiment, member, repository_queue_id, 'repository', int(member['replacement_count']) + 1, {'repository_queue_id': repository_queue_id}, depth.DOCKER_DEPTH_REPOSITORY_SKIP_REASON, now, ) self.db.conn.commit() def activate_targets(self, experiment_id, members, outcomes=()): now = scanner_db.utc_now_iso() outcomes = tuple(outcomes) outcome_member_ids = { members[(query_ordinal, repository_rank)] for query_ordinal, repository_rank, _outcome in outcomes } self.mark_members_skipped(experiment_id, outcome_member_ids) experiment = self.db.conn.execute( 'SELECT * FROM docker_depth_experiments WHERE id = ?', (experiment_id,), ).fetchone() target_count = 0 selection_count = 0 for query_ordinal, repository_rank, outcome in outcomes: member_id = members[(query_ordinal, repository_rank)] records = list(outcome.selection_records) self.db.conn.execute( '''UPDATE docker_depth_experiment_repositories SET work_state = 'resolved', candidate_distinct_graph_count = ?, selected_image_count = ?, last_error_code = NULL, resolved_at = ?, updated_at = ? WHERE id = ?''', ( outcome.candidate_distinct_graph_count, len(records), now, now, member_id, ), ) for record in records: normalized_target = scanner_db.normalize_target(record['target'], 'docker') queue = self.db.conn.execute( '''SELECT id FROM target_queue WHERE source = 'dockerhub' AND normalized_target = ?''', (normalized_target,), ).fetchone() if not queue: query = QUERIES[query_ordinal] queue = self.db.conn.execute( '''INSERT INTO target_queue( source, platform, query, target, normalized_target, status, created_at, updated_at ) VALUES ('dockerhub', 'docker', ?, ?, ?, 'pending', ?, ?) RETURNING id''', (query, record['target'], normalized_target, now, now), ).fetchone() manifest = self.db.conn.execute( 'SELECT id FROM docker_image_manifests WHERE target_queue_id = ?', (queue['id'],), ).fetchone() if not manifest: manifest = self.db.conn.execute( '''INSERT INTO docker_image_manifests( target_queue_id, source, repository, manifest_digest, manifest_media_type, config_digest, graph_sha256, manifest_size_bytes, layer_count, resolved_at, created_at ) VALUES (?, 'dockerhub', ?, ?, ?, ?, ?, ?, ?, ?, ?) RETURNING id''', ( queue['id'], record['repository'], record['manifest_digest'], record['manifest_media_type'], record['config_digest'], record['graph_sha256'], record['manifest_size_bytes'], record['layer_count'], now, now, ), ).fetchone() for layer in record['layer_metadata']: self.db.conn.execute( '''INSERT INTO docker_manifest_layers( manifest_id, position_from_base, position_from_top, layer_digest, media_type, layer_size_bytes, descriptor_sha256, created_at ) VALUES (?, ?, ?, ?, ?, ?, ?, ?)''', ( manifest['id'], layer['position_from_base'], layer['position_from_top'], layer['digest'], layer['media_type'], layer['size_bytes'], layer['descriptor_sha256'], now, ), ) wave, order = self.db._docker_depth_dispatch_position( query_ordinal, repository_rank, record['image_rank'], len(QUERIES), ) target = self.db.conn.execute( '''SELECT * FROM docker_depth_experiment_targets WHERE experiment_id = ? AND target_queue_id = ?''', (experiment_id, queue['id']), ).fetchone() if not target: target_count += 1 target = self.db.conn.execute( '''INSERT INTO docker_depth_experiment_targets( experiment_id, target_queue_id, manifest_id, counter_ordinal, state, dispatch_wave, dispatch_order, created_at, updated_at ) VALUES (?, ?, ?, ?, 'pending', ?, ?, ?, ?) RETURNING *''', ( experiment_id, queue['id'], manifest['id'], target_count, wave, order, now, now, ), ).fetchone() elif (wave, order) < ( int(target['dispatch_wave']), int(target['dispatch_order']) ): target = self.db.conn.execute( '''UPDATE docker_depth_experiment_targets SET dispatch_wave = ?, dispatch_order = ?, updated_at = ? WHERE id = ? RETURNING *''', (wave, order, now, target['id']), ).fetchone() self.db.conn.execute( '''INSERT INTO docker_depth_experiment_selections( experiment_id, query_ordinal, experiment_repository_id, experiment_target_id, image_rank, selection_reason, selection_evidence_sha256, graph_sha256, selected_at, created_at ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)''', ( experiment_id, query_ordinal, member_id, target['id'], record['image_rank'], record['selection_reason'], record['selection_evidence_sha256'], record['graph_sha256'], now, now, ), ) selection_count += 1 for query_ordinal in range(len(QUERIES)): deep = self.db.conn.execute( '''SELECT id FROM docker_depth_experiment_repositories WHERE experiment_id = ? AND query_ordinal = ? AND work_state = 'resolved' AND selected_image_count >= 1 ORDER BY candidate_distinct_graph_count DESC, repository_rank, repository_queue_id, id LIMIT 1''', (experiment_id, query_ordinal), ).fetchone() self.db.conn.execute( '''UPDATE docker_depth_experiment_repositories SET is_deep_probe = CASE WHEN id = ? THEN 1 ELSE 0 END WHERE experiment_id = ? AND query_ordinal = ?''', ( int(deep['id']) if deep else None, experiment_id, query_ordinal, ), ) target_count = self.db.conn.execute( '''SELECT COUNT(*) AS count FROM docker_depth_experiment_targets WHERE experiment_id = ?''', (experiment_id,), ).fetchone()['count'] self.db.conn.execute( '''UPDATE docker_depth_experiments SET state = 'active', target_count = ?, selection_count = ?, activated_at = ?, updated_at = ? WHERE id = ?''', (target_count, selection_count, now, now, experiment_id), ) experiment = self.db.conn.execute( 'SELECT * FROM docker_depth_experiments WHERE id = ?', (experiment_id,), ).fetchone() selection_sha256 = self.db._docker_depth_runtime_selection_sha256_locked( experiment, resolver_authority(), ) self.db.conn.execute( 'UPDATE docker_depth_experiments SET selection_sha256 = ? WHERE id = ?', (selection_sha256, experiment_id), ) self.db.conn.commit() def ready_ingester(self, supervisor='dispatch-supervisor'): now = scanner_db.utc_now_iso() self.db.conn.execute( '''INSERT OR REPLACE INTO pipeline_leases( worker_name, generation, lease_token, supervisor_instance_id, state, lease_expires_at, updated_at ) VALUES ('result_ingester', 1, 'ingester-token', ?, 'ready', '2999-01-01T00:00:00+00:00', ?)''', (supervisor, now), ) self.db.conn.commit() def claim_scan(self, index, authority=None, limits=None, max_attempts=3): return self.db.reserve_and_claim_target( 'dockerhub', 'docker', current_process_identity(), 'dispatch-supervisor', 1024, 2048, 1, 128, capacity_limits=limits or { 'bundle_items': 100, 'bundle_bytes': 1024 * 1024, 'projection_items': 100, 'projection_bytes': 1024 * 1024, 'keycheck_items': 100, 'keycheck_bytes': 1024 * 1024, }, max_attempts=max_attempts, reservation_token=f'dispatch-reservation-{index}', bundle_id=f'{index + 100:032x}', scan_event_id=f'{index + 200:032x}', docker_depth_authority=authority or resolver_authority(), final_cutover=True, ) def complete_scan_claim(self, claim): now = scanner_db.utc_now_iso() scan = self.db.conn.execute( '''INSERT INTO target_scans( scan_event_id, scan_event_hash, queue_id, claim_lease_token, queue_completion_applied, queue_completion_disposition, source, query, target, normalized_target, scan_type, status, result_reservation_id, created_at ) VALUES (?, ?, ?, ?, 1, 'applied', 'dockerhub', ?, ?, ?, 'docker', 'clean', ?, ?) RETURNING id''', ( claim['scan_event_id'], f"{claim['reservation_id']:064x}", claim['queue_id'], claim['claim_lease_token'], claim['query'], claim['target'], claim['normalized_target'], claim['reservation_id'], now, ), ).fetchone() self.db.conn.execute( '''UPDATE target_queue SET status = 'done', target_scan_id = ?, lease_owner = NULL, lease_token = NULL, claim_batch = NULL, leased_at = NULL, lease_expires_at = NULL, current_result_reservation_id = NULL, claim_event_id = NULL, completed_at = ?, updated_at = ? WHERE id = ?''', (scan['id'], now, now, claim['queue_id']), ) self.db.conn.execute( '''UPDATE result_reservations SET state = 'db_committed', updated_at = ? WHERE id = ?''', (now, claim['reservation_id']), ) self.db.conn.execute( '''UPDATE docker_depth_experiment_scan_bindings SET target_scan_id = ?, state = 'completed', scan_bound_at = ?, completed_at = ? WHERE reservation_id = ?''', (scan['id'], now, now, claim['reservation_id']), ) self.db.conn.execute( '''UPDATE docker_depth_experiment_targets SET state = 'done', terminal_at = ?, updated_at = ? WHERE id = ?''', (now, now, claim['experiment_target_id']), ) self.db.conn.commit() def add_owned_hold_event(self, experiment_id): now = scanner_db.utc_now_iso() queue = self.db.conn.execute( '''INSERT INTO target_queue( source, platform, query, target, normalized_target, status, resolver_state, created_at, updated_at ) VALUES ('dockerhub', 'docker', ?, 'owner/held-extra', 'owner/held-extra', 'cold', 'pending', ?, ?) RETURNING id''', (QUERIES[0], now, now), ).fetchone() entry = { 'queue_id': int(queue['id']), 'source': 'dockerhub', 'platform': 'docker', 'query': QUERIES[0], 'prior_status': 'deferred', 'prior_updated_at': now, } audit = self.db._target_queue_policy_audit_sha256( 'cold', HOLD_SHA256, entry, experiment_id, ) event = self.db.conn.execute( '''INSERT INTO target_queue_policy_events( queue_id, action, prior_status, next_status, source, platform, query, reason_code, config_sha256, policy_sha256, manifest_sha256, review_audit_sha256, experiment_id, prior_updated_at, created_at ) VALUES (?, 'cold', 'deferred', 'cold', 'dockerhub', 'docker', ?, ?, ?, ?, ?, ?, ?, ?, ?) RETURNING id''', ( queue['id'], QUERIES[0], depth.DOCKER_DEPTH_HOLD_REASON, experiment_config().config_hash, POLICY_SHA256, HOLD_SHA256, audit, experiment_id, now, now, ), ).fetchone() self.db.conn.commit() return int(queue['id']), int(event['id']), now def test_claims_are_fenced_fair_and_nonfinal_authority_holds(self): experiment_id, _members = self.seed_cohort() self.assertEqual(self.db.claim_docker_depth_experiment_resolutions( 'dockerhub', 1, 'owner', authority=resolver_authority(), final_cutover=True, ), []) self.use_postgres_shape() claims = self.db.claim_docker_depth_experiment_resolutions( 'dockerhub', 3, 'owner', authority=resolver_authority(), final_cutover=True, ) self.assertEqual( [(row['repository_rank'], row['query_ordinal'], row['stage']) for row in claims], [(1, 0, 'breadth'), (1, 1, 'breadth'), (1, 2, 'breadth')], ) before = self.db.conn.execute( "SELECT COUNT(*) AS count FROM docker_depth_experiment_repositories WHERE work_state = 'resolving'" ).fetchone()['count'] self.assertEqual(self.db.claim_docker_depth_experiment_resolutions( 'dockerhub', 1, 'owner', authority=resolver_authority(), final_cutover=False, ), []) after = self.db.conn.execute( "SELECT COUNT(*) AS count FROM docker_depth_experiment_repositories WHERE work_state = 'resolving'" ).fetchone()['count'] held = self.db.conn.execute( '''SELECT state, hold_reason_code FROM docker_depth_experiments WHERE id = ?''', (experiment_id,), ).fetchone() self.assertEqual((before, after), (3, 0)) self.assertEqual(tuple(held), ('held', 'final_cutover_unavailable')) def test_false_to_true_activation_preserves_frozen_config_authority(self): disabled = resolver_authority(False) enabled = resolver_authority(True) self.assertEqual(disabled['config_sha256'], enabled['config_sha256']) experiment_id, _members = self.seed_cohort() self.db.conn.execute( "UPDATE docker_depth_experiments SET state = 'holding' WHERE id = ?", (experiment_id,), ) self.db.conn.commit() frozen = self.db.conn.execute( 'SELECT config_sha256 FROM docker_depth_experiments WHERE id = ?', (experiment_id,), ).fetchone()['config_sha256'] self.assertEqual(frozen, disabled['config_sha256']) self.use_postgres_shape() claims = self.db.claim_docker_depth_experiment_resolutions( 'dockerhub', 1, 'activation-owner', authority=enabled, final_cutover=True, ) self.assertEqual(len(claims), 1) state = self.db.conn.execute( '''SELECT state, hold_reason_code FROM docker_depth_experiments WHERE id = ?''', (experiment_id,), ).fetchone() self.assertEqual(tuple(state), ('resolving', None)) def test_clean_stale_resolver_hold_resumes_on_the_next_claim(self): experiment_id, _members = self.seed_cohort() self.use_postgres_shape() first = self.claim_one() self.db.conn.execute( '''UPDATE docker_depth_experiment_repositories SET resolver_expires_at = '2000-01-01T00:00:00+00:00' WHERE id = ?''', (first['id'],), ) self.db.conn.commit() self.assertEqual(self.db.claim_docker_depth_experiment_resolutions( 'dockerhub', 1, 'second-owner', authority=resolver_authority(), final_cutover=True, ), []) held = self.db.conn.execute( '''SELECT state, hold_reason_code FROM docker_depth_experiments WHERE id = ?''', (experiment_id,), ).fetchone() cleaned = self.db.conn.execute( '''SELECT work_state, resolver_owner, resolver_token, resolver_expires_at, resolver_due_at, resolver_attempts FROM docker_depth_experiment_repositories WHERE id = ?''', (first['id'],), ).fetchone() self.assertEqual(tuple(held), ('held', 'stale_resolver_fence')) self.assertEqual(tuple(cleaned), ('pending', None, None, None, None, 1)) resumed = self.db.claim_docker_depth_experiment_resolutions( 'dockerhub', 1, 'second-owner', authority=resolver_authority(), final_cutover=True, ) self.assertEqual(len(resumed), 1) self.assertEqual(resumed[0]['id'], first['id']) self.assertEqual(resumed[0]['resolver_attempts'], 2) state = self.db.conn.execute( '''SELECT state, hold_reason_code, held_at FROM docker_depth_experiments WHERE id = ?''', (experiment_id,), ).fetchone() self.assertEqual(tuple(state), ('resolving', None, None)) def test_nonresolver_hold_is_never_automatically_resumed(self): experiment_id, _members = self.seed_cohort() now = scanner_db.utc_now_iso() self.db.conn.execute( '''UPDATE docker_depth_experiments SET state = 'held', hold_reason_code = 'config_hash_drift', held_at = ?, updated_at = ? WHERE id = ?''', (now, now, experiment_id), ) self.db.conn.commit() self.use_postgres_shape() claims = self.db.claim_docker_depth_experiment_resolutions( 'dockerhub', 1, 'owner', authority=resolver_authority(), final_cutover=True, ) self.assertEqual(claims, []) state = self.db.conn.execute( '''SELECT state, hold_reason_code, held_at FROM docker_depth_experiments WHERE id = ?''', (experiment_id,), ).fetchone() self.assertEqual(tuple(state), ('held', 'config_hash_drift', now)) def test_all_610_rank_one_memberships_activate_but_require_scan_evidence(self): experiment_id, members = self.seed_cohort() outcomes = [] for query_ordinal in range(61): for repository_rank in range(1, 11): repository = ( f'owner/q{query_ordinal:02d}-repository-{repository_rank:02d}' ) outcomes.append(( query_ordinal, repository_rank, rich_outcome( repository, ((9000 + query_ordinal * 10 + repository_rank,),), ), )) self.activate_targets(experiment_id, members, outcomes) self.db.conn.execute( '''UPDATE docker_depth_experiments SET state = 'resolving', selection_sha256 = NULL WHERE id = ?''', (experiment_id,), ) self.db.conn.commit() self.use_postgres_shape() result = self.db.refresh_docker_depth_experiment_state( resolver_authority(), final_cutover=True, ) self.assertEqual(result['status'], 'active') row = self.db.conn.execute( '''SELECT state, target_count, selection_count, selection_sha256 FROM docker_depth_experiments WHERE id = ?''', (experiment_id,), ).fetchone() self.assertEqual( (row['state'], row['target_count'], row['selection_count']), ('active', 610, 610), ) self.assertRegex(row['selection_sha256'], r'^[a-f0-9]{64}$') now = scanner_db.utc_now_iso() self.db.conn.execute( '''UPDATE docker_depth_experiment_targets SET state = 'done', terminal_at = ?, updated_at = ? WHERE experiment_id = ?''', (now, now, experiment_id), ) self.db.conn.execute( '''UPDATE target_queue SET status = 'done', completed_at = ?, updated_at = ? WHERE id IN ( SELECT target_queue_id FROM docker_depth_experiment_targets WHERE experiment_id = ? )''', (now, now, experiment_id), ) self.db.conn.commit() completed = self.db.refresh_docker_depth_experiment_state( resolver_authority(), final_cutover=True, ) self.assertEqual(completed['status'], 'held') self.assertEqual(completed['reason'], 'reservation_binding_drift') def test_released_experiment_does_not_intercept_ordinary_claims(self): experiment_id, _members = self.seed_cohort() now = scanner_db.utc_now_iso() digest = 'a' * 64 queue_id = self.db.conn.execute( '''INSERT INTO target_queue( source, platform, query, target, normalized_target, status, created_at, updated_at ) VALUES ('dockerhub', 'docker', 'ordinary', ?, ?, 'pending', ?, ?) RETURNING id''', (f'owner/ordinary@sha256:{digest}', f'owner/ordinary@sha256:{digest}', now, now), ).fetchone()['id'] self.db.conn.execute( '''UPDATE docker_depth_experiments SET state = 'released', released_at = ?, updated_at = ? WHERE id = ?''', (now, now, experiment_id), ) self.db.conn.commit() self.ready_ingester() self.use_postgres_shape() claim = self.claim_scan(900, authority=resolver_authority()) self.assertIsNotNone(claim) self.assertEqual(claim['queue_id'], queue_id) binding = self.db.conn.execute( '''SELECT COUNT(*) AS count FROM docker_depth_experiment_scan_bindings WHERE reservation_id = ?''', (claim['reservation_id'],), ).fetchone() self.assertEqual(binding['count'], 0) def test_reviewed_holding_activates_with_explicit_image_scarcity(self): experiment_id, members = self.seed_cohort() selected_member = members[(0, 1)] self.mark_members_skipped(experiment_id, (selected_member,)) self.db.conn.execute( "UPDATE docker_depth_experiments SET state = 'holding' WHERE id = ?", (experiment_id,), ) self.db.conn.commit() self.use_postgres_shape() claim = self.claim_one() self.assertEqual(claim['id'], selected_member) result = self.finish(claim, rich_outcome(claim['target'])) self.assertEqual(result['experiment_state'], 'active') state = self.db.conn.execute( '''SELECT state, hold_reason_code, released_at FROM docker_depth_experiments WHERE id = ?''', (experiment_id,), ).fetchone() self.assertEqual(tuple(state), ('active', None, None)) def test_rich_completion_persists_layers_then_stale_fence_is_read_only(self): experiment_id, _members = self.seed_cohort() self.use_postgres_shape() claim = self.claim_one() outcome = rich_outcome(claim['target'], ((101, 102, 101),)) result = self.finish(claim, outcome) self.assertEqual(result['status'], 'resolved') member = self.db.conn.execute( '''SELECT work_state, candidate_distinct_graph_count, selected_image_count, resolver_token FROM docker_depth_experiment_repositories WHERE id = ?''', (claim['id'],), ).fetchone() layers = self.db.conn.execute( '''SELECT layer_digest, position_from_base, position_from_top, media_type, layer_size_bytes, descriptor_sha256 FROM docker_manifest_layers ORDER BY position_from_base''' ).fetchall() experiment = self.db.conn.execute( '''SELECT target_count, selection_count FROM docker_depth_experiments WHERE id = ?''', (experiment_id,), ).fetchone() self.assertEqual(tuple(member), ('resolved', 1, 1, None)) self.assertEqual( [(row['layer_digest'], row['position_from_base'], row['position_from_top']) for row in layers], [(digest(101), 1, 3), (digest(102), 2, 2), (digest(101), 3, 1)], ) self.assertTrue(all(row['media_type'] and row['descriptor_sha256'] for row in layers)) self.assertEqual(tuple(experiment), (1, 1)) stale = self.db.finish_docker_depth_experiment_resolution( 'dockerhub', claim['id'], claim['resolver_generation'], claim['resolver_token'] + '-stale', outcome, complete=True, resolver_owner=claim['resolver_owner'], authority=resolver_authority(), final_cutover=True, ) self.assertEqual(stale, {'status': 'stale', 'committed': False}) state = self.db.conn.execute( '''SELECT state, hold_reason_code FROM docker_depth_experiments WHERE id = ?''', (experiment_id,), ).fetchone() self.assertEqual(tuple(state), ('resolving', None)) self.assertEqual(self.db.conn.execute( 'SELECT COUNT(*) AS count FROM docker_depth_experiment_targets' ).fetchone()['count'], 1) self.assertFalse(self.db.has_claimable_targets_v2( 'dockerhub', 'docker', max_attempts=3, )) def test_rank1_breadth_claims_resolved_target_while_cohort_is_resolving(self): selector = depth.DOCKER_RANK1_BREADTH_SELECTOR_VERSION experiment = experiment_config(selector_version=selector) authority = resolver_authority(selector_version=selector) experiment_id, _members = self.seed_cohort(experiment=experiment) self.use_postgres_shape() with mock.patch.object( depth, '_require_released_policy_history', return_value=True, ): claims = self.db.claim_docker_depth_experiment_resolutions( 'dockerhub', 1, 'breadth-owner', authority=authority, final_cutover=True, ) self.assertEqual(len(claims), 1) resolver_claim = claims[0] outcome = rich_outcome( resolver_claim['target'], ((601,),), selector_version=selector, ) resolved = self.db.finish_docker_depth_experiment_resolution( 'dockerhub', resolver_claim['id'], resolver_claim['resolver_generation'], resolver_claim['resolver_token'], outcome, complete=True, resolver_owner=resolver_claim['resolver_owner'], authority=authority, final_cutover=True, ) self.assertEqual( (resolved['status'], resolved['experiment_state']), ('resolved', 'resolving'), ) self.ready_ingester() scan_claim = self.claim_scan(901, authority=authority) self.assertIsNotNone(scan_claim) self.assertEqual(scan_claim['dispatch_wave'], 1) state = self.db.conn.execute( '''SELECT state FROM docker_depth_experiments WHERE id = ?''', (experiment_id,), ).fetchone() target = self.db.conn.execute( '''SELECT state FROM docker_depth_experiment_targets WHERE id = ?''', (scan_claim['experiment_target_id'],), ).fetchone() self.assertEqual((state['state'], target['state']), ('resolving', 'reserved')) def test_rank1_breadth_does_not_schedule_redundant_deep_resolution(self): selector = depth.DOCKER_RANK1_BREADTH_SELECTOR_VERSION experiment = experiment_config(selector_version=selector) experiment_id, members = self.seed_cohort(experiment=experiment) now = scanner_db.utc_now_iso() self.db.conn.execute( '''UPDATE docker_depth_experiment_repositories SET work_state = 'resolved', selected_image_count = 1, candidate_distinct_graph_count = 2, resolved_at = ?, updated_at = ? WHERE experiment_id = ? AND query_ordinal = 0''', (now, now, experiment_id), ) self.db.conn.commit() self.use_postgres_shape() experiment_row = self.db.conn.execute( 'SELECT * FROM docker_depth_experiments WHERE id = ?', (experiment_id,), ).fetchone() member = self.db.conn.execute( 'SELECT * FROM docker_depth_experiment_repositories WHERE id = ?', (members[(0, 10)],), ).fetchone() with mock.patch.object( self.db, '_docker_depth_terminal_repository_evidence_locked', return_value=(None, 'e' * 64), ): self.assertTrue(self.db._finalize_docker_depth_query_breadth_locked( experiment_row, member, now, )) query_members = self.db.conn.execute( '''SELECT work_state, is_deep_probe FROM docker_depth_experiment_repositories WHERE experiment_id = ? AND query_ordinal = 0''', (experiment_id,), ).fetchall() self.assertEqual( [(row['work_state'], row['is_deep_probe']) for row in query_members], [('resolved', 0)] * 10, ) def test_legacy_depth_does_not_claim_target_until_cohort_is_active(self): experiment_id, _members = self.seed_cohort() self.use_postgres_shape() resolver_claim = self.claim_one() resolved = self.finish( resolver_claim, rich_outcome(resolver_claim['target'], ((602,),)), ) self.assertEqual(resolved['experiment_state'], 'resolving') self.ready_ingester() scan_claim = self.claim_scan(902) self.assertIsNone(scan_claim) state = self.db.conn.execute( '''SELECT experiment.state, target.state AS target_state FROM docker_depth_experiments experiment JOIN docker_depth_experiment_targets target ON target.experiment_id = experiment.id WHERE experiment.id = ?''', (experiment_id,), ).fetchone() self.assertEqual((state['state'], state['target_state']), ('resolving', 'pending')) def test_resolver_renewal_is_fenced_and_remote_failures_skip_at_ceiling(self): experiment_id, _members = self.seed_cohort() self.use_postgres_shape() claim = self.claim_one() renewed = self.db.renew_docker_depth_experiment_resolution( 'dockerhub', claim['id'], claim['resolver_generation'], claim['resolver_token'], resolver_owner=claim['resolver_owner'], lease_seconds=600, authority=resolver_authority(), final_cutover=True, ) self.assertTrue(renewed['renewed']) stale = self.db.renew_docker_depth_experiment_resolution( 'dockerhub', claim['id'], claim['resolver_generation'], claim['resolver_token'] + '-stale', resolver_owner=claim['resolver_owner'], lease_seconds=600, authority=resolver_authority(), final_cutover=True, ) self.assertEqual(stale, {'status': 'stale', 'renewed': False}) current = self.db.conn.execute( '''SELECT resolver_token FROM docker_depth_experiment_repositories WHERE id = ?''', (claim['id'],), ).fetchone() self.assertEqual(current['resolver_token'], claim['resolver_token']) for expected in ('deferred', 'deferred', 'skipped'): result = self.finish( claim, None, complete=False, ) self.assertEqual(result['status'], expected) if expected != 'skipped': self.db.conn.execute( '''UPDATE docker_depth_experiment_repositories SET resolver_due_at = '2000-01-01T00:00:00+00:00' WHERE id = ?''', (claim['id'],), ) self.db.conn.commit() claim = self.claim_one() self.assertEqual( result['reason'], depth.DOCKER_DEPTH_REMOTE_UNAVAILABLE_SKIP_REASON, ) final = self.db.conn.execute( '''SELECT experiment.state, experiment.hold_reason_code, member.work_state, member.resolver_attempts FROM docker_depth_experiments experiment JOIN docker_depth_experiment_repositories member ON member.experiment_id = experiment.id WHERE experiment.id = ? AND member.id = ?''', (experiment_id, claim['id']), ).fetchone() self.assertEqual(tuple(final), ('resolving', None, 'skipped', 3)) skip = self.db.conn.execute( '''SELECT reason_code FROM docker_depth_experiment_candidate_skips WHERE experiment_repository_id = ?''', (claim['id'],), ).fetchone() self.assertEqual( skip['reason_code'], depth.DOCKER_DEPTH_REMOTE_UNAVAILABLE_SKIP_REASON, ) def test_existing_attempt_limit_hold_is_skipped_and_resolution_continues(self): experiment_id, members = self.seed_cohort() member_id = members[(0, 1)] self.hold_member_at_resolver_attempt_limit(experiment_id, member_id) self.use_postgres_shape() claimed = self.db.claim_docker_depth_experiment_resolutions( 'dockerhub', 1, 'resolver-worker', lease_seconds=300, authority=resolver_authority(), final_cutover=True, ) self.assertEqual(len(claimed), 1) self.assertNotEqual(claimed[0]['id'], member_id) recovered = self.db.conn.execute( '''SELECT experiment.state, experiment.hold_reason_code, member.work_state, member.last_error_code, member.resolver_attempts FROM docker_depth_experiments experiment JOIN docker_depth_experiment_repositories member ON member.experiment_id = experiment.id WHERE experiment.id = ? AND member.id = ?''', (experiment_id, member_id), ).fetchone() self.assertEqual(tuple(recovered), ( 'resolving', None, 'skipped', depth.DOCKER_DEPTH_REMOTE_UNAVAILABLE_SKIP_REASON, 3, )) def test_reviewed_attempt_limit_disposition_replaces_and_is_idempotent(self): experiment_id, members = self.seed_cohort() replacement_id, _repository = self.add_repository_candidate(experiment_id) self.hold_repository_candidate(experiment_id, replacement_id) member_id = members[(0, 1)] self.hold_member_at_resolver_attempt_limit(experiment_id, member_id) self.use_postgres_shape() experiment = experiment_config() manifest, manifest_sha256 = ( depth.generate_docker_depth_resolver_disposition_manifest( self.db, experiment, POLICY_SHA256, ) ) self.assertEqual(manifest['entries'][0]['outcome'], 'replaced') self.assertNotIn('normalized_target', json.dumps(manifest)) result = depth.apply_docker_depth_resolver_disposition_manifest( self.db, experiment, POLICY_SHA256, manifest, manifest_sha256, ) self.assertEqual( (result['applied'], result['duplicates'], result['outcome']), (1, 0, 'replaced'), ) member = self.db.conn.execute( '''SELECT work_state, resolver_attempts, replacement_repository_queue_id FROM docker_depth_experiment_repositories WHERE id = ?''', (member_id,), ).fetchone() self.assertEqual(tuple(member), ('pending', 0, replacement_id)) replay = depth.apply_docker_depth_resolver_disposition_manifest( self.db, experiment, POLICY_SHA256, manifest, manifest_sha256, ) self.assertEqual((replay['applied'], replay['duplicates']), (0, 1)) self.assertEqual(self.db.conn.execute( 'SELECT COUNT(*) AS count FROM docker_depth_resolver_dispositions' ).fetchone()['count'], 1) def test_reviewed_attempt_limit_disposition_records_terminal_scarcity(self): experiment_id, members = self.seed_cohort() member_id = members[(0, 1)] self.hold_member_at_resolver_attempt_limit(experiment_id, member_id) self.use_postgres_shape() experiment = experiment_config() manifest, manifest_sha256 = ( depth.generate_docker_depth_resolver_disposition_manifest( self.db, experiment, POLICY_SHA256, ) ) self.assertEqual(manifest['entries'][0]['outcome'], 'skipped') result = depth.apply_docker_depth_resolver_disposition_manifest( self.db, experiment, POLICY_SHA256, manifest, manifest_sha256, ) self.assertEqual((result['applied'], result['outcome']), (1, 'skipped')) member = self.db.conn.execute( '''SELECT work_state, resolver_attempts, selected_image_count, last_error_code FROM docker_depth_experiment_repositories WHERE id = ?''', (member_id,), ).fetchone() self.assertEqual(tuple(member), ( 'skipped', 3, 0, depth.DOCKER_DEPTH_REMOTE_UNAVAILABLE_SKIP_REASON, )) self.assertEqual(self.db.conn.execute( 'SELECT COUNT(*) AS count FROM docker_depth_experiment_targets' ).fetchone()['count'], 0) self.assertEqual(self.db.conn.execute( 'SELECT COUNT(*) AS count FROM docker_depth_resolver_dispositions' ).fetchone()['count'], 1) def test_shared_immutable_target_deduplicates_physical_capacity(self): experiment_id, _members = self.seed_cohort(shared_first=True) self.use_postgres_shape() first = self.claim_one() self.assertEqual(first['target'], 'owner/shared-repository') self.assertEqual(self.finish(first, rich_outcome(first['target']))['status'], 'resolved') second = self.claim_one() self.assertEqual((second['target'], second['query_ordinal']), ('owner/shared-repository', 1)) self.assertEqual(self.finish(second, rich_outcome(second['target']))['status'], 'resolved') counts = self.db.conn.execute( '''SELECT (SELECT target_count FROM docker_depth_experiments WHERE id = ?) AS targets, (SELECT selection_count FROM docker_depth_experiments WHERE id = ?) AS selections, (SELECT COUNT(*) FROM docker_image_manifests) AS manifests, (SELECT COUNT(*) FROM docker_depth_experiment_targets) AS physical''', (experiment_id, experiment_id), ).fetchone() self.assertEqual(tuple(counts), (1, 2, 1, 1)) def test_conflicting_manifest_layers_and_selector_evidence_roll_back(self): experiment_id, _members = self.seed_cohort(shared_first=True) self.use_postgres_shape() first = self.claim_one() self.assertEqual(self.finish(first, rich_outcome(first['target']))['status'], 'resolved') second = self.claim_one() conflicting_graph = rich_outcome(second['target'], ((999,),)) result = self.finish(second, conflicting_graph) self.assertEqual(result['status'], 'conflict_retry') self.db.conn.execute( '''UPDATE docker_depth_experiment_repositories SET resolver_due_at = '2000-01-01T00:00:00+00:00' WHERE id = ?''', (second['id'],), ) self.db.conn.commit() retried = self.claim_one() matching = rich_outcome(retried['target']) tampered_record = dict(matching.selection_records[0]) tampered_record['selection_evidence_sha256'] = 'f' * 64 tampered = SimpleNamespace( **{ **matching.__dict__, 'selection_records': (tampered_record,), } ) result = self.finish(retried, tampered) self.assertEqual(result['status'], 'conflict_retry') counts = self.db.conn.execute( '''SELECT (SELECT target_count FROM docker_depth_experiments WHERE id = ?) AS targets, (SELECT selection_count FROM docker_depth_experiments WHERE id = ?) AS selections, (SELECT COUNT(*) FROM docker_image_manifests) AS manifests, (SELECT COUNT(*) FROM docker_depth_experiment_selections) AS ranks''', (experiment_id, experiment_id), ).fetchone() self.assertEqual(tuple(counts), (1, 1, 1, 1)) self.db.conn.execute( '''UPDATE docker_depth_experiment_repositories SET resolver_due_at = '2000-01-01T00:00:00+00:00' WHERE id = ?''', (retried['id'],), ) self.db.conn.commit() final_claim = self.claim_one() final_result = self.finish(final_claim, tampered) self.assertEqual(final_result, { 'status': 'held', 'committed': True, 'reason': 'resolver_evidence_attempt_limit', }) held = self.db.conn.execute( '''SELECT state, hold_reason_code FROM docker_depth_experiments WHERE id = ?''', (experiment_id,), ).fetchone() self.assertEqual(tuple(held), ('held', 'resolver_evidence_attempt_limit')) def test_fresh_graph_counts_choose_one_deep_probe_after_all_breadth(self): experiment_id, members = self.seed_cohort() query_zero_members = { member_id for (query_ordinal, _rank), member_id in members.items() if query_ordinal == 0 } self.mark_members_skipped(experiment_id, query_zero_members) self.use_postgres_shape() for repository_rank in range(1, 11): claim = self.claim_one() self.assertEqual( (claim['query_ordinal'], claim['repository_rank'], claim['stage']), (0, repository_rank, 'breadth'), ) candidate_count = 5 if repository_rank == 3 else 2 outcome = rich_outcome( claim['target'], ((101 + repository_rank,),), candidate_count=candidate_count, ) self.assertEqual(self.finish(claim, outcome)['status'], 'resolved') deep_flags = self.db.conn.execute( '''SELECT repository_rank, planned_is_deep_probe, is_deep_probe, work_state, candidate_distinct_graph_count FROM docker_depth_experiment_repositories WHERE experiment_id = ? AND query_ordinal = 0 ORDER BY repository_rank''', (experiment_id,), ).fetchall() self.assertEqual( [row['repository_rank'] for row in deep_flags if row['planned_is_deep_probe']], [1], ) self.assertEqual( [row['repository_rank'] for row in deep_flags if row['is_deep_probe']], [3], ) deep_claim = self.claim_one() self.assertEqual( (deep_claim['id'], deep_claim['stage'], deep_claim['selection_limit']), (members[(0, 3)], 'deep', 10), ) deep_outcome = rich_outcome( deep_claim['target'], tuple((104,) if rank == 0 else (300 + rank,) for rank in range(5)), candidate_count=5, ) result = self.finish(deep_claim, deep_outcome) self.assertEqual((result['status'], result['selected_image_count']), ('resolved', 5)) dispatch = self.db.conn.execute( '''SELECT selection.image_rank, target.dispatch_wave, target.dispatch_order FROM docker_depth_experiment_selections selection JOIN docker_depth_experiment_targets target ON target.id = selection.experiment_target_id WHERE selection.experiment_repository_id = ? ORDER BY selection.image_rank''', (deep_claim['id'],), ).fetchall() self.assertEqual( [(row['image_rank'], row['dispatch_wave']) for row in dispatch], [(1, 1), (2, 2), (3, 2), (4, 3), (5, 3)], ) def test_terminal_candidate_without_replacement_records_image_scarcity(self): experiment_id, _members = self.seed_cohort() repository = 'owner/q00-repository-01' immutable = rich_outcome(repository).tags[0] now = scanner_db.utc_now_iso() self.db.conn.execute( '''INSERT INTO target_queue( source, platform, query, target, normalized_target, status, created_at, updated_at, completed_at ) VALUES ('dockerhub', 'docker', ?, ?, ?, 'done', ?, ?, ?)''', (QUERIES[0], immutable, immutable, now, now, now), ) self.db.conn.commit() self.use_postgres_shape() claim = self.claim_one() result = self.finish(claim, rich_outcome(repository)) self.assertEqual(result['status'], 'skipped') self.assertEqual(result['reason'], 'no_eligible_physical_target') counts = self.db.conn.execute( '''SELECT (SELECT target_count FROM docker_depth_experiments WHERE id = ?) AS targets, (SELECT selection_count FROM docker_depth_experiments WHERE id = ?) AS selections, (SELECT COUNT(*) FROM docker_image_manifests) AS manifests, (SELECT COUNT(*) FROM docker_depth_experiment_selections) AS ranks''', (experiment_id, experiment_id), ).fetchone() state = self.db.conn.execute( 'SELECT state, hold_reason_code FROM docker_depth_experiments WHERE id = ?', (experiment_id,), ).fetchone() self.assertEqual(tuple(counts), (0, 0, 0, 0)) self.assertEqual(tuple(state), ('resolving', None)) member = self.db.conn.execute( '''SELECT work_state, resolver_due_at, resolver_attempts, selected_image_count, last_error_code, resolved_at FROM docker_depth_experiment_repositories WHERE id = ?''', (claim['id'],), ).fetchone() self.assertEqual(tuple(member[:5]), ( 'skipped', None, 1, 0, 'no_eligible_physical_target', )) self.assertTrue(member['resolved_at']) skips = self.db.conn.execute( '''SELECT candidate_kind, reason_code FROM docker_depth_experiment_candidate_skips ORDER BY id''' ).fetchall() self.assertEqual( [tuple(row) for row in skips], [('image', 'immutable_target_terminal'), ('repository', 'no_eligible_physical_target')], ) preserved = self.db.conn.execute( 'SELECT status FROM target_queue WHERE normalized_target = ?', (immutable,), ).fetchone() self.assertEqual(preserved['status'], 'done') def test_terminal_candidate_uses_next_fresh_graph_as_rank_one(self): experiment_id, _members = self.seed_cohort() repository = 'owner/q00-repository-01' outcome = rich_outcome(repository, ((801,), (802,))) first_target = outcome.selection_records[0]['target'] now = scanner_db.utc_now_iso() self.db.conn.execute( '''INSERT INTO target_queue( source, platform, query, target, normalized_target, status, created_at, updated_at, completed_at ) VALUES ('dockerhub', 'docker', ?, ?, ?, 'done', ?, ?, ?)''', (QUERIES[0], first_target, first_target, now, now, now), ) self.db.conn.commit() self.use_postgres_shape() claim = self.claim_one() candidates = SimpleNamespace(**{ **outcome.__dict__, 'tags': (first_target,), 'selection_records': (outcome.selection_records[0],), 'candidate_records': outcome.selection_records, }) result = self.finish(claim, candidates) self.assertEqual(result['status'], 'resolved') selection = self.db.conn.execute( '''SELECT selection.image_rank, selection.selection_reason, queue.normalized_target FROM docker_depth_experiment_selections selection JOIN docker_depth_experiment_targets target ON target.id = selection.experiment_target_id JOIN target_queue queue ON queue.id = target.target_queue_id''' ).fetchone() self.assertEqual( tuple(selection), (1, 'replacement_after_exclusion', outcome.selection_records[1]['target']), ) skip = self.db.conn.execute( '''SELECT candidate_kind, reason_code, candidate_identity_sha256 FROM docker_depth_experiment_candidate_skips''' ).fetchone() self.assertEqual(tuple(skip[:2]), ('image', 'immutable_target_terminal')) self.assertRegex(skip['candidate_identity_sha256'], r'^[a-f0-9]{64}$') self.assertEqual(self.db.conn.execute( 'SELECT status FROM target_queue WHERE normalized_target = ?', (first_target,), ).fetchone()['status'], 'done') self.assertEqual(self.db.conn.execute( 'SELECT target_count FROM docker_depth_experiments WHERE id = ?', (experiment_id,), ).fetchone()['target_count'], 1) def test_cold_and_quarantined_candidates_are_skipped_without_reactivation(self): _experiment_id, _members = self.seed_cohort() repository = 'owner/q00-repository-01' outcome = rich_outcome(repository, ((811,), (812,), (813,))) now = scanner_db.utc_now_iso() for record, status in zip(outcome.selection_records[:2], ('cold', 'quarantined')): self.db.conn.execute( '''INSERT INTO target_queue( source, platform, query, target, normalized_target, status, created_at, updated_at ) VALUES ('dockerhub', 'docker', ?, ?, ?, ?, ?, ?)''', (QUERIES[0], record['target'], record['target'], status, now, now), ) self.db.conn.commit() self.use_postgres_shape() claim = self.claim_one() candidates = SimpleNamespace(**{ **outcome.__dict__, 'tags': (outcome.selection_records[0]['target'],), 'selection_records': (outcome.selection_records[0],), 'candidate_records': outcome.selection_records, }) self.assertEqual(self.finish(claim, candidates)['status'], 'resolved') states = self.db.conn.execute( '''SELECT status FROM target_queue WHERE normalized_target IN (?, ?) ORDER BY normalized_target''', tuple(record['target'] for record in outcome.selection_records[:2]), ).fetchall() self.assertEqual({row['status'] for row in states}, {'cold', 'quarantined'}) reasons = self.db.conn.execute( '''SELECT reason_code FROM docker_depth_experiment_candidate_skips ORDER BY id''' ).fetchall() self.assertEqual( {row['reason_code'] for row in reasons}, {'immutable_target_independently_cold', 'immutable_target_quarantined'}, ) def test_conclusive_zero_images_records_image_scarcity(self): experiment_id, _members = self.seed_cohort() self.use_postgres_shape() claim = self.claim_one() outcome = scanner.DockerTagResolutionOutcome( tags=(), status='empty', remote_attempted=True, selection_records=(), candidate_records=(), selector_version=depth.DOCKER_DEPTH_SELECTOR_VERSION, selector_hash=depth.canonical_selector_hash( depth.DOCKER_DEPTH_SELECTOR_VERSION ), candidate_distinct_graph_count=0, fresh_graph_evidence=True, cache_bypassed=True, ) result = self.finish(claim, outcome) self.assertEqual( (result['status'], result['reason']), ('skipped', 'no_eligible_physical_target'), ) row = self.db.conn.execute( '''SELECT state, hold_reason_code, target_count, selection_count FROM docker_depth_experiments WHERE id = ?''', (experiment_id,), ).fetchone() self.assertEqual(tuple(row), ( 'resolving', None, 0, 0, )) member = self.db.conn.execute( '''SELECT work_state, selected_image_count, last_error_code FROM docker_depth_experiment_repositories WHERE id = ?''', (claim['id'],), ).fetchone() self.assertEqual(tuple(member), ( 'skipped', 0, 'no_eligible_physical_target', )) def test_skipped_repository_evidence_drift_holds_experiment(self): experiment_id, _members = self.seed_cohort() self.use_postgres_shape() claim = self.claim_one() outcome = scanner.DockerTagResolutionOutcome( tags=(), status='empty', remote_attempted=True, selection_records=(), candidate_records=(), selector_version=depth.DOCKER_DEPTH_SELECTOR_VERSION, selector_hash=depth.canonical_selector_hash( depth.DOCKER_DEPTH_SELECTOR_VERSION ), candidate_distinct_graph_count=0, fresh_graph_evidence=True, cache_bypassed=True, ) self.assertEqual(self.finish(claim, outcome)['status'], 'skipped') self.db.conn.execute( '''UPDATE docker_depth_experiment_candidate_skips SET evidence_sha256 = ? WHERE experiment_repository_id = ? AND candidate_kind = 'repository' ''', ('f' * 64, claim['id']), ) self.db.conn.commit() result = self.db.refresh_docker_depth_experiment_state( resolver_authority(), final_cutover=True, ) self.assertEqual( (result['status'], result['reason']), ('held', 'repository_skip_evidence_drift'), ) row = self.db.conn.execute( 'SELECT state, hold_reason_code FROM docker_depth_experiments WHERE id = ?', (experiment_id,), ).fetchone() self.assertEqual(tuple(row), ('held', 'repository_skip_evidence_drift')) def test_conclusive_zero_replaces_repository_once_from_fresh_ranked_pool(self): experiment_id, _members = self.seed_cohort() replacement_id, replacement_repository = self.add_repository_candidate( experiment_id, ) queue = self.db.conn.execute( 'SELECT * FROM target_queue WHERE id = ?', (replacement_id,), ).fetchone() hold_entry = { 'queue_id': replacement_id, 'source': 'dockerhub', 'platform': 'docker', 'query': QUERIES[0], 'prior_status': 'deferred', 'prior_updated_at': str(queue['updated_at']), } audit = self.db._target_queue_policy_audit_sha256( 'cold', HOLD_SHA256, hold_entry, experiment_id, ) now = scanner_db.utc_now_iso() self.db.conn.execute( '''UPDATE target_queue SET status = 'cold', updated_at = ? WHERE id = ?''', (now, replacement_id), ) self.db.conn.execute( '''INSERT INTO target_queue_policy_events( queue_id, action, prior_status, next_status, source, platform, query, reason_code, config_sha256, policy_sha256, manifest_sha256, review_audit_sha256, experiment_id, prior_updated_at, created_at ) VALUES (?, 'cold', 'deferred', 'cold', 'dockerhub', 'docker', ?, ?, ?, ?, ?, ?, ?, ?, ?)''', ( replacement_id, QUERIES[0], depth.DOCKER_DEPTH_HOLD_REASON, experiment_config().config_hash, POLICY_SHA256, HOLD_SHA256, audit, experiment_id, hold_entry['prior_updated_at'], now, ), ) self.db.conn.commit() self.use_postgres_shape() first = self.claim_one() empty = scanner.DockerTagResolutionOutcome( tags=(), status='empty', remote_attempted=True, selector_version=depth.DOCKER_DEPTH_SELECTOR_VERSION, selector_hash=depth.canonical_selector_hash( depth.DOCKER_DEPTH_SELECTOR_VERSION ), candidate_distinct_graph_count=0, fresh_graph_evidence=True, cache_bypassed=True, ) result = self.finish(first, empty) self.assertEqual(result['status'], 'replaced') second = self.claim_one() self.assertEqual(second['repository_queue_id'], replacement_id) self.assertEqual(second['target'], replacement_repository) self.assertEqual(self.db.conn.execute( 'SELECT status FROM target_queue WHERE id = ?', (replacement_id,), ).fetchone()['status'], 'cold') member = self.db.conn.execute( '''SELECT replacement_repository_queue_id, replacement_count, replacement_evidence_sha256, resolver_attempts FROM docker_depth_experiment_repositories WHERE id = ?''', (first['id'],), ).fetchone() self.assertEqual( (member['replacement_repository_queue_id'], member['replacement_count']), (replacement_id, 1), ) self.assertRegex(member['replacement_evidence_sha256'], r'^[a-f0-9]{64}$') self.assertEqual(member['resolver_attempts'], 1) def test_sparse_query_selects_deep_probe_from_available_image_member(self): experiment_id, members = self.seed_cohort() self.shrink_query_cohort(experiment_id, 0, 1) self.use_postgres_shape() claim = self.claim_one() self.assertEqual(claim['id'], members[(0, 1)]) candidates = rich_outcome( 'owner/q00-repository-01', ((901,), (902,), (903,)), candidate_count=3, ) outcome = SimpleNamespace(**{ **candidates.__dict__, 'tags': (candidates.selection_records[0]['target'],), 'selection_records': (candidates.selection_records[0],), 'candidate_records': candidates.selection_records, }) result = self.finish(claim, outcome) self.assertEqual((result['status'], result['stage']), ('resolved', 'breadth')) member = self.db.conn.execute( '''SELECT work_state, is_deep_probe, candidate_distinct_graph_count, selected_image_count FROM docker_depth_experiment_repositories WHERE id = ?''', (claim['id'],), ).fetchone() self.assertEqual(tuple(member), ('pending', 1, 3, 1)) self.assertEqual(self.db.conn.execute( '''SELECT selected_repository_count FROM docker_depth_experiment_queries WHERE experiment_id = ? AND query_ordinal = 0''', (experiment_id,), ).fetchone()['selected_repository_count'], 1) def test_scan_dispatch_waits_for_completion_between_all_three_waves(self): experiment_id, members = self.seed_cohort() first = rich_outcome( 'owner/q00-repository-01', ((101,), (102,), (103,), (104,)), ) second = rich_outcome('owner/q01-repository-01', ((201,),)) self.activate_targets(experiment_id, members, ( (0, 1, first), (1, 1, second), )) self.ready_ingester() self.use_postgres_shape() wave_one = [self.claim_scan(index) for index in range(2)] self.assertIsNone(self.claim_scan(2)) for claim in wave_one: self.complete_scan_claim(claim) wave_two = [self.claim_scan(index) for index in range(3, 5)] self.assertIsNone(self.claim_scan(5)) for claim in wave_two: self.complete_scan_claim(claim) claims = wave_one + wave_two + [self.claim_scan(6)] self.assertEqual( [(claim['dispatch_wave'], claim['dispatch_order']) for claim in claims], [(1, 1), (1, 2), (2, 1), (2, 62), (3, 1)], ) self.assertEqual([claim['experiment_attempt'] for claim in claims], [1] * 5) identities = self.db.conn.execute( '''SELECT binding.id, binding.reservation_id, target.target_queue_id, reservation.queue_id FROM docker_depth_experiment_scan_bindings binding JOIN docker_depth_experiment_targets target ON target.id = binding.experiment_target_id JOIN result_reservations reservation ON reservation.id = binding.reservation_id ORDER BY binding.id''' ).fetchall() self.assertEqual(len(identities), 5) self.assertTrue(all( row['target_queue_id'] == row['queue_id'] for row in identities )) self.assertIsNone(self.claim_scan(7)) def test_shared_selection_uses_one_physical_scan_and_earliest_order(self): experiment_id, members = self.seed_cohort(shared_first=True) outcome = rich_outcome('owner/shared-repository', ((301,),)) self.activate_targets(experiment_id, members, ( (0, 1, outcome), (1, 1, outcome), )) self.ready_ingester() self.use_postgres_shape() claim = self.claim_scan(10) self.assertEqual((claim['dispatch_wave'], claim['dispatch_order']), (1, 1)) self.assertIsNone(self.claim_scan(11)) counts = self.db.conn.execute( '''SELECT (SELECT COUNT(*) FROM docker_depth_experiment_targets) AS targets, (SELECT COUNT(*) FROM docker_depth_experiment_selections) AS selections, (SELECT COUNT(*) FROM docker_depth_experiment_scan_bindings) AS bindings''' ).fetchone() self.assertEqual(tuple(counts), (1, 2, 1)) def test_refunded_attempt_requeues_only_the_same_bound_experiment_target(self): experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members, ( (0, 1, rich_outcome('owner/q00-repository-01', ((351,), (352,)))), )) self.ready_ingester() self.use_postgres_shape() first = self.claim_scan(12) self.assertTrue(self.db.refund_uncommitted_reservation( first['reservation_id'], current_process_identity(), 'synthetic pre-handoff failure', partial_absence_confirmed=True, )) second = self.claim_scan(13) self.assertEqual(second['experiment_target_id'], first['experiment_target_id']) self.assertEqual(second['experiment_attempt'], 2) self.assertIsNone(self.claim_scan(14)) attempts = self.db.conn.execute( '''SELECT binding.attempt, binding.state, reservation.state FROM docker_depth_experiment_scan_bindings binding JOIN result_reservations reservation ON reservation.id = binding.reservation_id ORDER BY binding.attempt''' ).fetchall() self.assertEqual( [tuple(row) for row in attempts], [(1, 'released', 'refunded'), (2, 'reserved', 'scanning')], ) self.complete_scan_claim(second) later_wave = self.claim_scan(15) self.assertEqual(later_wave['dispatch_wave'], 2) def test_expired_ready_recovery_advances_exact_binding_and_target_atomically(self): experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members, ( (0, 1, rich_outcome('owner/q00-repository-01', ((381,),))), )) self.ready_ingester() self.use_postgres_shape() claim = self.db.reserve_and_claim_target( 'dockerhub', 'docker', current_process_identity(), 'dispatch-supervisor', 64 * 1024, 2048, 1, 128, capacity_limits={ 'bundle_items': 100, 'bundle_bytes': 1024 * 1024, 'projection_items': 100, 'projection_bytes': 1024 * 1024, 'keycheck_items': 100, 'keycheck_bytes': 1024 * 1024, }, reservation_token='dispatch-reservation-18', bundle_id=f'{118:032x}', scan_event_id=f'{218:032x}', docker_depth_authority=resolver_authority(), final_cutover=True, ) bundle_root = os.path.join(self.temp.name, 'bundles') ensure_private_directory(bundle_root, reject_reparse=True) staged = scanner.stage_result_bundle( { 'scan_event_id': claim['scan_event_id'], 'target': claim['target'], 'scan_type': 'docker', 'timestamp': '2026-09-10T00:00:00+00:00', 'findings': [], 'errors': [], }, claim, bundle_root, {}, {'queue_status': 'done'}, candidate_max_items=1, candidate_max_bytes=128, ) metadata = ResultBundleReader( os.path.join(bundle_root, *staged.relative_path.split('/')), ).validate().as_dict() metadata['relative_path'] = staged.relative_path self.db.conn.execute( '''UPDATE result_reservations SET producer_lease_expires_at = '2000-01-01T00:00:00+00:00' WHERE id = ?''', (claim['reservation_id'],), ) self.db.conn.execute( '''UPDATE target_queue SET lease_expires_at = '2000-01-01T00:00:00+00:00' WHERE id = ?''', (claim['queue_id'],), ) self.db.conn.commit() with mock.patch.object( scanner_db, 'exact_process_identity_state', return_value='dead', ): self.assertTrue(self.db.recover_expired_result_bundle_ready( claim['reservation_id'], metadata, )) state = self.db.conn.execute( '''SELECT reservation.state AS reservation_state, binding.state AS binding_state, target.state AS target_state, binding.reservation_id, target.target_queue_id FROM result_reservations reservation JOIN docker_depth_experiment_scan_bindings binding ON binding.reservation_id = reservation.id JOIN docker_depth_experiment_targets target ON target.id = binding.experiment_target_id WHERE reservation.id = ?''', (claim['reservation_id'],), ).fetchone() self.assertEqual( tuple(state), ('ready', 'scanning', 'scanning', claim['reservation_id'], claim['queue_id']), ) def test_quarantine_retry_restores_exact_experiment_binding_and_resumes(self): experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members, ( (0, 1, rich_outcome('owner/q00-repository-01', ((391,),))), )) self.ready_ingester() self.use_postgres_shape() claim = self.db.reserve_and_claim_target( 'dockerhub', 'docker', current_process_identity(), 'dispatch-supervisor', 64 * 1024, 2048, 1, 128, capacity_limits={ 'bundle_items': 100, 'bundle_bytes': 1024 * 1024, 'projection_items': 100, 'projection_bytes': 1024 * 1024, 'keycheck_items': 100, 'keycheck_bytes': 1024 * 1024, }, reservation_token='dispatch-reservation-19', bundle_id=f'{119:032x}', scan_event_id=f'{219:032x}', docker_depth_authority=resolver_authority(), final_cutover=True, ) bundle_root = os.path.join(self.temp.name, 'retry-bundles') ensure_private_directory(bundle_root, reject_reparse=True) staged = scanner.stage_result_bundle( { 'scan_event_id': claim['scan_event_id'], 'target': claim['target'], 'scan_type': 'docker', 'timestamp': '2026-09-10T00:00:00+00:00', 'findings': [], 'errors': [], }, claim, bundle_root, {}, {'queue_status': 'done'}, candidate_max_items=1, candidate_max_bytes=128, ) metadata = ResultBundleReader( os.path.join(bundle_root, *staged.relative_path.split('/')), ).validate().as_dict() metadata['relative_path'] = staged.relative_path self.assertTrue(self.db.mark_result_bundle_ready( claim['reservation_id'], metadata, )) quarantine_id = self.db.quarantine_result_bundle( claim['reservation_id'], 'bundle_validation_failed', 'synthetic reviewed retry', payload_sha256=metadata['scan_event_hash'], byte_count=metadata['actual_bytes'], physical_confirmed=True, ) self.db.conn.execute( '''UPDATE docker_depth_experiments SET state = 'held', hold_reason_code = 'scan_target_held', held_at = ?, updated_at = ? WHERE id = ?''', (scanner_db.utc_now_iso(), scanner_db.utc_now_iso(), experiment_id), ) self.db.conn.commit() self.db.review_pipeline_quarantine( quarantine_id, 'bundle_validation_failed', metadata['scan_event_hash'], 'retry', '9' * 64, ) restored = self.db.conn.execute( '''SELECT reservation.state AS reservation_state, bundle.state AS bundle_state, queue.status AS queue_status, binding.state AS binding_state, target.state AS target_state, experiment.state AS experiment_state, experiment.hold_reason_code FROM result_reservations reservation JOIN result_bundles bundle ON bundle.reservation_id = reservation.id JOIN target_queue queue ON queue.id = reservation.queue_id JOIN docker_depth_experiment_scan_bindings binding ON binding.reservation_id = reservation.id JOIN docker_depth_experiment_targets target ON target.id = binding.experiment_target_id JOIN docker_depth_experiments experiment ON experiment.id = target.experiment_id WHERE reservation.id = ?''', (claim['reservation_id'],), ).fetchone() self.assertEqual( tuple(restored), ('scanning', 'ready', 'in_progress', 'reserved', 'reserved', 'active', None), ) self.assertTrue(self.db.mark_result_bundle_ready( claim['reservation_id'], metadata, )) recovered = self.db.conn.execute( '''SELECT reservation.state, binding.state, target.state FROM result_reservations reservation JOIN docker_depth_experiment_scan_bindings binding ON binding.reservation_id = reservation.id JOIN docker_depth_experiment_targets target ON target.id = binding.experiment_target_id WHERE reservation.id = ?''', (claim['reservation_id'],), ).fetchone() self.assertEqual(tuple(recovered), ('scanning', 'scanning', 'scanning')) def test_quarantine_rescan_preserves_attempt_and_requeues_experiment_target(self): experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members, ( (0, 1, rich_outcome('owner/q00-repository-01', ((392,),))), )) self.ready_ingester() self.use_postgres_shape() claim = self.claim_scan(20) from tests.test_docker_layer_scanning import docker_plan plan = docker_plan(coverage_state='skipped') repository, manifest_digest = claim['target'].rsplit('@', 1) plan['image'] = claim['target'] plan['repository'] = repository plan['manifest_digest'] = manifest_digest plan = scanner_db.validate_docker_layer_plan(plan) plan_bytes = scanner_db.canonical_docker_layer_plan_bytes(plan) self.db.conn.execute( '''UPDATE result_reservations SET docker_layer_plan_json = ?, docker_layer_plan_sha256 = ? WHERE id = ?''', ( json.dumps(plan, ensure_ascii=True, sort_keys=True, separators=(',', ':')), hashlib.sha256(plan_bytes).hexdigest(), claim['reservation_id'], ), ) self.db.conn.commit() quarantine_id = self.db.quarantine_result_bundle( claim['reservation_id'], 'bundle_validation_failed', 'synthetic reviewed rescan', physical_confirmed=True, ) self.db.conn.execute( '''UPDATE docker_depth_experiments SET state = 'draining', draining_at = ?, updated_at = ? WHERE id = ?''', (scanner_db.utc_now_iso(), scanner_db.utc_now_iso(), experiment_id), ) self.db.conn.commit() self.db.review_pipeline_quarantine( quarantine_id, 'bundle_validation_failed', '', 'rescan', '7' * 64, ) restored = self.db.conn.execute( '''SELECT reservation.state AS reservation_state, queue.status AS queue_status, queue.current_result_reservation_id, binding.state AS binding_state, target.state AS target_state, target.reservation_count, experiment.state AS experiment_state, experiment.draining_at FROM result_reservations reservation JOIN target_queue queue ON queue.id = reservation.queue_id JOIN docker_depth_experiment_scan_bindings binding ON binding.reservation_id = reservation.id JOIN docker_depth_experiment_targets target ON target.id = binding.experiment_target_id JOIN docker_depth_experiments experiment ON experiment.id = target.experiment_id WHERE reservation.id = ?''', (claim['reservation_id'],), ).fetchone() self.assertEqual( tuple(restored), ('quarantined', 'pending', None, 'quarantined', 'pending', 1, 'active', None), ) retried = self.claim_scan(21) self.assertEqual(retried['experiment_target_id'], claim['experiment_target_id']) self.assertEqual(retried['experiment_attempt'], 2) def test_atomic_binding_failure_rolls_back_queue_reservation_and_capacity(self): experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members, ( (0, 1, rich_outcome('owner/q00-repository-01', ((401,),))), )) self.ready_ingester() self.db.conn.execute( '''CREATE TRIGGER fail_experiment_binding BEFORE INSERT ON docker_depth_experiment_scan_bindings BEGIN SELECT RAISE(ABORT, 'synthetic binding failure'); END''' ) self.db.conn.commit() self.use_postgres_shape() with self.assertRaisesRegex(sqlite3.IntegrityError, 'synthetic binding failure'): self.claim_scan(20) state = self.db.conn.execute( '''SELECT target.state, target.reservation_count, queue.status, queue.current_result_reservation_id FROM docker_depth_experiment_targets target JOIN target_queue queue ON queue.id = target.target_queue_id''' ).fetchone() capacity = self.db.conn.execute( 'SELECT bundle_items, projection_items, keycheck_items FROM pipeline_capacity' ).fetchone() self.assertEqual(tuple(state), ('pending', 0, 'pending', None)) self.assertEqual(tuple(capacity), (0, 0, 0)) self.assertEqual(self.db.conn.execute( 'SELECT COUNT(*) AS count FROM result_reservations' ).fetchone()['count'], 0) def test_selection_mutation_holds_fail_closed(self): experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members, ( (0, 1, rich_outcome('owner/q00-repository-01', ((501,),))), )) self.ready_ingester() self.db.conn.execute( "UPDATE docker_depth_experiment_selections SET selection_evidence_sha256 = ?", ('f' * 64,), ) self.db.conn.commit() self.use_postgres_shape() self.assertIsNone(self.claim_scan(30)) held = self.db.conn.execute( 'SELECT state, hold_reason_code FROM docker_depth_experiments' ).fetchone() self.assertEqual(tuple(held), ('held', 'selection_mutation')) def test_candidate_graph_count_churn_does_not_mutate_selection(self): experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members, ( (0, 1, rich_outcome('owner/q00-repository-01', ((502,),))), )) self.ready_ingester() self.db.conn.execute( '''UPDATE docker_depth_experiment_repositories SET candidate_distinct_graph_count = 2 WHERE id = ?''', (members[(0, 1)],), ) self.db.conn.commit() self.use_postgres_shape() claim = self.claim_scan(31) self.assertIsNotNone(claim) state = self.db.conn.execute( 'SELECT state, hold_reason_code FROM docker_depth_experiments' ).fetchone() self.assertEqual(tuple(state), ('active', None)) def test_retryable_scan_returns_scanning_target_to_pending(self): experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members, ( (0, 1, rich_outcome('owner/q00-repository-01', ((503,),))), )) self.ready_ingester() self.use_postgres_shape() claim = self.claim_scan(32) reservation = self.db.conn.execute( 'SELECT * FROM result_reservations WHERE id = ?', (claim['reservation_id'],), ).fetchone() now = scanner_db.utc_now_iso() self.db._transition_docker_depth_binding_locked( reservation, 'scanning', 'scanning', now, ) self.db._transition_docker_depth_binding_locked( reservation, 'failed', 'pending', now, ) state = self.db.conn.execute( '''SELECT binding.state, target.state FROM docker_depth_experiment_scan_bindings binding JOIN docker_depth_experiment_targets target ON target.id = binding.experiment_target_id WHERE binding.reservation_id = ?''', (claim['reservation_id'],), ).fetchone() self.assertEqual(tuple(state), ('failed', 'pending')) def test_exhausted_deferred_target_is_failed_before_next_wave_claim(self): experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members, ( (0, 1, rich_outcome( 'owner/q00-repository-01', ((504,), (505,)), )), )) self.ready_ingester() self.use_postgres_shape() first = self.claim_scan(33) now = scanner_db.utc_now_iso() scan = self.db.conn.execute( '''INSERT INTO target_scans( scan_event_id, scan_event_hash, queue_id, claim_lease_token, queue_completion_applied, queue_completion_disposition, source, query, target, normalized_target, scan_type, status, result_reservation_id, created_at ) VALUES (?, ?, ?, ?, 1, 'applied', 'dockerhub', ?, ?, ?, 'docker', 'error', ?, ?) RETURNING id''', ( first['scan_event_id'], f"{first['reservation_id']:064x}", first['queue_id'], first['claim_lease_token'], first['query'], first['target'], first['normalized_target'], first['reservation_id'], now, ), ).fetchone() self.db.conn.execute( '''UPDATE target_queue SET status = 'deferred', attempts = 3, available_after = '2000-01-01T00:00:00+00:00', target_scan_id = ?, last_error = 'retryable layer failure', lease_owner = NULL, lease_token = NULL, claim_batch = NULL, leased_at = NULL, lease_expires_at = NULL, current_result_reservation_id = NULL, claim_event_id = NULL, completed_at = NULL, updated_at = ? WHERE id = ?''', (scan['id'], now, first['queue_id']), ) self.db.conn.execute( '''UPDATE result_reservations SET state = 'acknowledged', bundle_credit_released = 1, projection_credit_transferred = 1, candidate_credit_transferred = 1, released_at = ?, updated_at = ? WHERE id = ?''', (now, now, first['reservation_id']), ) self.db.conn.execute( '''UPDATE docker_depth_experiment_scan_bindings SET target_scan_id = ?, state = 'failed', scan_bound_at = ?, completed_at = ? WHERE reservation_id = ?''', (scan['id'], now, now, first['reservation_id']), ) self.db.conn.execute( '''UPDATE docker_depth_experiment_targets SET state = 'pending', terminal_at = NULL, updated_at = ? WHERE id = ?''', (now, first['experiment_target_id']), ) self.db.conn.execute( '''UPDATE pipeline_capacity SET bundle_items = 0, bundle_bytes = 0, projection_items = 0, projection_bytes = 0, keycheck_items = 0, keycheck_bytes = 0 WHERE id = 1''' ) self.db.conn.commit() second = self.claim_scan(34, max_attempts=3) self.assertEqual(second['dispatch_wave'], 2) recovered = self.db.conn.execute( '''SELECT target.state AS target_state, target.terminal_at, queue.status AS queue_status, queue.completed_at, queue.available_after, queue.attempts, binding.state AS binding_state FROM docker_depth_experiment_targets target JOIN target_queue queue ON queue.id = target.target_queue_id JOIN docker_depth_experiment_scan_bindings binding ON binding.experiment_target_id = target.id AND binding.attempt = target.reservation_count WHERE target.id = ?''', (first['experiment_target_id'],), ).fetchone() self.assertEqual( ( recovered['target_state'], recovered['queue_status'], recovered['available_after'], recovered['attempts'], recovered['binding_state'], ), ('failed', 'failed', None, 3, 'failed'), ) self.assertIsNotNone(recovered['terminal_at']) self.assertIsNotNone(recovered['completed_at']) def test_frozen_deep_probe_hash_detects_runtime_choice_drift(self): experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members, ( (0, 1, rich_outcome('owner/q00-repository-01', ((511,),))), )) self.db.conn.execute( '''UPDATE docker_depth_experiment_repositories SET is_deep_probe = CASE repository_rank WHEN 2 THEN 1 ELSE 0 END WHERE experiment_id = ? AND query_ordinal = 0''', (experiment_id,), ) self.db.conn.commit() self.use_postgres_shape() result = self.db.refresh_docker_depth_experiment_state( resolver_authority(), final_cutover=True, ) self.assertEqual( (result['status'], result['reason']), ('held', 'selection_hash_drift'), ) def test_authority_holds_on_tampered_owned_hold_event(self): experiment_id, _members = self.seed_cohort() _queue_id, event_id, _now = self.add_owned_hold_event(experiment_id) self.db.conn.execute( '''UPDATE target_queue_policy_events SET review_audit_sha256 = ? WHERE id = ?''', ('f' * 64, event_id), ) self.db.conn.commit() self.use_postgres_shape() result = self.db.refresh_docker_depth_experiment_state( resolver_authority(), final_cutover=True, ) self.assertEqual( (result['status'], result['reason']), ('held', 'target_history_drift'), ) def test_authority_holds_on_owned_event_reversal_before_release(self): experiment_id, _members = self.seed_cohort() queue_id, event_id, now = self.add_owned_hold_event(experiment_id) entry = { 'queue_id': queue_id, 'source': 'dockerhub', 'platform': 'docker', 'query': QUERIES[0], 'cold_event_id': event_id, 'restore_status': 'deferred', 'prior_updated_at': now, } audit = self.db._target_queue_policy_audit_sha256( 'reactivate', 'a' * 64, entry, experiment_id, ) self.db.conn.execute( '''INSERT INTO target_queue_policy_events( queue_id, action, prior_status, next_status, source, platform, query, reason_code, config_sha256, policy_sha256, manifest_sha256, review_audit_sha256, reverses_event_id, experiment_id, prior_updated_at, created_at ) VALUES (?, 'reactivate', 'cold', 'deferred', 'dockerhub', 'docker', ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)''', ( queue_id, QUERIES[0], depth.DOCKER_DEPTH_RELEASE_REASON, experiment_config().config_hash, POLICY_SHA256, 'a' * 64, audit, event_id, experiment_id, now, now, ), ) self.db.conn.commit() self.use_postgres_shape() result = self.db.refresh_docker_depth_experiment_state( resolver_authority(), final_cutover=True, ) self.assertEqual( (result['status'], result['reason']), ('held', 'target_history_drift'), ) def test_stale_reservation_holds_fail_closed(self): experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members, ( (0, 1, rich_outcome('owner/q00-repository-01', ((502,),))), )) self.ready_ingester() self.use_postgres_shape() self.assertIsNotNone(self.claim_scan(31)) self.db.conn.execute( '''UPDATE result_reservations SET producer_lease_expires_at = '2000-01-01T00:00:00+00:00' ''' ) self.db.conn.commit() self.assertIsNone(self.claim_scan(32)) held = self.db.conn.execute( 'SELECT state, hold_reason_code FROM docker_depth_experiments' ).fetchone() self.assertEqual(tuple(held), ('held', 'stale_scan_fence')) def test_config_drift_holds_fail_closed(self): experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members, ( (0, 1, rich_outcome('owner/q00-repository-01', ((601,),))), )) self.ready_ingester() self.use_postgres_shape() malformed = dict(resolver_authority(), selector_sha256='0' * 64) self.assertIsNone(self.claim_scan(39, authority=malformed)) stable = self.db.conn.execute( 'SELECT state, hold_reason_code FROM docker_depth_experiments' ).fetchone() self.assertEqual(tuple(stable), ('active', None)) drifted = dict(resolver_authority(), config_sha256='f' * 64) self.assertIsNone(self.claim_scan(40, authority=drifted)) held = self.db.conn.execute( 'SELECT state, hold_reason_code FROM docker_depth_experiments' ).fetchone() self.assertEqual(tuple(held), ('held', 'config_hash_drift')) replay = self.db.refresh_docker_depth_experiment_state( malformed, final_cutover=True, ) self.assertEqual( (replay['status'], replay['committed'], replay['reason']), ('invalid', False, 'authority_payload_invalid'), ) stable = self.db.conn.execute( 'SELECT state, hold_reason_code FROM docker_depth_experiments' ).fetchone() self.assertEqual(tuple(stable), ('held', 'config_hash_drift')) def test_capacity_backpressure_does_not_hold_and_recovers_legacy_hold(self): experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members, ( (0, 1, rich_outcome('owner/q00-repository-01', ((602,),))), )) self.ready_ingester() self.use_postgres_shape() limits = { 'bundle_items': 0, 'bundle_bytes': 0, 'projection_items': 0, 'projection_bytes': 0, 'keycheck_items': 0, 'keycheck_bytes': 0, } self.assertIsNone(self.claim_scan(41, limits=limits)) held = self.db.conn.execute( 'SELECT state, hold_reason_code FROM docker_depth_experiments' ).fetchone() self.assertEqual(tuple(held), ('active', None)) intent = self.db.conn.execute( '''SELECT state, resolution_detail FROM admission_intents WHERE reservation_token = 'dispatch-reservation-41' ''' ).fetchone() self.assertEqual(tuple(intent), ('aborted', 'pipeline_capacity_closed')) self.assertEqual(self.db.conn.execute( 'SELECT COUNT(*) AS count FROM result_reservations' ).fetchone()['count'], 0) now = scanner_db.utc_now_iso() self.db.conn.execute( '''UPDATE docker_depth_experiments SET state = 'held', hold_reason_code = 'pipeline_capacity_conflict', held_at = ?, updated_at = ? WHERE id = ?''', (now, now, experiment_id), ) self.db.conn.commit() self.assertIsNotNone(self.claim_scan(42)) recovered = self.db.conn.execute( 'SELECT state, hold_reason_code FROM docker_depth_experiments' ).fetchone() self.assertEqual(tuple(recovered), ('active', None)) def test_completion_accepts_explicit_scarcity_and_preserves_holds(self): experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members, ( (0, 1, rich_outcome('owner/q00-repository-01', ((651,),))), )) self.ready_ingester() now = scanner_db.utc_now_iso() queue = self.db.conn.execute( '''INSERT INTO target_queue( source, platform, query, target, normalized_target, status, created_at, updated_at ) VALUES ('dockerhub', 'docker', ?, 'owner/noncohort', 'owner/noncohort', 'cold', ?, ?) RETURNING id''', (QUERIES[0], now, now), ).fetchone() entry = { 'queue_id': int(queue['id']), 'source': 'dockerhub', 'platform': 'docker', 'query': QUERIES[0], 'prior_status': 'deferred', 'prior_updated_at': now, } audit = self.db._target_queue_policy_audit_sha256( 'cold', HOLD_SHA256, entry, experiment_id, ) self.db.conn.execute( '''INSERT INTO target_queue_policy_events( queue_id, action, prior_status, next_status, source, platform, query, reason_code, config_sha256, policy_sha256, manifest_sha256, review_audit_sha256, experiment_id, prior_updated_at, created_at ) VALUES (?, 'cold', 'deferred', 'cold', 'dockerhub', 'docker', ?, ?, ?, ?, ?, ?, ?, ?, ?)''', ( queue['id'], QUERIES[0], depth.DOCKER_DEPTH_HOLD_REASON, experiment_config().config_hash, POLICY_SHA256, HOLD_SHA256, audit, experiment_id, now, now, ), ) self.db.conn.commit() self.use_postgres_shape() claim = self.claim_scan(45) state = self.db.refresh_docker_depth_experiment_state( resolver_authority(), final_cutover=True, ) self.assertEqual(state['status'], 'active') scan = self.db.conn.execute( '''INSERT INTO target_scans( scan_event_id, scan_event_hash, queue_id, claim_lease_token, queue_completion_applied, queue_completion_disposition, source, query, target, normalized_target, scan_type, status, result_reservation_id, created_at ) VALUES (?, ?, ?, ?, 1, 'applied', 'dockerhub', ?, ?, ?, 'docker', 'clean', ?, ?) RETURNING id''', ( claim['scan_event_id'], 'a' * 64, claim['queue_id'], claim['claim_lease_token'], claim['query'], claim['target'], claim['normalized_target'], claim['reservation_id'], now, ), ).fetchone() self.db.conn.execute( '''UPDATE target_queue SET status = 'done', target_scan_id = ?, lease_owner = NULL, lease_token = NULL, claim_batch = NULL, leased_at = NULL, lease_expires_at = NULL, current_result_reservation_id = NULL, claim_event_id = NULL, completed_at = ?, updated_at = ? WHERE id = ?''', (scan['id'], now, now, claim['queue_id']), ) self.db.conn.execute( '''UPDATE result_reservations SET state = 'acknowledged', bundle_credit_released = 1, projection_credit_transferred = 1, candidate_credit_transferred = 1, released_at = ?, updated_at = ? WHERE id = ?''', (now, now, claim['reservation_id']), ) self.db.conn.execute( '''UPDATE docker_depth_experiment_scan_bindings SET target_scan_id = ?, state = 'completed', scan_bound_at = ?, completed_at = ? WHERE reservation_id = ?''', (scan['id'], now, now, claim['reservation_id']), ) self.db.conn.execute( '''UPDATE docker_depth_experiment_targets SET state = 'done', terminal_at = ?, updated_at = ? WHERE id = ?''', (now, now, claim['experiment_target_id']), ) self.db.conn.execute( '''UPDATE pipeline_capacity SET bundle_items = 0, bundle_bytes = 0, projection_items = 0, projection_bytes = 0, keycheck_items = 0, keycheck_bytes = 0 WHERE id = 1''' ) self.db.conn.commit() state = self.db.refresh_docker_depth_experiment_state( resolver_authority(), final_cutover=True, ) self.assertEqual(state['status'], 'completed') experiment = self.db.conn.execute( '''SELECT state, hold_manifest_sha256, released_at FROM docker_depth_experiments WHERE id = ?''', (experiment_id,), ).fetchone() events = self.db.conn.execute( '''SELECT action FROM target_queue_policy_events WHERE experiment_id = ? ORDER BY id''', (experiment_id,), ).fetchall() self.assertEqual(tuple(experiment), ('completed', HOLD_SHA256, None)) self.assertEqual([row['action'] for row in events], ['cold']) def test_disabled_authority_preserves_ordinary_claim_and_native_sqlite_fails_closed(self): with self.assertRaisesRegex(RuntimeError, 'require PostgreSQL'): self.db.reserve_and_claim_target( 'dockerhub', 'docker', current_process_identity(), 'dispatch-supervisor', 1024, 2048, 1, 128, reservation_token='native-sqlite', bundle_id='a' * 32, scan_event_id='b' * 32, docker_depth_authority=resolver_authority(), final_cutover=True, ) experiment_id, members = self.seed_cohort() self.activate_targets(experiment_id, members) now = scanner_db.utc_now_iso() self.db.conn.execute( '''INSERT INTO target_queue( source, platform, query, target, normalized_target, status, created_at, updated_at ) VALUES ('dockerhub', 'docker', ?, ?, ?, 'pending', ?, ?)''', (QUERIES[0], 'owner/ordinary@' + digest(701), 'owner/ordinary@' + digest(701), now, now), ) self.ready_ingester() self.use_postgres_shape() claim = self.claim_scan(50, authority=resolver_authority(False)) self.assertEqual(claim['target'], 'owner/ordinary@' + digest(701)) self.assertNotIn('experiment_target_id', claim) experiment = self.db.conn.execute( '''SELECT state, hold_reason_code FROM docker_depth_experiments WHERE id = ?''', (experiment_id,), ).fetchone() self.assertEqual(tuple(experiment), ('held', 'experiment_disabled')) def test_disabling_while_holding_transitions_fail_closed_to_held(self): experiment_id, _members = self.seed_cohort() self.db.conn.execute( "UPDATE docker_depth_experiments SET state = 'holding' WHERE id = ?", (experiment_id,), ) self.db.conn.commit() self.use_postgres_shape() result = self.db.refresh_docker_depth_experiment_state( resolver_authority(False), final_cutover=True, ) self.assertEqual( (result['status'], result['reason']), ('held', 'experiment_disabled'), ) self.assertEqual(self.db.conn.execute( 'SELECT state FROM docker_depth_experiments WHERE id = ?', (experiment_id,), ).fetchone()['state'], 'held') class DockerDepthResolverWorkerTests(unittest.TestCase): def test_fresh_zero_graph_pool_is_conclusive_without_limit_zero_selection(self): class Response: status_code = 200 @staticmethod def raise_for_status(): return None tag_payload = {'results': [ {'name': 'latest', 'digest': digest(1), 'last_updated': '2026-01-02T00:00:00Z'}, ]} frozen_selector_hash = depth.canonical_selector_hash( depth.DOCKER_DEPTH_SELECTOR_VERSION ) with mock.patch.object(scanner, 'get_dockerhub_tag_cache') as cache_get, \ mock.patch.object(scanner, 'put_dockerhub_tag_cache') as cache_put, \ mock.patch.object(scanner, 'dockerhub_tag_rate_limit_state', return_value={ 'active': False, 'retry_at': None, }), \ mock.patch.object(scanner, 'dockerhub_tags_response', return_value=Response()), \ mock.patch.object(scanner, '_bounded_docker_registry_json', return_value=tag_payload), \ mock.patch.object( scanner, 'resolve_docker_layer_graph', return_value=(None, None), ), \ mock.patch.object( scanner, 'select_docker_layer_graphs', side_effect=AssertionError('selector must not receive an empty pool'), ): outcome = scanner.fetch_dockerhub_tags( 'owner/empty-graph-pool', limit=10, retry_count=0, return_outcome=True, fresh_graph_evidence=True, ) cache_get.assert_not_called() cache_put.assert_not_called() self.assertEqual(outcome.status, 'unsupported') self.assertEqual(outcome.candidate_distinct_graph_count, 0) self.assertEqual(tuple(outcome.candidate_records), ()) self.assertEqual(tuple(outcome.selection_records), ()) self.assertEqual( depth.canonical_selector_hash(depth.DOCKER_DEPTH_SELECTOR_VERSION), frozen_selector_hash, ) def test_fresh_experiment_fetch_bypasses_lossy_cache_and_deduplicates_alias_graphs(self): class Response: status_code = 200 @staticmethod def raise_for_status(): return None tag_payload = {'results': [ {'name': 'latest', 'digest': digest(1), 'last_updated': '2026-01-02T00:00:00Z'}, {'name': 'alias', 'digest': digest(2), 'last_updated': '2026-01-01T00:00:00Z'}, ]} graph = { 'manifest_digest': digest(3), 'manifest_media_type': 'application/vnd.oci.image.manifest.v1+json', 'manifest_size_bytes': 1234, 'config_digest': digest(4), 'layers': (digest(5), digest(5)), 'layer_descriptors': ( {'digest': digest(5), 'media_type': 'application/vnd.oci.image.layer.v1.tar+gzip', 'size': 10}, {'digest': digest(5), 'media_type': 'application/vnd.oci.image.layer.v1.tar+gzip', 'size': 10}, ), } renew = mock.Mock(return_value=True) with mock.patch.object(scanner, 'get_dockerhub_tag_cache') as cache_get, \ mock.patch.object(scanner, 'put_dockerhub_tag_cache') as cache_put, \ mock.patch.object(scanner, 'dockerhub_tag_rate_limit_state', return_value={ 'active': False, 'retry_at': None, }), \ mock.patch.object(scanner, 'dockerhub_tags_response', return_value=Response()), \ mock.patch.object(scanner, '_bounded_docker_registry_json', return_value=tag_payload), \ mock.patch.object( scanner, 'resolve_docker_layer_graph', return_value=(graph, None), ) as resolve_graph: outcome = scanner.fetch_dockerhub_tags( 'owner/cache-bypass', limit=10, retry_count=0, return_outcome=True, fresh_graph_evidence=True, lease_renewal_callback=renew, ) cache_get.assert_not_called() cache_put.assert_not_called() self.assertTrue(outcome.cache_bypassed) self.assertTrue(outcome.fresh_graph_evidence) self.assertGreaterEqual(renew.call_count, 1) self.assertIs( resolve_graph.call_args.kwargs['lease_renewal_callback'], renew, ) self.assertEqual(outcome.candidate_distinct_graph_count, 1) self.assertEqual(len(outcome.selection_records), 1) self.assertEqual( [(layer['position_from_base'], layer['position_from_top']) for layer in outcome.selection_records[0]['layer_metadata']], [(1, 2), (2, 1)], ) def test_experiment_lane_bypasses_only_scan_queue_empty_gate(self): authority = resolver_authority() outcome = rich_outcome('owner/worker') class DB: conn = SimpleNamespace(is_postgres=True) def __init__(self): self.claimed = False self.finished = [] self.renewals = 0 def claim_docker_depth_experiment_resolutions(self, *_args, **_kwargs): if self.claimed: return [] self.claimed = True return [{ 'id': 1, 'target': 'owner/worker', 'selection_limit': 1, 'resolver_generation': 1, 'resolver_token': 'token', 'resolver_owner': 'owner', }] def renew_docker_depth_experiment_resolution(self, *_args, **_kwargs): self.renewals += 1 return {'status': 'renewed', 'renewed': True} def finish_docker_depth_experiment_resolution(self, *args, **kwargs): self.finished.append((args, kwargs)) return {'status': 'resolved', 'committed': True} @staticmethod def runtime_control_state(): return {'effective_discovery_paused': False} @staticmethod def has_claimable_targets_v2(*_args, **_kwargs): return True @staticmethod def claim_docker_resolutions(*_args, **_kwargs): raise AssertionError('ordinary resolver must remain behind the scan queue gate') last_error = '' args = SimpleNamespace( docker_depth_experiment_authority=authority, sync_file_queues=False, tag_resolve_limit=10, tag_retry_count=0, tag_retry_delay=0, docker_platform_filter_enabled=True, docker_platform_os='linux', docker_platform_arch='amd64', docker_platform_candidate_tags=20, target_retry_max_attempts=3, ) database = DB() with mock.patch.object( console_runner, 'fetch_dockerhub_tags', return_value=outcome, ) as fetch, mock.patch('builtins.print'): processed = console_runner.resolve_due_docker_queue_targets_if_scan_queue_empty( database, 'dockerhub', args, ) self.assertEqual(processed, 1) self.assertTrue(fetch.call_args.kwargs['fresh_graph_evidence']) self.assertIsNotNone(fetch.call_args.kwargs['lease_renewal_callback']) self.assertEqual(database.renewals, 2) self.assertTrue(database.finished) if __name__ == '__main__': unittest.main()