2835 lines
128 KiB
Python
2835 lines
128 KiB
Python
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()
|