3899 lines
174 KiB
Python
3899 lines
174 KiB
Python
import copy
|
|
from datetime import datetime, timezone
|
|
import hashlib
|
|
import inspect
|
|
import io
|
|
import json
|
|
import math
|
|
import os
|
|
import re
|
|
import textwrap
|
|
import tokenize
|
|
from dataclasses import dataclass
|
|
from typing import Optional
|
|
|
|
from target_identity import (
|
|
DOCKER_DIGEST_RE,
|
|
DOCKER_IMAGE_RE,
|
|
DOCKER_REVISION_RE,
|
|
DOCKER_TAG_TARGET_SCHEMA,
|
|
normalize_docker_digest,
|
|
parse_docker_target,
|
|
serialize_docker_tag_target,
|
|
validate_docker_image_reference,
|
|
)
|
|
|
|
|
|
DOCKER_DEPTH_SELECTOR_VERSION = 'docker-layer-graph-v2'
|
|
DOCKER_RANK1_BREADTH_SELECTOR_VERSION = 'docker-rank1-breadth-v1'
|
|
DOCKER_DEPTH_QUERY_COUNT = 61
|
|
DOCKER_DEPTH_ORDINARY_IMAGES_PER_REPOSITORY = 3
|
|
DOCKER_DEPTH_REPOSITORIES_PER_QUERY = 10
|
|
DOCKER_DEPTH_SHALLOW_IMAGES_PER_REPOSITORY = 1
|
|
DOCKER_DEPTH_DEEP_REPOSITORIES_PER_QUERY = 1
|
|
DOCKER_DEPTH_DEEP_IMAGES_PER_REPOSITORY = 10
|
|
DOCKER_DEPTH_MAX_UNIQUE_TARGETS = 1200
|
|
DOCKER_RANK1_BREADTH_REPOSITORIES_PER_QUERY = 39
|
|
DOCKER_RANK1_BREADTH_MAX_UNIQUE_TARGETS = 2000
|
|
DOCKERHUB_DISCOVERY_MAX_PAGES = 30
|
|
DOCKERHUB_DISCOVERY_MAX_PER_PAGE = 100
|
|
DOCKERHUB_DISCOVERY_ALGORITHM_VERSION = 1
|
|
DOCKER_DEPTH_REVIEW_MANIFEST_MAX_ROWS = 250000
|
|
DOCKER_DEPTH_RESOLVER_MAX_ATTEMPTS = 3
|
|
DOCKER_DEPTH_RESOLVER_RETRY_SECONDS = 300
|
|
DOCKER_DEPTH_RESOLVER_RETRY_MAX_SECONDS = 3600
|
|
DOCKER_DEPTH_COLLECTION_GENERATION = 'docker-depth-provenance-v1'
|
|
DOCKER_DEPTH_COHORT_PLAN_TYPE = 'truf-docker-depth-cohort-plan-v2'
|
|
DOCKER_DEPTH_COHORT_MANIFEST_TYPE = 'truf-docker-depth-cohort-review-v2'
|
|
DOCKER_DEPTH_HOLD_MANIFEST_TYPE = 'truf-docker-depth-hold-review-v1'
|
|
DOCKER_DEPTH_REACTIVATION_MANIFEST_TYPE = 'truf-docker-depth-reactivation-review-v1'
|
|
DOCKER_DEPTH_RESOLVER_REFUND_MANIFEST_TYPE = (
|
|
'truf-docker-depth-resolver-attempt-refund-v1'
|
|
)
|
|
DOCKER_DEPTH_RESOLVER_REFUND_KIND = 'zero_graph_limit_v1'
|
|
DOCKER_DEPTH_RESOLVER_REFUND_ATTEMPTS = 2
|
|
DOCKER_DEPTH_RESOLVER_REFUND_LOG_MAX_BYTES = 64 * 1024 * 1024
|
|
DOCKER_DEPTH_RESOLVER_REFUND_OLD_ERROR = (
|
|
'Docker replacement candidate pool must contain 1 through 100 graphs'
|
|
)
|
|
DOCKER_DEPTH_HOLD_REASON = 'docker_depth_experiment_hold'
|
|
DOCKER_DEPTH_DYNAMIC_HOLD_REASON = 'docker_depth_experiment_dynamic_hold'
|
|
DOCKER_DEPTH_RELEASE_REASON = 'docker_depth_experiment_reviewed_release'
|
|
DOCKER_DEPTH_REPOSITORY_SKIP_REASON = 'no_eligible_physical_target'
|
|
DOCKER_DEPTH_REMOTE_UNAVAILABLE_SKIP_REASON = 'remote_unavailable_after_attempt_limit'
|
|
DOCKER_DEPTH_RESOLVER_DISPOSITION_MANIFEST_TYPE = (
|
|
'truf-docker-depth-resolver-disposition-v1'
|
|
)
|
|
DOCKER_DEPTH_RESOLVER_DISPOSITION_KIND = (
|
|
'resolver_attempt_limit_replace_or_skip_v1'
|
|
)
|
|
DOCKER_DEPTH_HOLD_ACTIVE_STATES = frozenset({
|
|
'holding', 'resolving', 'active', 'draining', 'completed', 'held',
|
|
})
|
|
|
|
_DOCKER_EXPERIMENT_PROFILES = {
|
|
DOCKER_DEPTH_SELECTOR_VERSION: {
|
|
'query_count': DOCKER_DEPTH_QUERY_COUNT,
|
|
'repositories_per_query': DOCKER_DEPTH_REPOSITORIES_PER_QUERY,
|
|
'shallow_images_per_repository': DOCKER_DEPTH_SHALLOW_IMAGES_PER_REPOSITORY,
|
|
'deep_repositories_per_query': DOCKER_DEPTH_DEEP_REPOSITORIES_PER_QUERY,
|
|
'images_per_repository': DOCKER_DEPTH_DEEP_IMAGES_PER_REPOSITORY,
|
|
'target_limit': DOCKER_DEPTH_MAX_UNIQUE_TARGETS,
|
|
'theoretical_max_targets': DOCKER_DEPTH_QUERY_COUNT * (
|
|
DOCKER_DEPTH_REPOSITORIES_PER_QUERY
|
|
+ DOCKER_DEPTH_DEEP_IMAGES_PER_REPOSITORY - 1
|
|
),
|
|
},
|
|
DOCKER_RANK1_BREADTH_SELECTOR_VERSION: {
|
|
'query_count': DOCKER_DEPTH_QUERY_COUNT,
|
|
'repositories_per_query': DOCKER_RANK1_BREADTH_REPOSITORIES_PER_QUERY,
|
|
'shallow_images_per_repository': 1,
|
|
'deep_repositories_per_query': 1,
|
|
'images_per_repository': 1,
|
|
'target_limit': DOCKER_RANK1_BREADTH_MAX_UNIQUE_TARGETS,
|
|
'theoretical_max_targets': DOCKER_RANK1_BREADTH_MAX_UNIQUE_TARGETS,
|
|
},
|
|
}
|
|
|
|
|
|
def reviewed_docker_experiment_profile(selector_version):
|
|
try:
|
|
return dict(_DOCKER_EXPERIMENT_PROFILES[selector_version])
|
|
except KeyError:
|
|
raise ValueError('Docker experiment selector_version is not reviewed') from None
|
|
|
|
DOCKER_DEPTH_EXPERIMENT_KEYS = frozenset({
|
|
'experiment_key',
|
|
'enabled',
|
|
'queries',
|
|
'repositories_per_query',
|
|
'shallow_images_per_repository',
|
|
'deep_repositories_per_query',
|
|
'deep_images_per_repository',
|
|
'target_limit',
|
|
'selector_version',
|
|
})
|
|
|
|
_IDENTIFIER_RE = re.compile(r'^[a-z0-9](?:[a-z0-9._-]{0,126}[a-z0-9])?$')
|
|
_PLATFORM_COMPONENT_RE = re.compile(r'^[a-z0-9][a-z0-9._-]{0,63}$')
|
|
_POSTGRES_URL_RE = re.compile(r'^postgres(?:ql)?://', re.IGNORECASE)
|
|
_SELECTOR_DISTINCT_GRAPH = 'ordered-layer-digests'
|
|
_SELECTOR_FIRST_THREE_STEPS = (
|
|
'newest_distinct_graph',
|
|
'maximum_marginal_layer_novelty',
|
|
'oldest_distinct_graph',
|
|
)
|
|
_SELECTOR_LATER_SCORE = 'maximum_marginal_layer_novelty'
|
|
_SELECTOR_LATER_TIE_BREAKS = (
|
|
'maximum_minimum_temporal_distance',
|
|
'newest_update',
|
|
'normalized_target',
|
|
'ordered_graph_sha256',
|
|
)
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class DockerDepthExperimentConfig:
|
|
experiment_key: str
|
|
enabled: bool
|
|
collection_generation: str
|
|
queries: tuple
|
|
repositories_per_query: int
|
|
shallow_images_per_repository: int
|
|
deep_repositories_per_query: int
|
|
deep_images_per_repository: int
|
|
target_limit: int
|
|
selector_version: str
|
|
theoretical_max_targets: int
|
|
ordered_query_hash: str
|
|
config_hash: str
|
|
selector_hash: str
|
|
|
|
@property
|
|
def ordered_query_sha256(self):
|
|
return self.ordered_query_hash
|
|
|
|
@property
|
|
def ordered_queries_sha256(self):
|
|
return self.ordered_query_hash
|
|
|
|
@property
|
|
def config_sha256(self):
|
|
return self.config_hash
|
|
|
|
@property
|
|
def selector_sha256(self):
|
|
return self.selector_hash
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class ValidatedDockerDepthConfig:
|
|
normalized_config: dict
|
|
docker_images_per_repository: int
|
|
experiment: Optional[DockerDepthExperimentConfig] = None
|
|
ordered_query_hash: str = ''
|
|
config_hash: str = ''
|
|
selector_hash: str = ''
|
|
|
|
@property
|
|
def config(self):
|
|
return copy.deepcopy(self.normalized_config)
|
|
|
|
@property
|
|
def ordered_query_sha256(self):
|
|
return self.ordered_query_hash
|
|
|
|
@property
|
|
def ordered_queries_sha256(self):
|
|
return self.ordered_query_hash
|
|
|
|
@property
|
|
def config_sha256(self):
|
|
return self.config_hash
|
|
|
|
@property
|
|
def selector_sha256(self):
|
|
return self.selector_hash
|
|
|
|
|
|
def _canonical_sha256(value):
|
|
payload = json.dumps(
|
|
value,
|
|
ensure_ascii=True,
|
|
allow_nan=False,
|
|
sort_keys=True,
|
|
separators=(',', ':'),
|
|
).encode('utf-8')
|
|
return hashlib.sha256(payload).hexdigest()
|
|
|
|
|
|
def canonical_ordered_query_hash(queries):
|
|
return _canonical_sha256(list(queries))
|
|
|
|
|
|
def canonical_dockerhub_discovery_policy(pages, per_page, sort_by, sort_order):
|
|
if isinstance(pages, bool) or not isinstance(pages, int) or not 1 <= pages <= DOCKERHUB_DISCOVERY_MAX_PAGES:
|
|
raise ValueError(
|
|
f'DockerHub discovery pages must be an integer from 1 through {DOCKERHUB_DISCOVERY_MAX_PAGES}'
|
|
)
|
|
if (
|
|
isinstance(per_page, bool)
|
|
or not isinstance(per_page, int)
|
|
or not 1 <= per_page <= DOCKERHUB_DISCOVERY_MAX_PER_PAGE
|
|
):
|
|
raise ValueError(
|
|
'DockerHub discovery per_page must be an integer from 1 through '
|
|
f'{DOCKERHUB_DISCOVERY_MAX_PER_PAGE}'
|
|
)
|
|
if not isinstance(sort_by, str) or not sort_by or sort_by != sort_by.strip():
|
|
raise ValueError('DockerHub discovery sort_by must be a non-empty canonical string')
|
|
if sort_order not in ('asc', 'desc'):
|
|
raise ValueError("DockerHub discovery sort_order must be 'asc' or 'desc'")
|
|
payload = {
|
|
'algorithm_version': DOCKERHUB_DISCOVERY_ALGORITHM_VERSION,
|
|
'page_hard_cap': DOCKERHUB_DISCOVERY_MAX_PAGES,
|
|
'pages': pages,
|
|
'per_page': per_page,
|
|
'per_page_hard_cap': DOCKERHUB_DISCOVERY_MAX_PER_PAGE,
|
|
'sort_by': sort_by,
|
|
'sort_order': sort_order,
|
|
}
|
|
return {**payload, 'policy_sha256': _canonical_sha256(payload)}
|
|
|
|
|
|
def canonical_selector_hash(selector_version):
|
|
return _canonical_sha256({
|
|
'identity': {
|
|
'dependencies': {
|
|
'docker_digest_pattern': DOCKER_DIGEST_RE.pattern,
|
|
'docker_image_pattern': DOCKER_IMAGE_RE.pattern,
|
|
'docker_revision_pattern': DOCKER_REVISION_RE.pattern,
|
|
'docker_target_schema': DOCKER_TAG_TARGET_SCHEMA,
|
|
},
|
|
'distinct_graph': _SELECTOR_DISTINCT_GRAPH,
|
|
'first_three': _SELECTOR_FIRST_THREE_STEPS,
|
|
'implementation_sha256': _selector_implementation_sha256(),
|
|
'later_score': _SELECTOR_LATER_SCORE,
|
|
'later_tie_breaks': _SELECTOR_LATER_TIE_BREAKS,
|
|
},
|
|
'version': selector_version,
|
|
})
|
|
|
|
|
|
def canonical_docker_layer_graph_hash(layers):
|
|
return _canonical_sha256(list(layers))
|
|
|
|
|
|
def canonical_docker_descriptor_hash(digest, media_type, size_bytes):
|
|
return _canonical_sha256({
|
|
'digest': str(digest),
|
|
'media_type': str(media_type),
|
|
'size_bytes': int(size_bytes),
|
|
})
|
|
|
|
|
|
def canonical_docker_depth_selection_evidence_hash(evidence):
|
|
normalized = dict(evidence)
|
|
# Repository tag churn may change this audit metric after rank-one selection.
|
|
normalized.pop('candidate_distinct_graph_count', None)
|
|
return _canonical_sha256(normalized)
|
|
|
|
|
|
def validate_docker_images_per_repository(value):
|
|
if isinstance(value, bool) or not isinstance(value, int):
|
|
raise ValueError('docker_images_per_repository must be an integer from 1 through 10')
|
|
if value < 1 or value > 10:
|
|
raise ValueError('docker_images_per_repository must be an integer from 1 through 10')
|
|
return value
|
|
|
|
|
|
def _selector_graph_layers(candidate):
|
|
layers = candidate.get('layers')
|
|
if not isinstance(layers, (list, tuple)) or not layers:
|
|
return None
|
|
graph = []
|
|
for layer in layers:
|
|
digest = normalize_docker_digest(
|
|
layer.get('digest') if isinstance(layer, dict) else layer
|
|
)
|
|
if not digest:
|
|
return None
|
|
graph.append(digest)
|
|
return tuple(graph)
|
|
|
|
|
|
def _selector_graph_timestamp(value):
|
|
if isinstance(value, bool):
|
|
return float('-inf')
|
|
try:
|
|
timestamp = float(value)
|
|
return timestamp if math.isfinite(timestamp) else float('-inf')
|
|
except (TypeError, ValueError, OverflowError):
|
|
pass
|
|
if isinstance(value, str):
|
|
try:
|
|
parsed = datetime.fromisoformat(value.strip().replace('Z', '+00:00'))
|
|
if parsed.tzinfo is None:
|
|
parsed = parsed.replace(tzinfo=timezone.utc)
|
|
return parsed.timestamp()
|
|
except (ValueError, OverflowError, OSError):
|
|
pass
|
|
return float('-inf')
|
|
|
|
|
|
def _selector_graph_target(candidate):
|
|
target = str(candidate.get('target') or '').strip()
|
|
try:
|
|
return parse_docker_target(target)['target'].lower()
|
|
except (TypeError, ValueError):
|
|
return target.lower()
|
|
|
|
|
|
def _selector_graph_selection_record(candidate, rank, reason, marginal_layers):
|
|
record = {
|
|
key: value for key, value in candidate.items()
|
|
if not key.startswith('_selector_') and key != 'layer_descriptors'
|
|
}
|
|
graph = candidate['_selector_graph']
|
|
layer_count = len(graph)
|
|
descriptors = candidate.get('layer_descriptors')
|
|
if isinstance(descriptors, (tuple, list)) and len(descriptors) == layer_count:
|
|
layer_metadata = tuple({
|
|
'digest': descriptor['digest'],
|
|
'media_type': descriptor['media_type'],
|
|
'size_bytes': descriptor['size'],
|
|
'descriptor_sha256': canonical_docker_descriptor_hash(
|
|
descriptor['digest'], descriptor['media_type'], descriptor['size'],
|
|
),
|
|
'position_from_base': position,
|
|
'position_from_top': layer_count - position + 1,
|
|
} for position, descriptor in enumerate(descriptors, 1))
|
|
else:
|
|
layer_metadata = tuple({
|
|
'digest': digest,
|
|
'position_from_base': position,
|
|
'position_from_top': layer_count - position + 1,
|
|
} for position, digest in enumerate(graph, 1))
|
|
record.update({
|
|
'rank': rank,
|
|
'image_rank': rank,
|
|
'reason': reason,
|
|
'selection_reason': reason,
|
|
'graph': graph,
|
|
'graph_hash': candidate['_selector_graph_hash'],
|
|
'graph_sha256': candidate['_selector_graph_hash'],
|
|
'layers': graph,
|
|
'layer_count': layer_count,
|
|
'marginal_layer_count': marginal_layers,
|
|
'layer_metadata': layer_metadata,
|
|
})
|
|
return record
|
|
|
|
|
|
def select_docker_layer_graphs(candidates, limit=1, *, replacement_pool=False):
|
|
if replacement_pool:
|
|
if isinstance(limit, bool) or not isinstance(limit, int) or not 1 <= limit <= 100:
|
|
raise ValueError('Docker replacement candidate pool must contain 1 through 100 graphs')
|
|
else:
|
|
limit = validate_docker_images_per_repository(limit)
|
|
|
|
def source_index(candidate):
|
|
value = candidate.get('source_index', 0)
|
|
if isinstance(value, bool):
|
|
return 0
|
|
try:
|
|
return int(value)
|
|
except (TypeError, ValueError, OverflowError):
|
|
return 0
|
|
|
|
def order_key(candidate):
|
|
return (
|
|
-candidate['_selector_updated_at'],
|
|
source_index(candidate),
|
|
_selector_graph_target(candidate),
|
|
)
|
|
|
|
valid = []
|
|
for item in candidates or ():
|
|
if not isinstance(item, dict):
|
|
continue
|
|
candidate = dict(item)
|
|
graph = _selector_graph_layers(candidate)
|
|
if graph is None:
|
|
continue
|
|
candidate['_selector_graph'] = graph
|
|
candidate['_selector_graph_hash'] = canonical_docker_layer_graph_hash(graph)
|
|
candidate['_selector_updated_at'] = _selector_graph_timestamp(candidate.get('updated_at'))
|
|
valid.append(candidate)
|
|
|
|
distinct = []
|
|
seen_graphs = set()
|
|
for candidate in sorted(valid, key=order_key):
|
|
graph = candidate['_selector_graph']
|
|
if graph in seen_graphs:
|
|
continue
|
|
seen_graphs.add(graph)
|
|
distinct.append(candidate)
|
|
if not distinct:
|
|
return []
|
|
|
|
selected = []
|
|
first = distinct.pop(0)
|
|
covered_layers = set(first['_selector_graph'])
|
|
selected.append((first, _SELECTOR_FIRST_THREE_STEPS[0], len(covered_layers)))
|
|
|
|
if limit >= 2 and distinct:
|
|
index = max(
|
|
range(len(distinct)),
|
|
key=lambda position: (
|
|
len(set(distinct[position]['_selector_graph']) - covered_layers),
|
|
-position,
|
|
),
|
|
)
|
|
candidate = distinct.pop(index)
|
|
marginal = len(set(candidate['_selector_graph']) - covered_layers)
|
|
selected.append((candidate, _SELECTOR_FIRST_THREE_STEPS[1], marginal))
|
|
covered_layers.update(candidate['_selector_graph'])
|
|
|
|
if limit >= 3 and distinct:
|
|
candidate = distinct.pop()
|
|
marginal = len(set(candidate['_selector_graph']) - covered_layers)
|
|
selected.append((candidate, _SELECTOR_FIRST_THREE_STEPS[2], marginal))
|
|
covered_layers.update(candidate['_selector_graph'])
|
|
|
|
while len(selected) < limit and distinct:
|
|
selected_times = [
|
|
candidate['_selector_updated_at'] for candidate, _reason, _marginal in selected
|
|
if math.isfinite(candidate['_selector_updated_at'])
|
|
]
|
|
|
|
def later_key(candidate):
|
|
updated_at = candidate['_selector_updated_at']
|
|
temporal_distance = (
|
|
min(abs(updated_at - selected_at) for selected_at in selected_times)
|
|
if math.isfinite(updated_at) and selected_times
|
|
else float('-inf')
|
|
)
|
|
return (
|
|
-len(set(candidate['_selector_graph']) - covered_layers),
|
|
-temporal_distance,
|
|
-updated_at,
|
|
_selector_graph_target(candidate),
|
|
candidate['_selector_graph_hash'],
|
|
)
|
|
|
|
candidate = min(distinct, key=later_key)
|
|
distinct.remove(candidate)
|
|
marginal = len(set(candidate['_selector_graph']) - covered_layers)
|
|
selected.append((candidate, _SELECTOR_LATER_SCORE, marginal))
|
|
covered_layers.update(candidate['_selector_graph'])
|
|
|
|
return [
|
|
_selector_graph_selection_record(candidate, rank, reason, marginal)
|
|
for rank, (candidate, reason, marginal) in enumerate(selected, 1)
|
|
]
|
|
|
|
|
|
def _canonical_selector_source(function):
|
|
try:
|
|
source = textwrap.dedent(inspect.getsource(function))
|
|
except (OSError, TypeError) as exc:
|
|
raise RuntimeError('Docker depth selector source authority is unavailable') from exc
|
|
canonical = []
|
|
ignored = {
|
|
tokenize.COMMENT, tokenize.ENCODING, tokenize.ENDMARKER, tokenize.NL,
|
|
}
|
|
structural = {tokenize.INDENT, tokenize.DEDENT, tokenize.NEWLINE}
|
|
for token in tokenize.generate_tokens(io.StringIO(source).readline):
|
|
if token.type in ignored:
|
|
continue
|
|
if token.type in structural:
|
|
canonical.append(tokenize.tok_name[token.type])
|
|
else:
|
|
canonical.append(f'{tokenize.tok_name[token.type]}:{token.string}')
|
|
return '\n'.join(canonical)
|
|
|
|
|
|
def _selector_implementation_sha256(source_functions=None):
|
|
if source_functions is None:
|
|
source_functions = (
|
|
_canonical_sha256,
|
|
canonical_docker_layer_graph_hash,
|
|
canonical_docker_descriptor_hash,
|
|
validate_docker_images_per_repository,
|
|
normalize_docker_digest,
|
|
validate_docker_image_reference,
|
|
serialize_docker_tag_target,
|
|
parse_docker_target,
|
|
_selector_graph_layers,
|
|
_selector_graph_timestamp,
|
|
_selector_graph_target,
|
|
_selector_graph_selection_record,
|
|
select_docker_layer_graphs,
|
|
)
|
|
payload = '\nFUNCTION\n'.join(
|
|
_canonical_selector_source(function) for function in source_functions
|
|
).encode('utf-8')
|
|
return hashlib.sha256(payload).hexdigest()
|
|
|
|
|
|
def _strict_integer(mapping, key, minimum=1, maximum=None):
|
|
value = mapping.get(key)
|
|
if isinstance(value, bool) or not isinstance(value, int):
|
|
raise ValueError(f'docker_depth_experiment.{key} must be an integer')
|
|
if value < minimum or (maximum is not None and value > maximum):
|
|
upper = f' through {maximum}' if maximum is not None else ''
|
|
raise ValueError(
|
|
f'docker_depth_experiment.{key} must be from {minimum}{upper}'
|
|
)
|
|
return value
|
|
|
|
|
|
def _strict_identifier(value, name):
|
|
if not isinstance(value, str) or not _IDENTIFIER_RE.fullmatch(value):
|
|
raise ValueError(
|
|
f'docker_depth_experiment.{name} must be a lowercase stable identifier'
|
|
)
|
|
return value
|
|
|
|
|
|
def _strict_config_integer(value, name, minimum=0, maximum=None):
|
|
if isinstance(value, bool) or not isinstance(value, int):
|
|
raise ValueError(f'{name} must be an integer')
|
|
if value < minimum or (maximum is not None and value > maximum):
|
|
upper = f' through {maximum}' if maximum is not None else ' or greater'
|
|
raise ValueError(f'{name} must be from {minimum}{upper}')
|
|
return value
|
|
|
|
|
|
def validate_dockerhub_discovery_policies(source, queries=None):
|
|
"""Return strict ordered effective query policies without runtime side effects."""
|
|
if not isinstance(source, dict):
|
|
raise ValueError('sources.dockerhub must be a mapping')
|
|
configured_queries = source.get('queries') if queries is None else list(queries)
|
|
if not isinstance(configured_queries, list):
|
|
raise ValueError('sources.dockerhub.queries must be an ordered list')
|
|
configured_queries = tuple(configured_queries)
|
|
if any(
|
|
not isinstance(query, str)
|
|
or not query
|
|
or query != query.strip()
|
|
for query in configured_queries
|
|
):
|
|
raise ValueError('sources.dockerhub.queries contains a malformed query')
|
|
if len(set(configured_queries)) != len(configured_queries):
|
|
raise ValueError('sources.dockerhub.queries must be unique')
|
|
|
|
pages = _strict_config_integer(
|
|
source.get('pages', 1), 'sources.dockerhub.pages', 1,
|
|
DOCKERHUB_DISCOVERY_MAX_PAGES,
|
|
)
|
|
per_page = _strict_config_integer(
|
|
source.get('per_page', 50), 'sources.dockerhub.per_page', 1,
|
|
DOCKERHUB_DISCOVERY_MAX_PER_PAGE,
|
|
)
|
|
max_targets = _strict_config_integer(
|
|
source.get('max_targets', 0), 'sources.dockerhub.max_targets', 0,
|
|
)
|
|
sort_by = source.get('docker_sort_by', source.get('sort_by', 'updated_at'))
|
|
sort_order = source.get('sort_order', 'desc')
|
|
canonical_dockerhub_discovery_policy(pages, per_page, sort_by, sort_order)
|
|
|
|
overrides = source.get('query_overrides', {})
|
|
if not isinstance(overrides, dict):
|
|
raise ValueError('sources.dockerhub.query_overrides must be a mapping')
|
|
unknown_queries = [query for query in overrides if query not in configured_queries]
|
|
if unknown_queries:
|
|
names = ', '.join(sorted(repr(query) for query in unknown_queries))
|
|
raise ValueError(
|
|
'sources.dockerhub.query_overrides contains unconfigured queries: ' + names
|
|
)
|
|
|
|
effective = []
|
|
for query in configured_queries:
|
|
override = overrides.get(query, {})
|
|
if not isinstance(override, dict):
|
|
raise ValueError(
|
|
f'sources.dockerhub.query override for {query!r} must be a mapping'
|
|
)
|
|
unknown = set(override) - {'pages', 'per_page', 'max_targets'}
|
|
if unknown:
|
|
raise ValueError(
|
|
f'sources.dockerhub query override for {query!r} has unsupported keys: '
|
|
+ ', '.join(sorted(repr(key) for key in unknown))
|
|
)
|
|
effective_pages = _strict_config_integer(
|
|
override.get('pages', pages),
|
|
f'sources.dockerhub query override pages for {query!r}', 1,
|
|
DOCKERHUB_DISCOVERY_MAX_PAGES,
|
|
)
|
|
effective_per_page = _strict_config_integer(
|
|
override.get('per_page', per_page),
|
|
f'sources.dockerhub query override per_page for {query!r}', 1,
|
|
DOCKERHUB_DISCOVERY_MAX_PER_PAGE,
|
|
)
|
|
effective_max_targets = _strict_config_integer(
|
|
override.get('max_targets', max_targets),
|
|
f'sources.dockerhub query override max_targets for {query!r}', 0,
|
|
)
|
|
policy = canonical_dockerhub_discovery_policy(
|
|
effective_pages, effective_per_page, sort_by, sort_order,
|
|
)
|
|
effective.append({
|
|
'query': query,
|
|
**policy,
|
|
'max_targets': effective_max_targets,
|
|
})
|
|
return tuple(effective)
|
|
|
|
|
|
def _docker_source(config):
|
|
sources = config.get('sources')
|
|
if sources is None:
|
|
sources = {}
|
|
if not isinstance(sources, dict):
|
|
raise ValueError('sources must be a mapping')
|
|
source = sources.get('dockerhub')
|
|
if source is None:
|
|
source = {}
|
|
if not isinstance(source, dict):
|
|
raise ValueError('sources.dockerhub must be a mapping')
|
|
return source
|
|
|
|
|
|
def validate_docker_depth_config(
|
|
config, *, managed_postgres=None, final_cutover=None,
|
|
):
|
|
"""Validate and copy Docker depth configuration without runtime side effects."""
|
|
if not isinstance(config, dict):
|
|
raise ValueError('configuration root must be a mapping')
|
|
normalized = copy.deepcopy(config)
|
|
global_config = normalized.get('global')
|
|
if global_config is None:
|
|
global_config = {}
|
|
if not isinstance(global_config, dict):
|
|
raise ValueError('global must be a mapping')
|
|
source = _docker_source(normalized)
|
|
|
|
configured_depths = []
|
|
if 'docker_images_per_repository' in global_config:
|
|
configured_depths.append((global_config, 'global'))
|
|
if 'docker_images_per_repository' in source:
|
|
configured_depths.append((source, 'sources.dockerhub'))
|
|
for owner, _name in configured_depths:
|
|
validate_docker_images_per_repository(owner['docker_images_per_repository'])
|
|
image_depth = source.get(
|
|
'docker_images_per_repository',
|
|
global_config.get('docker_images_per_repository', 1),
|
|
)
|
|
image_depth = validate_docker_images_per_repository(image_depth)
|
|
|
|
experiment_mapping = source.get('docker_depth_experiment')
|
|
if experiment_mapping is None:
|
|
return ValidatedDockerDepthConfig(normalized, image_depth)
|
|
if not isinstance(experiment_mapping, dict):
|
|
raise ValueError('docker_depth_experiment must be a mapping')
|
|
|
|
unknown = sorted(set(experiment_mapping) - DOCKER_DEPTH_EXPERIMENT_KEYS)
|
|
if unknown:
|
|
raise ValueError(
|
|
'docker_depth_experiment has unsupported keys: ' + ', '.join(unknown)
|
|
)
|
|
missing = sorted(DOCKER_DEPTH_EXPERIMENT_KEYS - set(experiment_mapping))
|
|
if missing:
|
|
raise ValueError(
|
|
'docker_depth_experiment is missing required keys: ' + ', '.join(missing)
|
|
)
|
|
|
|
experiment_key = _strict_identifier(
|
|
experiment_mapping['experiment_key'], 'experiment_key',
|
|
)
|
|
enabled = experiment_mapping['enabled']
|
|
if not isinstance(enabled, bool):
|
|
raise ValueError('docker_depth_experiment.enabled must be a boolean')
|
|
selector_version = _strict_identifier(
|
|
experiment_mapping['selector_version'], 'selector_version',
|
|
)
|
|
profile = reviewed_docker_experiment_profile(selector_version)
|
|
|
|
queries_value = experiment_mapping['queries']
|
|
if not isinstance(queries_value, list):
|
|
raise ValueError('docker_depth_experiment.queries must be an ordered list')
|
|
queries = tuple(queries_value)
|
|
if len(queries) != DOCKER_DEPTH_QUERY_COUNT:
|
|
raise ValueError(
|
|
f'docker_depth_experiment.queries must contain exactly {DOCKER_DEPTH_QUERY_COUNT} queries'
|
|
)
|
|
if any(
|
|
not isinstance(query, str)
|
|
or not query
|
|
or query != query.strip()
|
|
for query in queries
|
|
):
|
|
raise ValueError('docker_depth_experiment.queries contains a malformed query')
|
|
if len(set(queries)) != len(queries):
|
|
raise ValueError('docker_depth_experiment.queries must be unique')
|
|
source_queries = source.get('queries')
|
|
if not isinstance(source_queries, list) or tuple(source_queries) != queries:
|
|
raise ValueError(
|
|
'docker_depth_experiment.queries must exactly match ordered sources.dockerhub.queries'
|
|
)
|
|
|
|
repositories_per_query = _strict_integer(
|
|
experiment_mapping, 'repositories_per_query', maximum=100,
|
|
)
|
|
shallow_images = _strict_integer(
|
|
experiment_mapping, 'shallow_images_per_repository', maximum=10,
|
|
)
|
|
deep_repositories = _strict_integer(
|
|
experiment_mapping, 'deep_repositories_per_query', maximum=100,
|
|
)
|
|
deep_images = _strict_integer(
|
|
experiment_mapping, 'deep_images_per_repository', maximum=10,
|
|
)
|
|
target_limit = _strict_integer(
|
|
experiment_mapping, 'target_limit', maximum=1000000,
|
|
)
|
|
if deep_repositories > repositories_per_query:
|
|
raise ValueError(
|
|
'docker_depth_experiment deep repositories exceed repositories_per_query'
|
|
)
|
|
if shallow_images > deep_images:
|
|
raise ValueError(
|
|
'docker_depth_experiment shallow image depth exceeds deep image depth'
|
|
)
|
|
theoretical_max = len(queries) * (
|
|
repositories_per_query * shallow_images
|
|
+ deep_repositories * (deep_images - shallow_images)
|
|
)
|
|
if (
|
|
selector_version == DOCKER_DEPTH_SELECTOR_VERSION
|
|
and theoretical_max > target_limit
|
|
):
|
|
raise ValueError(
|
|
'docker_depth_experiment theoretical target maximum exceeds target_limit'
|
|
)
|
|
|
|
fixed_limits = {
|
|
'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'],
|
|
}
|
|
for key, expected in fixed_limits.items():
|
|
if experiment_mapping[key] != expected:
|
|
raise ValueError(
|
|
f'docker_depth_experiment.{key} must equal {expected}'
|
|
)
|
|
theoretical_max = profile['theoretical_max_targets']
|
|
if image_depth != DOCKER_DEPTH_ORDINARY_IMAGES_PER_REPOSITORY:
|
|
raise ValueError(
|
|
'docker_images_per_repository must equal the reviewed ordinary '
|
|
f'resolver depth {DOCKER_DEPTH_ORDINARY_IMAGES_PER_REPOSITORY}'
|
|
)
|
|
|
|
effective_policies = validate_dockerhub_discovery_policies(source, queries)
|
|
if len({policy['policy_sha256'] for policy in effective_policies}) != 1:
|
|
raise ValueError(
|
|
'docker_depth_experiment requires one consistent pages/per_page '
|
|
'discovery policy across all queries'
|
|
)
|
|
|
|
configured_platform_filters = []
|
|
if 'docker_platform_filter_enabled' in global_config:
|
|
configured_platform_filters.append(global_config['docker_platform_filter_enabled'])
|
|
if 'docker_platform_filter_enabled' in source:
|
|
configured_platform_filters.append(source['docker_platform_filter_enabled'])
|
|
if any(not isinstance(value, bool) for value in configured_platform_filters):
|
|
raise ValueError('docker_platform_filter_enabled must be a boolean')
|
|
platform_filter_enabled = source.get(
|
|
'docker_platform_filter_enabled',
|
|
global_config.get('docker_platform_filter_enabled', True),
|
|
)
|
|
|
|
def platform_component(key, default):
|
|
configured = []
|
|
if key in global_config:
|
|
configured.append(global_config[key])
|
|
if key in source:
|
|
configured.append(source[key])
|
|
for value in configured:
|
|
if (
|
|
not isinstance(value, str)
|
|
or not _PLATFORM_COMPONENT_RE.fullmatch(value)
|
|
or value != value.lower()
|
|
):
|
|
raise ValueError(f'{key} must be a lowercase Docker platform identifier')
|
|
return source.get(key, global_config.get(key, default))
|
|
|
|
platform_os = platform_component('docker_platform_os', 'linux')
|
|
platform_arch = platform_component('docker_platform_arch', 'amd64')
|
|
if platform_filter_enabled and (platform_os, platform_arch) != ('linux', 'amd64'):
|
|
raise ValueError(
|
|
'docker_depth_experiment platform filtering supports only linux/amd64'
|
|
)
|
|
|
|
configured_candidate_counts = []
|
|
if 'docker_platform_candidate_tags' in global_config:
|
|
configured_candidate_counts.append((
|
|
global_config['docker_platform_candidate_tags'],
|
|
'global.docker_platform_candidate_tags',
|
|
))
|
|
if 'docker_platform_candidate_tags' in source:
|
|
configured_candidate_counts.append((
|
|
source['docker_platform_candidate_tags'],
|
|
'sources.dockerhub.docker_platform_candidate_tags',
|
|
))
|
|
for value, name in configured_candidate_counts:
|
|
_strict_config_integer(value, name, 1, 100)
|
|
candidate_tags = source.get(
|
|
'docker_platform_candidate_tags',
|
|
global_config.get('docker_platform_candidate_tags', 20),
|
|
)
|
|
candidate_tags = _strict_config_integer(
|
|
candidate_tags, 'docker_platform_candidate_tags', 1, 100,
|
|
)
|
|
if candidate_tags < deep_images:
|
|
raise ValueError(
|
|
'docker_platform_candidate_tags must be at least the configured deep image depth'
|
|
)
|
|
|
|
for key in (
|
|
'docker_repository_refresh_interval_sec',
|
|
'docker_repository_refresh_max_per_cycle',
|
|
):
|
|
for owner, name in ((global_config, 'global'), (source, 'sources.dockerhub')):
|
|
if key in owner:
|
|
_strict_config_integer(owner[key], f'{name}.{key}')
|
|
refresh_max = source.get(
|
|
'docker_repository_refresh_max_per_cycle',
|
|
global_config.get('docker_repository_refresh_max_per_cycle', 0),
|
|
)
|
|
if enabled and refresh_max != 0:
|
|
raise ValueError(
|
|
'docker_depth_experiment requires '
|
|
'docker_repository_refresh_max_per_cycle=0'
|
|
)
|
|
|
|
database_url = global_config.get('database_url')
|
|
if database_url not in (None, ''):
|
|
if not isinstance(database_url, str) or not _POSTGRES_URL_RE.match(database_url):
|
|
raise ValueError(
|
|
'docker_depth_experiment collection requires a PostgreSQL database_url when configured'
|
|
)
|
|
if managed_postgres is None:
|
|
managed_postgres = True
|
|
if managed_postgres is not None and not isinstance(managed_postgres, bool):
|
|
raise ValueError('managed_postgres validation authority must be boolean or None')
|
|
configured_final_cutover = source.get(
|
|
'sync_file_queues', global_config.get('sync_file_queues', True),
|
|
) is False
|
|
if final_cutover is None:
|
|
final_cutover = configured_final_cutover
|
|
elif not isinstance(final_cutover, bool):
|
|
raise ValueError('final_cutover validation authority must be boolean or None')
|
|
if managed_postgres is False:
|
|
raise ValueError(
|
|
'docker_depth_experiment collection requires managed PostgreSQL'
|
|
)
|
|
if not configured_final_cutover or final_cutover is not True:
|
|
raise ValueError(
|
|
'docker_depth_experiment collection requires PostgreSQL final cutover'
|
|
)
|
|
if source.get('mode') != 'search':
|
|
raise ValueError(
|
|
'docker_depth_experiment collection requires Docker Hub search mode'
|
|
)
|
|
if source.get('require_digest') is not True:
|
|
raise ValueError(
|
|
'docker_depth_experiment collection requires require_digest=true'
|
|
)
|
|
|
|
canonical_experiment = {
|
|
'experiment_key': experiment_key,
|
|
'enabled': enabled,
|
|
'queries': list(queries),
|
|
'repositories_per_query': repositories_per_query,
|
|
'shallow_images_per_repository': shallow_images,
|
|
'deep_repositories_per_query': deep_repositories,
|
|
'deep_images_per_repository': deep_images,
|
|
'target_limit': target_limit,
|
|
'selector_version': selector_version,
|
|
}
|
|
source['docker_depth_experiment'] = canonical_experiment
|
|
ordered_query_hash = canonical_ordered_query_hash(queries)
|
|
selector_hash = canonical_selector_hash(selector_version)
|
|
semantic_experiment = {
|
|
key: value for key, value in canonical_experiment.items()
|
|
if key != 'enabled'
|
|
}
|
|
config_hash = _canonical_sha256({
|
|
'collection_generation': DOCKER_DEPTH_COLLECTION_GENERATION,
|
|
'discovery_policies': list(effective_policies),
|
|
'experiment': semantic_experiment,
|
|
'mode': source.get('mode'),
|
|
'ordinary_images_per_repository': image_depth,
|
|
'platform': {
|
|
'architecture': platform_arch,
|
|
'candidate_tags': candidate_tags,
|
|
'filter_enabled': platform_filter_enabled,
|
|
'os': platform_os,
|
|
},
|
|
'require_digest': source.get('require_digest'),
|
|
'schema': 'docker-depth-experiment-v3',
|
|
})
|
|
experiment = DockerDepthExperimentConfig(
|
|
experiment_key=experiment_key,
|
|
enabled=enabled,
|
|
collection_generation=DOCKER_DEPTH_COLLECTION_GENERATION,
|
|
queries=queries,
|
|
repositories_per_query=repositories_per_query,
|
|
shallow_images_per_repository=shallow_images,
|
|
deep_repositories_per_query=deep_repositories,
|
|
deep_images_per_repository=deep_images,
|
|
target_limit=target_limit,
|
|
selector_version=selector_version,
|
|
theoretical_max_targets=theoretical_max,
|
|
ordered_query_hash=ordered_query_hash,
|
|
config_hash=config_hash,
|
|
selector_hash=selector_hash,
|
|
)
|
|
return ValidatedDockerDepthConfig(
|
|
normalized_config=normalized,
|
|
docker_images_per_repository=image_depth,
|
|
experiment=experiment,
|
|
ordered_query_hash=ordered_query_hash,
|
|
config_hash=config_hash,
|
|
selector_hash=selector_hash,
|
|
)
|
|
|
|
|
|
def _strict_positive_int(value, name, maximum=None):
|
|
if isinstance(value, bool) or not isinstance(value, int) or value < 1:
|
|
raise ValueError(f'{name} must be a positive integer')
|
|
if maximum is not None and value > maximum:
|
|
raise ValueError(f'{name} exceeds its {maximum} bound')
|
|
return value
|
|
|
|
|
|
def _strict_nonnegative_int(value, name):
|
|
if isinstance(value, bool) or not isinstance(value, int) or value < 0:
|
|
raise ValueError(f'{name} must be a non-negative integer')
|
|
return value
|
|
|
|
|
|
def _valid_sha256(value):
|
|
return bool(re.fullmatch(r'[a-f0-9]{64}', str(value or '')))
|
|
|
|
|
|
def _docker_depth_authority(experiment, provenance_policy_sha256):
|
|
if experiment is None:
|
|
raise ValueError('Docker depth experiment authority is unavailable')
|
|
queries = tuple(getattr(experiment, 'queries', ()))
|
|
if (
|
|
len(queries) != DOCKER_DEPTH_QUERY_COUNT
|
|
or len(set(queries)) != DOCKER_DEPTH_QUERY_COUNT
|
|
or any(not isinstance(query, str) or not query or query != query.strip()
|
|
for query in queries)
|
|
):
|
|
raise ValueError('Docker depth planning requires all 61 unique ordered queries')
|
|
ordered_queries_sha256 = str(getattr(experiment, 'ordered_query_hash', '') or '')
|
|
if ordered_queries_sha256 != canonical_ordered_query_hash(queries):
|
|
raise ValueError('Docker depth ordered-query hash conflicts')
|
|
selector_version = str(getattr(experiment, 'selector_version', '') or '')
|
|
selector_sha256 = str(getattr(experiment, 'selector_hash', '') or '')
|
|
if (
|
|
selector_version not in _DOCKER_EXPERIMENT_PROFILES
|
|
or selector_sha256 != canonical_selector_hash(selector_version)
|
|
):
|
|
raise ValueError('Docker depth selector authority conflicts')
|
|
config_sha256 = str(getattr(experiment, 'config_hash', '') or '')
|
|
provenance_policy_sha256 = str(provenance_policy_sha256 or '')
|
|
if not _valid_sha256(config_sha256) or not _valid_sha256(provenance_policy_sha256):
|
|
raise ValueError('Docker depth configuration or provenance policy hash is invalid')
|
|
profile = reviewed_docker_experiment_profile(selector_version)
|
|
fixed = (
|
|
(getattr(experiment, 'repositories_per_query', None),
|
|
profile['repositories_per_query']),
|
|
(getattr(experiment, 'shallow_images_per_repository', None),
|
|
profile['shallow_images_per_repository']),
|
|
(getattr(experiment, 'deep_repositories_per_query', None),
|
|
profile['deep_repositories_per_query']),
|
|
(getattr(experiment, 'deep_images_per_repository', None),
|
|
profile['images_per_repository']),
|
|
(getattr(experiment, 'target_limit', None), profile['target_limit']),
|
|
)
|
|
if any(actual != expected for actual, expected in fixed):
|
|
raise ValueError('Docker depth planning limits conflict with the reviewed experiment')
|
|
theoretical_max = profile['theoretical_max_targets']
|
|
if getattr(experiment, 'theoretical_max_targets', None) != theoretical_max:
|
|
raise ValueError('Docker depth theoretical target capacity conflicts')
|
|
experiment_key = _strict_identifier(
|
|
getattr(experiment, 'experiment_key', None), 'experiment_key',
|
|
)
|
|
collection_generation = str(
|
|
getattr(experiment, 'collection_generation', '') or ''
|
|
)
|
|
if collection_generation != DOCKER_DEPTH_COLLECTION_GENERATION:
|
|
raise ValueError('Docker depth collection generation authority conflicts')
|
|
return {
|
|
'experiment_key': experiment_key,
|
|
'source': 'dockerhub',
|
|
'collection_generation': collection_generation,
|
|
'queries': queries,
|
|
'config_sha256': config_sha256,
|
|
'ordered_queries_sha256': ordered_queries_sha256,
|
|
'selector_version': selector_version,
|
|
'selector_sha256': selector_sha256,
|
|
'provenance_policy_sha256': provenance_policy_sha256,
|
|
'query_count': DOCKER_DEPTH_QUERY_COUNT,
|
|
'repositories_per_query': profile['repositories_per_query'],
|
|
'images_per_repository': profile['images_per_repository'],
|
|
'target_limit': profile['target_limit'],
|
|
'theoretical_max_targets': theoretical_max,
|
|
}
|
|
|
|
|
|
def docker_depth_resolver_authority(experiment, provenance_policy_sha256):
|
|
authority = _docker_depth_authority(experiment, provenance_policy_sha256)
|
|
enabled = getattr(experiment, 'enabled', None)
|
|
if not isinstance(enabled, bool):
|
|
raise ValueError('Docker depth resolver enabled authority is invalid')
|
|
return {**authority, 'enabled': enabled}
|
|
|
|
|
|
def _cohort_plan_document(authority, planned_queries):
|
|
round_robin = []
|
|
by_ordinal = {
|
|
item['query_ordinal']: item for item in planned_queries
|
|
}
|
|
for repository_rank in range(1, authority['repositories_per_query'] + 1):
|
|
for query_ordinal in range(authority['query_count']):
|
|
query_plan = by_ordinal[query_ordinal]
|
|
if repository_rank > len(query_plan['repositories']):
|
|
continue
|
|
repository = query_plan['repositories'][repository_rank - 1]
|
|
round_robin.append({
|
|
'query_ordinal': query_ordinal,
|
|
'repository_rank': repository_rank,
|
|
'repository_queue_id': repository['repository_queue_id'],
|
|
})
|
|
return {
|
|
'schema': 2,
|
|
'type': DOCKER_DEPTH_COHORT_PLAN_TYPE,
|
|
'experiment_key': authority['experiment_key'],
|
|
'source': authority['source'],
|
|
'collection_generation': authority['collection_generation'],
|
|
'config_sha256': authority['config_sha256'],
|
|
'ordered_queries_sha256': authority['ordered_queries_sha256'],
|
|
'selector_version': authority['selector_version'],
|
|
'selector_sha256': authority['selector_sha256'],
|
|
'provenance_policy_sha256': authority['provenance_policy_sha256'],
|
|
'query_count': authority['query_count'],
|
|
'repositories_per_query': authority['repositories_per_query'],
|
|
'images_per_repository': authority['images_per_repository'],
|
|
'target_limit': authority['target_limit'],
|
|
'theoretical_max_targets': authority['theoretical_max_targets'],
|
|
'queries': planned_queries,
|
|
'repository_round_robin': round_robin,
|
|
}
|
|
|
|
|
|
def build_docker_depth_cohort_plan(
|
|
experiment, provenance_policy_sha256, candidates_by_query,
|
|
):
|
|
"""Build the immutable cohort from already fenced fresh-candidate rows."""
|
|
authority = _docker_depth_authority(experiment, provenance_policy_sha256)
|
|
if not isinstance(candidates_by_query, dict) or set(candidates_by_query) != set(
|
|
authority['queries']
|
|
):
|
|
raise ValueError('Docker depth candidates must cover exactly all 61 ordered queries')
|
|
normalized_by_query = []
|
|
for query_ordinal, query in enumerate(authority['queries']):
|
|
values = candidates_by_query[query]
|
|
if isinstance(values, (str, bytes)):
|
|
raise ValueError('Docker depth query candidates are invalid')
|
|
try:
|
|
values = list(values)
|
|
except TypeError:
|
|
raise ValueError('Docker depth query candidates are invalid') from None
|
|
normalized = []
|
|
seen = set()
|
|
for raw in values:
|
|
if not isinstance(raw, dict):
|
|
raw = dict(raw)
|
|
queue_id = _strict_positive_int(
|
|
raw.get('repository_queue_id'), 'repository_queue_id',
|
|
)
|
|
if queue_id in seen:
|
|
raise ValueError('Docker depth query candidates contain a duplicate repository')
|
|
seen.add(queue_id)
|
|
normalized.append({
|
|
'repository_queue_id': queue_id,
|
|
'eligibility_page_id': _strict_positive_int(
|
|
raw.get('eligibility_page_id'), 'eligibility_page_id',
|
|
),
|
|
'best_search_rank': _strict_positive_int(
|
|
raw.get('best_search_rank'), 'best_search_rank',
|
|
),
|
|
'valid_distinct_graph_count': _strict_nonnegative_int(
|
|
raw.get('valid_distinct_graph_count', 0),
|
|
'valid_distinct_graph_count',
|
|
),
|
|
})
|
|
normalized.sort(key=lambda item: (
|
|
item['best_search_rank'], item['repository_queue_id'],
|
|
))
|
|
normalized_by_query.append(normalized)
|
|
|
|
if authority['selector_version'] == DOCKER_DEPTH_SELECTOR_VERSION:
|
|
selected_by_query = [
|
|
values[:DOCKER_DEPTH_REPOSITORIES_PER_QUERY]
|
|
for values in normalized_by_query
|
|
]
|
|
else:
|
|
selected_by_query = [[] for _query in authority['queries']]
|
|
cursors = [0] * len(selected_by_query)
|
|
seen_queue_ids = set()
|
|
selected_count = 0
|
|
while selected_count < authority['target_limit']:
|
|
progressed = False
|
|
for ordinal, values in enumerate(normalized_by_query):
|
|
if len(selected_by_query[ordinal]) >= authority['repositories_per_query']:
|
|
continue
|
|
while (
|
|
cursors[ordinal] < len(values)
|
|
and values[cursors[ordinal]]['repository_queue_id'] in seen_queue_ids
|
|
):
|
|
cursors[ordinal] += 1
|
|
if cursors[ordinal] >= len(values):
|
|
continue
|
|
candidate = values[cursors[ordinal]]
|
|
cursors[ordinal] += 1
|
|
selected_by_query[ordinal].append(candidate)
|
|
seen_queue_ids.add(candidate['repository_queue_id'])
|
|
selected_count += 1
|
|
progressed = True
|
|
if selected_count == authority['target_limit']:
|
|
break
|
|
if not progressed:
|
|
break
|
|
if selected_count != authority['target_limit']:
|
|
raise ValueError('Docker rank1 breadth cohort cannot fill its exact target limit')
|
|
|
|
planned_queries = []
|
|
for query_ordinal, query in enumerate(authority['queries']):
|
|
selected = selected_by_query[query_ordinal]
|
|
deep = min(selected, key=lambda item: (
|
|
-item['valid_distinct_graph_count'],
|
|
item['best_search_rank'],
|
|
item['repository_queue_id'],
|
|
)) if selected else None
|
|
repositories = []
|
|
for repository_rank, item in enumerate(selected, 1):
|
|
repositories.append({
|
|
'repository_queue_id': item['repository_queue_id'],
|
|
'eligibility_page_id': item['eligibility_page_id'],
|
|
'repository_rank': repository_rank,
|
|
'is_deep_probe': bool(
|
|
deep and item['repository_queue_id'] == deep['repository_queue_id']
|
|
),
|
|
})
|
|
planned_queries.append({
|
|
'query_ordinal': query_ordinal,
|
|
'query': query,
|
|
'query_sha256': _canonical_sha256(query),
|
|
'selected_repository_count': len(repositories),
|
|
'repositories': repositories,
|
|
})
|
|
if authority['selector_version'] == DOCKER_RANK1_BREADTH_SELECTOR_VERSION:
|
|
queue_ids = [
|
|
repository['repository_queue_id']
|
|
for query in planned_queries for repository in query['repositories']
|
|
]
|
|
if len(queue_ids) != len(set(queue_ids)):
|
|
raise ValueError(
|
|
'Docker rank1 breadth cohort contains a duplicate physical repository'
|
|
)
|
|
return _cohort_plan_document(authority, planned_queries)
|
|
|
|
|
|
def canonical_docker_depth_plan_hash(plan):
|
|
if not isinstance(plan, dict) or plan.get('type') != DOCKER_DEPTH_COHORT_PLAN_TYPE:
|
|
raise ValueError('Docker depth cohort plan type is invalid')
|
|
return _canonical_sha256(plan)
|
|
|
|
|
|
def _validate_docker_depth_cohort_plan(
|
|
manifest, experiment=None, provenance_policy_sha256=None,
|
|
):
|
|
"""Validate the immutable plan carried by a reviewed cohort manifest."""
|
|
expected_keys = {
|
|
'schema', 'type', 'experiment_key', 'source', 'collection_generation',
|
|
'config_sha256', 'ordered_queries_sha256', 'selector_version',
|
|
'selector_sha256', 'provenance_policy_sha256', 'query_count',
|
|
'repositories_per_query', 'images_per_repository', 'target_limit',
|
|
'theoretical_max_targets', 'queries', 'repository_round_robin',
|
|
}
|
|
if (
|
|
not isinstance(manifest, dict)
|
|
or set(manifest) != expected_keys
|
|
or type(manifest.get('schema')) is not int
|
|
or manifest.get('schema') != 2
|
|
or manifest.get('type') != DOCKER_DEPTH_COHORT_PLAN_TYPE
|
|
):
|
|
raise ValueError('Docker depth cohort manifest shape is invalid')
|
|
if (
|
|
manifest.get('source') != 'dockerhub'
|
|
or manifest.get('collection_generation') != DOCKER_DEPTH_COLLECTION_GENERATION
|
|
or manifest.get('selector_version') not in _DOCKER_EXPERIMENT_PROFILES
|
|
or manifest.get('selector_sha256')
|
|
!= canonical_selector_hash(manifest.get('selector_version'))
|
|
or any(not _valid_sha256(manifest.get(name)) for name in (
|
|
'config_sha256', 'ordered_queries_sha256', 'selector_sha256',
|
|
'provenance_policy_sha256',
|
|
))
|
|
):
|
|
raise ValueError('Docker depth cohort manifest authority is invalid')
|
|
_strict_identifier(manifest.get('experiment_key'), 'experiment_key')
|
|
profile = reviewed_docker_experiment_profile(manifest['selector_version'])
|
|
fixed = {
|
|
'query_count': profile['query_count'],
|
|
'repositories_per_query': profile['repositories_per_query'],
|
|
'images_per_repository': profile['images_per_repository'],
|
|
'target_limit': profile['target_limit'],
|
|
'theoretical_max_targets': profile['theoretical_max_targets'],
|
|
}
|
|
if any(type(manifest.get(name)) is not int or manifest[name] != value
|
|
for name, value in fixed.items()):
|
|
raise ValueError('Docker depth cohort manifest limits conflict')
|
|
raw_queries = manifest.get('queries')
|
|
if not isinstance(raw_queries, list) or len(raw_queries) != fixed['query_count']:
|
|
raise ValueError('Docker depth cohort manifest query coverage is incomplete')
|
|
planned_queries = []
|
|
queries = []
|
|
global_queue_ids = set()
|
|
for ordinal, raw_query in enumerate(raw_queries):
|
|
if not isinstance(raw_query, dict) or set(raw_query) != {
|
|
'query_ordinal', 'query', 'query_sha256',
|
|
'selected_repository_count', 'repositories',
|
|
}:
|
|
raise ValueError('Docker depth cohort manifest query shape is invalid')
|
|
query = raw_query.get('query')
|
|
if (
|
|
type(raw_query.get('query_ordinal')) is not int
|
|
or raw_query['query_ordinal'] != ordinal
|
|
or not isinstance(query, str)
|
|
or not query
|
|
or query != query.strip()
|
|
or raw_query.get('query_sha256') != _canonical_sha256(query)
|
|
):
|
|
raise ValueError('Docker depth cohort manifest query identity conflicts')
|
|
raw_repositories = raw_query.get('repositories')
|
|
selected_repository_count = raw_query.get('selected_repository_count')
|
|
if (
|
|
not isinstance(raw_repositories, list)
|
|
or type(selected_repository_count) is not int
|
|
or not 0 <= selected_repository_count <= fixed['repositories_per_query']
|
|
or len(raw_repositories) != selected_repository_count
|
|
):
|
|
raise ValueError('Docker depth cohort manifest membership count conflicts')
|
|
repositories = []
|
|
queue_ids = set()
|
|
for rank, raw_repository in enumerate(raw_repositories, 1):
|
|
if not isinstance(raw_repository, dict) or set(raw_repository) != {
|
|
'repository_queue_id', 'eligibility_page_id', 'repository_rank',
|
|
'is_deep_probe',
|
|
}:
|
|
raise ValueError('Docker depth cohort manifest membership shape is invalid')
|
|
queue_id = _strict_positive_int(
|
|
raw_repository.get('repository_queue_id'), 'repository_queue_id',
|
|
)
|
|
if queue_id in queue_ids:
|
|
raise ValueError('Docker depth cohort manifest membership is duplicated')
|
|
queue_ids.add(queue_id)
|
|
if (
|
|
manifest['selector_version'] == DOCKER_RANK1_BREADTH_SELECTOR_VERSION
|
|
and queue_id in global_queue_ids
|
|
):
|
|
raise ValueError('Docker rank1 breadth cohort physical membership is duplicated')
|
|
global_queue_ids.add(queue_id)
|
|
if (
|
|
type(raw_repository.get('repository_rank')) is not int
|
|
or raw_repository['repository_rank'] != rank
|
|
or type(raw_repository.get('is_deep_probe')) is not bool
|
|
):
|
|
raise ValueError('Docker depth cohort manifest rank identity conflicts')
|
|
repositories.append({
|
|
'repository_queue_id': queue_id,
|
|
'eligibility_page_id': _strict_positive_int(
|
|
raw_repository.get('eligibility_page_id'), 'eligibility_page_id',
|
|
),
|
|
'repository_rank': rank,
|
|
'is_deep_probe': raw_repository['is_deep_probe'],
|
|
})
|
|
expected_deep_count = 1 if repositories else 0
|
|
if sum(int(item['is_deep_probe']) for item in repositories) != expected_deep_count:
|
|
raise ValueError('Docker depth cohort manifest deep membership conflicts')
|
|
queries.append(query)
|
|
planned_queries.append({
|
|
'query_ordinal': ordinal,
|
|
'query': query,
|
|
'query_sha256': _canonical_sha256(query),
|
|
'selected_repository_count': selected_repository_count,
|
|
'repositories': repositories,
|
|
})
|
|
if len(set(queries)) != fixed['query_count']:
|
|
raise ValueError('Docker depth cohort manifest queries are duplicated')
|
|
if (
|
|
manifest['selector_version'] == DOCKER_RANK1_BREADTH_SELECTOR_VERSION
|
|
and len(global_queue_ids) != fixed['target_limit']
|
|
):
|
|
raise ValueError('Docker rank1 breadth cohort target count conflicts')
|
|
if manifest['ordered_queries_sha256'] != canonical_ordered_query_hash(queries):
|
|
raise ValueError('Docker depth cohort manifest ordered-query hash conflicts')
|
|
authority = {
|
|
'experiment_key': manifest['experiment_key'],
|
|
'source': manifest['source'],
|
|
'collection_generation': manifest['collection_generation'],
|
|
'queries': tuple(queries),
|
|
'config_sha256': manifest['config_sha256'],
|
|
'ordered_queries_sha256': manifest['ordered_queries_sha256'],
|
|
'selector_version': manifest['selector_version'],
|
|
'selector_sha256': manifest['selector_sha256'],
|
|
'provenance_policy_sha256': manifest['provenance_policy_sha256'],
|
|
**fixed,
|
|
}
|
|
if experiment is not None:
|
|
expected_authority = _docker_depth_authority(
|
|
experiment,
|
|
provenance_policy_sha256 or manifest['provenance_policy_sha256'],
|
|
)
|
|
if authority != expected_authority:
|
|
raise ValueError('Docker depth cohort manifest authority drifted')
|
|
canonical = _cohort_plan_document(authority, planned_queries)
|
|
if canonical != manifest:
|
|
raise ValueError('Docker depth cohort manifest is not canonical')
|
|
normalized = copy.deepcopy(canonical)
|
|
return normalized, canonical_docker_depth_plan_hash(normalized)
|
|
|
|
|
|
def _cohort_provenance_snapshot(authority, candidates_by_query):
|
|
if not isinstance(candidates_by_query, dict) or set(candidates_by_query) != set(
|
|
authority['queries']
|
|
):
|
|
raise ValueError('Docker depth candidates must cover exactly all 61 ordered queries')
|
|
queries = []
|
|
for query_ordinal, query in enumerate(authority['queries']):
|
|
values = candidates_by_query[query]
|
|
if isinstance(values, (str, bytes)):
|
|
raise ValueError('Docker depth query candidates are invalid')
|
|
try:
|
|
values = list(values)
|
|
except TypeError:
|
|
raise ValueError('Docker depth query candidates are invalid') from None
|
|
normalized = []
|
|
seen = set()
|
|
for raw in values:
|
|
if not isinstance(raw, dict):
|
|
raw = dict(raw)
|
|
queue_id = _strict_positive_int(
|
|
raw.get('repository_queue_id'), 'repository_queue_id',
|
|
)
|
|
if queue_id in seen:
|
|
raise ValueError('Docker depth query candidates contain a duplicate repository')
|
|
seen.add(queue_id)
|
|
normalized.append({
|
|
'repository_queue_id': queue_id,
|
|
'eligibility_page_id': _strict_positive_int(
|
|
raw.get('eligibility_page_id'), 'eligibility_page_id',
|
|
),
|
|
'best_search_rank': _strict_positive_int(
|
|
raw.get('best_search_rank'), 'best_search_rank',
|
|
),
|
|
'valid_distinct_graph_count': _strict_nonnegative_int(
|
|
raw.get('valid_distinct_graph_count', 0),
|
|
'valid_distinct_graph_count',
|
|
),
|
|
})
|
|
normalized.sort(key=lambda item: (
|
|
item['best_search_rank'], item['repository_queue_id'],
|
|
))
|
|
selected = normalized[:authority['repositories_per_query']]
|
|
queries.append({
|
|
'query_ordinal': query_ordinal,
|
|
'query_sha256': _canonical_sha256(query),
|
|
'selected_repository_count': len(selected),
|
|
'repositories': selected,
|
|
})
|
|
return {
|
|
'schema': 2,
|
|
'type': 'truf-docker-depth-cohort-provenance-snapshot-v2',
|
|
'collection_generation': authority['collection_generation'],
|
|
'config_sha256': authority['config_sha256'],
|
|
'ordered_queries_sha256': authority['ordered_queries_sha256'],
|
|
'selector_sha256': authority['selector_sha256'],
|
|
'provenance_policy_sha256': authority['provenance_policy_sha256'],
|
|
'queries': queries,
|
|
}
|
|
|
|
|
|
def build_docker_depth_cohort_manifest(
|
|
experiment, provenance_policy_sha256, candidates_by_query,
|
|
):
|
|
authority = _docker_depth_authority(experiment, provenance_policy_sha256)
|
|
plan = build_docker_depth_cohort_plan(
|
|
experiment, provenance_policy_sha256, candidates_by_query,
|
|
)
|
|
snapshot = _cohort_provenance_snapshot(authority, candidates_by_query)
|
|
manifest = {
|
|
'schema': 2,
|
|
'type': DOCKER_DEPTH_COHORT_MANIFEST_TYPE,
|
|
'version': 2,
|
|
'experiment_key': authority['experiment_key'],
|
|
'source': authority['source'],
|
|
'collection_generation': authority['collection_generation'],
|
|
'collection_generation_sha256': _canonical_sha256(
|
|
authority['collection_generation']
|
|
),
|
|
'config_sha256': authority['config_sha256'],
|
|
'ordered_queries_sha256': authority['ordered_queries_sha256'],
|
|
'selector_version': authority['selector_version'],
|
|
'selector_sha256': authority['selector_sha256'],
|
|
'provenance_policy_sha256': authority['provenance_policy_sha256'],
|
|
'provenance_snapshot_sha256': _canonical_sha256(snapshot),
|
|
'query_count': authority['query_count'],
|
|
'repository_count': sum(
|
|
item['selected_repository_count'] for item in plan['queries']
|
|
),
|
|
'plan_sha256': canonical_docker_depth_plan_hash(plan),
|
|
'plan': plan,
|
|
}
|
|
return validate_docker_depth_cohort_manifest(
|
|
manifest, experiment, provenance_policy_sha256,
|
|
)[0]
|
|
|
|
|
|
def validate_docker_depth_cohort_manifest(
|
|
manifest, experiment=None, provenance_policy_sha256=None,
|
|
):
|
|
expected_keys = {
|
|
'schema', 'type', 'version', 'experiment_key', 'source',
|
|
'collection_generation', 'collection_generation_sha256',
|
|
'config_sha256', 'ordered_queries_sha256', 'selector_version',
|
|
'selector_sha256', 'provenance_policy_sha256',
|
|
'provenance_snapshot_sha256', 'query_count', 'repository_count',
|
|
'plan_sha256', 'plan',
|
|
}
|
|
if (
|
|
not isinstance(manifest, dict)
|
|
or set(manifest) != expected_keys
|
|
or type(manifest.get('schema')) is not int
|
|
or manifest.get('schema') != 2
|
|
or type(manifest.get('version')) is not int
|
|
or manifest.get('version') != 2
|
|
or manifest.get('type') != DOCKER_DEPTH_COHORT_MANIFEST_TYPE
|
|
):
|
|
raise ValueError('Docker depth cohort review manifest shape is invalid')
|
|
plan, plan_sha256 = _validate_docker_depth_cohort_plan(
|
|
manifest.get('plan'), experiment, provenance_policy_sha256,
|
|
)
|
|
for name in (
|
|
'collection_generation_sha256', 'config_sha256',
|
|
'ordered_queries_sha256', 'selector_sha256',
|
|
'provenance_policy_sha256', 'provenance_snapshot_sha256', 'plan_sha256',
|
|
):
|
|
if not _valid_sha256(manifest.get(name)):
|
|
raise ValueError('Docker depth cohort review manifest contains an invalid hash')
|
|
if (
|
|
manifest['experiment_key'] != plan['experiment_key']
|
|
or manifest['source'] != plan['source']
|
|
or manifest['collection_generation'] != plan['collection_generation']
|
|
or manifest['collection_generation_sha256']
|
|
!= _canonical_sha256(plan['collection_generation'])
|
|
or manifest['config_sha256'] != plan['config_sha256']
|
|
or manifest['ordered_queries_sha256'] != plan['ordered_queries_sha256']
|
|
or manifest['selector_version'] != plan['selector_version']
|
|
or manifest['selector_sha256'] != plan['selector_sha256']
|
|
or manifest['provenance_policy_sha256'] != plan['provenance_policy_sha256']
|
|
or type(manifest.get('query_count')) is not int
|
|
or manifest['query_count'] != plan['query_count']
|
|
or type(manifest.get('repository_count')) is not int
|
|
or manifest['repository_count']
|
|
!= sum(item['selected_repository_count'] for item in plan['queries'])
|
|
or manifest['plan_sha256'] != plan_sha256
|
|
):
|
|
raise ValueError('Docker depth cohort review manifest evidence conflicts')
|
|
normalized = copy.deepcopy(manifest)
|
|
normalized['plan'] = plan
|
|
if normalized != manifest:
|
|
raise ValueError('Docker depth cohort review manifest is not canonical')
|
|
return normalized, _canonical_sha256(normalized)
|
|
|
|
|
|
def _require_postgres_experiment_db(db, operation):
|
|
conn = getattr(db, 'conn', None)
|
|
if not conn or not getattr(conn, 'is_postgres', False):
|
|
raise RuntimeError(f'{operation} requires managed PostgreSQL')
|
|
return conn
|
|
|
|
|
|
def _experiment_identity_matches(row, authority):
|
|
expected = {
|
|
'experiment_key': authority['experiment_key'],
|
|
'source': authority['source'],
|
|
'collection_generation': authority['collection_generation'],
|
|
'config_sha256': authority['config_sha256'],
|
|
'ordered_queries_sha256': authority['ordered_queries_sha256'],
|
|
'selector_version': authority['selector_version'],
|
|
'selector_sha256': authority['selector_sha256'],
|
|
'provenance_policy_sha256': authority['provenance_policy_sha256'],
|
|
'query_count': authority['query_count'],
|
|
'repositories_per_query': authority['repositories_per_query'],
|
|
'images_per_repository': authority['images_per_repository'],
|
|
'target_limit': authority['target_limit'],
|
|
}
|
|
return all(
|
|
(int(row[name]) if isinstance(value, int) else str(row[name])) == value
|
|
for name, value in expected.items()
|
|
)
|
|
|
|
|
|
def _stored_cohort_plan(conn, experiment_row, authority):
|
|
experiment_id = int(experiment_row['id'])
|
|
query_rows = conn.execute(
|
|
'''SELECT query_ordinal, query, query_sha256,
|
|
required_repository_count, selected_repository_count
|
|
FROM docker_depth_experiment_queries
|
|
WHERE experiment_id = ? ORDER BY query_ordinal''',
|
|
(experiment_id,),
|
|
).fetchall()
|
|
repository_rows = conn.execute(
|
|
'''SELECT query_ordinal, repository_queue_id, eligibility_page_id,
|
|
repository_rank, planned_is_deep_probe
|
|
FROM docker_depth_experiment_repositories
|
|
WHERE experiment_id = ? ORDER BY query_ordinal, repository_rank''',
|
|
(experiment_id,),
|
|
).fetchall()
|
|
expected_repository_count = sum(
|
|
int(row['selected_repository_count']) for row in query_rows
|
|
) if len(query_rows) == authority['query_count'] else -1
|
|
if (
|
|
len(query_rows) != authority['query_count']
|
|
or len(repository_rows) != expected_repository_count
|
|
):
|
|
raise RuntimeError('Docker depth persisted cohort is incomplete')
|
|
repositories_by_query = {
|
|
ordinal: [] for ordinal in range(authority['query_count'])
|
|
}
|
|
for row in repository_rows:
|
|
ordinal = int(row['query_ordinal'])
|
|
if ordinal not in repositories_by_query:
|
|
raise RuntimeError('Docker depth persisted repository ordinal is invalid')
|
|
repositories_by_query[ordinal].append({
|
|
'repository_queue_id': int(row['repository_queue_id']),
|
|
'eligibility_page_id': int(row['eligibility_page_id']),
|
|
'repository_rank': int(row['repository_rank']),
|
|
'is_deep_probe': bool(row['planned_is_deep_probe']),
|
|
})
|
|
planned_queries = []
|
|
for ordinal, row in enumerate(query_rows):
|
|
repositories = repositories_by_query[ordinal]
|
|
selected_repository_count = int(row['selected_repository_count'])
|
|
if (
|
|
int(row['query_ordinal']) != ordinal
|
|
or str(row['query']) != authority['queries'][ordinal]
|
|
or str(row['query_sha256']) != _canonical_sha256(authority['queries'][ordinal])
|
|
or int(row['required_repository_count']) != authority['repositories_per_query']
|
|
or not 0 <= selected_repository_count <= authority['repositories_per_query']
|
|
or len(repositories) != selected_repository_count
|
|
or [item['repository_rank'] for item in repositories]
|
|
!= list(range(1, selected_repository_count + 1))
|
|
or sum(int(item['is_deep_probe']) for item in repositories)
|
|
!= (1 if selected_repository_count else 0)
|
|
):
|
|
raise RuntimeError('Docker depth persisted cohort identity conflicts')
|
|
planned_queries.append({
|
|
'query_ordinal': ordinal,
|
|
'query': authority['queries'][ordinal],
|
|
'query_sha256': str(row['query_sha256']),
|
|
'selected_repository_count': selected_repository_count,
|
|
'repositories': repositories,
|
|
})
|
|
return _cohort_plan_document(authority, planned_queries)
|
|
|
|
|
|
def _fresh_cohort_candidates(conn, authority, include_experiment_id=None):
|
|
candidates_by_query = {}
|
|
coverage = []
|
|
breadth_profile = (
|
|
authority['selector_version'] == DOCKER_RANK1_BREADTH_SELECTOR_VERSION
|
|
)
|
|
candidate_limit = 3000 if breadth_profile else authority['repositories_per_query']
|
|
if breadth_profile:
|
|
policy_filter = '''AND NOT EXISTS (
|
|
SELECT 1
|
|
FROM target_queue_policy_events event
|
|
LEFT JOIN docker_depth_experiments prior
|
|
ON prior.id = event.experiment_id
|
|
WHERE event.queue_id = queue.id
|
|
AND (
|
|
prior.id IS NULL OR prior.state <> 'released'
|
|
OR event.source <> queue.source
|
|
OR event.platform <> queue.platform
|
|
OR event.query <> queue.query
|
|
OR event.config_sha256 <> prior.config_sha256
|
|
OR event.policy_sha256 <> prior.provenance_policy_sha256
|
|
OR (event.action = 'cold' AND NOT EXISTS (
|
|
SELECT 1 FROM target_queue_policy_events reverse_event
|
|
WHERE reverse_event.reverses_event_id = event.id
|
|
AND reverse_event.action = 'reactivate'
|
|
AND reverse_event.experiment_id = event.experiment_id
|
|
AND reverse_event.queue_id = event.queue_id
|
|
))
|
|
OR (event.action = 'reactivate' AND NOT EXISTS (
|
|
SELECT 1 FROM target_queue_policy_events cold_event
|
|
WHERE cold_event.id = event.reverses_event_id
|
|
AND cold_event.action = 'cold'
|
|
AND cold_event.experiment_id = event.experiment_id
|
|
AND cold_event.queue_id = event.queue_id
|
|
))
|
|
)
|
|
)'''
|
|
else:
|
|
policy_filter = '''AND NOT EXISTS (
|
|
SELECT 1 FROM target_queue_policy_events event
|
|
WHERE event.queue_id = queue.id
|
|
)'''
|
|
if include_experiment_id is None:
|
|
member_filter = '''AND NOT EXISTS (
|
|
SELECT 1 FROM docker_depth_experiment_repositories member
|
|
WHERE member.repository_queue_id = queue.id
|
|
)'''
|
|
member_params = ()
|
|
else:
|
|
member_filter = '''AND NOT EXISTS (
|
|
SELECT 1 FROM docker_depth_experiment_repositories member
|
|
WHERE member.repository_queue_id = queue.id
|
|
AND member.experiment_id <> ?
|
|
)'''
|
|
member_params = (int(include_experiment_id),)
|
|
for query_ordinal, query in enumerate(authority['queries']):
|
|
candidates = conn.execute(
|
|
f'''SELECT provenance.repository_queue_id,
|
|
MIN(observation.search_rank) AS best_search_rank,
|
|
MIN(observation.page_id) AS eligibility_page_id,
|
|
COALESCE((
|
|
SELECT COUNT(DISTINCT manifest.graph_sha256)
|
|
FROM docker_image_manifests manifest
|
|
WHERE manifest.source = queue.source
|
|
AND manifest.repository = queue.normalized_target
|
|
), 0) AS valid_distinct_graph_count
|
|
FROM docker_repository_query_provenance provenance
|
|
JOIN docker_repository_query_observations observation
|
|
ON observation.source = provenance.source
|
|
AND observation.query = provenance.query
|
|
AND observation.repository_queue_id = provenance.repository_queue_id
|
|
JOIN docker_discovery_pages page ON page.id = observation.page_id
|
|
JOIN docker_discovery_passes discovery_pass ON discovery_pass.id = page.pass_id
|
|
JOIN target_queue queue ON queue.id = provenance.repository_queue_id
|
|
WHERE provenance.source = ? AND provenance.query = ?
|
|
AND provenance.provenance_kind = 'fresh_page'
|
|
AND provenance.fresh_coverage_eligible = 1
|
|
AND provenance.fresh_complete_observation_count > 0
|
|
AND page.query_ordinal = ? AND page.query = ?
|
|
AND discovery_pass.source = ? AND discovery_pass.pass_kind = 'deep'
|
|
AND discovery_pass.collection_generation = ?
|
|
AND discovery_pass.policy_sha256 = ?
|
|
AND discovery_pass.ordered_queries_sha256 = ?
|
|
AND discovery_pass.expected_query_count = ?
|
|
AND discovery_pass.state = 'complete'
|
|
AND queue.source = ? AND queue.platform = 'docker'
|
|
AND queue.status IN ('pending','deferred')
|
|
AND queue.target_scan_id IS NULL
|
|
AND queue.target NOT LIKE '%@%' AND queue.normalized_target NOT LIKE '%@%'
|
|
AND queue.lease_owner IS NULL AND queue.lease_token IS NULL
|
|
AND queue.claim_batch IS NULL AND queue.leased_at IS NULL
|
|
AND queue.lease_expires_at IS NULL
|
|
AND queue.current_result_reservation_id IS NULL
|
|
AND queue.claim_event_id IS NULL AND queue.resolver_token IS NULL
|
|
AND COALESCE(queue.resolver_state, '') <> 'resolving'
|
|
AND NOT EXISTS (SELECT 1 FROM target_scans scan WHERE scan.queue_id = queue.id)
|
|
{policy_filter}
|
|
AND NOT EXISTS (
|
|
SELECT 1 FROM result_reservations reservation
|
|
WHERE reservation.queue_id = queue.id
|
|
AND reservation.state IN ('scanning','ready','ingesting','db_committed')
|
|
)
|
|
AND NOT EXISTS (
|
|
SELECT 1 FROM result_reservations reservation
|
|
JOIN pipeline_quarantine quarantine ON quarantine.reservation_id = reservation.id
|
|
WHERE reservation.queue_id = queue.id
|
|
)
|
|
AND NOT EXISTS (
|
|
SELECT 1 FROM docker_image_manifests manifest
|
|
WHERE manifest.target_queue_id = queue.id
|
|
)
|
|
AND NOT EXISTS (
|
|
SELECT 1 FROM docker_image_blob_coverage coverage
|
|
JOIN docker_content_blobs blob
|
|
ON blob.digest = coverage.blob_digest
|
|
AND blob.coverage_policy_sha256 = coverage.coverage_policy_sha256
|
|
WHERE coverage.queue_id = queue.id
|
|
AND blob.state IN ('leased','submitted')
|
|
)
|
|
{member_filter}
|
|
GROUP BY provenance.repository_queue_id, queue.source, queue.normalized_target
|
|
ORDER BY MIN(observation.search_rank), provenance.repository_queue_id
|
|
LIMIT ?''',
|
|
(
|
|
authority['source'], query, query_ordinal, query,
|
|
authority['source'], authority['collection_generation'],
|
|
authority['provenance_policy_sha256'],
|
|
authority['ordered_queries_sha256'], authority['query_count'],
|
|
authority['source'], *member_params,
|
|
candidate_limit,
|
|
),
|
|
).fetchall()
|
|
normalized = [dict(candidate) for candidate in candidates]
|
|
candidates_by_query[query] = normalized
|
|
coverage.append({
|
|
'query_ordinal': query_ordinal,
|
|
'query_sha256': _canonical_sha256(query),
|
|
'eligible_slots': len(normalized),
|
|
'required_slots': authority['repositories_per_query'],
|
|
'eligible_slots_sha256': _canonical_sha256(normalized),
|
|
})
|
|
return candidates_by_query, coverage
|
|
|
|
|
|
def _require_complete_docker_depth_collection_pass(conn, authority):
|
|
row = conn.execute(
|
|
'''SELECT id FROM docker_discovery_passes
|
|
WHERE source = ? AND pass_kind = 'deep'
|
|
AND collection_generation = ? AND policy_sha256 = ?
|
|
AND ordered_queries_sha256 = ? AND expected_query_count = ?
|
|
AND completed_query_count = expected_query_count
|
|
AND state = 'complete' AND completed_at IS NOT NULL
|
|
ORDER BY id LIMIT 1''',
|
|
(
|
|
authority['source'], authority['collection_generation'],
|
|
authority['provenance_policy_sha256'],
|
|
authority['ordered_queries_sha256'], authority['query_count'],
|
|
),
|
|
).fetchone()
|
|
if not row:
|
|
raise RuntimeError(
|
|
'Docker depth cohort requires a complete fresh deep discovery pass'
|
|
)
|
|
return int(row['id'])
|
|
|
|
|
|
def summarize_docker_depth_fresh_coverage(
|
|
db, experiment, provenance_policy_sha256,
|
|
):
|
|
"""Return a query-name-free summary of the current 61-by-10 eligibility gate."""
|
|
authority = _docker_depth_authority(experiment, provenance_policy_sha256)
|
|
conn = _require_postgres_experiment_db(db, 'Docker depth coverage status')
|
|
candidates_by_query, coverage = _fresh_cohort_candidates(conn, authority)
|
|
complete = sum(
|
|
len(candidates_by_query[query]) >= authority['repositories_per_query']
|
|
for query in authority['queries']
|
|
)
|
|
eligible_slot_histogram = {
|
|
str(slot_count): sum(
|
|
len(candidates_by_query[query]) == slot_count
|
|
for query in authority['queries']
|
|
)
|
|
for slot_count in range(authority['repositories_per_query'] + 1)
|
|
}
|
|
return {
|
|
'query_count': authority['query_count'],
|
|
'complete_query_count': complete,
|
|
'incomplete_query_count': authority['query_count'] - complete,
|
|
'eligible_slot_histogram': eligible_slot_histogram,
|
|
'required_repository_count': (
|
|
authority['query_count'] * authority['repositories_per_query']
|
|
),
|
|
'eligible_repository_slots': sum(
|
|
len(candidates_by_query[query]) for query in authority['queries']
|
|
),
|
|
'coverage_sha256': _canonical_sha256({
|
|
'collection_generation': authority['collection_generation'],
|
|
'config_sha256': authority['config_sha256'],
|
|
'ordered_queries_sha256': authority['ordered_queries_sha256'],
|
|
'provenance_policy_sha256': authority['provenance_policy_sha256'],
|
|
'queries': coverage,
|
|
'schema': 1,
|
|
'selector_sha256': authority['selector_sha256'],
|
|
'type': 'truf-docker-depth-fresh-coverage-v1',
|
|
}),
|
|
}
|
|
|
|
|
|
def _cohort_partial_count(conn, experiment_id):
|
|
return int(conn.execute(
|
|
'''SELECT
|
|
(SELECT COUNT(*) FROM docker_depth_experiment_queries
|
|
WHERE experiment_id = ?) +
|
|
(SELECT COUNT(*) FROM docker_depth_experiment_repositories
|
|
WHERE experiment_id = ?) AS count''',
|
|
(experiment_id, experiment_id),
|
|
).fetchone()['count'])
|
|
|
|
|
|
def _experiment_fence_active(row):
|
|
return any(
|
|
row[name] is not None
|
|
for name in ('fence_owner', 'fence_token', 'fence_expires_at')
|
|
)
|
|
|
|
|
|
def generate_docker_depth_cohort_manifest(
|
|
db, experiment, provenance_policy_sha256,
|
|
):
|
|
"""Build a deterministic cohort manifest without changing database state."""
|
|
authority = _docker_depth_authority(experiment, provenance_policy_sha256)
|
|
conn = _require_postgres_experiment_db(db, 'Docker depth cohort review')
|
|
try:
|
|
row = conn.execute(
|
|
'SELECT * FROM docker_depth_experiments WHERE experiment_key = ?',
|
|
(authority['experiment_key'],),
|
|
).fetchone()
|
|
if row:
|
|
if not _experiment_identity_matches(row, authority):
|
|
raise RuntimeError('Docker depth experiment authority hash drifted')
|
|
if _experiment_fence_active(row):
|
|
raise RuntimeError('Docker depth experiment has an active authority fence')
|
|
if row['plan_sha256'] is not None:
|
|
raise RuntimeError('Docker depth cohort is already persisted')
|
|
if str(row['state']) != 'collecting' or _cohort_partial_count(conn, row['id']):
|
|
raise RuntimeError('Docker depth collecting experiment is not empty')
|
|
_require_complete_docker_depth_collection_pass(conn, authority)
|
|
candidates_by_query, _coverage = _fresh_cohort_candidates(conn, authority)
|
|
manifest = build_docker_depth_cohort_manifest(
|
|
experiment, provenance_policy_sha256, candidates_by_query,
|
|
)
|
|
_require_planned_cohort_current(conn, manifest['plan'], False)
|
|
manifest, manifest_sha256 = validate_docker_depth_cohort_manifest(
|
|
manifest, experiment, provenance_policy_sha256,
|
|
)
|
|
conn.commit()
|
|
return manifest, manifest_sha256
|
|
except Exception:
|
|
conn.rollback()
|
|
raise
|
|
|
|
|
|
def _lock_cohort_manifest_rows(conn, manifest):
|
|
queue_ids = sorted({
|
|
repository['repository_queue_id']
|
|
for query_plan in manifest['plan']['queries']
|
|
for repository in query_plan['repositories']
|
|
})
|
|
if not queue_ids:
|
|
return
|
|
placeholders = ','.join('?' for _ in queue_ids)
|
|
rows = conn.execute(
|
|
f'''SELECT id FROM target_queue WHERE id IN ({placeholders})
|
|
ORDER BY id FOR UPDATE''',
|
|
tuple(queue_ids),
|
|
).fetchall()
|
|
if [int(row['id']) for row in rows] != queue_ids:
|
|
raise RuntimeError('Docker depth reviewed cohort contains a missing repository')
|
|
|
|
|
|
def _policy_event_audit_sha256(action, manifest_sha256, entry, experiment_id):
|
|
return _canonical_sha256({
|
|
'action': str(action),
|
|
'manifest_sha256': str(manifest_sha256),
|
|
'entry': dict(entry),
|
|
'experiment_id': int(experiment_id),
|
|
})
|
|
|
|
|
|
def _require_released_policy_history(conn, queue_ids):
|
|
grouped = {int(queue_id): [] for queue_id in queue_ids}
|
|
values = sorted(grouped)
|
|
for offset in range(0, len(values), 500):
|
|
chunk = values[offset:offset + 500]
|
|
placeholders = ','.join('?' for _ in chunk)
|
|
rows = conn.execute(
|
|
f'''SELECT event.*, queue.status AS queue_status,
|
|
queue.source AS queue_source,
|
|
queue.platform AS queue_platform,
|
|
queue.query AS queue_query,
|
|
prior.state AS experiment_state,
|
|
prior.config_sha256 AS experiment_config_sha256,
|
|
prior.provenance_policy_sha256 AS experiment_policy_sha256,
|
|
prior.hold_manifest_sha256 AS experiment_hold_manifest_sha256
|
|
FROM target_queue_policy_events event
|
|
JOIN target_queue queue ON queue.id = event.queue_id
|
|
LEFT JOIN docker_depth_experiments prior
|
|
ON prior.id = event.experiment_id
|
|
WHERE event.queue_id IN ({placeholders})
|
|
ORDER BY event.queue_id, event.id''',
|
|
tuple(chunk),
|
|
).fetchall()
|
|
for row in rows:
|
|
grouped[int(row['queue_id'])].append(row)
|
|
|
|
for queue_id, events in grouped.items():
|
|
if not events:
|
|
continue
|
|
by_id = {int(event['id']): event for event in events}
|
|
cold_events = [event for event in events if event['action'] == 'cold']
|
|
reverse_events = [event for event in events if event['action'] == 'reactivate']
|
|
if len(cold_events) != len(reverse_events):
|
|
raise RuntimeError('Docker breadth candidate policy history is incomplete')
|
|
used_reverse_ids = set()
|
|
for cold in cold_events:
|
|
reverses = [
|
|
event for event in reverse_events
|
|
if int(event['reverses_event_id'] or 0) == int(cold['id'])
|
|
]
|
|
if len(reverses) != 1:
|
|
raise RuntimeError('Docker breadth candidate policy history is ambiguous')
|
|
reverse = reverses[0]
|
|
used_reverse_ids.add(int(reverse['id']))
|
|
experiment_id = int(cold['experiment_id'] or 0)
|
|
shared_conflict = any(
|
|
event['experiment_id'] is None
|
|
or int(event['experiment_id']) != experiment_id
|
|
or int(event['queue_id']) != queue_id
|
|
or str(event['source']) != str(event['queue_source'])
|
|
or str(event['platform']) != str(event['queue_platform'])
|
|
or str(event['query']) != str(event['queue_query'])
|
|
or str(event['experiment_state']) != 'released'
|
|
or str(event['config_sha256'])
|
|
!= str(event['experiment_config_sha256'])
|
|
or str(event['policy_sha256'])
|
|
!= str(event['experiment_policy_sha256'])
|
|
for event in (cold, reverse)
|
|
)
|
|
if (
|
|
shared_conflict
|
|
or cold['reverses_event_id'] is not None
|
|
or str(cold['reason_code']) not in (
|
|
DOCKER_DEPTH_HOLD_REASON, DOCKER_DEPTH_DYNAMIC_HOLD_REASON,
|
|
)
|
|
or str(cold['manifest_sha256'])
|
|
!= str(cold['experiment_hold_manifest_sha256'])
|
|
or str(cold['prior_status']) not in ('pending', 'deferred')
|
|
or str(cold['next_status']) != 'cold'
|
|
or str(reverse['reason_code']) != DOCKER_DEPTH_RELEASE_REASON
|
|
or str(reverse['prior_status']) != 'cold'
|
|
or str(reverse['next_status']) != str(cold['prior_status'])
|
|
or str(reverse['queue_status']) != str(reverse['next_status'])
|
|
):
|
|
raise RuntimeError('Docker breadth candidate policy history conflicts')
|
|
cold_entry = {
|
|
'queue_id': queue_id,
|
|
'source': str(cold['source']),
|
|
'platform': str(cold['platform']),
|
|
'query': str(cold['query']),
|
|
'prior_status': str(cold['prior_status']),
|
|
'prior_updated_at': str(cold['prior_updated_at']),
|
|
}
|
|
reverse_entry = {
|
|
'queue_id': queue_id,
|
|
'source': str(reverse['source']),
|
|
'platform': str(reverse['platform']),
|
|
'query': str(reverse['query']),
|
|
'cold_event_id': int(cold['id']),
|
|
'restore_status': str(reverse['next_status']),
|
|
'prior_updated_at': str(reverse['prior_updated_at']),
|
|
}
|
|
if (
|
|
str(cold['review_audit_sha256'])
|
|
!= _policy_event_audit_sha256(
|
|
'cold', cold['manifest_sha256'], cold_entry, experiment_id,
|
|
)
|
|
or str(reverse['review_audit_sha256'])
|
|
!= _policy_event_audit_sha256(
|
|
'reactivate', reverse['manifest_sha256'], reverse_entry,
|
|
experiment_id,
|
|
)
|
|
):
|
|
raise RuntimeError('Docker breadth candidate policy audit conflicts')
|
|
if len(used_reverse_ids) != len(reverse_events) or any(
|
|
int(event['reverses_event_id'] or 0) not in by_id
|
|
for event in reverse_events
|
|
):
|
|
raise RuntimeError('Docker breadth candidate policy reversal conflicts')
|
|
|
|
|
|
def _require_planned_cohort_current(conn, manifest, lock_rows):
|
|
queue_ids = sorted({
|
|
repository['repository_queue_id']
|
|
for query_plan in manifest['queries']
|
|
for repository in query_plan['repositories']
|
|
})
|
|
if not queue_ids:
|
|
return
|
|
placeholders = ','.join('?' for _ in queue_ids)
|
|
lock_suffix = ' FOR UPDATE OF queue' if lock_rows else ''
|
|
rows = conn.execute(
|
|
f'''SELECT queue.id
|
|
FROM target_queue queue
|
|
WHERE queue.id IN ({placeholders})
|
|
AND queue.source = 'dockerhub' AND queue.platform = 'docker'
|
|
AND queue.status IN ('pending','deferred')
|
|
AND queue.target_scan_id IS NULL
|
|
AND queue.target NOT LIKE '%@%' AND queue.normalized_target NOT LIKE '%@%'
|
|
AND queue.lease_owner IS NULL AND queue.lease_token IS NULL
|
|
AND queue.claim_batch IS NULL AND queue.leased_at IS NULL
|
|
AND queue.lease_expires_at IS NULL
|
|
AND queue.current_result_reservation_id IS NULL
|
|
AND queue.claim_event_id IS NULL AND queue.resolver_token IS NULL
|
|
AND COALESCE(queue.resolver_state, '') <> 'resolving'
|
|
AND NOT EXISTS (SELECT 1 FROM target_scans scan WHERE scan.queue_id = queue.id)
|
|
AND NOT EXISTS (
|
|
SELECT 1 FROM result_reservations reservation
|
|
WHERE reservation.queue_id = queue.id
|
|
AND reservation.state IN ('scanning','ready','ingesting','db_committed')
|
|
)
|
|
AND NOT EXISTS (
|
|
SELECT 1 FROM result_reservations reservation
|
|
JOIN pipeline_quarantine quarantine
|
|
ON quarantine.reservation_id = reservation.id
|
|
WHERE reservation.queue_id = queue.id
|
|
)
|
|
AND NOT EXISTS (
|
|
SELECT 1 FROM docker_image_manifests image_manifest
|
|
WHERE image_manifest.target_queue_id = queue.id
|
|
)
|
|
AND NOT EXISTS (
|
|
SELECT 1 FROM docker_image_blob_coverage coverage
|
|
JOIN docker_content_blobs blob
|
|
ON blob.digest = coverage.blob_digest
|
|
AND blob.coverage_policy_sha256 = coverage.coverage_policy_sha256
|
|
WHERE coverage.queue_id = queue.id
|
|
AND blob.state IN ('leased','submitted')
|
|
)
|
|
ORDER BY queue.id{lock_suffix}''',
|
|
tuple(queue_ids),
|
|
).fetchall()
|
|
if [int(row['id']) for row in rows] != queue_ids:
|
|
raise RuntimeError('Docker depth cohort selection drifted after review')
|
|
_require_released_policy_history(conn, queue_ids)
|
|
|
|
|
|
def apply_docker_depth_cohort_manifest(
|
|
db, experiment, provenance_policy_sha256, manifest, manifest_sha256,
|
|
):
|
|
"""Atomically persist only the exact currently eligible reviewed cohort."""
|
|
authority = _docker_depth_authority(experiment, provenance_policy_sha256)
|
|
manifest, expected_sha256 = validate_docker_depth_cohort_manifest(
|
|
manifest, experiment, provenance_policy_sha256,
|
|
)
|
|
if str(manifest_sha256 or '') != expected_sha256:
|
|
raise ValueError('Docker depth cohort manifest hash conflicts')
|
|
plan = manifest['plan']
|
|
plan_sha256 = manifest['plan_sha256']
|
|
conn = _require_postgres_experiment_db(db, 'Docker depth cohort application')
|
|
now = datetime.now(timezone.utc).isoformat(timespec='seconds')
|
|
try:
|
|
_require_complete_docker_depth_collection_pass(conn, authority)
|
|
conn.execute(
|
|
'''INSERT INTO docker_depth_experiments(
|
|
experiment_key, source, state, collection_generation,
|
|
config_sha256, ordered_queries_sha256, selector_version,
|
|
selector_sha256, provenance_policy_sha256, query_count,
|
|
repositories_per_query, images_per_repository, target_limit,
|
|
created_at, updated_at
|
|
) VALUES (?, ?, 'collecting', ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
|
ON CONFLICT(experiment_key) DO NOTHING''',
|
|
(
|
|
authority['experiment_key'], authority['source'],
|
|
authority['collection_generation'], authority['config_sha256'],
|
|
authority['ordered_queries_sha256'], authority['selector_version'],
|
|
authority['selector_sha256'], authority['provenance_policy_sha256'],
|
|
authority['query_count'], authority['repositories_per_query'],
|
|
authority['images_per_repository'], authority['target_limit'], now, now,
|
|
),
|
|
)
|
|
row = conn.execute(
|
|
'''SELECT * FROM docker_depth_experiments
|
|
WHERE experiment_key = ? FOR UPDATE''',
|
|
(authority['experiment_key'],),
|
|
).fetchone()
|
|
if not row or not _experiment_identity_matches(row, authority):
|
|
raise RuntimeError('Docker depth experiment authority is absent or drifted')
|
|
if _experiment_fence_active(row):
|
|
raise RuntimeError('Docker depth experiment has an active authority fence')
|
|
if row['plan_sha256'] is not None:
|
|
stored = _stored_cohort_plan(conn, row, authority)
|
|
stored_sha256 = canonical_docker_depth_plan_hash(stored)
|
|
if (
|
|
stored != plan
|
|
or stored_sha256 != plan_sha256
|
|
or str(row['plan_sha256']) != plan_sha256
|
|
or str(row['state']) not in ('planned', 'holding')
|
|
):
|
|
raise RuntimeError('Docker depth reviewed cohort authority conflicts')
|
|
_require_planned_cohort_current(conn, stored, True)
|
|
candidates_by_query, _coverage = _fresh_cohort_candidates(
|
|
conn, authority, include_experiment_id=int(row['id']),
|
|
)
|
|
current = build_docker_depth_cohort_manifest(
|
|
experiment, provenance_policy_sha256, candidates_by_query,
|
|
)
|
|
if current != manifest or _canonical_sha256(current) != expected_sha256:
|
|
raise RuntimeError('Docker depth cohort provenance drifted after review')
|
|
conn.commit()
|
|
return {
|
|
'experiment_id': int(row['id']),
|
|
'state': str(row['state']),
|
|
'planning_allowed': True,
|
|
'planned': False,
|
|
'plan': stored,
|
|
'plan_sha256': stored_sha256,
|
|
}
|
|
if str(row['state']) != 'collecting':
|
|
raise RuntimeError('Docker depth experiment is not in collecting state')
|
|
if _cohort_partial_count(conn, row['id']):
|
|
raise RuntimeError('Docker depth collecting experiment has a partial cohort')
|
|
_lock_cohort_manifest_rows(conn, manifest)
|
|
_require_planned_cohort_current(conn, manifest['plan'], True)
|
|
candidates_by_query, _coverage = _fresh_cohort_candidates(conn, authority)
|
|
current = build_docker_depth_cohort_manifest(
|
|
experiment, provenance_policy_sha256, candidates_by_query,
|
|
)
|
|
current, current_sha256 = validate_docker_depth_cohort_manifest(
|
|
current, experiment, provenance_policy_sha256,
|
|
)
|
|
if current != manifest or current_sha256 != expected_sha256:
|
|
raise RuntimeError('Docker depth cohort selection drifted after review')
|
|
experiment_id = int(row['id'])
|
|
for query_plan in plan['queries']:
|
|
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 (?, ?, ?, ?, ?, ?, ?, ?)''',
|
|
(
|
|
experiment_id, authority['source'], query_plan['query_ordinal'],
|
|
query_plan['query'], query_plan['query_sha256'],
|
|
authority['repositories_per_query'],
|
|
query_plan['selected_repository_count'], now,
|
|
),
|
|
)
|
|
for repository in query_plan['repositories']:
|
|
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 (?, ?, ?, ?, ?, ?, ?, ?, ?, 'pending', ?, ?)''',
|
|
(
|
|
experiment_id, query_plan['query_ordinal'], authority['source'],
|
|
query_plan['query'], repository['repository_queue_id'],
|
|
repository['eligibility_page_id'], repository['repository_rank'],
|
|
int(repository['is_deep_probe']), int(repository['is_deep_probe']),
|
|
now, now,
|
|
),
|
|
)
|
|
cursor = conn.execute(
|
|
'''UPDATE docker_depth_experiments
|
|
SET state = 'planned', plan_sha256 = ?, planned_at = ?, updated_at = ?
|
|
WHERE id = ? AND state = 'collecting' AND plan_sha256 IS NULL
|
|
AND target_count = 0 AND selection_count = 0''',
|
|
(plan_sha256, now, now, experiment_id),
|
|
)
|
|
if int(cursor.rowcount or 0) != 1:
|
|
raise RuntimeError('Docker depth cohort plan lost its authority fence')
|
|
conn.commit()
|
|
return {
|
|
'experiment_id': experiment_id,
|
|
'state': 'planned',
|
|
'planning_allowed': True,
|
|
'planned': True,
|
|
'plan': plan,
|
|
'plan_sha256': plan_sha256,
|
|
}
|
|
except Exception:
|
|
conn.rollback()
|
|
raise
|
|
|
|
|
|
def plan_docker_depth_experiment(
|
|
db, experiment, provenance_policy_sha256, *, manifest=None,
|
|
manifest_sha256=None,
|
|
):
|
|
"""Compatibility name that cannot bypass reviewed cohort approval."""
|
|
if manifest is None or manifest_sha256 is None:
|
|
raise RuntimeError('Docker depth planning requires a reviewed cohort manifest')
|
|
return apply_docker_depth_cohort_manifest(
|
|
db, experiment, provenance_policy_sha256, manifest, manifest_sha256,
|
|
)
|
|
|
|
|
|
def _experiment_row(conn, authority, lock_rows):
|
|
lock_suffix = ' FOR UPDATE' if lock_rows else ''
|
|
row = conn.execute(
|
|
f'''SELECT * FROM docker_depth_experiments
|
|
WHERE experiment_key = ?{lock_suffix}''',
|
|
(authority['experiment_key'],),
|
|
).fetchone()
|
|
if not row or not _experiment_identity_matches(row, authority):
|
|
raise RuntimeError('Docker depth experiment authority is absent or drifted')
|
|
return row
|
|
|
|
|
|
def _locked_experiment_row(conn, authority):
|
|
return _experiment_row(conn, authority, True)
|
|
|
|
|
|
def _persisted_experiment_drift_reason(db, row, authority, now):
|
|
checker = getattr(db, '_docker_depth_persisted_drift_reason_locked', None)
|
|
if not callable(checker):
|
|
raise RuntimeError('Docker depth persisted authority validator is unavailable')
|
|
return checker(row, authority, now)
|
|
|
|
|
|
def _require_persisted_experiment_authority(db, conn, row, authority, now):
|
|
holder = getattr(db, '_hold_docker_depth_experiment_locked', None)
|
|
if not callable(holder):
|
|
raise RuntimeError('Docker depth persisted authority validator is unavailable')
|
|
reason = _persisted_experiment_drift_reason(db, row, authority, now)
|
|
if reason:
|
|
holder(row, reason, now)
|
|
conn.commit()
|
|
raise RuntimeError(f'Docker depth persisted authority drifted: {reason}')
|
|
|
|
|
|
def _noncohort_hold_snapshot(conn, experiment_id, source, maximum, lock_rows):
|
|
lock_suffix = ' FOR UPDATE OF queue' if lock_rows else ''
|
|
rows = conn.execute(
|
|
f'''SELECT queue.id, queue.source, queue.platform, queue.query,
|
|
queue.status, queue.updated_at, queue.target_scan_id,
|
|
queue.resolver_state, queue.lease_owner, queue.lease_token,
|
|
queue.claim_batch, queue.leased_at, queue.lease_expires_at,
|
|
queue.current_result_reservation_id, queue.claim_event_id,
|
|
queue.resolver_token,
|
|
EXISTS (
|
|
SELECT 1 FROM result_reservations reservation
|
|
WHERE reservation.queue_id = queue.id
|
|
AND reservation.state IN ('scanning','ready','ingesting','db_committed')
|
|
) AS active_reservation,
|
|
EXISTS (
|
|
SELECT 1 FROM docker_image_blob_coverage coverage
|
|
JOIN docker_content_blobs blob
|
|
ON blob.digest = coverage.blob_digest
|
|
AND blob.coverage_policy_sha256 = coverage.coverage_policy_sha256
|
|
WHERE coverage.queue_id = queue.id
|
|
AND blob.state IN ('leased','submitted')
|
|
) AS active_docker_blob,
|
|
EXISTS (
|
|
SELECT 1 FROM result_reservations reservation
|
|
JOIN pipeline_quarantine quarantine
|
|
ON quarantine.reservation_id = reservation.id
|
|
WHERE reservation.queue_id = queue.id
|
|
) AS quarantined_reservation,
|
|
EXISTS (
|
|
SELECT 1 FROM target_scans scan WHERE scan.queue_id = queue.id
|
|
) AS prior_scan,
|
|
EXISTS (
|
|
SELECT 1 FROM target_queue_policy_events event
|
|
WHERE event.queue_id = queue.id
|
|
) AS prior_policy_event
|
|
FROM target_queue queue
|
|
WHERE queue.source = ? AND queue.platform = 'docker'
|
|
AND queue.target NOT LIKE '%@%' AND queue.normalized_target NOT LIKE '%@%'
|
|
AND queue.status NOT IN ('done','failed')
|
|
AND NOT EXISTS (
|
|
SELECT 1 FROM docker_depth_experiment_repositories member
|
|
WHERE member.experiment_id = ?
|
|
AND member.repository_queue_id = queue.id
|
|
)
|
|
ORDER BY queue.id LIMIT ?{lock_suffix}''',
|
|
(source, experiment_id, maximum + 1),
|
|
).fetchall()
|
|
if len(rows) > maximum:
|
|
raise RuntimeError(
|
|
f'Docker depth hold selection exceeds its reviewed {maximum}-row bound'
|
|
)
|
|
entries = []
|
|
conflicts = []
|
|
fence_fields = (
|
|
'lease_owner', 'lease_token', 'claim_batch', 'leased_at', 'lease_expires_at',
|
|
'current_result_reservation_id', 'claim_event_id', 'resolver_token',
|
|
)
|
|
for row in rows:
|
|
status = str(row['status'] or '')
|
|
query = str(row['query'] or '')
|
|
reason = None
|
|
if status == 'cold':
|
|
reason = 'independently_cold'
|
|
elif status in ('done', 'failed'):
|
|
reason = 'terminal_status'
|
|
elif status == 'quarantined' or bool(row['quarantined_reservation']):
|
|
reason = 'quarantined'
|
|
elif row['target_scan_id'] is not None or bool(row['prior_scan']):
|
|
reason = 'previously_scanned'
|
|
elif status not in ('pending', 'deferred'):
|
|
reason = 'ineligible_status'
|
|
elif any(row[field] is not None for field in fence_fields):
|
|
reason = 'active_claim_fence'
|
|
elif str(row['resolver_state'] or '') == 'resolving':
|
|
reason = 'active_resolver_fence'
|
|
elif bool(row['active_reservation']) or bool(row['active_docker_blob']):
|
|
reason = 'active_content_fence'
|
|
elif bool(row['prior_policy_event']):
|
|
reason = 'unrelated_policy_event'
|
|
elif not query:
|
|
reason = 'missing_query_identity'
|
|
if reason is None:
|
|
entries.append({
|
|
'queue_id': int(row['id']),
|
|
'source': str(row['source']),
|
|
'platform': str(row['platform']),
|
|
'query': query,
|
|
'prior_status': status,
|
|
'prior_updated_at': str(row['updated_at']),
|
|
})
|
|
else:
|
|
conflicts.append({
|
|
'queue_id': int(row['id']),
|
|
'source': str(row['source']),
|
|
'platform': str(row['platform']),
|
|
'query': query,
|
|
'status': status,
|
|
'prior_updated_at': str(row['updated_at']),
|
|
'reason': reason,
|
|
})
|
|
return entries, conflicts
|
|
|
|
|
|
def _hold_manifest_document(authority, experiment_row, entries, conflicts):
|
|
return {
|
|
'schema': 1,
|
|
'type': DOCKER_DEPTH_HOLD_MANIFEST_TYPE,
|
|
'experiment_id': int(experiment_row['id']),
|
|
'experiment_key': authority['experiment_key'],
|
|
'source': authority['source'],
|
|
'config_sha256': authority['config_sha256'],
|
|
'ordered_queries_sha256': authority['ordered_queries_sha256'],
|
|
'selector_sha256': authority['selector_sha256'],
|
|
'provenance_policy_sha256': authority['provenance_policy_sha256'],
|
|
'plan_sha256': str(experiment_row['plan_sha256']),
|
|
'reason_code': DOCKER_DEPTH_HOLD_REASON,
|
|
'entry_count': len(entries),
|
|
'conflict_count': len(conflicts),
|
|
'selection_sha256': _canonical_sha256(entries),
|
|
'conflicts_sha256': _canonical_sha256(conflicts),
|
|
'entries': entries,
|
|
'conflicts': conflicts,
|
|
}
|
|
|
|
|
|
def _normalize_hold_entries(entries, maximum):
|
|
if not isinstance(entries, list) or len(entries) > maximum:
|
|
raise ValueError('Docker depth hold manifest entries exceed their bound')
|
|
required = {
|
|
'queue_id', 'source', 'platform', 'query', 'prior_status', 'prior_updated_at',
|
|
}
|
|
normalized = []
|
|
seen = set()
|
|
for raw in entries:
|
|
if not isinstance(raw, dict) or set(raw) != required:
|
|
raise ValueError('Docker depth hold manifest entry shape is invalid')
|
|
queue_id = _strict_positive_int(raw['queue_id'], 'queue_id')
|
|
if queue_id in seen:
|
|
raise ValueError('Docker depth hold manifest contains a duplicate queue row')
|
|
seen.add(queue_id)
|
|
item = {
|
|
'queue_id': queue_id,
|
|
'source': str(raw['source'] or ''),
|
|
'platform': str(raw['platform'] or ''),
|
|
'query': str(raw['query'] or ''),
|
|
'prior_status': str(raw['prior_status'] or ''),
|
|
'prior_updated_at': str(raw['prior_updated_at'] or ''),
|
|
}
|
|
if (
|
|
item['source'] != 'dockerhub'
|
|
or item['platform'] != 'docker'
|
|
or not item['query']
|
|
or item['prior_status'] not in ('pending', 'deferred')
|
|
or not item['prior_updated_at']
|
|
):
|
|
raise ValueError('Docker depth hold manifest entry identity is invalid')
|
|
normalized.append(item)
|
|
normalized.sort(key=lambda item: item['queue_id'])
|
|
return normalized
|
|
|
|
|
|
def _normalize_hold_conflicts(conflicts, maximum):
|
|
if not isinstance(conflicts, list) or len(conflicts) > maximum:
|
|
raise ValueError('Docker depth hold manifest conflicts exceed their bound')
|
|
required = {
|
|
'queue_id', 'source', 'platform', 'query', 'status', 'prior_updated_at', 'reason',
|
|
}
|
|
allowed_reasons = {
|
|
'independently_cold', 'terminal_status', 'quarantined', 'previously_scanned',
|
|
'ineligible_status', 'active_claim_fence', 'active_resolver_fence',
|
|
'active_content_fence', 'unrelated_policy_event', 'missing_query_identity',
|
|
}
|
|
normalized = []
|
|
seen = set()
|
|
for raw in conflicts:
|
|
if not isinstance(raw, dict) or set(raw) != required:
|
|
raise ValueError('Docker depth hold conflict shape is invalid')
|
|
queue_id = _strict_positive_int(raw['queue_id'], 'queue_id')
|
|
if queue_id in seen:
|
|
raise ValueError('Docker depth hold conflicts contain a duplicate queue row')
|
|
seen.add(queue_id)
|
|
item = {
|
|
'queue_id': queue_id,
|
|
'source': str(raw['source'] or ''),
|
|
'platform': str(raw['platform'] or ''),
|
|
'query': str(raw['query'] or ''),
|
|
'status': str(raw['status'] or ''),
|
|
'prior_updated_at': str(raw['prior_updated_at'] or ''),
|
|
'reason': str(raw['reason'] or ''),
|
|
}
|
|
if (
|
|
item['source'] != 'dockerhub'
|
|
or item['platform'] != 'docker'
|
|
or not item['prior_updated_at']
|
|
or item['reason'] not in allowed_reasons
|
|
):
|
|
raise ValueError('Docker depth hold conflict identity is invalid')
|
|
normalized.append(item)
|
|
normalized.sort(key=lambda item: item['queue_id'])
|
|
return normalized
|
|
|
|
|
|
def validate_docker_depth_hold_manifest(
|
|
manifest, experiment=None, provenance_policy_sha256=None,
|
|
max_rows=DOCKER_DEPTH_REVIEW_MANIFEST_MAX_ROWS,
|
|
):
|
|
maximum = min(
|
|
DOCKER_DEPTH_REVIEW_MANIFEST_MAX_ROWS,
|
|
_strict_positive_int(max_rows, 'max_rows'),
|
|
)
|
|
expected_keys = {
|
|
'schema', 'type', 'experiment_id', 'experiment_key', 'source',
|
|
'config_sha256', 'ordered_queries_sha256', 'selector_sha256',
|
|
'provenance_policy_sha256', 'plan_sha256', 'reason_code', 'entry_count',
|
|
'conflict_count', 'selection_sha256', 'conflicts_sha256', 'entries', 'conflicts',
|
|
}
|
|
if (
|
|
not isinstance(manifest, dict)
|
|
or set(manifest) != expected_keys
|
|
or type(manifest.get('schema')) is not int
|
|
or manifest.get('schema') != 1
|
|
or manifest.get('type') != DOCKER_DEPTH_HOLD_MANIFEST_TYPE
|
|
):
|
|
raise ValueError('Docker depth hold manifest shape is invalid')
|
|
entries = _normalize_hold_entries(manifest['entries'], maximum)
|
|
conflicts = _normalize_hold_conflicts(manifest['conflicts'], maximum)
|
|
if entries != manifest['entries'] or conflicts != manifest['conflicts']:
|
|
raise ValueError('Docker depth hold manifest rows are not canonical')
|
|
if len(entries) + len(conflicts) > maximum:
|
|
raise ValueError('Docker depth hold manifest selection exceeds its bound')
|
|
if set(item['queue_id'] for item in entries).intersection(
|
|
item['queue_id'] for item in conflicts
|
|
):
|
|
raise ValueError('Docker depth hold manifest entry and conflict sets overlap')
|
|
hashes = (
|
|
'config_sha256', 'ordered_queries_sha256', 'selector_sha256',
|
|
'provenance_policy_sha256', 'plan_sha256', 'selection_sha256',
|
|
'conflicts_sha256',
|
|
)
|
|
if any(not _valid_sha256(manifest.get(name)) for name in hashes):
|
|
raise ValueError('Docker depth hold manifest contains an invalid hash')
|
|
entry_count = _strict_nonnegative_int(manifest['entry_count'], 'entry_count')
|
|
conflict_count = _strict_nonnegative_int(
|
|
manifest['conflict_count'], 'conflict_count',
|
|
)
|
|
_strict_identifier(manifest['experiment_key'], 'experiment_key')
|
|
if (
|
|
_strict_positive_int(manifest['experiment_id'], 'experiment_id') <= 0
|
|
or manifest['source'] != 'dockerhub'
|
|
or manifest['reason_code'] != DOCKER_DEPTH_HOLD_REASON
|
|
or entry_count != len(entries)
|
|
or conflict_count != len(conflicts)
|
|
or manifest['selection_sha256'] != _canonical_sha256(entries)
|
|
or manifest['conflicts_sha256'] != _canonical_sha256(conflicts)
|
|
):
|
|
raise ValueError('Docker depth hold manifest evidence conflicts')
|
|
if experiment is not None:
|
|
authority = _docker_depth_authority(
|
|
experiment,
|
|
provenance_policy_sha256 or manifest['provenance_policy_sha256'],
|
|
)
|
|
for manifest_name, authority_name in (
|
|
('experiment_key', 'experiment_key'),
|
|
('source', 'source'),
|
|
('config_sha256', 'config_sha256'),
|
|
('ordered_queries_sha256', 'ordered_queries_sha256'),
|
|
('selector_sha256', 'selector_sha256'),
|
|
('provenance_policy_sha256', 'provenance_policy_sha256'),
|
|
):
|
|
if manifest[manifest_name] != authority[authority_name]:
|
|
raise ValueError('Docker depth hold manifest authority drifted')
|
|
normalized = copy.deepcopy(manifest)
|
|
return normalized, _canonical_sha256(normalized)
|
|
|
|
|
|
def generate_docker_depth_hold_manifest(
|
|
db, experiment, provenance_policy_sha256,
|
|
max_rows=DOCKER_DEPTH_REVIEW_MANIFEST_MAX_ROWS,
|
|
):
|
|
authority = _docker_depth_authority(experiment, provenance_policy_sha256)
|
|
conn = _require_postgres_experiment_db(db, 'Docker depth hold planning')
|
|
maximum = min(
|
|
DOCKER_DEPTH_REVIEW_MANIFEST_MAX_ROWS,
|
|
_strict_positive_int(max_rows, 'max_rows'),
|
|
)
|
|
try:
|
|
row = _experiment_row(conn, authority, False)
|
|
if str(row['state']) != 'planned' or not _valid_sha256(row['plan_sha256']):
|
|
raise RuntimeError('Docker depth hold planning requires an immutable planned cohort')
|
|
if _experiment_fence_active(row):
|
|
raise RuntimeError('Docker depth experiment has an active authority fence')
|
|
plan = _stored_cohort_plan(conn, row, authority)
|
|
if canonical_docker_depth_plan_hash(plan) != str(row['plan_sha256']):
|
|
raise RuntimeError('Docker depth hold planning found plan hash drift')
|
|
_require_planned_cohort_current(conn, plan, False)
|
|
entries, conflicts = _noncohort_hold_snapshot(
|
|
conn, int(row['id']), authority['source'], maximum, False,
|
|
)
|
|
manifest = _hold_manifest_document(authority, row, entries, conflicts)
|
|
manifest, manifest_sha256 = validate_docker_depth_hold_manifest(
|
|
manifest, experiment, provenance_policy_sha256, maximum,
|
|
)
|
|
conn.commit()
|
|
return manifest, manifest_sha256
|
|
except Exception:
|
|
conn.rollback()
|
|
raise
|
|
|
|
|
|
def apply_docker_depth_hold_manifest(
|
|
db, experiment, provenance_policy_sha256, manifest, manifest_sha256,
|
|
max_rows=DOCKER_DEPTH_REVIEW_MANIFEST_MAX_ROWS,
|
|
):
|
|
authority = _docker_depth_authority(experiment, provenance_policy_sha256)
|
|
manifest, expected_sha256 = validate_docker_depth_hold_manifest(
|
|
manifest, experiment, provenance_policy_sha256, max_rows,
|
|
)
|
|
if str(manifest_sha256 or '') != expected_sha256:
|
|
raise ValueError('Docker depth hold manifest hash conflicts')
|
|
conn = _require_postgres_experiment_db(db, 'Docker depth hold application')
|
|
maximum = min(
|
|
DOCKER_DEPTH_REVIEW_MANIFEST_MAX_ROWS,
|
|
_strict_positive_int(max_rows, 'max_rows'),
|
|
)
|
|
now = datetime.now(timezone.utc).isoformat(timespec='seconds')
|
|
try:
|
|
row = _locked_experiment_row(conn, authority)
|
|
if _experiment_fence_active(row):
|
|
raise RuntimeError('Docker depth experiment has an active authority fence')
|
|
if int(row['id']) != manifest['experiment_id']:
|
|
raise RuntimeError('Docker depth hold manifest experiment identity drifted')
|
|
if str(row['plan_sha256'] or '') != manifest['plan_sha256']:
|
|
raise RuntimeError('Docker depth hold manifest plan identity drifted')
|
|
plan = _stored_cohort_plan(conn, row, authority)
|
|
if canonical_docker_depth_plan_hash(plan) != manifest['plan_sha256']:
|
|
raise RuntimeError('Docker depth hold manifest plan identity drifted')
|
|
_require_planned_cohort_current(conn, plan, True)
|
|
stored_manifest_sha256 = str(row['hold_manifest_sha256'] or '')
|
|
if stored_manifest_sha256:
|
|
_require_persisted_experiment_authority(
|
|
db, conn, row, authority, now,
|
|
)
|
|
if (
|
|
stored_manifest_sha256 != expected_sha256
|
|
or str(row['state']) != 'holding'
|
|
):
|
|
raise RuntimeError('Docker depth reviewed hold authority conflicts')
|
|
reviewed_rows = conn.execute(
|
|
'''SELECT queue_id, reason_code, manifest_sha256
|
|
FROM target_queue_policy_events
|
|
WHERE experiment_id = ? AND action = 'cold'
|
|
ORDER BY queue_id, id''',
|
|
(row['id'],),
|
|
).fetchall()
|
|
if (
|
|
[int(item['queue_id']) for item in reviewed_rows]
|
|
!= [item['queue_id'] for item in manifest['entries']]
|
|
or any(
|
|
str(item['reason_code']) != DOCKER_DEPTH_HOLD_REASON
|
|
or str(item['manifest_sha256']) != expected_sha256
|
|
for item in reviewed_rows
|
|
)
|
|
):
|
|
raise RuntimeError('Docker depth reviewed hold set changed')
|
|
else:
|
|
if str(row['state']) != 'planned':
|
|
raise RuntimeError('Docker depth experiment is not ready for reviewed holds')
|
|
entries, conflicts = _noncohort_hold_snapshot(
|
|
conn, int(row['id']), authority['source'], maximum, True,
|
|
)
|
|
current = _hold_manifest_document(authority, row, entries, conflicts)
|
|
if current != manifest:
|
|
raise RuntimeError('Docker depth hold selection drifted after review')
|
|
entries = db._normalize_target_queue_policy_entries(
|
|
manifest['entries'], 'cold', maximum,
|
|
) if manifest['entries'] else []
|
|
result = db._cold_target_queue_rows_locked(
|
|
entries,
|
|
reason_code=DOCKER_DEPTH_HOLD_REASON,
|
|
config_sha256=authority['config_sha256'],
|
|
policy_sha256=authority['provenance_policy_sha256'],
|
|
manifest_sha256=expected_sha256,
|
|
experiment_id=int(row['id']),
|
|
now=now,
|
|
)
|
|
if not stored_manifest_sha256:
|
|
cursor = conn.execute(
|
|
'''UPDATE docker_depth_experiments
|
|
SET state = 'holding', hold_manifest_sha256 = ?,
|
|
activated_at = COALESCE(activated_at, ?), updated_at = ?
|
|
WHERE id = ? AND state = 'planned'
|
|
AND plan_sha256 = ? AND hold_manifest_sha256 IS NULL''',
|
|
(expected_sha256, now, now, row['id'], manifest['plan_sha256']),
|
|
)
|
|
if int(cursor.rowcount or 0) != 1:
|
|
raise RuntimeError('Docker depth hold activation lost its authority fence')
|
|
conn.commit()
|
|
return {
|
|
**result,
|
|
'experiment_id': int(row['id']),
|
|
'state': 'holding' if not stored_manifest_sha256 else str(row['state']),
|
|
'conflicts': manifest['conflict_count'],
|
|
}
|
|
except Exception:
|
|
conn.rollback()
|
|
raise
|
|
|
|
|
|
def _experiment_reactivation_entries(conn, experiment_id, maximum, lock_rows):
|
|
lock_suffix = ' FOR UPDATE OF queue, cold_event' if lock_rows else ''
|
|
rows = conn.execute(
|
|
f'''SELECT queue.*, cold_event.id AS cold_event_id,
|
|
cold_event.prior_status AS restore_status,
|
|
reverse_event.id AS reverse_event_id
|
|
FROM target_queue_policy_events cold_event
|
|
JOIN target_queue queue ON queue.id = cold_event.queue_id
|
|
LEFT JOIN target_queue_policy_events reverse_event
|
|
ON reverse_event.reverses_event_id = cold_event.id
|
|
WHERE cold_event.experiment_id = ? AND cold_event.action = 'cold'
|
|
AND reverse_event.id IS NULL
|
|
ORDER BY queue.id, cold_event.id LIMIT ?{lock_suffix}''',
|
|
(experiment_id, maximum + 1),
|
|
).fetchall()
|
|
if len(rows) > maximum:
|
|
raise RuntimeError(
|
|
f'Docker depth reactivation selection exceeds its reviewed {maximum}-row bound'
|
|
)
|
|
entries = []
|
|
seen = set()
|
|
for row in rows:
|
|
queue_id = int(row['id'])
|
|
if queue_id in seen or row['reverse_event_id'] is not None or row['status'] != 'cold':
|
|
raise RuntimeError('Docker depth reactivation evidence conflicts')
|
|
seen.add(queue_id)
|
|
entries.append({
|
|
'queue_id': queue_id,
|
|
'source': str(row['source']),
|
|
'platform': str(row['platform']),
|
|
'query': str(row['query']),
|
|
'cold_event_id': int(row['cold_event_id']),
|
|
'restore_status': str(row['restore_status']),
|
|
'prior_updated_at': str(row['updated_at']),
|
|
})
|
|
return entries
|
|
|
|
|
|
def _reactivation_manifest_document(authority, experiment_row, entries):
|
|
return {
|
|
'schema': 1,
|
|
'type': DOCKER_DEPTH_REACTIVATION_MANIFEST_TYPE,
|
|
'experiment_id': int(experiment_row['id']),
|
|
'experiment_key': authority['experiment_key'],
|
|
'source': authority['source'],
|
|
'config_sha256': authority['config_sha256'],
|
|
'ordered_queries_sha256': authority['ordered_queries_sha256'],
|
|
'selector_sha256': authority['selector_sha256'],
|
|
'provenance_policy_sha256': authority['provenance_policy_sha256'],
|
|
'plan_sha256': str(experiment_row['plan_sha256']),
|
|
'hold_manifest_sha256': str(experiment_row['hold_manifest_sha256']),
|
|
'reason_code': DOCKER_DEPTH_RELEASE_REASON,
|
|
'entry_count': len(entries),
|
|
'selection_sha256': _canonical_sha256(entries),
|
|
'entries': entries,
|
|
}
|
|
|
|
|
|
def _normalize_reactivation_entries(entries, maximum):
|
|
if not isinstance(entries, list) or len(entries) > maximum:
|
|
raise ValueError('Docker depth reactivation entries exceed their bound')
|
|
required = {
|
|
'queue_id', 'source', 'platform', 'query', 'cold_event_id',
|
|
'restore_status', 'prior_updated_at',
|
|
}
|
|
normalized = []
|
|
queue_ids = set()
|
|
event_ids = set()
|
|
for raw in entries:
|
|
if not isinstance(raw, dict) or set(raw) != required:
|
|
raise ValueError('Docker depth reactivation entry shape is invalid')
|
|
queue_id = _strict_positive_int(raw['queue_id'], 'queue_id')
|
|
event_id = _strict_positive_int(raw['cold_event_id'], 'cold_event_id')
|
|
if queue_id in queue_ids or event_id in event_ids:
|
|
raise ValueError('Docker depth reactivation identity is duplicated')
|
|
queue_ids.add(queue_id)
|
|
event_ids.add(event_id)
|
|
item = {
|
|
'queue_id': queue_id,
|
|
'source': str(raw['source'] or ''),
|
|
'platform': str(raw['platform'] or ''),
|
|
'query': str(raw['query'] or ''),
|
|
'cold_event_id': event_id,
|
|
'restore_status': str(raw['restore_status'] or ''),
|
|
'prior_updated_at': str(raw['prior_updated_at'] or ''),
|
|
}
|
|
if (
|
|
item['source'] != 'dockerhub'
|
|
or item['platform'] != 'docker'
|
|
or not item['query']
|
|
or item['restore_status'] not in ('pending', 'deferred')
|
|
or not item['prior_updated_at']
|
|
):
|
|
raise ValueError('Docker depth reactivation entry identity is invalid')
|
|
normalized.append(item)
|
|
normalized.sort(key=lambda item: item['queue_id'])
|
|
return normalized
|
|
|
|
|
|
def validate_docker_depth_reactivation_manifest(
|
|
manifest, experiment=None, provenance_policy_sha256=None,
|
|
max_rows=DOCKER_DEPTH_REVIEW_MANIFEST_MAX_ROWS,
|
|
):
|
|
maximum = min(
|
|
DOCKER_DEPTH_REVIEW_MANIFEST_MAX_ROWS,
|
|
_strict_positive_int(max_rows, 'max_rows'),
|
|
)
|
|
expected_keys = {
|
|
'schema', 'type', 'experiment_id', 'experiment_key', 'source',
|
|
'config_sha256', 'ordered_queries_sha256', 'selector_sha256',
|
|
'provenance_policy_sha256', 'plan_sha256', 'hold_manifest_sha256',
|
|
'reason_code', 'entry_count', 'selection_sha256', 'entries',
|
|
}
|
|
if (
|
|
not isinstance(manifest, dict)
|
|
or set(manifest) != expected_keys
|
|
or type(manifest.get('schema')) is not int
|
|
or manifest.get('schema') != 1
|
|
or manifest.get('type') != DOCKER_DEPTH_REACTIVATION_MANIFEST_TYPE
|
|
):
|
|
raise ValueError('Docker depth reactivation manifest shape is invalid')
|
|
entries = _normalize_reactivation_entries(manifest['entries'], maximum)
|
|
if entries != manifest['entries']:
|
|
raise ValueError('Docker depth reactivation entries are not canonical')
|
|
for name in (
|
|
'config_sha256', 'ordered_queries_sha256', 'selector_sha256',
|
|
'provenance_policy_sha256', 'plan_sha256', 'hold_manifest_sha256',
|
|
'selection_sha256',
|
|
):
|
|
if not _valid_sha256(manifest.get(name)):
|
|
raise ValueError('Docker depth reactivation manifest contains an invalid hash')
|
|
entry_count = _strict_nonnegative_int(manifest['entry_count'], 'entry_count')
|
|
_strict_identifier(manifest['experiment_key'], 'experiment_key')
|
|
if (
|
|
_strict_positive_int(manifest['experiment_id'], 'experiment_id') <= 0
|
|
or manifest['source'] != 'dockerhub'
|
|
or manifest['reason_code'] != DOCKER_DEPTH_RELEASE_REASON
|
|
or entry_count != len(entries)
|
|
or manifest['selection_sha256'] != _canonical_sha256(entries)
|
|
):
|
|
raise ValueError('Docker depth reactivation evidence conflicts')
|
|
if experiment is not None:
|
|
authority = _docker_depth_authority(
|
|
experiment,
|
|
provenance_policy_sha256 or manifest['provenance_policy_sha256'],
|
|
)
|
|
for manifest_name, authority_name in (
|
|
('experiment_key', 'experiment_key'), ('source', 'source'),
|
|
('config_sha256', 'config_sha256'),
|
|
('ordered_queries_sha256', 'ordered_queries_sha256'),
|
|
('selector_sha256', 'selector_sha256'),
|
|
('provenance_policy_sha256', 'provenance_policy_sha256'),
|
|
):
|
|
if manifest[manifest_name] != authority[authority_name]:
|
|
raise ValueError('Docker depth reactivation authority drifted')
|
|
normalized = copy.deepcopy(manifest)
|
|
return normalized, _canonical_sha256(normalized)
|
|
|
|
|
|
def generate_docker_depth_reactivation_manifest(
|
|
db, experiment, provenance_policy_sha256,
|
|
max_rows=DOCKER_DEPTH_REVIEW_MANIFEST_MAX_ROWS,
|
|
):
|
|
authority = _docker_depth_authority(experiment, provenance_policy_sha256)
|
|
conn = _require_postgres_experiment_db(db, 'Docker depth reactivation planning')
|
|
maximum = min(
|
|
DOCKER_DEPTH_REVIEW_MANIFEST_MAX_ROWS,
|
|
_strict_positive_int(max_rows, 'max_rows'),
|
|
)
|
|
try:
|
|
row = _experiment_row(conn, authority, False)
|
|
if (
|
|
str(row['state']) != 'completed'
|
|
or not _valid_sha256(row['plan_sha256'])
|
|
or not _valid_sha256(row['hold_manifest_sha256'])
|
|
):
|
|
raise RuntimeError('Docker depth reactivation requires a completed experiment')
|
|
now = datetime.now(timezone.utc).isoformat(timespec='seconds')
|
|
reason = _persisted_experiment_drift_reason(db, row, authority, now)
|
|
if reason:
|
|
raise RuntimeError(f'Docker depth persisted authority drifted: {reason}')
|
|
entries = _experiment_reactivation_entries(
|
|
conn, int(row['id']), maximum, False,
|
|
)
|
|
for entry in entries:
|
|
queue = conn.execute(
|
|
'SELECT * FROM target_queue WHERE id = ?',
|
|
(entry['queue_id'],),
|
|
).fetchone()
|
|
db._require_target_queue_policy_unfenced(queue, lock_rows=False)
|
|
manifest = _reactivation_manifest_document(authority, row, entries)
|
|
manifest, manifest_sha256 = validate_docker_depth_reactivation_manifest(
|
|
manifest, experiment, provenance_policy_sha256, maximum,
|
|
)
|
|
conn.commit()
|
|
return manifest, manifest_sha256
|
|
except Exception:
|
|
conn.rollback()
|
|
raise
|
|
|
|
|
|
def _require_released_reactivation_replay(
|
|
conn, experiment_id, manifest, manifest_sha256,
|
|
):
|
|
rows = conn.execute(
|
|
'''SELECT cold_event.id AS cold_event_id, cold_event.queue_id,
|
|
cold_event.reason_code AS cold_reason_code,
|
|
cold_event.config_sha256 AS cold_config_sha256,
|
|
cold_event.policy_sha256 AS cold_policy_sha256,
|
|
cold_event.manifest_sha256 AS cold_manifest_sha256,
|
|
reverse_event.id AS reverse_event_id,
|
|
reverse_event.action AS reverse_action,
|
|
reverse_event.experiment_id AS reverse_experiment_id,
|
|
reverse_event.reason_code AS reverse_reason_code,
|
|
reverse_event.config_sha256 AS reverse_config_sha256,
|
|
reverse_event.policy_sha256 AS reverse_policy_sha256,
|
|
reverse_event.manifest_sha256 AS reverse_manifest_sha256,
|
|
queue.status AS queue_status
|
|
FROM target_queue_policy_events cold_event
|
|
JOIN target_queue queue ON queue.id = cold_event.queue_id
|
|
LEFT JOIN target_queue_policy_events reverse_event
|
|
ON reverse_event.reverses_event_id = cold_event.id
|
|
WHERE cold_event.experiment_id = ? AND cold_event.action = 'cold'
|
|
ORDER BY cold_event.queue_id, cold_event.id''',
|
|
(experiment_id,),
|
|
).fetchall()
|
|
expected = [
|
|
(entry['queue_id'], entry['cold_event_id']) for entry in manifest['entries']
|
|
]
|
|
actual = [(int(row['queue_id']), int(row['cold_event_id'])) for row in rows]
|
|
if actual != expected or any(
|
|
str(row['cold_reason_code']) not in (
|
|
DOCKER_DEPTH_HOLD_REASON, DOCKER_DEPTH_DYNAMIC_HOLD_REASON,
|
|
)
|
|
or str(row['cold_config_sha256']) != manifest['config_sha256']
|
|
or str(row['cold_policy_sha256']) != manifest['provenance_policy_sha256']
|
|
or str(row['cold_manifest_sha256']) != manifest['hold_manifest_sha256']
|
|
or row['reverse_event_id'] is None
|
|
or str(row['reverse_action']) != 'reactivate'
|
|
or row['reverse_experiment_id'] is None
|
|
or int(row['reverse_experiment_id']) != experiment_id
|
|
or str(row['reverse_reason_code']) != DOCKER_DEPTH_RELEASE_REASON
|
|
or str(row['reverse_config_sha256']) != manifest['config_sha256']
|
|
or str(row['reverse_policy_sha256']) != manifest['provenance_policy_sha256']
|
|
or str(row['reverse_manifest_sha256']) != manifest_sha256
|
|
or str(row['queue_status']) != manifest['entries'][index]['restore_status']
|
|
for index, row in enumerate(rows)
|
|
):
|
|
raise RuntimeError('Docker depth released hold set changed')
|
|
|
|
|
|
def apply_docker_depth_reactivation_manifest(
|
|
db, experiment, provenance_policy_sha256, manifest, manifest_sha256,
|
|
max_rows=DOCKER_DEPTH_REVIEW_MANIFEST_MAX_ROWS,
|
|
):
|
|
authority = _docker_depth_authority(experiment, provenance_policy_sha256)
|
|
manifest, expected_sha256 = validate_docker_depth_reactivation_manifest(
|
|
manifest, experiment, provenance_policy_sha256, max_rows,
|
|
)
|
|
if str(manifest_sha256 or '') != expected_sha256:
|
|
raise ValueError('Docker depth reactivation manifest hash conflicts')
|
|
conn = _require_postgres_experiment_db(db, 'Docker depth reactivation')
|
|
maximum = min(
|
|
DOCKER_DEPTH_REVIEW_MANIFEST_MAX_ROWS,
|
|
_strict_positive_int(max_rows, 'max_rows'),
|
|
)
|
|
now = datetime.now(timezone.utc).isoformat(timespec='seconds')
|
|
try:
|
|
row = _locked_experiment_row(conn, authority)
|
|
if (
|
|
int(row['id']) != manifest['experiment_id']
|
|
or str(row['plan_sha256'] or '') != manifest['plan_sha256']
|
|
or str(row['hold_manifest_sha256'] or '') != manifest['hold_manifest_sha256']
|
|
):
|
|
raise RuntimeError('Docker depth reactivation experiment identity drifted')
|
|
plan = _stored_cohort_plan(conn, row, authority)
|
|
if canonical_docker_depth_plan_hash(plan) != manifest['plan_sha256']:
|
|
raise RuntimeError('Docker depth reactivation cohort identity drifted')
|
|
released = str(row['state']) == 'released'
|
|
if released:
|
|
_require_released_reactivation_replay(
|
|
conn, int(row['id']), manifest, expected_sha256,
|
|
)
|
|
else:
|
|
if str(row['state']) != 'completed':
|
|
raise RuntimeError('Docker depth reactivation requires a completed experiment')
|
|
_require_persisted_experiment_authority(
|
|
db, conn, row, authority, now,
|
|
)
|
|
current_entries = _experiment_reactivation_entries(
|
|
conn, int(row['id']), maximum, True,
|
|
)
|
|
current = _reactivation_manifest_document(authority, row, current_entries)
|
|
if current != manifest:
|
|
raise RuntimeError('Docker depth reactivation selection drifted after review')
|
|
entries = db._normalize_target_queue_policy_entries(
|
|
manifest['entries'], 'reactivate', maximum,
|
|
) if manifest['entries'] else []
|
|
result = db._reactivate_target_queue_rows_locked(
|
|
entries,
|
|
reason_code=DOCKER_DEPTH_RELEASE_REASON,
|
|
config_sha256=authority['config_sha256'],
|
|
policy_sha256=authority['provenance_policy_sha256'],
|
|
manifest_sha256=expected_sha256,
|
|
experiment_id=int(row['id']),
|
|
now=now,
|
|
)
|
|
if not released:
|
|
remaining = conn.execute(
|
|
'''SELECT COUNT(*) AS count
|
|
FROM target_queue_policy_events cold_event
|
|
LEFT JOIN target_queue_policy_events reverse_event
|
|
ON reverse_event.reverses_event_id = cold_event.id
|
|
WHERE cold_event.experiment_id = ? AND cold_event.action = 'cold'
|
|
AND reverse_event.id IS NULL''',
|
|
(row['id'],),
|
|
).fetchone()['count']
|
|
if int(remaining):
|
|
raise RuntimeError('Docker depth reactivation left owned cold events unreversed')
|
|
cursor = conn.execute(
|
|
'''UPDATE docker_depth_experiments
|
|
SET state = 'released', released_at = COALESCE(released_at, ?),
|
|
updated_at = ?
|
|
WHERE id = ? AND state = 'completed'
|
|
AND hold_manifest_sha256 = ?''',
|
|
(now, now, row['id'], manifest['hold_manifest_sha256']),
|
|
)
|
|
if int(cursor.rowcount or 0) != 1:
|
|
raise RuntimeError('Docker depth release lost its authority fence')
|
|
conn.commit()
|
|
return {
|
|
**result,
|
|
'experiment_id': int(row['id']),
|
|
'state': 'released',
|
|
}
|
|
except Exception:
|
|
conn.rollback()
|
|
raise
|
|
|
|
|
|
def _docker_depth_resolver_refund_text_sha256(value):
|
|
return hashlib.sha256(str(value).encode('utf-8')).hexdigest()
|
|
|
|
|
|
def _docker_depth_resolver_refund_log(path):
|
|
absolute = os.path.abspath(os.fspath(path))
|
|
size = os.path.getsize(absolute)
|
|
if size < 1 or size > DOCKER_DEPTH_RESOLVER_REFUND_LOG_MAX_BYTES:
|
|
raise RuntimeError('Docker depth resolver refund log size is invalid')
|
|
with open(absolute, 'rb') as handle:
|
|
payload = handle.read(DOCKER_DEPTH_RESOLVER_REFUND_LOG_MAX_BYTES + 1)
|
|
if len(payload) != size or len(payload) > DOCKER_DEPTH_RESOLVER_REFUND_LOG_MAX_BYTES:
|
|
raise RuntimeError('Docker depth resolver refund log changed while reading')
|
|
return payload, hashlib.sha256(payload).hexdigest()
|
|
|
|
|
|
def _docker_depth_resolver_refund_entry(row, log_lines):
|
|
target = str(row['normalized_target'] or '')
|
|
if not target:
|
|
raise RuntimeError('Docker depth resolver refund target identity is absent')
|
|
exact_error = (
|
|
f'Unable to fetch Docker Hub tags for {target}: '
|
|
f'{DOCKER_DEPTH_RESOLVER_REFUND_OLD_ERROR}'
|
|
)
|
|
bug_count = sum(exact_error in line for line in log_lines)
|
|
if bug_count != DOCKER_DEPTH_RESOLVER_REFUND_ATTEMPTS:
|
|
return None
|
|
work_state = str(row['work_state'] or '')
|
|
attempts = int(row['resolver_attempts'])
|
|
last_error = str(row['last_error_code'] or '')
|
|
expected_shape = (
|
|
work_state == 'pending'
|
|
and attempts == DOCKER_DEPTH_RESOLVER_REFUND_ATTEMPTS
|
|
and last_error == DOCKER_DEPTH_RESOLVER_REFUND_OLD_ERROR
|
|
) or (
|
|
work_state == 'held'
|
|
and attempts == DOCKER_DEPTH_RESOLVER_REFUND_ATTEMPTS + 1
|
|
and last_error == 'resolver_attempt_limit'
|
|
)
|
|
if not expected_shape:
|
|
raise RuntimeError('Docker depth bug-tainted resolver state is not refundable')
|
|
entry = {
|
|
'experiment_repository_id': int(row['id']),
|
|
'repository_queue_id': int(row['effective_repository_queue_id']),
|
|
'query_ordinal': int(row['query_ordinal']),
|
|
'repository_rank': int(row['repository_rank']),
|
|
'recovery_kind': DOCKER_DEPTH_RESOLVER_REFUND_KIND,
|
|
'prior_work_state': work_state,
|
|
'next_work_state': 'pending',
|
|
'prior_resolver_attempts': attempts,
|
|
'refund_attempts': DOCKER_DEPTH_RESOLVER_REFUND_ATTEMPTS,
|
|
'next_resolver_attempts': attempts - DOCKER_DEPTH_RESOLVER_REFUND_ATTEMPTS,
|
|
'confirmed_bug_event_count': bug_count,
|
|
'prior_error_code_sha256': _docker_depth_resolver_refund_text_sha256(last_error),
|
|
'target_identity_sha256': _docker_depth_resolver_refund_text_sha256(target),
|
|
}
|
|
entry['entry_evidence_sha256'] = _canonical_sha256(entry)
|
|
return entry
|
|
|
|
|
|
def validate_docker_depth_resolver_refund_manifest(
|
|
manifest, experiment=None, provenance_policy_sha256=None,
|
|
):
|
|
expected_keys = {
|
|
'schema', 'type', 'version', 'experiment_id', 'experiment_key', 'source',
|
|
'config_sha256', 'ordered_queries_sha256', 'selector_sha256',
|
|
'provenance_policy_sha256', 'plan_sha256', 'hold_reason_code',
|
|
'recovery_kind', 'refund_attempts', 'evidence_log_name',
|
|
'evidence_log_sha256', 'entry_count', 'held_entry_count',
|
|
'selection_sha256', 'entries',
|
|
}
|
|
if (
|
|
not isinstance(manifest, dict)
|
|
or set(manifest) != expected_keys
|
|
or type(manifest.get('schema')) is not int
|
|
or manifest.get('schema') != 1
|
|
or type(manifest.get('version')) is not int
|
|
or manifest.get('version') != 1
|
|
or manifest.get('type') != DOCKER_DEPTH_RESOLVER_REFUND_MANIFEST_TYPE
|
|
):
|
|
raise ValueError('Docker depth resolver refund manifest shape is invalid')
|
|
_strict_identifier(manifest.get('experiment_key'), 'experiment_key')
|
|
entries = manifest.get('entries')
|
|
if not isinstance(entries, list) or not entries:
|
|
raise ValueError('Docker depth resolver refund manifest entries are invalid')
|
|
normalized_entries = []
|
|
member_ids = set()
|
|
for raw in entries:
|
|
expected_entry_keys = {
|
|
'experiment_repository_id', 'repository_queue_id', 'query_ordinal',
|
|
'repository_rank', 'recovery_kind', 'prior_work_state',
|
|
'next_work_state', 'prior_resolver_attempts', 'refund_attempts',
|
|
'next_resolver_attempts', 'confirmed_bug_event_count',
|
|
'prior_error_code_sha256', 'target_identity_sha256',
|
|
'entry_evidence_sha256',
|
|
}
|
|
if not isinstance(raw, dict) or set(raw) != expected_entry_keys:
|
|
raise ValueError('Docker depth resolver refund entry shape is invalid')
|
|
item = copy.deepcopy(raw)
|
|
member_id = _strict_positive_int(
|
|
item['experiment_repository_id'], 'experiment_repository_id',
|
|
)
|
|
if member_id in member_ids:
|
|
raise ValueError('Docker depth resolver refund entry is duplicated')
|
|
member_ids.add(member_id)
|
|
_strict_positive_int(item['repository_queue_id'], 'repository_queue_id')
|
|
_strict_nonnegative_int(item['query_ordinal'], 'query_ordinal')
|
|
_strict_positive_int(item['repository_rank'], 'repository_rank')
|
|
prior_attempts = _strict_positive_int(
|
|
item['prior_resolver_attempts'], 'prior_resolver_attempts',
|
|
)
|
|
next_attempts = _strict_nonnegative_int(
|
|
item['next_resolver_attempts'], 'next_resolver_attempts',
|
|
)
|
|
if (
|
|
item['recovery_kind'] != DOCKER_DEPTH_RESOLVER_REFUND_KIND
|
|
or item['prior_work_state'] not in ('pending', 'held')
|
|
or item['next_work_state'] != 'pending'
|
|
or item['refund_attempts'] != DOCKER_DEPTH_RESOLVER_REFUND_ATTEMPTS
|
|
or item['confirmed_bug_event_count'] != DOCKER_DEPTH_RESOLVER_REFUND_ATTEMPTS
|
|
or next_attempts != prior_attempts - DOCKER_DEPTH_RESOLVER_REFUND_ATTEMPTS
|
|
or next_attempts < 0
|
|
or any(not _valid_sha256(item.get(name)) for name in (
|
|
'prior_error_code_sha256', 'target_identity_sha256',
|
|
'entry_evidence_sha256',
|
|
))
|
|
):
|
|
raise ValueError('Docker depth resolver refund entry evidence is invalid')
|
|
evidence = dict(item)
|
|
evidence.pop('entry_evidence_sha256')
|
|
if item['entry_evidence_sha256'] != _canonical_sha256(evidence):
|
|
raise ValueError('Docker depth resolver refund entry hash conflicts')
|
|
normalized_entries.append(item)
|
|
normalized_entries.sort(key=lambda item: item['experiment_repository_id'])
|
|
if normalized_entries != entries:
|
|
raise ValueError('Docker depth resolver refund entries are not canonical')
|
|
hashes = (
|
|
'config_sha256', 'ordered_queries_sha256', 'selector_sha256',
|
|
'provenance_policy_sha256', 'plan_sha256', 'evidence_log_sha256',
|
|
'selection_sha256',
|
|
)
|
|
if any(not _valid_sha256(manifest.get(name)) for name in hashes):
|
|
raise ValueError('Docker depth resolver refund manifest hash is invalid')
|
|
if (
|
|
_strict_positive_int(manifest.get('experiment_id'), 'experiment_id') <= 0
|
|
or manifest.get('source') != 'dockerhub'
|
|
or manifest.get('hold_reason_code') != 'resolver_attempt_limit'
|
|
or manifest.get('recovery_kind') != DOCKER_DEPTH_RESOLVER_REFUND_KIND
|
|
or manifest.get('refund_attempts') != DOCKER_DEPTH_RESOLVER_REFUND_ATTEMPTS
|
|
or manifest.get('evidence_log_name') != 'dockerhub.log'
|
|
or manifest.get('entry_count') != len(normalized_entries)
|
|
or manifest.get('held_entry_count')
|
|
!= sum(item['prior_work_state'] == 'held' for item in normalized_entries)
|
|
or manifest.get('held_entry_count') != 1
|
|
or manifest.get('selection_sha256') != _canonical_sha256(normalized_entries)
|
|
):
|
|
raise ValueError('Docker depth resolver refund manifest evidence conflicts')
|
|
if experiment is not None:
|
|
authority = _docker_depth_authority(
|
|
experiment,
|
|
provenance_policy_sha256 or manifest['provenance_policy_sha256'],
|
|
)
|
|
for manifest_name, authority_name in (
|
|
('experiment_key', 'experiment_key'), ('source', 'source'),
|
|
('config_sha256', 'config_sha256'),
|
|
('ordered_queries_sha256', 'ordered_queries_sha256'),
|
|
('selector_sha256', 'selector_sha256'),
|
|
('provenance_policy_sha256', 'provenance_policy_sha256'),
|
|
):
|
|
if manifest[manifest_name] != authority[authority_name]:
|
|
raise ValueError('Docker depth resolver refund authority drifted')
|
|
normalized = copy.deepcopy(manifest)
|
|
normalized['entries'] = normalized_entries
|
|
if normalized != manifest:
|
|
raise ValueError('Docker depth resolver refund manifest is not canonical')
|
|
return normalized, _canonical_sha256(normalized)
|
|
|
|
|
|
def generate_docker_depth_resolver_refund_manifest(
|
|
db, experiment, provenance_policy_sha256, evidence_log_path,
|
|
):
|
|
authority = _docker_depth_authority(experiment, provenance_policy_sha256)
|
|
conn = _require_postgres_experiment_db(db, 'Docker depth resolver refund planning')
|
|
log_payload, log_sha256 = _docker_depth_resolver_refund_log(evidence_log_path)
|
|
log_lines = log_payload.decode('utf-8', 'replace').splitlines()
|
|
try:
|
|
row = _experiment_row(conn, authority, False)
|
|
if (
|
|
str(row['state']) != 'held'
|
|
or str(row['hold_reason_code'] or '') != 'resolver_attempt_limit'
|
|
or not _valid_sha256(row['plan_sha256'])
|
|
or _experiment_fence_active(row)
|
|
):
|
|
raise RuntimeError('Docker depth resolver refund requires an attempt-limit hold')
|
|
members = conn.execute(
|
|
'''SELECT member.*, queue.id AS effective_repository_queue_id,
|
|
queue.normalized_target
|
|
FROM docker_depth_experiment_repositories member
|
|
JOIN target_queue queue
|
|
ON queue.id = COALESCE(
|
|
member.replacement_repository_queue_id,
|
|
member.repository_queue_id
|
|
)
|
|
WHERE member.experiment_id = ? AND member.resolver_attempts >= ?
|
|
ORDER BY member.id''',
|
|
(row['id'], DOCKER_DEPTH_RESOLVER_REFUND_ATTEMPTS),
|
|
).fetchall()
|
|
entries = []
|
|
for member in members:
|
|
entry = _docker_depth_resolver_refund_entry(member, log_lines)
|
|
if entry is not None:
|
|
entries.append(entry)
|
|
if not entries:
|
|
raise RuntimeError('No confirmed bug-tainted resolver attempts were found')
|
|
manifest = {
|
|
'schema': 1,
|
|
'type': DOCKER_DEPTH_RESOLVER_REFUND_MANIFEST_TYPE,
|
|
'version': 1,
|
|
'experiment_id': int(row['id']),
|
|
'experiment_key': authority['experiment_key'],
|
|
'source': authority['source'],
|
|
'config_sha256': authority['config_sha256'],
|
|
'ordered_queries_sha256': authority['ordered_queries_sha256'],
|
|
'selector_sha256': authority['selector_sha256'],
|
|
'provenance_policy_sha256': authority['provenance_policy_sha256'],
|
|
'plan_sha256': str(row['plan_sha256']),
|
|
'hold_reason_code': 'resolver_attempt_limit',
|
|
'recovery_kind': DOCKER_DEPTH_RESOLVER_REFUND_KIND,
|
|
'refund_attempts': DOCKER_DEPTH_RESOLVER_REFUND_ATTEMPTS,
|
|
'evidence_log_name': 'dockerhub.log',
|
|
'evidence_log_sha256': log_sha256,
|
|
'entry_count': len(entries),
|
|
'held_entry_count': sum(
|
|
item['prior_work_state'] == 'held' for item in entries
|
|
),
|
|
'selection_sha256': _canonical_sha256(entries),
|
|
'entries': entries,
|
|
}
|
|
normalized, manifest_sha256 = validate_docker_depth_resolver_refund_manifest(
|
|
manifest, experiment, provenance_policy_sha256,
|
|
)
|
|
conn.commit()
|
|
return normalized, manifest_sha256
|
|
except Exception:
|
|
conn.rollback()
|
|
raise
|
|
|
|
|
|
def apply_docker_depth_resolver_refund_manifest(
|
|
db, experiment, provenance_policy_sha256, manifest, manifest_sha256,
|
|
evidence_log_path,
|
|
):
|
|
authority = _docker_depth_authority(experiment, provenance_policy_sha256)
|
|
manifest, expected_sha256 = validate_docker_depth_resolver_refund_manifest(
|
|
manifest, experiment, provenance_policy_sha256,
|
|
)
|
|
if str(manifest_sha256 or '') != expected_sha256:
|
|
raise ValueError('Docker depth resolver refund manifest hash conflicts')
|
|
_, current_log_sha256 = _docker_depth_resolver_refund_log(evidence_log_path)
|
|
if current_log_sha256 != manifest['evidence_log_sha256']:
|
|
raise RuntimeError('Docker depth resolver refund evidence log drifted')
|
|
conn = _require_postgres_experiment_db(db, 'Docker depth resolver refund')
|
|
now = datetime.now(timezone.utc).isoformat(timespec='seconds')
|
|
try:
|
|
row = _locked_experiment_row(conn, authority)
|
|
if (
|
|
int(row['id']) != manifest['experiment_id']
|
|
or str(row['plan_sha256'] or '') != manifest['plan_sha256']
|
|
or _experiment_fence_active(row)
|
|
):
|
|
raise RuntimeError('Docker depth resolver refund experiment identity drifted')
|
|
existing = conn.execute(
|
|
'''SELECT * FROM docker_depth_resolver_attempt_refunds
|
|
WHERE experiment_id = ? AND recovery_kind = ? ORDER BY id FOR UPDATE''',
|
|
(row['id'], DOCKER_DEPTH_RESOLVER_REFUND_KIND),
|
|
).fetchall()
|
|
if existing:
|
|
expected_entries = {
|
|
item['experiment_repository_id']: item for item in manifest['entries']
|
|
}
|
|
if (
|
|
len(existing) != len(expected_entries)
|
|
or any(
|
|
int(audit['experiment_repository_id']) not in expected_entries
|
|
or str(audit['manifest_sha256']) != expected_sha256
|
|
or str(audit['entry_evidence_sha256'])
|
|
!= expected_entries[int(audit['experiment_repository_id'])][
|
|
'entry_evidence_sha256'
|
|
]
|
|
for audit in existing
|
|
)
|
|
):
|
|
raise RuntimeError('Docker depth resolver refund audit conflicts')
|
|
conn.commit()
|
|
return {
|
|
'experiment_id': int(row['id']),
|
|
'refunded': 0,
|
|
'duplicates': len(existing),
|
|
'state': str(row['state']),
|
|
}
|
|
if (
|
|
str(row['state']) != 'held'
|
|
or str(row['hold_reason_code'] or '') != 'resolver_attempt_limit'
|
|
):
|
|
raise RuntimeError('Docker depth resolver refund requires an attempt-limit hold')
|
|
member_ids = [item['experiment_repository_id'] for item in manifest['entries']]
|
|
placeholders = ','.join('?' for _ in member_ids)
|
|
members = conn.execute(
|
|
f'''SELECT member.*, queue.id AS effective_repository_queue_id,
|
|
queue.normalized_target
|
|
FROM docker_depth_experiment_repositories member
|
|
JOIN target_queue queue
|
|
ON queue.id = COALESCE(
|
|
member.replacement_repository_queue_id,
|
|
member.repository_queue_id
|
|
)
|
|
WHERE member.experiment_id = ? AND member.id IN ({placeholders})
|
|
ORDER BY member.id FOR UPDATE OF member, queue''',
|
|
(row['id'], *member_ids),
|
|
).fetchall()
|
|
if len(members) != len(member_ids):
|
|
raise RuntimeError('Docker depth resolver refund membership set drifted')
|
|
entries_by_id = {
|
|
item['experiment_repository_id']: item for item in manifest['entries']
|
|
}
|
|
for member in members:
|
|
member_id = int(member['id'])
|
|
entry = entries_by_id.get(member_id)
|
|
last_error = str(member['last_error_code'] or '')
|
|
target = str(member['normalized_target'] or '')
|
|
if (
|
|
entry is None
|
|
or int(member['effective_repository_queue_id'])
|
|
!= entry['repository_queue_id']
|
|
or int(member['query_ordinal']) != entry['query_ordinal']
|
|
or int(member['repository_rank']) != entry['repository_rank']
|
|
or str(member['work_state']) != entry['prior_work_state']
|
|
or int(member['resolver_attempts']) != entry['prior_resolver_attempts']
|
|
or _docker_depth_resolver_refund_text_sha256(last_error)
|
|
!= entry['prior_error_code_sha256']
|
|
or _docker_depth_resolver_refund_text_sha256(target)
|
|
!= entry['target_identity_sha256']
|
|
or any(member[name] is not None for name in (
|
|
'resolver_owner', 'resolver_token', 'resolver_expires_at',
|
|
))
|
|
):
|
|
raise RuntimeError('Docker depth resolver refund membership evidence drifted')
|
|
cursor = conn.execute(
|
|
'''UPDATE docker_depth_experiment_repositories
|
|
SET work_state = 'pending', resolver_attempts = ?,
|
|
resolver_due_at = NULL, last_error_code = NULL,
|
|
updated_at = ?
|
|
WHERE id = ? AND experiment_id = ? AND work_state = ?
|
|
AND resolver_attempts = ? AND resolver_owner IS NULL
|
|
AND resolver_token IS NULL AND resolver_expires_at IS NULL''',
|
|
(
|
|
entry['next_resolver_attempts'], now, member_id, row['id'],
|
|
entry['prior_work_state'], entry['prior_resolver_attempts'],
|
|
),
|
|
)
|
|
if int(cursor.rowcount or 0) != 1:
|
|
raise RuntimeError('Docker depth resolver refund lost its member fence')
|
|
conn.execute(
|
|
'''INSERT INTO docker_depth_resolver_attempt_refunds(
|
|
experiment_id, experiment_repository_id, repository_queue_id,
|
|
query_ordinal, repository_rank, recovery_kind, manifest_sha256,
|
|
entry_evidence_sha256, log_sha256, target_identity_sha256,
|
|
prior_error_code_sha256, prior_work_state, next_work_state,
|
|
prior_resolver_attempts, refund_attempts, next_resolver_attempts,
|
|
confirmed_bug_event_count, applied_at, created_at
|
|
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)''',
|
|
(
|
|
row['id'], member_id, entry['repository_queue_id'],
|
|
entry['query_ordinal'], entry['repository_rank'],
|
|
DOCKER_DEPTH_RESOLVER_REFUND_KIND, expected_sha256,
|
|
entry['entry_evidence_sha256'], manifest['evidence_log_sha256'],
|
|
entry['target_identity_sha256'], entry['prior_error_code_sha256'],
|
|
entry['prior_work_state'], entry['next_work_state'],
|
|
entry['prior_resolver_attempts'], entry['refund_attempts'],
|
|
entry['next_resolver_attempts'], entry['confirmed_bug_event_count'],
|
|
now, now,
|
|
),
|
|
)
|
|
cursor = conn.execute(
|
|
'''UPDATE docker_depth_experiments
|
|
SET state = 'resolving', hold_reason_code = NULL, held_at = NULL,
|
|
updated_at = ?
|
|
WHERE id = ? AND state = 'held'
|
|
AND hold_reason_code = 'resolver_attempt_limit'
|
|
AND fence_owner IS NULL AND fence_token IS NULL
|
|
AND fence_expires_at IS NULL''',
|
|
(now, row['id']),
|
|
)
|
|
if int(cursor.rowcount or 0) != 1:
|
|
raise RuntimeError('Docker depth resolver refund lost its experiment fence')
|
|
refreshed = conn.execute(
|
|
'SELECT * FROM docker_depth_experiments WHERE id = ? FOR UPDATE',
|
|
(row['id'],),
|
|
).fetchone()
|
|
drift = _persisted_experiment_drift_reason(db, refreshed, authority, now)
|
|
if drift:
|
|
raise RuntimeError(f'Docker depth resolver refund left authority drift: {drift}')
|
|
conn.commit()
|
|
return {
|
|
'experiment_id': int(row['id']),
|
|
'refunded': len(manifest['entries']),
|
|
'duplicates': 0,
|
|
'state': 'resolving',
|
|
}
|
|
except Exception:
|
|
conn.rollback()
|
|
raise
|
|
|
|
|
|
def _docker_depth_resolver_disposition_entry(member, candidates):
|
|
target = str(member['normalized_target'] or '')
|
|
if not target:
|
|
raise RuntimeError('Docker depth resolver disposition target identity is absent')
|
|
if (
|
|
str(member['work_state'] or '') != 'held'
|
|
or str(member['last_error_code'] or '') != 'resolver_attempt_limit'
|
|
or int(member['resolver_attempts']) < DOCKER_DEPTH_RESOLVER_MAX_ATTEMPTS
|
|
or int(member['selected_image_count']) != 0
|
|
or any(member[name] is not None for name in (
|
|
'resolver_owner', 'resolver_token', 'resolver_expires_at',
|
|
'resolver_due_at',
|
|
))
|
|
):
|
|
raise RuntimeError('Docker depth attempt-limit disposition state is invalid')
|
|
replacement = next(
|
|
(candidate for candidate in candidates if not candidate['conflict_reason']),
|
|
None,
|
|
)
|
|
outcome = 'replaced' if replacement else 'skipped'
|
|
attempts = int(member['resolver_attempts'])
|
|
entry = {
|
|
'experiment_repository_id': int(member['id']),
|
|
'repository_queue_id': int(member['effective_repository_queue_id']),
|
|
'query_ordinal': int(member['query_ordinal']),
|
|
'repository_rank': int(member['repository_rank']),
|
|
'disposition_kind': DOCKER_DEPTH_RESOLVER_DISPOSITION_KIND,
|
|
'outcome': outcome,
|
|
'prior_work_state': 'held',
|
|
'next_work_state': 'pending' if replacement else 'skipped',
|
|
'prior_resolver_attempts': attempts,
|
|
'next_resolver_attempts': 0 if replacement else attempts,
|
|
'prior_error_code_sha256': _docker_depth_resolver_refund_text_sha256(
|
|
member['last_error_code']
|
|
),
|
|
'target_identity_sha256': _docker_depth_resolver_refund_text_sha256(target),
|
|
'candidate_snapshot_sha256': _canonical_sha256(candidates),
|
|
'inspected_candidate_count': len(candidates),
|
|
'replacement_repository_queue_id': (
|
|
replacement['repository_queue_id'] if replacement else None
|
|
),
|
|
'replacement_eligibility_page_id': (
|
|
replacement['eligibility_page_id'] if replacement else None
|
|
),
|
|
'replacement_best_search_rank': (
|
|
replacement['best_search_rank'] if replacement else None
|
|
),
|
|
'replacement_target_identity_sha256': (
|
|
replacement['target_identity_sha256'] if replacement else None
|
|
),
|
|
'terminal_reason': (
|
|
'' if replacement else DOCKER_DEPTH_REMOTE_UNAVAILABLE_SKIP_REASON
|
|
),
|
|
}
|
|
entry['entry_evidence_sha256'] = _canonical_sha256(entry)
|
|
return entry
|
|
|
|
|
|
def validate_docker_depth_resolver_disposition_manifest(
|
|
manifest, experiment=None, provenance_policy_sha256=None,
|
|
):
|
|
expected_keys = {
|
|
'schema', 'type', 'version', 'experiment_id', 'experiment_key', 'source',
|
|
'config_sha256', 'ordered_queries_sha256', 'selector_sha256',
|
|
'provenance_policy_sha256', 'plan_sha256', 'hold_reason_code',
|
|
'disposition_kind', 'entry_count', 'selection_sha256', 'entries',
|
|
}
|
|
if (
|
|
not isinstance(manifest, dict)
|
|
or set(manifest) != expected_keys
|
|
or type(manifest.get('schema')) is not int
|
|
or manifest.get('schema') != 1
|
|
or type(manifest.get('version')) is not int
|
|
or manifest.get('version') != 1
|
|
or manifest.get('type') != DOCKER_DEPTH_RESOLVER_DISPOSITION_MANIFEST_TYPE
|
|
):
|
|
raise ValueError('Docker depth resolver disposition manifest shape is invalid')
|
|
_strict_identifier(manifest.get('experiment_key'), 'experiment_key')
|
|
entries = manifest.get('entries')
|
|
if not isinstance(entries, list) or len(entries) != 1:
|
|
raise ValueError('Docker depth resolver disposition requires one reviewed entry')
|
|
raw = entries[0]
|
|
expected_entry_keys = {
|
|
'experiment_repository_id', 'repository_queue_id', 'query_ordinal',
|
|
'repository_rank', 'disposition_kind', 'outcome', 'prior_work_state',
|
|
'next_work_state', 'prior_resolver_attempts', 'next_resolver_attempts',
|
|
'prior_error_code_sha256', 'target_identity_sha256',
|
|
'candidate_snapshot_sha256', 'inspected_candidate_count',
|
|
'replacement_repository_queue_id', 'replacement_eligibility_page_id',
|
|
'replacement_best_search_rank', 'replacement_target_identity_sha256',
|
|
'terminal_reason', 'entry_evidence_sha256',
|
|
}
|
|
if not isinstance(raw, dict) or set(raw) != expected_entry_keys:
|
|
raise ValueError('Docker depth resolver disposition entry shape is invalid')
|
|
entry = copy.deepcopy(raw)
|
|
_strict_positive_int(entry['experiment_repository_id'], 'experiment_repository_id')
|
|
_strict_positive_int(entry['repository_queue_id'], 'repository_queue_id')
|
|
_strict_nonnegative_int(entry['query_ordinal'], 'query_ordinal')
|
|
_strict_positive_int(entry['repository_rank'], 'repository_rank')
|
|
attempts = _strict_positive_int(
|
|
entry['prior_resolver_attempts'], 'prior_resolver_attempts',
|
|
)
|
|
next_attempts = _strict_nonnegative_int(
|
|
entry['next_resolver_attempts'], 'next_resolver_attempts',
|
|
)
|
|
inspected = _strict_nonnegative_int(
|
|
entry['inspected_candidate_count'], 'inspected_candidate_count',
|
|
)
|
|
if inspected > 3000 or attempts < DOCKER_DEPTH_RESOLVER_MAX_ATTEMPTS:
|
|
raise ValueError('Docker depth resolver disposition range is invalid')
|
|
outcome = entry['outcome']
|
|
replacement_fields = (
|
|
'replacement_repository_queue_id', 'replacement_eligibility_page_id',
|
|
'replacement_best_search_rank', 'replacement_target_identity_sha256',
|
|
)
|
|
if outcome == 'replaced':
|
|
for name in replacement_fields[:3]:
|
|
_strict_positive_int(entry[name], name)
|
|
if (
|
|
entry['next_work_state'] != 'pending'
|
|
or next_attempts != 0
|
|
or entry['terminal_reason'] != ''
|
|
or not _valid_sha256(entry['replacement_target_identity_sha256'])
|
|
or inspected < 1
|
|
):
|
|
raise ValueError('Docker depth replacement disposition is invalid')
|
|
elif outcome == 'skipped':
|
|
if (
|
|
any(entry[name] is not None for name in replacement_fields)
|
|
or entry['next_work_state'] != 'skipped'
|
|
or next_attempts != attempts
|
|
or entry['terminal_reason'] != DOCKER_DEPTH_REMOTE_UNAVAILABLE_SKIP_REASON
|
|
):
|
|
raise ValueError('Docker depth skip disposition is invalid')
|
|
else:
|
|
raise ValueError('Docker depth resolver disposition outcome is invalid')
|
|
if (
|
|
entry['disposition_kind'] != DOCKER_DEPTH_RESOLVER_DISPOSITION_KIND
|
|
or entry['prior_work_state'] != 'held'
|
|
or any(not _valid_sha256(entry.get(name)) for name in (
|
|
'prior_error_code_sha256', 'target_identity_sha256',
|
|
'candidate_snapshot_sha256', 'entry_evidence_sha256',
|
|
))
|
|
):
|
|
raise ValueError('Docker depth resolver disposition evidence is invalid')
|
|
evidence = dict(entry)
|
|
evidence.pop('entry_evidence_sha256')
|
|
if entry['entry_evidence_sha256'] != _canonical_sha256(evidence):
|
|
raise ValueError('Docker depth resolver disposition entry hash conflicts')
|
|
hashes = (
|
|
'config_sha256', 'ordered_queries_sha256', 'selector_sha256',
|
|
'provenance_policy_sha256', 'plan_sha256', 'selection_sha256',
|
|
)
|
|
if any(not _valid_sha256(manifest.get(name)) for name in hashes):
|
|
raise ValueError('Docker depth resolver disposition manifest hash is invalid')
|
|
if (
|
|
_strict_positive_int(manifest.get('experiment_id'), 'experiment_id') <= 0
|
|
or manifest.get('source') != 'dockerhub'
|
|
or manifest.get('hold_reason_code') != 'resolver_attempt_limit'
|
|
or manifest.get('disposition_kind') != DOCKER_DEPTH_RESOLVER_DISPOSITION_KIND
|
|
or manifest.get('entry_count') != 1
|
|
or manifest.get('selection_sha256') != _canonical_sha256(entries)
|
|
):
|
|
raise ValueError('Docker depth resolver disposition authority is invalid')
|
|
if experiment is not None:
|
|
authority = _docker_depth_authority(
|
|
experiment,
|
|
provenance_policy_sha256 or manifest['provenance_policy_sha256'],
|
|
)
|
|
for manifest_name, authority_name in (
|
|
('experiment_key', 'experiment_key'), ('source', 'source'),
|
|
('config_sha256', 'config_sha256'),
|
|
('ordered_queries_sha256', 'ordered_queries_sha256'),
|
|
('selector_sha256', 'selector_sha256'),
|
|
('provenance_policy_sha256', 'provenance_policy_sha256'),
|
|
):
|
|
if manifest[manifest_name] != authority[authority_name]:
|
|
raise ValueError('Docker depth resolver disposition authority drifted')
|
|
normalized = copy.deepcopy(manifest)
|
|
normalized['entries'] = [entry]
|
|
if normalized != manifest:
|
|
raise ValueError('Docker depth resolver disposition manifest is not canonical')
|
|
return normalized, _canonical_sha256(normalized)
|
|
|
|
|
|
def generate_docker_depth_resolver_disposition_manifest(
|
|
db, experiment, provenance_policy_sha256,
|
|
):
|
|
authority = _docker_depth_authority(experiment, provenance_policy_sha256)
|
|
conn = _require_postgres_experiment_db(
|
|
db, 'Docker depth resolver disposition planning',
|
|
)
|
|
try:
|
|
row = _experiment_row(conn, authority, False)
|
|
if (
|
|
str(row['state']) != 'held'
|
|
or str(row['hold_reason_code'] or '') != 'resolver_attempt_limit'
|
|
or not _valid_sha256(row['plan_sha256'])
|
|
or _experiment_fence_active(row)
|
|
):
|
|
raise RuntimeError(
|
|
'Docker depth resolver disposition requires an attempt-limit hold'
|
|
)
|
|
members = conn.execute(
|
|
'''SELECT member.*, queue.id AS effective_repository_queue_id,
|
|
queue.normalized_target
|
|
FROM docker_depth_experiment_repositories member
|
|
JOIN target_queue queue ON queue.id = COALESCE(
|
|
member.replacement_repository_queue_id,
|
|
member.repository_queue_id
|
|
)
|
|
WHERE member.experiment_id = ? AND member.work_state = 'held'
|
|
AND member.last_error_code = 'resolver_attempt_limit'
|
|
ORDER BY member.id''',
|
|
(row['id'],),
|
|
).fetchall()
|
|
if len(members) != 1:
|
|
raise RuntimeError(
|
|
'Docker depth resolver disposition requires one held membership'
|
|
)
|
|
member = members[0]
|
|
candidates = db._docker_depth_repository_replacement_candidates(
|
|
row, member, authority, lock=False,
|
|
)
|
|
entry = _docker_depth_resolver_disposition_entry(member, candidates)
|
|
manifest = {
|
|
'schema': 1,
|
|
'type': DOCKER_DEPTH_RESOLVER_DISPOSITION_MANIFEST_TYPE,
|
|
'version': 1,
|
|
'experiment_id': int(row['id']),
|
|
'experiment_key': authority['experiment_key'],
|
|
'source': authority['source'],
|
|
'config_sha256': authority['config_sha256'],
|
|
'ordered_queries_sha256': authority['ordered_queries_sha256'],
|
|
'selector_sha256': authority['selector_sha256'],
|
|
'provenance_policy_sha256': authority['provenance_policy_sha256'],
|
|
'plan_sha256': str(row['plan_sha256']),
|
|
'hold_reason_code': 'resolver_attempt_limit',
|
|
'disposition_kind': DOCKER_DEPTH_RESOLVER_DISPOSITION_KIND,
|
|
'entry_count': 1,
|
|
'selection_sha256': _canonical_sha256([entry]),
|
|
'entries': [entry],
|
|
}
|
|
normalized, manifest_sha256 = (
|
|
validate_docker_depth_resolver_disposition_manifest(
|
|
manifest, experiment, provenance_policy_sha256,
|
|
)
|
|
)
|
|
conn.commit()
|
|
return normalized, manifest_sha256
|
|
except Exception:
|
|
conn.rollback()
|
|
raise
|
|
|
|
|
|
def apply_docker_depth_resolver_disposition_manifest(
|
|
db, experiment, provenance_policy_sha256, manifest, manifest_sha256,
|
|
):
|
|
authority = _docker_depth_authority(experiment, provenance_policy_sha256)
|
|
manifest, expected_sha256 = validate_docker_depth_resolver_disposition_manifest(
|
|
manifest, experiment, provenance_policy_sha256,
|
|
)
|
|
if str(manifest_sha256 or '') != expected_sha256:
|
|
raise ValueError('Docker depth resolver disposition manifest hash conflicts')
|
|
conn = _require_postgres_experiment_db(db, 'Docker depth resolver disposition')
|
|
now = datetime.now(timezone.utc).isoformat(timespec='seconds')
|
|
entry = manifest['entries'][0]
|
|
try:
|
|
row = _locked_experiment_row(conn, authority)
|
|
if (
|
|
int(row['id']) != manifest['experiment_id']
|
|
or str(row['plan_sha256'] or '') != manifest['plan_sha256']
|
|
or _experiment_fence_active(row)
|
|
):
|
|
raise RuntimeError('Docker depth resolver disposition identity drifted')
|
|
existing = conn.execute(
|
|
'''SELECT * FROM docker_depth_resolver_dispositions
|
|
WHERE experiment_repository_id = ? AND disposition_kind = ?
|
|
FOR UPDATE''',
|
|
(
|
|
entry['experiment_repository_id'],
|
|
DOCKER_DEPTH_RESOLVER_DISPOSITION_KIND,
|
|
),
|
|
).fetchone()
|
|
if existing:
|
|
if (
|
|
int(existing['experiment_id']) != int(row['id'])
|
|
or str(existing['manifest_sha256']) != expected_sha256
|
|
or str(existing['entry_evidence_sha256'])
|
|
!= entry['entry_evidence_sha256']
|
|
or str(existing['outcome']) != entry['outcome']
|
|
):
|
|
raise RuntimeError('Docker depth resolver disposition audit conflicts')
|
|
conn.commit()
|
|
return {
|
|
'experiment_id': int(row['id']), 'applied': 0, 'duplicates': 1,
|
|
'outcome': str(existing['outcome']), 'state': str(row['state']),
|
|
}
|
|
if (
|
|
str(row['state']) != 'held'
|
|
or str(row['hold_reason_code'] or '') != 'resolver_attempt_limit'
|
|
):
|
|
raise RuntimeError(
|
|
'Docker depth resolver disposition requires an attempt-limit hold'
|
|
)
|
|
member = conn.execute(
|
|
'''SELECT member.*, queue.id AS effective_repository_queue_id,
|
|
queue.normalized_target
|
|
FROM docker_depth_experiment_repositories member
|
|
JOIN target_queue queue ON queue.id = COALESCE(
|
|
member.replacement_repository_queue_id,
|
|
member.repository_queue_id
|
|
)
|
|
WHERE member.id = ? AND member.experiment_id = ?
|
|
FOR UPDATE OF member, queue''',
|
|
(entry['experiment_repository_id'], row['id']),
|
|
).fetchone()
|
|
if not member:
|
|
raise RuntimeError('Docker depth resolver disposition membership is absent')
|
|
target = str(member['normalized_target'] or '')
|
|
if (
|
|
int(member['effective_repository_queue_id']) != entry['repository_queue_id']
|
|
or int(member['query_ordinal']) != entry['query_ordinal']
|
|
or int(member['repository_rank']) != entry['repository_rank']
|
|
or str(member['work_state']) != entry['prior_work_state']
|
|
or int(member['resolver_attempts']) != entry['prior_resolver_attempts']
|
|
or str(member['last_error_code'] or '') != 'resolver_attempt_limit'
|
|
or int(member['selected_image_count']) != 0
|
|
or _docker_depth_resolver_refund_text_sha256(target)
|
|
!= entry['target_identity_sha256']
|
|
or _docker_depth_resolver_refund_text_sha256(member['last_error_code'])
|
|
!= entry['prior_error_code_sha256']
|
|
or any(member[name] is not None for name in (
|
|
'resolver_owner', 'resolver_token', 'resolver_expires_at',
|
|
'resolver_due_at',
|
|
))
|
|
):
|
|
raise RuntimeError('Docker depth resolver disposition evidence drifted')
|
|
if conn.execute(
|
|
'''SELECT 1 FROM docker_depth_experiment_selections
|
|
WHERE experiment_repository_id = ? LIMIT 1''',
|
|
(member['id'],),
|
|
).fetchone():
|
|
raise RuntimeError('Docker depth resolver disposition cannot replace selected work')
|
|
candidates = db._docker_depth_repository_replacement_candidates(
|
|
row, member, authority, lock=True,
|
|
)
|
|
current_entry = _docker_depth_resolver_disposition_entry(member, candidates)
|
|
if current_entry != entry:
|
|
raise RuntimeError('Docker depth resolver disposition selection drifted after review')
|
|
for candidate in candidates:
|
|
if not candidate['conflict_reason']:
|
|
break
|
|
db._record_docker_depth_candidate_skip_locked(
|
|
row, member, candidate['repository_queue_id'], 'repository',
|
|
candidate['best_search_rank'],
|
|
{'repository_queue_id': candidate['repository_queue_id']},
|
|
candidate['conflict_reason'], now,
|
|
)
|
|
current_queue_id = int(member['effective_repository_queue_id'])
|
|
db._record_docker_depth_candidate_skip_locked(
|
|
row, member, current_queue_id, 'repository',
|
|
int(member['replacement_count']) + 1,
|
|
{'repository_queue_id': current_queue_id},
|
|
DOCKER_DEPTH_REMOTE_UNAVAILABLE_SKIP_REASON, now,
|
|
)
|
|
replacement = next(
|
|
(candidate for candidate in candidates if not candidate['conflict_reason']),
|
|
None,
|
|
)
|
|
if replacement:
|
|
previous_hash = str(member['replacement_evidence_sha256'] or '')
|
|
replacement_document = {
|
|
'schema': 1,
|
|
'type': 'docker-depth-repository-replacement-v1',
|
|
'experiment_id': int(row['id']),
|
|
'experiment_repository_id': int(member['id']),
|
|
'replacement_number': int(member['replacement_count']) + 1,
|
|
'from_repository_queue_id': current_queue_id,
|
|
'to_repository_queue_id': replacement['repository_queue_id'],
|
|
'eligibility_page_id': replacement['eligibility_page_id'],
|
|
'best_search_rank': replacement['best_search_rank'],
|
|
'previous_evidence_sha256': previous_hash,
|
|
'reason_code': DOCKER_DEPTH_REMOTE_UNAVAILABLE_SKIP_REASON,
|
|
}
|
|
replacement_sha256 = _canonical_sha256(replacement_document)
|
|
cursor = conn.execute(
|
|
'''UPDATE docker_depth_experiment_repositories
|
|
SET replacement_repository_queue_id = ?,
|
|
replacement_eligibility_page_id = ?,
|
|
replacement_count = replacement_count + 1,
|
|
replacement_evidence_sha256 = ?, work_state = 'pending',
|
|
resolver_attempts = 0, resolver_due_at = NULL,
|
|
last_error_code = 'repository_candidate_replaced',
|
|
resolved_at = NULL, updated_at = ?
|
|
WHERE id = ? AND experiment_id = ? AND work_state = 'held'
|
|
AND resolver_attempts = ? AND resolver_owner IS NULL
|
|
AND resolver_token IS NULL AND resolver_expires_at IS NULL''',
|
|
(
|
|
replacement['repository_queue_id'],
|
|
replacement['eligibility_page_id'], replacement_sha256, now,
|
|
member['id'], row['id'], entry['prior_resolver_attempts'],
|
|
),
|
|
)
|
|
else:
|
|
cursor = conn.execute(
|
|
'''UPDATE docker_depth_experiment_repositories
|
|
SET work_state = 'skipped', resolver_due_at = NULL,
|
|
last_error_code = ?, selected_image_count = 0,
|
|
resolved_at = ?, updated_at = ?
|
|
WHERE id = ? AND experiment_id = ? AND work_state = 'held'
|
|
AND resolver_attempts = ? AND resolver_owner IS NULL
|
|
AND resolver_token IS NULL AND resolver_expires_at IS NULL''',
|
|
(
|
|
DOCKER_DEPTH_REMOTE_UNAVAILABLE_SKIP_REASON, now, now,
|
|
member['id'], row['id'], entry['prior_resolver_attempts'],
|
|
),
|
|
)
|
|
if int(cursor.rowcount or 0) != 1:
|
|
raise RuntimeError('Docker depth resolver disposition lost its member fence')
|
|
resume = conn.execute(
|
|
'''UPDATE docker_depth_experiments
|
|
SET state = 'resolving', hold_reason_code = NULL, held_at = NULL,
|
|
updated_at = ?
|
|
WHERE id = ? AND state = 'held'
|
|
AND hold_reason_code = 'resolver_attempt_limit'
|
|
AND fence_owner IS NULL AND fence_token IS NULL
|
|
AND fence_expires_at IS NULL''',
|
|
(now, row['id']),
|
|
)
|
|
if int(resume.rowcount or 0) != 1:
|
|
raise RuntimeError('Docker depth resolver disposition lost its experiment fence')
|
|
if not replacement:
|
|
refreshed_member = conn.execute(
|
|
'SELECT * FROM docker_depth_experiment_repositories WHERE id = ?',
|
|
(member['id'],),
|
|
).fetchone()
|
|
db._finalize_docker_depth_query_breadth_locked(row, refreshed_member, now)
|
|
refreshed = conn.execute(
|
|
'SELECT * FROM docker_depth_experiments WHERE id = ? FOR UPDATE',
|
|
(row['id'],),
|
|
).fetchone()
|
|
refreshed, selection_reason = db._freeze_docker_depth_runtime_selection_locked(
|
|
refreshed, authority, now,
|
|
)
|
|
if selection_reason:
|
|
raise RuntimeError(
|
|
f'Docker depth resolver disposition selection drift: {selection_reason}'
|
|
)
|
|
drift = _persisted_experiment_drift_reason(db, refreshed, authority, now)
|
|
if drift:
|
|
raise RuntimeError(
|
|
f'Docker depth resolver disposition left authority drift: {drift}'
|
|
)
|
|
refreshed = db._advance_docker_depth_experiment_state_locked(refreshed, now)
|
|
conn.execute(
|
|
'''INSERT INTO docker_depth_resolver_dispositions(
|
|
experiment_id, experiment_repository_id,
|
|
prior_repository_queue_id, replacement_repository_queue_id,
|
|
replacement_eligibility_page_id, replacement_best_search_rank,
|
|
query_ordinal, repository_rank, disposition_kind, outcome,
|
|
manifest_sha256, entry_evidence_sha256,
|
|
candidate_snapshot_sha256, prior_target_identity_sha256,
|
|
replacement_target_identity_sha256, prior_error_code_sha256,
|
|
prior_work_state, next_work_state, prior_resolver_attempts,
|
|
next_resolver_attempts, applied_at, created_at
|
|
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)''',
|
|
(
|
|
row['id'], member['id'], entry['repository_queue_id'],
|
|
entry['replacement_repository_queue_id'],
|
|
entry['replacement_eligibility_page_id'],
|
|
entry['replacement_best_search_rank'], entry['query_ordinal'],
|
|
entry['repository_rank'], DOCKER_DEPTH_RESOLVER_DISPOSITION_KIND,
|
|
entry['outcome'], expected_sha256, entry['entry_evidence_sha256'],
|
|
entry['candidate_snapshot_sha256'], entry['target_identity_sha256'],
|
|
entry['replacement_target_identity_sha256'],
|
|
entry['prior_error_code_sha256'], entry['prior_work_state'],
|
|
entry['next_work_state'], entry['prior_resolver_attempts'],
|
|
entry['next_resolver_attempts'], now, now,
|
|
),
|
|
)
|
|
conn.commit()
|
|
return {
|
|
'experiment_id': int(row['id']), 'applied': 1, 'duplicates': 0,
|
|
'outcome': entry['outcome'], 'state': str(refreshed['state']),
|
|
}
|
|
except Exception:
|
|
conn.rollback()
|
|
raise
|