Files
truf-server/tests/test_dockerhub_search_pagination.py
2026-09-30 20:30:56 +03:00

750 lines
29 KiB
Python

from datetime import datetime
from pathlib import Path
import concurrent.futures
import copy
import sys
import threading
import time
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
import scanner
class Response:
def __init__(self, status_code=200, payload=None, headers=None):
self.status_code = status_code
self.payload = payload if payload is not None else {}
self.headers = headers or {}
self.text = ''
self.content = b''
def raise_for_status(self):
if self.status_code >= 400:
raise RuntimeError(f'HTTP {self.status_code}')
def json(self):
return self.payload
def close(self):
return None
def manager_with_accounts(count=1):
manager = scanner.DockerTokenManager()
manager.accounts = [
scanner.DockerAccount(
f'account-{index}', f'user-{index}', f'fixture-secret-{index}', '',
)
for index in range(count)
]
manager.explicit_pool = True
return manager
class DockerHubSearchAuthTests(unittest.TestCase):
def test_search_uses_hub_bearer_and_bounded_request_attempts(self):
manager = manager_with_accounts()
responses = [
Response(payload={'access_token': 'fixture-bearer', 'expires_in': 600}),
Response(payload={'count': 0, 'results': []}),
]
with mock.patch.object(scanner, 'docker_token_manager', manager), \
mock.patch.object(scanner, 'api_request', side_effect=responses) as request, \
mock.patch.object(
manager, 'report_success', wraps=manager.report_success,
) as report_success:
response = scanner.dockerhub_search_response(
'https://hub.docker.com/v2/search/repositories',
{'query': 'fixture', 'page': 3, 'page_size': 100},
)
self.assertEqual(response.status_code, 200)
self.assertEqual(request.call_count, 2)
search_call = request.call_args_list[1]
self.assertEqual(
search_call.kwargs['headers']['Authorization'],
'Bearer fixture-bearer',
)
self.assertEqual(search_call.kwargs['max_retries'], 1)
self.assertTrue(any(
call.args[1] == 'hub_search'
for call in report_success.call_args_list
))
self.assertNotIn(('account-0', 'hub_search'), manager.cooldown_until)
self.assertNotIn(('account-0', 'hub_tags'), manager.cooldown_until)
def test_explicit_empty_pool_fails_without_anonymous_request(self):
manager = scanner.DockerTokenManager()
manager.explicit_pool = True
with mock.patch.object(scanner, 'docker_token_manager', manager), \
mock.patch.object(scanner, 'api_request') as request:
with self.assertRaisesRegex(
scanner.DockerRemoteAccessError, 'authentication is unavailable',
):
scanner.dockerhub_search_response(
'https://hub.docker.com/v2/search/repositories',
{'query': 'fixture', 'page': 1, 'page_size': 100},
)
request.assert_not_called()
def test_search_refreshes_once_after_401(self):
manager = manager_with_accounts()
responses = [
Response(payload={'access_token': 'old-bearer'}),
Response(status_code=401),
Response(payload={'access_token': 'new-bearer'}),
Response(payload={'count': 0, 'results': []}),
]
with mock.patch.object(scanner, 'docker_token_manager', manager), \
mock.patch.object(scanner, 'api_request', side_effect=responses) as request:
scanner.dockerhub_search_response(
'https://hub.docker.com/v2/search/repositories',
{'query': 'fixture', 'page': 1, 'page_size': 100},
)
self.assertEqual(request.call_count, 4)
self.assertEqual(
request.call_args_list[1].kwargs['headers']['Authorization'],
'Bearer old-bearer',
)
self.assertEqual(
request.call_args_list[3].kwargs['headers']['Authorization'],
'Bearer new-bearer',
)
def test_search_rotates_and_cools_account_after_403_or_429(self):
for status_code, expected_category in ((403, 'auth_forbidden'), (429, 'rate_limit')):
with self.subTest(status_code=status_code):
manager = manager_with_accounts(2)
responses = [
Response(payload={'access_token': 'bearer-a'}),
Response(status_code=status_code),
Response(payload={'access_token': 'bearer-b'}),
Response(payload={'count': 0, 'results': []}),
]
with mock.patch.object(scanner, 'docker_token_manager', manager), \
mock.patch.object(scanner, 'api_request', side_effect=responses) as request:
scanner.dockerhub_search_response(
'https://hub.docker.com/v2/search/repositories',
{'query': 'fixture', 'page': 1, 'page_size': 100},
)
self.assertEqual(request.call_count, 4)
self.assertEqual(
manager.cooldown_categories[('account-0', 'hub_search')],
expected_category,
)
self.assertEqual(
request.call_args_list[3].kwargs['headers']['Authorization'],
'Bearer bearer-b',
)
def test_transient_page_failure_gets_exactly_two_attempts(self):
manager = manager_with_accounts()
manager.cache_hub_token('account-0', 'fixture-bearer', 600)
timeout = scanner.requests.exceptions.ReadTimeout('fixture timeout')
response = Response(payload={'count': 0, 'results': []})
with mock.patch.object(scanner, 'docker_token_manager', manager), \
mock.patch.object(scanner.scan_config, 'api_proxy_enabled', False), \
mock.patch.object(
scanner.requests, 'request', side_effect=[timeout, response],
) as request, \
mock.patch.object(scanner, '_wait_or_raise_scan_slot_fatal'):
returned = scanner.dockerhub_search_response(
'https://hub.docker.com/v2/search/repositories',
{'query': 'fixture', 'page': 1, 'page_size': 100},
)
self.assertIs(returned, response)
self.assertEqual(request.call_count, 2)
def test_page_attempt_budget_does_not_reset_after_account_rotation(self):
manager = manager_with_accounts(3)
for index in range(3):
manager.cache_hub_token(f'account-{index}', f'bearer-{index}', 600)
with mock.patch.object(scanner, 'docker_token_manager', manager), \
mock.patch.object(
scanner, 'api_request',
side_effect=[
Response(status_code=403),
scanner.ApiRequestError('fixture transient failure'),
],
) as request:
with self.assertRaises(scanner.ApiRequestError):
scanner.dockerhub_search_response(
'https://hub.docker.com/v2/search/repositories',
{'query': 'fixture', 'page': 1, 'page_size': 100},
)
self.assertEqual(request.call_count, 2)
self.assertEqual(
request.call_args_list[1].kwargs['headers']['Authorization'],
'Bearer bearer-1',
)
def test_search_pool_exhaustion_does_not_set_tag_backoff(self):
manager = manager_with_accounts(2)
for index in range(2):
manager.cache_hub_token(f'account-{index}', f'bearer-{index}', 600)
with mock.patch.object(scanner, 'docker_token_manager', manager), \
mock.patch.object(
scanner, 'api_request',
side_effect=[Response(status_code=429), Response(status_code=429)],
), \
mock.patch.object(scanner, 'put_dockerhub_tags_rate_limit') as tag_backoff:
with self.assertRaisesRegex(
scanner.DockerRemoteAccessError, 'accounts are rate-limited',
):
scanner.dockerhub_search_response(
'https://hub.docker.com/v2/search/repositories',
{'query': 'fixture', 'page': 1, 'page_size': 100},
)
tag_backoff.assert_not_called()
self.assertTrue(manager.all_unavailable('hub_search'))
self.assertFalse(manager.all_unavailable('hub_tags'))
def test_concurrent_pages_share_token_acquisition(self):
manager = manager_with_accounts()
def token_response(*_args, **_kwargs):
time.sleep(0.02)
return Response(payload={'access_token': 'shared-bearer'})
with mock.patch.object(scanner, 'docker_token_manager', manager), \
mock.patch.object(
scanner, 'api_request', side_effect=token_response,
) as request, \
concurrent.futures.ThreadPoolExecutor(max_workers=2) as executor:
futures = [
executor.submit(
scanner._docker_hub_access_token,
manager.accounts[0],
False,
'hub_search',
)
for _ in range(2)
]
tokens = [future.result() for future in futures]
self.assertEqual(tokens, ['shared-bearer', 'shared-bearer'])
self.assertEqual(request.call_count, 1)
def test_search_event_persists_without_disabling_tag_account(self):
source_config = {
'auth_pool': 'dockerhub_main', 'rate_limit_cooldown': 300,
}
secrets = {'auth_pools': {'dockerhub_main': [{
'name': 'account-0', 'username': 'user-0', 'token': 'fixture-secret',
}]}}
state = {'sources': {'dockerhub': console_runner.default_source_state()}}
events = [{
'name': 'account-0',
'endpoint': 'hub_search',
'category': 'rate_limit',
'reset_at': '2099-01-01T00:00:00+00:00',
'message': 'Docker hub_search HTTP 429',
}]
with mock.patch.object(
console_runner, 'drain_docker_auth_events', return_value=events,
):
console_runner.persist_docker_auth_events(
'dockerhub', source_config, state, secrets,
)
source_state = state['sources']['dockerhub']
self.assertEqual(source_state['auth_status'].get('account-0'), {})
self.assertTrue(console_runner.auth_entry_is_available(
'dockerhub', secrets['auth_pools']['dockerhub_main'][0], state,
))
endpoint_state = source_state['auth_endpoint_status']['hub_search']['account-0']
self.assertEqual(endpoint_state['disabled_reason'], 'rate_limit')
self.assertNotIn('fixture-secret', repr(state))
manager = manager_with_accounts()
manager.restore_endpoint_cooldowns(source_state['auth_endpoint_status'])
self.assertTrue(manager.all_unavailable('hub_search'))
self.assertFalse(manager.all_unavailable('hub_tags'))
def test_auth_invalid_remains_global_after_late_endpoint_events(self):
manager = manager_with_accounts()
account = manager.accounts[0]
manager.report_http_status(
account, 'hub_search', 401, Response(status_code=401), 'auth_invalid',
)
manager.report_http_status(
account, 'hub_tags', 429, Response(status_code=429), 'rate_limit',
)
manager.report_success(account, 'registry')
self.assertTrue(manager.all_unavailable('hub_search'))
self.assertTrue(manager.all_unavailable('hub_tags'))
self.assertTrue(manager.all_unavailable('registry'))
self.assertEqual(
manager.cooldown_categories[('account-0', 'hub_search')],
'auth_invalid',
)
self.assertEqual(
[(event['endpoint'], event['category']) for event in manager.drain_status_events()],
[('hub_search', 'auth_invalid')],
)
def test_auth_invalid_account_is_excluded_from_cli_config_rotation(self):
manager = scanner.DockerTokenManager()
manager.accounts = [
scanner.DockerAccount('account-0', 'user-0', 'secret-0', 'config-0'),
scanner.DockerAccount('account-1', 'user-1', 'secret-1', 'config-1'),
]
manager.report_http_status(
manager.accounts[0], 'hub_search', 401,
Response(status_code=401), 'auth_invalid',
)
self.assertEqual(manager.get_next_config(), 'config-1')
self.assertEqual(manager.get_next_config(), 'config-1')
manager.report_http_status(
manager.accounts[1], 'registry', 401,
Response(status_code=401), 'auth_invalid',
)
self.assertIsNone(manager.get_next_config())
def test_persisted_auth_invalid_cannot_be_downgraded(self):
source_config = {
'auth_pool': 'dockerhub_main', 'rate_limit_cooldown': 300,
}
secrets = {'auth_pools': {'dockerhub_main': [{
'name': 'account-0', 'username': 'user-0', 'token': 'fixture-secret',
}]}}
state = {'sources': {'dockerhub': console_runner.default_source_state()}}
events = [
{
'name': 'account-0', 'endpoint': 'hub_search',
'category': 'auth_invalid', 'reset_at': 'manual',
'message': 'Docker hub_search HTTP 401',
},
{
'name': 'account-0', 'endpoint': 'hub_tags',
'category': 'rate_limit',
'reset_at': '2099-01-01T00:00:00+00:00',
'message': 'Docker hub_tags HTTP 429',
},
]
with mock.patch.object(
console_runner, 'drain_docker_auth_events', return_value=events,
):
console_runner.persist_docker_auth_events(
'dockerhub', source_config, state, secrets,
)
status = state['sources']['dockerhub']['auth_status']['account-0']
self.assertEqual(status['status'], 'dead')
self.assertEqual(status['disabled_until'], 'manual')
self.assertEqual(status['disabled_reason'], 'auth_invalid')
self.assertNotIn('fixture-secret', repr(state))
class DockerHubSearchPageTests(unittest.TestCase):
def test_page_helper_returns_validated_structured_data(self):
response = Response(payload={
'count': '26',
'results': [{
'repo_name': 'synthetic/repository-a',
'last_updated': '2026-01-02T03:04:05Z',
'url': 'raw-url-marker',
'token': 'raw-token-marker',
}],
})
with mock.patch.object(
scanner, 'dockerhub_search_response', return_value=response,
) as search, mock.patch.object(
scanner.time, 'perf_counter', side_effect=[10.0, 10.25],
):
result = scanner.fetch_dockerhub_search_page(
'synthetic-query', 2, per_page=25, sort_by='name',
sort_order='asc', request_timeout=9,
)
self.assertEqual(result, {
'page': 2,
'repositories': [{
'repo_name': 'synthetic/repository-a',
'last_updated': '2026-01-02T03:04:05Z',
}],
'total_count': 26,
'elapsed': 0.25,
})
search.assert_called_once_with(
'https://hub.docker.com/v2/search/repositories',
{
'query': 'synthetic-query',
'page': 2,
'page_size': 25,
'sort': 'name',
'order': 'asc',
},
request_timeout=9,
)
self.assertNotIn('raw-url-marker', repr(result))
self.assertNotIn('raw-token-marker', repr(result))
def test_page_helper_rejects_invalid_counts(self):
for count in (None, -1, True, 1.5, float('inf'), 'not-a-count'):
with self.subTest(count=count), mock.patch.object(
scanner, 'dockerhub_search_response',
return_value=Response(payload={'count': count, 'results': []}),
):
with self.assertRaisesRegex(
scanner.DockerHubDiscoveryTransportError,
'invalid result count',
):
scanner.fetch_dockerhub_search_page('synthetic-query', 1)
def test_page_helper_rejects_count_below_absolute_result_bound(self):
response = Response(payload={
'count': 25,
'results': [{'repo_name': 'synthetic/repository-a'}],
})
with mock.patch.object(
scanner, 'dockerhub_search_response', return_value=response,
), self.assertRaisesRegex(
scanner.DockerHubDiscoveryTransportError,
'incoherent pagination evidence',
):
scanner.fetch_dockerhub_search_page(
'synthetic-query', 2, per_page=25,
)
def test_page_helper_rejects_malformed_payloads(self):
payloads = (
[],
{'count': 1},
{'count': 1, 'results': {}},
{'count': 1, 'results': [None]},
{'count': 1, 'results': [{}]},
{'count': 1, 'results': [{'repo_name': ''}]},
{'count': 1, 'results': [{'repo_name': ' '}]},
{'count': 1, 'results': [{'repo_name': 42}]},
)
for payload in payloads:
with self.subTest(payload=payload), mock.patch.object(
scanner, 'dockerhub_search_response',
return_value=Response(payload=payload),
):
with self.assertRaisesRegex(
scanner.DockerHubDiscoveryTransportError,
'invalid payload',
):
scanner.fetch_dockerhub_search_page('synthetic-query', 1)
def test_page_helper_deduplicates_names_in_first_seen_order(self):
response = Response(payload={
'count': 3,
'results': [
{
'repo_name': 'synthetic/repository-a',
'last_updated': '2026-01-03T00:00:00Z',
},
{
'repo_name': 'synthetic/repository-b',
'last_modified': '2026-01-02T00:00:00Z',
},
{
'repo_name': 'synthetic/repository-a',
'last_updated': '2026-01-01T00:00:00Z',
},
],
})
with mock.patch.object(
scanner, 'dockerhub_search_response', return_value=response,
):
result = scanner.fetch_dockerhub_search_page('synthetic-query', 1)
self.assertEqual(result['repositories'], [
{
'repo_name': 'synthetic/repository-a',
'last_updated': '2026-01-03T00:00:00Z',
},
{
'repo_name': 'synthetic/repository-b',
'last_modified': '2026-01-02T00:00:00Z',
},
])
def test_page_helper_preserves_explicit_pool_fail_closed(self):
manager = scanner.DockerTokenManager()
manager.explicit_pool = True
with mock.patch.object(scanner, 'docker_token_manager', manager), \
mock.patch.object(scanner, 'api_request') as request:
with self.assertRaisesRegex(
scanner.DockerHubDiscoveryTransportError,
'failed after bounded attempts',
):
scanner.fetch_dockerhub_search_page('synthetic-query', 1)
request.assert_not_called()
def test_page_helper_preserves_two_get_authenticated_budget(self):
manager = manager_with_accounts()
manager.cache_hub_token('account-0', 'synthetic-bearer', 600)
with mock.patch.object(scanner, 'docker_token_manager', manager), \
mock.patch.object(
scanner, 'api_request',
side_effect=[
scanner.ApiRequestError('raw-url-marker'),
scanner.ApiRequestError('raw-token-marker'),
],
) as request, \
mock.patch.object(scanner, '_wait_or_raise_scan_slot_fatal'):
with self.assertRaises(
scanner.DockerHubDiscoveryTransportError,
) as raised:
scanner.fetch_dockerhub_search_page('synthetic-query', 1)
self.assertEqual(request.call_count, 2)
self.assertTrue(all(
call.kwargs['max_retries'] == 1
for call in request.call_args_list
))
self.assertNotIn('raw-url-marker', str(raised.exception))
self.assertNotIn('raw-token-marker', str(raised.exception))
class DockerHubSearchPaginationTests(unittest.TestCase):
@staticmethod
def response(page, count):
return Response(payload={
'count': count,
'results': [{'repo_name': f'owner/repo-{page}'}],
})
def test_search_caps_at_thirty_pages_and_preserves_page_order(self):
requested_pages = []
lock = threading.Lock()
def search(_url, params, request_timeout=15):
page = int(params['page'])
with lock:
requested_pages.append(page)
if page == 2:
time.sleep(0.02)
return self.response(page, 4000)
with mock.patch.object(scanner, 'dockerhub_search_response', side_effect=search):
repos = scanner.fetch_dockerhub_images(
'fixture', pages=40, per_page=100,
fetch_workers=8, resolve_tags=False,
)
self.assertEqual(requested_pages[0], 1)
self.assertEqual(sorted(requested_pages), list(range(1, 31)))
self.assertEqual(
repos,
[f'owner/repo-{page}' for page in range(1, 31)],
)
def test_page_one_count_limits_expected_request_range(self):
requested_pages = []
def search(_url, params, request_timeout=15):
page = int(params['page'])
requested_pages.append(page)
return self.response(page, 250)
with mock.patch.object(scanner, 'dockerhub_search_response', side_effect=search):
repos = scanner.fetch_dockerhub_images(
'fixture', pages=30, per_page=100, resolve_tags=False,
)
self.assertEqual(sorted(requested_pages), [1, 2, 3])
self.assertEqual(repos, [
'owner/repo-1', 'owner/repo-2', 'owner/repo-3',
])
def test_invalid_page_one_count_fails_closed(self):
for payload in (
{'results': []},
{'count': -1, 'results': []},
):
with self.subTest(payload=payload), mock.patch.object(
scanner, 'dockerhub_search_response',
return_value=Response(payload=payload),
):
with self.assertRaises(scanner.DockerHubDiscoveryTransportError):
scanner.fetch_dockerhub_images(
'fixture', pages=30, per_page=100, resolve_tags=False,
)
def test_expected_page_failure_rejects_all_partial_results(self):
def search(_url, params, request_timeout=15):
page = int(params['page'])
if page == 2:
raise scanner.ApiRequestError('bounded fixture failure')
return self.response(page, 300)
with mock.patch.object(scanner, 'dockerhub_search_response', side_effect=search), \
mock.patch.object(scanner, 'fetch_dockerhub_tags') as tags:
with self.assertRaisesRegex(
scanner.DockerHubDiscoveryTransportError,
'pagination incomplete',
):
scanner.fetch_dockerhub_images(
'fixture', pages=3, per_page=100, resolve_tags=True,
)
tags.assert_not_called()
def test_recent_search_uses_authenticated_response_path(self):
response = Response(payload={'count': 0, 'results': [], 'next': None})
with mock.patch.object(
scanner, 'dockerhub_search_response', return_value=response,
) as search:
repos = scanner.fetch_recent_dockerhub_images(
'fixture', datetime(2026, 1, 1), pages=30,
per_page=100, resolve_tags=False,
)
self.assertEqual(repos, [])
self.assertEqual(search.call_count, 1)
self.assertEqual(search.call_args.args[1]['page'], 1)
self.assertEqual(search.call_args.kwargs['request_timeout'], 30)
def test_recent_search_preserves_page_update_metadata(self):
response = Response(payload={
'count': 1,
'results': [{
'repo_name': 'synthetic/repository-a',
'last_updated': '2025-12-31T00:00:00Z',
}],
})
with mock.patch.object(
scanner, 'dockerhub_search_response', return_value=response,
), mock.patch.object(scanner, 'fetch_dockerhub_last_updated') as metadata:
repos = scanner.fetch_recent_dockerhub_images(
'synthetic-query', datetime(2026, 1, 1), pages=1,
per_page=100, resolve_tags=False,
)
self.assertEqual(repos, [])
metadata.assert_not_called()
class DockerHubSearchCycleTests(unittest.TestCase):
class Connection:
is_postgres = True
def __init__(self, events):
self.events = events
def rollback(self):
self.events.append('rollback')
class DB:
def __init__(self, events):
self.conn = DockerHubSearchCycleTests.Connection(events)
self.events = events
self.finished = []
def start_source_cycle(self, *args, **kwargs):
return 7
def finish_source_cycle(self, *args, **kwargs):
self.events.append('finish')
self.finished.append((args, kwargs))
@staticmethod
def state():
source_state = console_runner.default_source_state()
return {'sources': {'dockerhub': source_state}}
@staticmethod
def patches(args, result):
run_cycle = (
mock.patch.object(console_runner, 'run_cycle', side_effect=result)
if isinstance(result, BaseException)
else mock.patch.object(console_runner, 'run_cycle', return_value=result)
)
return (
mock.patch.object(console_runner, 'save_state'),
mock.patch.object(
console_runner, 'select_auth_entry',
return_value={'name': 'docker-1', 'token': 'fixture'},
),
mock.patch.object(console_runner, 'refresh_auth_summary'),
mock.patch.object(
console_runner, 'build_args_from_source_config',
return_value=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={}),
run_cycle,
)
def run_source(self, state, db, run_result):
args = SimpleNamespace(platform='docker', mode='search')
patches = self.patches(args, run_result)
with patches[0], patches[1], patches[2], patches[3], \
patches[4], patches[5], patches[6], patches[7]:
return console_runner.run_configured_source(
'dockerhub',
{'global': {}, 'sources': {'dockerhub': {
'queries': ['first', 'second'], 'mode': 'search',
}}},
state, 'state.json', {}, db, run_id=3,
)
def test_failed_pagination_finishes_cycle_without_advancing(self):
events = []
db = self.DB(events)
state = self.state()
metrics = self.run_source(
state, db,
scanner.DockerHubDiscoveryTransportError('bounded failure'),
)
self.assertEqual(events, ['rollback', 'finish'])
self.assertTrue(metrics['discovery_transport_failed'])
self.assertEqual(state['sources']['dockerhub']['query_index'], 0)
self.assertEqual(state['sources']['dockerhub']['last_status'], 'failed')
self.assertEqual(db.finished[0][0][1], 'failed')
def test_successful_cycle_keeps_existing_query_advance(self):
events = []
db = self.DB(events)
state = self.state()
metrics = self.run_source(
state, db, {'fetched_count': 0, 'scanned_count': 0},
)
self.assertEqual(metrics['fetched_count'], 0)
self.assertEqual(state['sources']['dockerhub']['query_index'], 1)
self.assertEqual(state['sources']['dockerhub']['last_status'], 'completed')
if __name__ == '__main__':
unittest.main()