from datetime import datetime, timedelta, timezone from pathlib import Path import os import sys from types import SimpleNamespace import unittest from unittest import mock ROOT = Path(__file__).resolve().parents[1] APP_DIR = ROOT / 'app' sys.path.insert(0, str(APP_DIR)) import console_runner def page(number, names, total_count=300): return { 'page': number, 'repositories': [{'repo_name': name} for name in names], 'total_count': total_count, 'elapsed': 0.01, } def discovery_args(*, pages=3, per_page=100, deep=False, query='fixture-query'): args = SimpleNamespace( platform='docker', mode='search', query=query, pages=pages, per_page=per_page, docker_sort_by='updated_at', sort_order='desc', fetch_timeout=15, ) policy = console_runner.dockerhub_discovery_policy(args) args.dockerhub_discovery_deep = deep args.dockerhub_discovery_pass_kind = 'deep' if deep else 'ordinary' args.dockerhub_discovery_policy_sha256 = policy['policy_sha256'] return args class Connection: is_postgres = True def __init__(self, events=None): self.events = events if events is not None else [] def rollback(self): self.events.append(('rollback',)) class PageDB: def __init__(self, known=(), enqueue_error=False, paused=False): self.conn = Connection() self.known = set(known) self.admissions = [] self.retries = [] self.enqueue_error = enqueue_error self.paused = paused def runtime_control_state(self): return {'effective_discovery_paused': self.paused} def persist_dockerhub_discovery_page(self, source, query, repositories, **kwargs): normalized = frozenset( str(repository['repo_name']).strip().lower() for repository in repositories ) preexisting = normalized & self.known self.known.update(normalized) self.admissions.append((tuple(sorted(normalized)), dict(kwargs))) complete = bool(kwargs.get('complete')) progress = kwargs.get('next_page') is not None return { 'attempted_count': len(repositories), 'normalized_count': len(normalized), 'preexisting_count': len(preexisting), 'inserted_count': len(normalized - preexisting), 'normalized_repositories': normalized, 'preexisting_repositories': frozenset(preexisting), 'retry_progress_count': int(progress), 'retry_completed_count': int(complete), } def enqueue_discovery_retry(self, *args, **kwargs): if self.enqueue_error: raise RuntimeError('raw database detail') self.retries.append((args, kwargs)) return { 'id': len(self.retries), 'status': 'pending', 'inserted_count': 1, 'coalesced_count': 0, } class DockerHubDeepStateTests(unittest.TestCase): def test_new_policy_and_72_hour_boundary_are_deep_due(self): now = datetime(2026, 1, 1, tzinfo=timezone.utc) policies = console_runner.configured_dockerhub_discovery_policies({ 'queries': ['first'], 'pages': 30, 'per_page': 100, }) state = console_runner.default_source_state() annotation = console_runner.prepare_dockerhub_discovery_state( state, policies, 'first', now, ) self.assertTrue(annotation['deep']) console_runner.mark_dockerhub_deep_dispatched( state, 'first', annotation['policy_sha256'], now, ) self.assertFalse(console_runner.prepare_dockerhub_discovery_state( state, policies, 'first', now + timedelta(hours=71, minutes=59), )['deep']) self.assertTrue(console_runner.prepare_dockerhub_discovery_state( state, policies, 'first', now + timedelta(hours=72), )['deep']) changed = console_runner.configured_dockerhub_discovery_policies({ 'queries': ['first'], 'pages': 29, 'per_page': 100, }) self.assertTrue(console_runner.prepare_dockerhub_discovery_state( state, changed, 'first', now + timedelta(hours=1), )['deep']) def test_appended_query_is_due_and_removed_or_malformed_state_is_pruned(self): now = datetime(2026, 1, 1, tzinfo=timezone.utc) base = console_runner.configured_dockerhub_discovery_policies({ 'queries': ['first'], 'pages': 30, 'per_page': 100, }) state = console_runner.default_source_state() console_runner.mark_dockerhub_deep_dispatched( state, 'first', base['first']['policy_sha256'], now, ) state[console_runner.DOCKERHUB_DISCOVERY_STATE_KEY]['deep_by_query']['removed'] = { 'policy_sha256': 'a' * 64, 'last_dispatched_at': now.isoformat(timespec='seconds'), } expanded = console_runner.configured_dockerhub_discovery_policies({ 'queries': ['first', 'second'], 'pages': 30, 'per_page': 100, }) second = console_runner.prepare_dockerhub_discovery_state( state, expanded, 'second', now + timedelta(hours=1), ) self.assertTrue(second['deep']) records = state[console_runner.DOCKERHUB_DISCOVERY_STATE_KEY]['deep_by_query'] self.assertIn('first', records) self.assertNotIn('removed', records) state[console_runner.DOCKERHUB_DISCOVERY_STATE_KEY] = { 'schema': 1, 'deep_by_query': {'first': {'policy_sha256': 'malformed'}}, } self.assertTrue(console_runner.prepare_dockerhub_discovery_state( state, base, 'first', now + timedelta(hours=1), )['deep']) self.assertEqual( state[console_runner.DOCKERHUB_DISCOVERY_STATE_KEY]['deep_by_query'], {}, ) def test_policy_hash_uses_effective_search_policy_not_query_list_shape(self): first = console_runner.configured_dockerhub_discovery_policies({ 'queries': ['one', 'two'], 'pages': 30, 'per_page': 100, }) reordered = console_runner.configured_dockerhub_discovery_policies({ 'queries': ['three', 'two', 'one'], 'pages': 30, 'per_page': 100, }) changed_sort = console_runner.configured_dockerhub_discovery_policies({ 'queries': ['one'], 'pages': 30, 'per_page': 100, 'docker_sort_by': 'name', }) self.assertEqual(first['one']['pages'], 30) self.assertEqual(first['one']['per_page'], 100) self.assertEqual( first['one']['policy_sha256'], reordered['one']['policy_sha256'], ) self.assertNotEqual( first['one']['policy_sha256'], changed_sort['one']['policy_sha256'], ) for key, value in (('pages', 31), ('per_page', 101), ('pages', '30')): with self.subTest(key=key, value=value), self.assertRaises(ValueError): console_runner.configured_dockerhub_discovery_policies({ 'queries': ['one'], key: value, }) def test_incomplete_pass_forces_deep_until_durably_cleared(self): now = datetime(2026, 1, 1, tzinfo=timezone.utc) policies = console_runner.configured_dockerhub_discovery_policies({ 'queries': ['first'], 'pages': 30, 'per_page': 100, }) state = console_runner.default_source_state() policy = policies['first']['policy_sha256'] console_runner.mark_dockerhub_deep_dispatched(state, 'first', policy, now) self.assertFalse(console_runner.prepare_dockerhub_discovery_state( state, policies, 'first', now + timedelta(hours=1), )['deep']) console_runner.mark_dockerhub_discovery_incomplete(state, 'first', policy) self.assertTrue(console_runner.prepare_dockerhub_discovery_state( state, policies, 'first', now + timedelta(hours=1), )['deep']) console_runner.clear_dockerhub_discovery_incomplete(state, 'first') self.assertFalse(console_runner.prepare_dockerhub_discovery_state( state, policies, 'first', now + timedelta(hours=1), )['deep']) class DockerHubIncrementalPageTests(unittest.TestCase): def test_first_two_preexisting_pages_stop_before_page_three(self): db = PageDB({'owner/one', 'owner/two'}) requested = [] def fetch(_query, requested_page, **_kwargs): requested.append(requested_page) return page(requested_page, [f'owner/{"one" if requested_page == 1 else "two"}']) with mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=fetch, ): metrics = console_runner.run_dockerhub_incremental_discovery( discovery_args(), db, ) self.assertEqual(requested, [1, 2]) self.assertTrue(metrics['discovery_stopped_on_preexisting']) self.assertEqual(metrics['discovery_pages_fetched'], 2) def test_current_pass_duplicate_prevents_known_page_stop(self): db = PageDB({'owner/known'}) responses = { 1: page(1, ['owner/new']), 2: page(2, ['owner/new']), 3: page(3, ['owner/known']), } with mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=lambda _query, number, **_kwargs: responses[number], ) as fetch: metrics = console_runner.run_dockerhub_incremental_discovery( discovery_args(), db, ) self.assertEqual([call.args[1] for call in fetch.call_args_list], [1, 2, 3]) self.assertFalse(metrics['discovery_stopped_on_preexisting']) self.assertEqual(metrics['queued_new_count'], 1) def test_deep_pass_bypasses_only_preexisting_stop(self): db = PageDB({'owner/one', 'owner/two', 'owner/three'}) with mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=lambda _query, number, **_kwargs: page( number, [f'owner/{("one", "two", "three")[number - 1]}'], ), ) as fetch: metrics = console_runner.run_dockerhub_incremental_discovery( discovery_args(deep=True), db, ) self.assertEqual([call.args[1] for call in fetch.call_args_list], [1, 2, 3]) self.assertTrue(metrics['deep_dispatch_durable']) self.assertEqual(metrics['discovery_pass_kind'], 'deep') def test_page_one_retry_is_durable_query_work(self): db = PageDB() failure = console_runner.DockerHubDiscoveryTransportError( 'safe failure', category='network', retryable=True, ) with mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=failure, ): metrics = console_runner.run_dockerhub_incremental_discovery( discovery_args(), db, ) self.assertEqual(metrics['cycle_status'], 'completed_with_retries') self.assertEqual(metrics['discovery_retry_enqueued_count'], 1) args, kwargs = db.retries[0] self.assertEqual(args[4], 'query') self.assertEqual((kwargs['page_start'], kwargs['page_end']), (1, 3)) def test_later_page_gap_keeps_admissions_and_continues_after_remote_attempt(self): db = PageDB() failure = console_runner.DockerHubDiscoveryTransportError( 'safe failure', category='network', remote_attempted=True, retryable=True, ) def fetch(_query, number, **_kwargs): if number == 2: raise failure return page(number, [f'owner/page-{number}']) with mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=fetch, ): metrics = console_runner.run_dockerhub_incremental_discovery( discovery_args(), db, ) self.assertEqual(metrics['cycle_status'], 'completed_with_retries') self.assertEqual([item[0] for item in db.admissions], [ ('owner/page-1',), ('owner/page-3',), ]) self.assertEqual(db.retries[0][0][4], 'page') self.assertEqual(db.retries[0][1]['page_start'], 2) def test_unattempted_tail_is_one_range_and_stops(self): db = PageDB() failure = console_runner.DockerHubDiscoveryTransportError( 'safe failure', category='provider_cooldown', retry_at='2026-01-01T00:00:00+00:00', remote_attempted=False, retryable=True, ) requested = [] def fetch(_query, number, **_kwargs): requested.append(number) if number == 2: raise failure return page(number, [f'owner/page-{number}'], total_count=500) with mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=fetch, ): console_runner.run_dockerhub_incremental_discovery( discovery_args(pages=5), db, ) self.assertEqual(requested, [1, 2]) self.assertEqual(db.retries[0][0][4], 'range') self.assertEqual( (db.retries[0][1]['page_start'], db.retries[0][1]['page_end']), (2, 5), ) def test_delegation_failure_and_invalid_payload_are_hard_failures(self): retryable = console_runner.DockerHubDiscoveryTransportError( 'safe failure', category='network', retryable=True, ) db = PageDB(enqueue_error=True) with mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=retryable, ), self.assertRaisesRegex(RuntimeError, 'delegation failed'): console_runner.run_dockerhub_incremental_discovery(discovery_args(), db) invalid = console_runner.DockerHubDiscoveryTransportError( 'safe invalid payload', category='invalid_payload', retryable=False, ) db = PageDB() with mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=invalid, ), self.assertRaises(console_runner.DockerHubDiscoveryTransportError): console_runner.run_dockerhub_incremental_discovery(discovery_args(), db) self.assertEqual(db.retries, []) def test_underreported_total_count_cannot_complete_a_later_page(self): db = PageDB() responses = { 1: page(1, ['owner/first'], total_count=200), 2: page(2, ['owner/second'], total_count=1), } with mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=lambda _query, number, **_kwargs: responses[number], ), self.assertRaises(console_runner.DockerHubDiscoveryTransportError): console_runner.run_dockerhub_incremental_discovery( discovery_args(pages=2), db, ) self.assertEqual([entry[0] for entry in db.admissions], [('owner/first',)]) self.assertEqual(db.retries, []) class RetryDB(PageDB): def __init__(self, claim): super().__init__() self.claim = claim self.claim_calls = [] self.finished = [] self.updated = [] self.held = [] self.renewed = [] def claim_discovery_retries(self, *args, **kwargs): self.claim_calls.append((args, kwargs)) return [dict(self.claim)] if self.claim else [] def finish_discovery_retry(self, *args, **kwargs): self.finished.append((args, kwargs)) return {'deleted_count': 1} def renew_discovery_retry_lease(self, *args, **kwargs): self.renewed.append((args, kwargs)) return {'status': 'leased'} def update_discovery_retry(self, *args, **kwargs): self.updated.append((args, kwargs)) return {'status': 'pending'} def hold_discovery_retry(self, *args, **kwargs): self.held.append((args, kwargs)) return {'status': 'held'} class DockerHubRetryLaneTests(unittest.TestCase): @staticmethod def policy(query='retry-query', pages=3): args = discovery_args(query=query, pages=pages) return args, {query: console_runner.dockerhub_discovery_policy(args)} @staticmethod def claim( policy, *, work_kind='range', page_start=2, page_end=3, next_page=2, source_cycle_id=None, ): return { 'id': 7, 'source': 'dockerhub', 'query': 'retry-query', 'policy_sha256': policy['retry-query']['policy_sha256'], 'pass_kind': 'ordinary', 'work_kind': work_kind, 'page_start': page_start, 'page_end': page_end, 'next_page': next_page, 'lease_owner': 'owner', 'lease_token': 'token', 'source_cycle_id': source_cycle_id, } def test_retry_success_persists_progress_then_completes_atomically(self): args, policies = self.policy() db = RetryDB(self.claim(policies)) with mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=lambda _query, number, **_kwargs: page( number, [f'owner/retry-{number}'], ), ): metrics = console_runner.process_dockerhub_discovery_retry( args, db, 'dockerhub', policies, ) self.assertEqual(len(db.claim_calls), 1) self.assertEqual(db.claim_calls[0][1], {'limit': 1, 'lease_seconds': 300}) self.assertEqual(len(db.renewed), 2) self.assertEqual(len(db.admissions), 2) self.assertEqual(db.admissions[0][1], { 'retry_id': 7, 'lease_owner': 'owner', 'lease_token': 'token', 'next_page': 3, 'complete': False, }) self.assertEqual(db.admissions[1][1], { 'retry_id': 7, 'lease_owner': 'owner', 'lease_token': 'token', 'next_page': None, 'complete': True, }) self.assertEqual(db.finished, []) self.assertEqual(metrics['discovery_retry_completed_count'], 1) def test_retry_uses_claim_origin_instead_of_unrelated_current_cycle(self): args, policies = self.policy() args.query = 'current-query' args.dockerhub_discovery_ordered_queries = ('current-query', 'retry-query') args.dockerhub_discovery_ordered_query_hash = 'd' * 64 args.dockerhub_discovery_query_count = 2 args.dockerhub_discovery_cycle_id = 91 db = RetryDB(self.claim( policies, work_kind='page', page_start=2, page_end=2, next_page=2, source_cycle_id=37, )) with mock.patch.object( console_runner, 'fetch_dockerhub_search_page', return_value=page(2, ['owner/origin']), ): console_runner.process_dockerhub_discovery_retry( args, db, 'dockerhub', policies, ) observation = db.admissions[0][1]['observation'] self.assertEqual(observation['cycle_id'], 37) self.assertNotEqual(observation['cycle_id'], args.dockerhub_discovery_cycle_id) self.assertEqual(observation['query_ordinal'], 1) self.assertEqual(db.admissions[0][1]['retry_id'], 7) def test_retry_failure_updates_with_refund_only_before_remote_attempt(self): args, policies = self.policy() db = RetryDB(self.claim( policies, work_kind='page', page_start=2, page_end=2, next_page=2, )) failure = console_runner.DockerHubDiscoveryTransportError( 'safe failure', category='provider_cooldown', retry_at='2026-01-01T00:00:00+00:00', remote_attempted=False, retryable=True, ) with mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=failure, ): metrics = console_runner.process_dockerhub_discovery_retry( args, db, 'dockerhub', policies, ) self.assertEqual(metrics['discovery_retry_deferred_count'], 1) self.assertEqual(db.updated[0][0][:4], (7, 'owner', 'token', 'provider_cooldown')) self.assertEqual(db.updated[0][1], { 'retry_at': '2026-01-01T00:00:00+00:00', 'refund_attempt': True, 'next_page': 2, }) def test_retry_does_not_refund_after_an_earlier_remote_page(self): args, policies = self.policy() db = RetryDB(self.claim(policies)) failure = console_runner.DockerHubDiscoveryTransportError( 'safe failure', category='provider_cooldown', retry_at='2026-01-01T00:00:00+00:00', remote_attempted=False, retryable=True, ) def fetch(_query, number, **_kwargs): if number == 3: raise failure return page(number, [f'owner/retry-{number}']) with mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=fetch, ): console_runner.process_dockerhub_discovery_retry( args, db, 'dockerhub', policies, ) self.assertEqual(db.updated[0][1], { 'retry_at': None, 'refund_attempt': False, 'next_page': 3, }) def test_pause_after_retry_claim_defers_without_provider_io_or_attempt_charge(self): args, policies = self.policy() db = RetryDB(self.claim( policies, work_kind='page', page_start=2, page_end=2, next_page=2, )) db.paused = True with mock.patch.object(console_runner, 'fetch_dockerhub_search_page') as fetch: metrics = console_runner.process_dockerhub_discovery_retry( args, db, 'dockerhub', policies, ) fetch.assert_not_called() self.assertEqual(metrics['discovery_retry_deferred_count'], 1) self.assertEqual(db.updated[0][0][:4], ( 7, 'owner', 'token', 'provider_cooldown', )) self.assertEqual(db.updated[0][1], { 'retry_at': None, 'refund_attempt': True, 'next_page': 2, }) class SourceDB: def __init__(self, events=None, state=None, claim_enabled=False): self.events = events if events is not None else [] self.conn = Connection(self.events) self.finished = [] self.state = state self.claim_enabled = claim_enabled def start_source_cycle(self, *_args, **_kwargs): return 7 def runtime_control_state(self): return {'effective_discovery_paused': False} def dockerhub_discovery_generation_complete(self, *_args, **_kwargs): return True def finish_source_cycle(self, *args, **kwargs): self.finished.append((args, kwargs)) def claim_discovery_retries(self, *_args, **_kwargs): if self.claim_enabled: self.events.append(('claim', self.state['sources']['dockerhub']['query_index'])) return [] def persist_dockerhub_discovery_page(self, *_args, **_kwargs): return {} def finish_discovery_retry(self, *_args, **_kwargs): return {} def renew_discovery_retry_lease(self, *_args, **_kwargs): return {} def update_discovery_retry(self, *_args, **_kwargs): return {} def hold_discovery_retry(self, *_args, **_kwargs): return {} class DockerHubConfiguredSourceTests(unittest.TestCase): @staticmethod def config(): return { 'global': {}, 'sources': {'dockerhub': { 'queries': ['first', 'second'], 'mode': 'search', 'pages': 3, 'per_page': 100, }}, } @staticmethod def args(): return discovery_args(query='first') def run_source(self, state, db, run_cycle, save_state=lambda *_args: None): with mock.patch.object(console_runner, 'save_state', side_effect=save_state), \ mock.patch.object(console_runner, 'select_auth_entry', return_value=None), \ mock.patch.object(console_runner, 'refresh_auth_summary'), \ mock.patch.object( console_runner, 'build_args_from_source_config', return_value=self.args(), ), \ mock.patch.object(console_runner, 'configure_source_auth'), \ mock.patch.object(console_runner, 'persist_docker_auth_events'), \ mock.patch.object(console_runner, 'queue_counts_for_args', return_value={}), \ mock.patch.object(console_runner, 'run_cycle', side_effect=run_cycle): return console_runner.run_configured_source( 'dockerhub', self.config(), state, 'state.json', {}, db, run_id=3, ) def test_completed_with_retries_advances_and_records_durable_deep_dispatch(self): state = {'sources': {'dockerhub': console_runner.default_source_state()}} metrics = self.run_source(state, SourceDB(), lambda *_args, **_kwargs: { 'cycle_status': 'completed_with_retries', 'deep_dispatch_durable': True, 'source_failure_count': 0, 'scanned_count': 0, }) source_state = state['sources']['dockerhub'] self.assertEqual(metrics['cycle_status'], 'completed_with_retries') self.assertEqual(source_state['last_status'], 'completed_with_retries') self.assertEqual(source_state['query_index'], 1) self.assertIn( 'first', source_state[console_runner.DOCKERHUB_DISCOVERY_STATE_KEY]['deep_by_query'], ) self.assertNotIn( 'first', source_state[console_runner.DOCKERHUB_DISCOVERY_STATE_KEY]['incomplete_by_query'], ) def test_returned_status_controls_cursor_advancement(self): for status, expected_index in ( ('failed', 0), ('source_failed', 0), ('backlog_only', 0), ('query_invalid', 1), ): with self.subTest(status=status): state = { 'sources': {'dockerhub': console_runner.default_source_state()}, } metrics = self.run_source( state, SourceDB(), lambda *_args, **_kwargs: { 'cycle_status': status, 'source_failure_count': 0, 'scanned_count': 0, }, ) self.assertEqual(metrics['cycle_status'], status) self.assertEqual(state['sources']['dockerhub']['last_status'], status) self.assertEqual( state['sources']['dockerhub']['query_index'], expected_index, ) def test_delegation_failure_retains_cursor(self): state = {'sources': {'dockerhub': console_runner.default_source_state()}} db = SourceDB() db.enqueue_discovery_retry = mock.Mock(side_effect=RuntimeError('raw detail')) failure = console_runner.DockerHubDiscoveryTransportError( 'safe failure', category='network', retryable=True, ) def run(args, *_args, **_kwargs): with mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=failure, ): return console_runner.run_dockerhub_incremental_discovery(args, db) with self.assertRaisesRegex(RuntimeError, 'delegation failed'): self.run_source(state, db, run) self.assertEqual(state['sources']['dockerhub']['query_index'], 0) source_state = state['sources']['dockerhub'] self.assertIn( 'first', source_state[console_runner.DOCKERHUB_DISCOVERY_STATE_KEY]['incomplete_by_query'], ) policies = console_runner.configured_dockerhub_discovery_policies( self.config()['sources']['dockerhub'], ) self.assertTrue(console_runner.prepare_dockerhub_discovery_state( source_state, policies, 'first', )['deep']) def test_auth_retry_reapplies_discovery_annotations(self): state = {'sources': {'dockerhub': console_runner.default_source_state()}} db = SourceDB() built_args = [self.args(), self.args()] seen = [] experiment = SimpleNamespace( enabled=True, queries=('first', 'second'), ordered_query_hash='e' * 64, collection_generation='fixture-generation', ) authority = {'enabled': True, 'selector_version': 'fixture-selector'} def run(args, *_args, **_kwargs): seen.append(( args.dockerhub_discovery_deep, args.dockerhub_discovery_pass_kind, args.dockerhub_discovery_policy_sha256, args.docker_depth_collection_only, args.dockerhub_discovery_ordered_queries, args.dockerhub_discovery_ordered_query_hash, args.dockerhub_discovery_query_count, args.dockerhub_discovery_collection_generation, args.dockerhub_discovery_cycle_id, args.docker_depth_experiment_authority, )) if len(seen) == 1: raise console_runner.RateLimitError( 'dockerhub', 'safe rate limit', category='rate_limit', auth_related=True, ) return { 'cycle_status': 'completed', 'deep_dispatch_durable': True, 'source_failure_count': 0, 'scanned_count': 0, } with mock.patch.object(console_runner, 'save_state'), \ mock.patch.object( console_runner, 'select_auth_entry', side_effect=[{'name': 'one'}, {'name': 'two'}], ), \ mock.patch.object(console_runner, 'refresh_auth_summary'), \ mock.patch.object( console_runner, 'build_args_from_source_config', side_effect=built_args, ), \ mock.patch.object(console_runner, 'configure_source_auth'), \ mock.patch.object(console_runner, 'persist_docker_auth_events'), \ mock.patch.object(console_runner, 'queue_counts_for_args', return_value={}), \ mock.patch.object( console_runner, 'validate_docker_depth_config', return_value=SimpleNamespace(experiment=experiment), ), \ mock.patch.object( console_runner, 'docker_depth_resolver_authority', return_value=authority, ), \ mock.patch.object(console_runner, 'run_cycle', side_effect=run): console_runner.run_configured_source( 'dockerhub', self.config(), state, 'state.json', {}, db, run_id=3, ) self.assertEqual(len(seen), 2) self.assertEqual(seen[0], seen[1]) self.assertEqual(seen[1][1], 'deep') self.assertFalse(seen[1][3]) self.assertEqual(seen[1][4], ('first', 'second')) self.assertEqual(seen[1][5], 'e' * 64) self.assertEqual(seen[1][6], 2) self.assertEqual(seen[1][7], 'fixture-generation') self.assertEqual(seen[1][8], 7) self.assertEqual(seen[1][9], authority) def test_main_state_is_saved_before_one_retry_claim(self): state = {'sources': {'dockerhub': console_runner.default_source_state()}} events = [] db = SourceDB(events, state, claim_enabled=True) def save(_path, saved_state): events.append(('save', saved_state['sources']['dockerhub']['query_index'])) self.run_source(state, db, lambda *_args, **_kwargs: { 'cycle_status': 'completed', 'source_failure_count': 0, 'scanned_count': 0, }, save_state=save) claims = [index for index, event in enumerate(events) if event[0] == 'claim'] self.assertEqual(len(claims), 1) self.assertEqual(events[claims[0] - 1], ('save', 1)) class DockerHubResolverGateTests(unittest.TestCase): class DB(PageDB): url = 'postgresql://fixture' last_error = '' def __init__(self): super().__init__({'owner/one', 'owner/two'}) self.finished = [] def require_runtime_safety_schema(self): return True def require_final_cutover(self): return True def has_claimable_targets_v2(self, *_args, **_kwargs): raise AssertionError('discovery role must not inspect scan backlog') def enqueue_targets(self, _source, _platform, _query, targets, **kwargs): return len(targets) + len(kwargs.get('unresolved_targets') or []) def target_queue_counts(self, _source): return {} def claim_docker_resolutions(self, *_args, **_kwargs): return [] def finish_source_cycle(self, *args, **kwargs): self.finished.append((args, kwargs)) class Lease: def __init__(self): self.released = False def release(self): self.released = True def test_resolver_gate_runs_once_after_incremental_pagination(self): args = discovery_args() args.__dict__.update({ 'sync_file_queues': False, 'docker_depth_collection_only': False, }) db = self.DB() order = [] with mock.patch.object( console_runner, 'resolve_due_docker_experiment_targets', side_effect=lambda *_args, **_kwargs: (order.append('depth') or (0, False)), ) as experiment_resolver, \ mock.patch.object( console_runner, 'resolve_due_docker_queue_targets', side_effect=lambda *_args, **_kwargs: (order.append('ordinary') or 0), ) as resolver, \ mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=lambda _query, number, **_kwargs: ( order.append(f'page-{number}') or page(number, [f'owner/{"one" if number == 1 else "two"}']) ), ) as fetch: metrics = console_runner.run_discovery_cycle(args, db, 1, 2, 'dockerhub') self.assertEqual([call.args[1] for call in fetch.call_args_list], [1, 2]) experiment_resolver.assert_called_once() resolver.assert_called_once() self.assertEqual(order, ['depth', 'page-1', 'page-2', 'ordinary']) self.assertEqual(metrics['cycle_status'], 'completed') self.assertEqual(metrics['scan_requested_count'], 0) self.assertEqual(db.finished[0][0][1], 'completed') def test_disabled_experiment_collection_skips_resolver_and_scan_admission(self): args = discovery_args() args.__dict__.update({ 'sync_file_queues': False, 'docker_depth_collection_only': True, }) db = self.DB() with mock.patch.object( console_runner, 'resolve_due_docker_experiment_targets', ) as experiment_resolver, \ mock.patch.object( console_runner, 'resolve_due_docker_queue_targets', ) as ordinary_resolver, \ mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=lambda _query, number, **_kwargs: page( number, [f'owner/{"one" if number == 1 else "two"}'], ), ): metrics = console_runner.run_discovery_cycle(args, db, 1, 2, 'dockerhub') experiment_resolver.assert_not_called() ordinary_resolver.assert_not_called() self.assertTrue(metrics['docker_depth_collection_only']) self.assertEqual(metrics['scan_requested_count'], 0) self.assertEqual(metrics['cycle_status'], 'completed') self.assertEqual(db.finished[0][0][1], 'completed') def test_midpass_pause_preserves_committed_page_metrics(self): args = discovery_args(pages=2) args.__dict__.update({ 'sync_file_queues': False, 'docker_depth_collection_only': False, }) db = self.DB() original_admit = db.persist_dockerhub_discovery_page def admit(*admit_args, **admit_kwargs): report = original_admit(*admit_args, **admit_kwargs) db.paused = True return report db.persist_dockerhub_discovery_page = admit with mock.patch.object( console_runner, 'resolve_due_docker_experiment_targets', return_value=(0, False), ), mock.patch.object( console_runner, 'resolve_due_docker_queue_targets', ) as resolver, mock.patch.object( console_runner, 'fetch_dockerhub_search_page', side_effect=lambda _query, number, **_kwargs: page( number, [f'owner/page-{number}'], total_count=200, ), ) as fetch: metrics = console_runner.run_discovery_cycle( args, db, 1, 2, 'dockerhub', ) self.assertEqual(fetch.call_count, 1) resolver.assert_not_called() self.assertEqual(metrics['cycle_status'], 'paused') self.assertEqual(metrics['fetched_count'], 1) self.assertEqual(metrics['queued_new_count'], 1) self.assertEqual(metrics['discovery_pages_fetched'], 1) self.assertEqual(len(db.finished), 1) self.assertEqual(db.finished[0][0][1], 'paused') def test_pause_after_ordinary_resolver_claim_refunds_without_provider_io(self): args = discovery_args() db = self.DB() db.paused = True row = {'id': 9, 'target': 'owner/repository', 'resolver_token': 'token'} db.claim_docker_resolutions = mock.Mock(return_value=[row]) db.finish_docker_resolution = mock.Mock(return_value=True) with mock.patch.object(console_runner, 'fetch_dockerhub_tags') as fetch, \ self.assertRaises(console_runner.DiscoveryPausedError): console_runner.resolve_due_docker_queue_targets(db, 'dockerhub', args) fetch.assert_not_called() self.assertFalse(db.finish_docker_resolution.call_args.kwargs["claim_attempt_consumed"]) def test_pause_after_depth_resolver_claim_refunds_without_provider_io(self): args = discovery_args() args.sync_file_queues = False args.docker_depth_experiment_authority = { 'enabled': True, 'selector_version': 'fixture-selector', } db = self.DB() db.paused = True row = { 'id': 10, 'target': 'owner/repository', 'selection_limit': 1, 'resolver_generation': 3, 'resolver_token': 'token', 'resolver_owner': 'owner', } db.claim_docker_depth_experiment_resolutions = mock.Mock(return_value=[row]) db.renew_docker_depth_experiment_resolution = mock.Mock() db.finish_docker_depth_experiment_resolution = mock.Mock(return_value={ 'committed': True, 'status': 'deferred', }) with mock.patch.object(console_runner, 'fetch_dockerhub_tags') as fetch, \ self.assertRaises(console_runner.DiscoveryPausedError): console_runner.resolve_due_docker_experiment_targets( db, 'dockerhub', args, ) fetch.assert_not_called() db.renew_docker_depth_experiment_resolution.assert_not_called() self.assertFalse( db.finish_docker_depth_experiment_resolution.call_args.kwargs[ 'claim_attempt_consumed' ] ) if __name__ == '__main__': unittest.main()