from __future__ import annotations import asyncio import logging import stat from concurrent.futures import ThreadPoolExecutor from pathlib import Path from threading import Barrier from typing import Any from unittest import mock import pytest import yaml import theHarvester.lib.core as core_module from theHarvester.lib.core import ( CONFIG_DIRS, DATA_DIR, AsyncFetcher, Core, FetcherResponse, ProxyUnavailableError, ResponseStreamError, ) from theHarvester.lib.output import configure_logging from theHarvester.lib.source_catalog import SOURCE_SPECS, ActivityClass, resolve_sources @pytest.fixture(autouse=True) def mock_environ(monkeypatch, tmp_path: Path): monkeypatch.setenv("HOME", str(tmp_path)) def test_email_capability_expands_to_email_sources() -> None: assert resolve_sources("emails") == [ "apis-guru", "baidu", "brave", "censys", "dehashed", "duckduckgo", "github-code", "gitlab", "hibpverified", "hudsonrock", "hunter", "intelx", "leaklookup", "mojeek", "rocketreach", "sherlockeye", "tomba", "windvane", "yahoo", "zoomeye", ] def test_capabilities_and_explicit_sources_form_a_union() -> None: assert resolve_sources("certspotter, urls") == [ "apis-guru", "bevigil", "builtwith", "certspotter", "gitlab", "intelx", "rocketreach", "urlscan", "zoomeye", ] def test_multiple_capabilities_form_a_union() -> None: assert resolve_sources("asns,people") == [ "criminalip", "onyphe", "urlscan", "zoomeye", ] def test_breach_capability_includes_every_matching_source() -> None: assert resolve_sources('breaches') == ['haveibeenpwned', 'hibpverified', 'leaklookup'] def test_named_source_can_be_combined_with_a_capability() -> None: assert resolve_sources('breaches,hibpverified') == ['haveibeenpwned', 'hibpverified', 'leaklookup'] def test_core_supported_engines_compatibility_uses_the_catalog() -> None: assert Core.get_supportedengines() == sorted(SOURCE_SPECS) def test_core_source_selection_compatibility_uses_the_catalog() -> None: assert Core.expand_source_selection('breaches') == resolve_sources('breaches') def test_all_selects_only_passive_catalog_sources() -> None: assert resolve_sources("ALL") == sorted( spec.name for spec in SOURCE_SPECS.values() if spec.activity is ActivityClass.PASSIVE ) assert { name: spec.activity for name, spec in SOURCE_SPECS.items() if spec.activity is not ActivityClass.PASSIVE } == { "criminalip": ActivityClass.DIRECT, "pentesttools": ActivityClass.DNS, "shodan": ActivityClass.DNS, "shodanInternetDB": ActivityClass.DNS, "subdomainfinderc99": ActivityClass.DNS, } @pytest.mark.parametrize( 'source', ['criminalip', 'pentesttools', 'shodan', 'shodanInternetDB', 'subdomainfinderc99'], ) def test_non_passive_sources_run_only_when_explicitly_selected(source: str) -> None: assert source not in resolve_sources('all') assert resolve_sources(source) == [source] def mock_read_text(mocked: dict[Path, str | Exception]): read_text = Path.read_text def _read_text(self: Path, *args, **kwargs): if result := mocked.get(self): if isinstance(result, Exception): raise result return result return read_text(self, *args, **kwargs) return _read_text @pytest.mark.parametrize( ("name", "contents", "expected"), [ ("api-keys", "apikeys: {}", {}), ("proxies", "http: [localhost:8080]", {"http": ["http://localhost:8080"], "socks5": []}), ], ) @pytest.mark.parametrize("dir", CONFIG_DIRS) def test_read_config_searches_config_dirs( name: str, contents: str, expected: Any, dir: Path, caplog ): caplog.set_level(logging.INFO, logger=core_module.__name__) file = dir.expanduser() / f"{name}.yaml" config_files = [d.expanduser() / file.name for d in CONFIG_DIRS] side_effect = mock_read_text( {f: contents if f == file else FileNotFoundError() for f in config_files} ) with mock.patch("pathlib.Path.read_text", autospec=True, side_effect=side_effect): got = Core.api_keys() if name == "api-keys" else Core.proxy_list() assert got == expected assert f"Read {file.name} from {file}" in caplog.messages @pytest.mark.parametrize("name", ("api-keys", "proxies")) def test_read_config_copies_default_to_home(name: str, capsys): configure_logging(verbose=False) file = Path(f"~/.theHarvester/{name}.yaml").expanduser() config_files = [d.expanduser() / file.name for d in CONFIG_DIRS] side_effect = mock_read_text({f: FileNotFoundError() for f in config_files}) with mock.patch("pathlib.Path.read_text", autospec=True, side_effect=side_effect): got = Core.api_keys() if name == "api-keys" else Core.proxy_list() default = yaml.safe_load((DATA_DIR / file.name).read_text()) expected = ( default["apikeys"] if name == "api-keys" else { "http": [f"http://{h}" for h in default["http"]] if default.get("http") else [], "socks5": [f"socks5://{h}" for h in default["socks5"]] if default.get("socks5") else [], } ) assert got == expected assert f"Created default {file.name} at {file}" in capsys.readouterr().out assert file.exists() assert stat.S_IMODE(file.parent.stat().st_mode) == 0o700 assert stat.S_IMODE(file.stat().st_mode) == 0o600 def test_read_config_concurrent_first_use_publishes_one_complete_file(monkeypatch: pytest.MonkeyPatch) -> None: destination = CONFIG_DIRS[0].expanduser() / 'api-keys.yaml' config_files = {directory.expanduser() / destination.name for directory in CONFIG_DIRS} publication_barrier = Barrier(2) link = core_module.os.link read_text = Path.read_text def isolated_read_text(path: Path, *args, **kwargs) -> str: if path == destination and path.exists(): return read_text(path, *args, **kwargs) if path in config_files: raise FileNotFoundError return read_text(path, *args, **kwargs) def synchronized_link(source: str | Path, target: str | Path) -> None: publication_barrier.wait(timeout=5) link(source, target) monkeypatch.setattr(core_module.os, 'link', synchronized_link) monkeypatch.setattr(Path, 'read_text', isolated_read_text) with ThreadPoolExecutor(max_workers=2) as executor: configs = list(executor.map(lambda _index: Core._read_config('api-keys.yaml'), range(2))) assert configs == [(DATA_DIR / 'api-keys.yaml').read_text()] * 2 assert destination.read_text() == configs[0] assert stat.S_IMODE(destination.stat().st_mode) == 0o600 assert list(destination.parent.glob(f'.{destination.name}.*')) == [] _DEFAULT_JSON = object() class DummyContent: def __init__(self, chunks: tuple[bytes, ...], error: BaseException | None = None): self.chunks = chunks self.error = error async def iter_any(self): for chunk in self.chunks: yield chunk if self.error is not None: raise self.error class DummyResponse: def __init__( self, text_value: str = 'response-text', json_value: Any = _DEFAULT_JSON, status: int = 200, headers: dict[str, str] | None = None, chunks: tuple[bytes, ...] | None = None, stream_error: BaseException | None = None, ): self.text_value = text_value self.json_value = {'ok': True} if json_value is _DEFAULT_JSON else json_value self.status = status self.headers = headers or {} self.content = DummyContent(chunks if chunks is not None else (text_value.encode(),), stream_error) async def __aenter__(self): return self async def __aexit__(self, exc_type, exc, tb): return False async def text(self): return self.text_value async def json(self): if isinstance(self.json_value, Exception): raise self.json_value return self.json_value class DummySession: instances: list[DummySession] = [] def __init__( self, *, headers: dict[str, str] | None = None, timeout: Any = None, connector: Any = None, proxy: str | None = None, cookie_jar: Any = None, ) -> None: self.headers = headers self.timeout = timeout self.connector = connector self.proxy = proxy self.cookie_jar = cookie_jar self.closed = False self.requests: list[tuple[str, str, dict[str, Any]]] = [] DummySession.instances.append(self) async def __aenter__(self): return self async def __aexit__(self, exc_type, exc, tb): await self.close() return False def request(self, method: str, url: str, **kwargs): self.requests.append((method, url, kwargs)) return DummyResponse() def get(self, url: str, **kwargs): self.requests.append(('GET', url, kwargs)) return DummyResponse() def post(self, url: str, **kwargs): self.requests.append(('POST', url, kwargs)) return DummyResponse(json_value={'posted': True}) async def close(self): self.closed = True def reset_dummy_sessions() -> None: DummySession.instances.clear() def install_stream_response( monkeypatch, *, chunks: tuple[bytes, ...], status: int = 200, headers: dict[str, str] | None = None, error: BaseException | None = None, ) -> None: reset_dummy_sessions() monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(core_module.aiohttp, 'TCPConnector', lambda *, ssl=None: object()) monkeypatch.setattr(core_module.ssl, 'create_default_context', lambda cafile=None: 'ssl-context') monkeypatch.setattr(core_module.certifi, 'where', lambda: '/tmp/cacert.pem') def request_stream(self, method: str, url: str, **kwargs): self.requests.append((method, url, kwargs)) return DummyResponse(status=status, headers=headers, chunks=chunks, stream_error=error) monkeypatch.setattr(DummySession, 'request', request_stream) async def collect_default_stream(lines: list[str]) -> None: async with AsyncFetcher.stream_records( 'https://provider.example/stream', framing='ndjson', ) as response: async for record in response: lines.append(record) async def fake_sleep(_seconds: float) -> None: return None def test_api_keys_yaml_is_in_sync_with_core_accessors(): required = core_module.Core._API_KEY_FIELDS assert required, "No API-key references were detected in `Core`" config = yaml.safe_load((DATA_DIR / "api-keys.yaml").read_text(encoding="utf-8")) apikeys = config["apikeys"] missing_providers = sorted(set(required) - set(apikeys)) assert not missing_providers, f"Missing providers in api-keys.yaml: {missing_providers}" missing_fields: dict[str, list[str]] = {} for provider, fields in required.items(): for field in sorted(fields): if field not in apikeys[provider]: missing_fields.setdefault(provider, []).append(field) assert not missing_fields, f"Missing fields in api-keys.yaml: {missing_fields}" def test_user_agent_policy_separates_provider_and_browser_identities() -> None: assert Core.get_user_agent() == ( f'theHarvester/{core_module.__version__} (+https://github.com/laramies/theHarvester)' ) assert Core.get_browser_user_agent() == ( 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) ' 'AppleWebKit/537.36 (KHTML, like Gecko) ' 'Chrome/151.0.0.0 Safari/537.36' ) @pytest.mark.parametrize( ("accessor_name", "expected"), [ ("bevigil_key", "bevigil-key"), ("censys_key", ("censys-token", "censys-org")), ("fofa_key", ("fofa-key", "fofa-email")), ("routeviews_key", "routeviews-key"), ("tomba_key", ("tomba-key", "tomba-secret")), ], ) def test_api_key_accessors_read_configured_values(monkeypatch, accessor_name: str, expected: Any): monkeypatch.setattr( Core, 'api_keys', staticmethod( lambda: { 'bevigil': {'key': 'bevigil-key'}, 'censys': {'token': 'censys-token', 'organization_id': 'censys-org'}, 'fofa': {'key': 'fofa-key', 'email': 'fofa-email'}, 'routeviews': {'key': 'routeviews-key'}, 'tomba': {'key': 'tomba-key', 'secret': 'tomba-secret'}, } ), ) accessor = getattr(Core, accessor_name) assert accessor() == expected @pytest.mark.parametrize('configured_value', [None, '', ' ', 10]) def test_routeviews_key_ignores_missing_or_invalid_optional_credentials(monkeypatch, configured_value: object) -> None: monkeypatch.setattr(Core, 'api_keys', staticmethod(lambda: {'routeviews': {'key': configured_value}})) assert Core.routeviews_key() is None def test_routeviews_key_strips_configuration_whitespace(monkeypatch) -> None: monkeypatch.setattr(Core, 'api_keys', staticmethod(lambda: {'routeviews': {'key': ' routeviews-key '}})) assert Core.routeviews_key() == 'routeviews-key' @pytest.mark.asyncio async def test_fetch_creates_session_with_default_headers(monkeypatch) -> None: async def fail_if_fetch_sleeps(seconds: float) -> None: raise AssertionError(f'fetch delayed reading a ready response by {seconds} seconds') reset_dummy_sessions() monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(core_module.ssl, 'create_default_context', lambda cafile=None: 'ssl-context') monkeypatch.setattr(core_module.certifi, 'where', lambda: '/tmp/cacert.pem') monkeypatch.setattr(core_module.asyncio, 'sleep', fail_if_fetch_sleeps) monkeypatch.setattr(Core, 'get_user_agent', staticmethod(lambda: 'test-agent')) result = await AsyncFetcher.fetch(url='https://example.com', follow_redirects=False) assert result == 'response-text' assert len(DummySession.instances) == 1 session = DummySession.instances[0] assert session.headers == {'User-Agent': 'test-agent'} assert session.closed is True assert session.requests == [ ('GET', 'https://example.com', {'ssl': 'ssl-context', 'allow_redirects': False}) ] @pytest.mark.asyncio async def test_fetch_reused_session_uses_a_stable_explicit_ssl_policy(monkeypatch) -> None: session = DummySession() monkeypatch.setattr(AsyncFetcher, '_ssl_context', staticmethod(lambda _verify=True: object())) await AsyncFetcher.fetch(session=session, url='https://example.com/one', verify=True) await AsyncFetcher.fetch(session=session, url='https://example.com/two', verify=True) ssl_policies = [request_options['ssl'] for _method, _url, request_options in session.requests] assert ssl_policies == [True, True] assert ssl_policies[0] is ssl_policies[1] @pytest.mark.asyncio async def test_open_session_owns_one_proxy_and_cookie_policy(monkeypatch: pytest.MonkeyPatch) -> None: reset_dummy_sessions() cookie_jar = core_module.aiohttp.DummyCookieJar() monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(AsyncFetcher, '_ssl_context', staticmethod(lambda: 'ssl-context')) async def fake_connector(*_args: object) -> str: return 'provider-connector' monkeypatch.setattr(AsyncFetcher, '_create_connector', fake_connector) async with AsyncFetcher.open_session( headers={'Accept': 'application/json'}, proxy='http://proxy.example:8080', request_timeout=45, cookie_jar=cookie_jar, ) as session: assert session.closed is False assert session.closed is True assert session.headers['Accept'] == 'application/json' assert session.proxy == 'http://proxy.example:8080' assert session.connector == 'provider-connector' assert session.cookie_jar is cookie_jar assert session.timeout.total == 45 @pytest.mark.asyncio async def test_open_session_fails_closed_when_proxy_mode_has_no_proxy(monkeypatch: pytest.MonkeyPatch) -> None: monkeypatch.setattr(AsyncFetcher, '_proxy_list', {'http': [], 'socks5': []}) with pytest.raises(ProxyUnavailableError, match='proxy-unavailable'): async with AsyncFetcher.open_session(proxy=True): pytest.fail('a direct session must not open when proxy mode is required') @pytest.mark.asyncio async def test_open_session_preserves_project_ca_and_aiohttp_default_timeout( monkeypatch: pytest.MonkeyPatch, ) -> None: reset_dummy_sessions() connector_options: list[tuple[str | None, str | None, object]] = [] monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(AsyncFetcher, '_ssl_context', staticmethod(lambda: 'ssl-context')) async def fake_connector( proxy_url: str | None, proxy_type: str | None, ssl_context: object, ) -> str: connector_options.append((proxy_url, proxy_type, ssl_context)) return 'direct-provider-connector' monkeypatch.setattr(AsyncFetcher, '_create_connector', fake_connector) async with AsyncFetcher.open_session() as session: assert session.timeout is None assert connector_options == [(None, None, 'ssl-context')] assert session.connector == 'direct-provider-connector' @pytest.mark.asyncio async def test_open_session_finishes_close_and_preserves_the_first_cancellation( monkeypatch: pytest.MonkeyPatch, ) -> None: close_started = asyncio.Event() release_close = asyncio.Event() cancelled = asyncio.CancelledError('provider-cancelled') class BlockingCloseSession: closed = False async def close(self) -> None: close_started.set() await release_close.wait() self.closed = True session = BlockingCloseSession() async def fake_build_session(*_args: object, **_kwargs: object) -> BlockingCloseSession: return session monkeypatch.setattr(AsyncFetcher, '_build_session', fake_build_session) async def use_session() -> None: async with AsyncFetcher.open_session(): raise cancelled task = asyncio.create_task(use_session()) await close_started.wait() task.cancel('second-cancellation') release_close.set() with pytest.raises(asyncio.CancelledError) as raised: await task assert raised.value is cancelled assert session.closed is True def test_default_headers_add_project_identity_without_mutating_caller_headers(monkeypatch) -> None: monkeypatch.setattr(Core, 'get_user_agent', staticmethod(lambda: 'test-agent')) supplied = {'Accept': 'application/json'} headers = AsyncFetcher._default_headers(supplied) assert headers == {'Accept': 'application/json', 'User-Agent': 'test-agent'} assert supplied == {'Accept': 'application/json'} @pytest.mark.parametrize('header_name', ['User-Agent', 'User-agent', 'user-agent']) def test_default_headers_preserve_explicit_user_agent_case_insensitively(monkeypatch, header_name: str) -> None: monkeypatch.setattr(Core, 'get_user_agent', staticmethod(lambda: 'default-agent')) headers = AsyncFetcher._default_headers({header_name: 'caller-agent', 'Accept': 'application/json'}) assert headers == {header_name: 'caller-agent', 'Accept': 'application/json'} @pytest.mark.asyncio async def test_fetch_can_include_buffered_response_metadata(monkeypatch) -> None: reset_dummy_sessions() monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(core_module.ssl, 'create_default_context', lambda cafile=None: 'ssl-context') monkeypatch.setattr(core_module.certifi, 'where', lambda: '/tmp/cacert.pem') monkeypatch.setattr(core_module.asyncio, 'sleep', fake_sleep) def request_with_metadata(self, method: str, url: str, **kwargs): self.requests.append((method, url, kwargs)) return DummyResponse( text_value='rate limited', status=429, headers={'Retry-After': '60', 'X-RateLimit-Remaining': '0'}, ) monkeypatch.setattr(DummySession, 'request', request_with_metadata) result = await AsyncFetcher.fetch(url='https://example.com', include_metadata=True) assert result == FetcherResponse( body='rate limited', status=429, headers={'retry-after': '60', 'x-ratelimit-remaining': '0'}, ) @pytest.mark.asyncio async def test_fetch_rejects_a_buffered_response_over_the_explicit_byte_limit(monkeypatch) -> None: reset_dummy_sessions() monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(core_module.ssl, 'create_default_context', lambda cafile=None: 'ssl-context') monkeypatch.setattr(core_module.certifi, 'where', lambda: '/tmp/cacert.pem') def oversized_response(self, method: str, url: str, **kwargs): self.requests.append((method, url, kwargs)) return DummyResponse(status=206, headers={'Content-Type': 'application/json', 'X-Evidence': 'kept'}) monkeypatch.setattr(DummySession, 'request', oversized_response) with pytest.raises(ResponseStreamError, match='response-limit') as raised: await AsyncFetcher.fetch(url='https://example.com', include_metadata=True, response_byte_limit=4) assert raised.value.status == 206 assert raised.value.headers == {'content-type': 'application/json', 'x-evidence': 'kept'} assert DummySession.instances[0].closed is True @pytest.mark.asyncio async def test_limited_metadata_fetch_preserves_non_json_response_text(monkeypatch) -> None: reset_dummy_sessions() monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(core_module.ssl, 'create_default_context', lambda cafile=None: 'ssl-context') monkeypatch.setattr(core_module.certifi, 'where', lambda: '/tmp/cacert.pem') result = await AsyncFetcher.fetch( url='https://example.com', json=True, include_metadata=True, response_byte_limit=100, ) assert result == FetcherResponse(body='response-text', status=200, headers={}) @pytest.mark.asyncio async def test_stream_records_yields_fragmented_utf8_records_and_response_metadata(monkeypatch) -> None: install_stream_response( monkeypatch, chunks=(b'first\ncaf\xc3', b'\xa9\nlast'), status=206, headers={'X-Stream': 'ready'}, ) async with AsyncFetcher.stream_records( 'https://provider.example/stream', framing='ndjson', params={'q': 'example.com'}, follow_redirects=False, ) as response: lines = [line async for line in response] assert response.status == 206 assert response.headers == {'x-stream': 'ready'} assert lines == ['first', 'café', 'last'] assert DummySession.instances[0].closed is True assert DummySession.instances[0].requests == [ ( 'GET', 'https://provider.example/stream', { 'ssl': 'ssl-context', 'allow_redirects': False, 'params': {'q': 'example.com'}, }, ) ] @pytest.mark.asyncio async def test_stream_records_frames_fragmented_sse_events(monkeypatch) -> None: install_stream_response( monkeypatch, chunks=( b'event: matches\r\n', b'data: [{"type":"content"}]\r', b'\n\r\n', b'event: done\n', b'data: {}\n\n', ), ) async with AsyncFetcher.stream_records( 'https://provider.example/stream', framing='sse', ) as response: records = [record async for record in response] assert records == [ 'event: matches\ndata: [{"type":"content"}]', 'event: done\ndata: {}', ] @pytest.mark.asyncio async def test_stream_records_normalizes_tls_setup_failure(monkeypatch) -> None: def fail_tls_setup(*_args, **_kwargs): raise core_module.ssl.SSLError('TLS unavailable') monkeypatch.setattr( AsyncFetcher, '_ssl_context', staticmethod(fail_tls_setup), ) with pytest.raises(core_module.ResponseStreamError) as raised: await collect_default_stream([]) assert raised.value.reason == 'transport-error' @pytest.mark.asyncio async def test_stream_records_does_not_rewrite_consumer_exceptions(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=(b'ok\n',)) with pytest.raises(PermissionError, match='provider denied access'): async with AsyncFetcher.stream_records( 'https://provider.example/stream', framing='ndjson', ): raise PermissionError('provider denied access') assert DummySession.instances[0].closed is True @pytest.mark.asyncio async def test_stream_records_uses_resolved_http_proxy_without_redirects(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=(b'ok\n',)) async def fake_connector(*_args, **_kwargs): return 'proxy-connector' monkeypatch.setattr(AsyncFetcher, '_create_connector', fake_connector) async with AsyncFetcher.stream_records( 'https://provider.example/stream', framing='ndjson', proxy='http://proxy.example:8080', ) as response: assert [record async for record in response] == ['ok'] session = DummySession.instances[0] assert session.connector == 'proxy-connector' assert session.requests == [ ( 'GET', 'https://provider.example/stream', { 'ssl': 'ssl-context', 'allow_redirects': False, 'proxy': 'http://proxy.example:8080', }, ) ] @pytest.mark.asyncio async def test_stream_records_rejects_second_iteration(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=(b'one\ntwo\n',)) async with AsyncFetcher.stream_records( 'https://provider.example/stream', framing='ndjson', ) as response: iterator = response.__aiter__() assert await anext(iterator) == 'one' with pytest.raises(RuntimeError, match='only be consumed once'): response.__aiter__() await iterator.aclose() @pytest.mark.asyncio async def test_stream_records_rejects_incomplete_sse_after_complete_prefix(monkeypatch) -> None: install_stream_response( monkeypatch, chunks=(b'event: progress\ndata: {}\n\nevent: matches\ndata: []',), ) records: list[str] = [] with pytest.raises(core_module.ResponseStreamError) as raised: async with AsyncFetcher.stream_records( 'https://provider.example/stream', framing='sse', ) as response: async for record in response: records.append(record) assert records == ['event: progress\ndata: {}'] assert raised.value.reason == 'invalid-response' @pytest.mark.asyncio async def test_stream_records_preserves_complete_prefix_before_response_limit(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=(b'one\n', b'two\n', b'x')) monkeypatch.setattr(core_module, 'MAX_PROVIDER_STREAM_BYTES', 8) lines: list[str] = [] with pytest.raises(core_module.ResponseStreamError) as raised: await collect_default_stream(lines) assert lines == ['one', 'two'] assert raised.value.reason == 'response-limit' assert DummySession.instances[0].closed is True @pytest.mark.asyncio async def test_stream_records_accepts_final_record_ending_exactly_at_response_limit(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=(b'one\ntwo',)) monkeypatch.setattr(core_module, 'MAX_PROVIDER_STREAM_BYTES', 7) lines: list[str] = [] await collect_default_stream(lines) assert lines == ['one', 'two'] @pytest.mark.asyncio async def test_stream_records_rejects_oversized_record_after_complete_prefix(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=(b'ok\n123', b'456\n')) monkeypatch.setattr(core_module, 'MAX_STREAM_RECORD_BYTES', 5) lines: list[str] = [] with pytest.raises(core_module.ResponseStreamError) as raised: await collect_default_stream(lines) assert lines == ['ok'] assert raised.value.reason == 'response-limit' @pytest.mark.asyncio async def test_stream_records_accepts_record_at_limit_with_split_crlf(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=(b'abc\r', b'\n')) monkeypatch.setattr(core_module, 'MAX_STREAM_RECORD_BYTES', 3) records: list[str] = [] await collect_default_stream(records) assert records == ['abc'] @pytest.mark.asyncio async def test_stream_records_bounds_complete_sse_event_across_many_lines(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=(b'a\nb\nc\n\n',)) monkeypatch.setattr(core_module, 'MAX_STREAM_RECORD_BYTES', 5) async with AsyncFetcher.stream_records( 'https://provider.example/stream', framing='sse', ) as response: records = [record async for record in response] assert records == ['a\nb\nc'] @pytest.mark.asyncio async def test_stream_records_reports_invalid_utf8_after_complete_prefix(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=(b'ok\n', b'\xff\n')) lines: list[str] = [] with pytest.raises(core_module.ResponseStreamError) as raised: await collect_default_stream(lines) assert lines == ['ok'] assert raised.value.reason == 'invalid-response' @pytest.mark.asyncio async def test_stream_records_reports_transport_failure_after_complete_prefix(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=(b'ok\n',), error=TimeoutError()) lines: list[str] = [] with pytest.raises(core_module.ResponseStreamError) as raised: await collect_default_stream(lines) assert lines == ['ok'] assert raised.value.reason == 'transport-error' assert DummySession.instances[0].closed is True @pytest.mark.asyncio async def test_stream_records_reports_connection_failure_and_closes_session(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=()) def failed_request(self, method: str, url: str, **kwargs): self.requests.append((method, url, kwargs)) raise core_module.aiohttp.ClientConnectionError('unavailable') monkeypatch.setattr(DummySession, 'request', failed_request) with pytest.raises(core_module.ResponseStreamError) as raised: await collect_default_stream([]) assert raised.value.reason == 'transport-error' assert DummySession.instances[0].closed is True @pytest.mark.asyncio async def test_stream_records_reports_session_creation_failure(monkeypatch) -> None: async def failed_build_session(*_args: object, **_kwargs: object): raise core_module.aiohttp.ClientConnectionError('unavailable') monkeypatch.setattr(AsyncFetcher, '_build_session', failed_build_session) with pytest.raises(core_module.ResponseStreamError) as raised: await collect_default_stream([]) assert raised.value.reason == 'transport-error' @pytest.mark.asyncio async def test_stream_records_propagates_cancellation_and_closes_session(monkeypatch) -> None: cancelled = asyncio.CancelledError() install_stream_response(monkeypatch, chunks=(), error=cancelled) with pytest.raises(asyncio.CancelledError) as raised: await collect_default_stream([]) assert raised.value is cancelled assert DummySession.instances[0].closed is True @pytest.mark.asyncio async def test_fetch_json_reads_bounded_fragmented_utf8_without_redirects(monkeypatch) -> None: install_stream_response( monkeypatch, chunks=(b'{"name":"caf\xc3', b'\xa9"}'), headers={'X-Provider': 'ready'}, ) result = await AsyncFetcher.fetch_json( 'https://provider.example/data', params={'asn': '64500'}, headers={'Api-Key': 'secret'}, request_timeout=30, ) assert result == FetcherResponse( body={'name': 'café'}, status=200, headers={'x-provider': 'ready'}, ) assert DummySession.instances[0].closed is True assert DummySession.instances[0].requests == [ ( 'GET', 'https://provider.example/data', { 'ssl': 'ssl-context', 'allow_redirects': False, 'params': {'asn': '64500'}, }, ) ] @pytest.mark.asyncio async def test_fetch_text_preserves_non_success_metadata_without_redirects(monkeypatch) -> None: install_stream_response( monkeypatch, chunks=(b'not ', b'found'), status=404, headers={'Location': 'https://outside.example/path'}, ) result = await AsyncFetcher.fetch_text( 'https://app.example.test', request_timeout=None, response_byte_limit=32, ) assert result == FetcherResponse( body='not found', status=404, headers={'location': 'https://outside.example/path'}, ) assert DummySession.instances[0].closed is True assert DummySession.instances[0].requests[0][2]['allow_redirects'] is False assert isinstance(DummySession.instances[0].timeout, core_module.aiohttp.ClientTimeout) assert DummySession.instances[0].timeout.total is None @pytest.mark.asyncio async def test_fetch_text_reuses_caller_owned_session_without_closing_it( monkeypatch: pytest.MonkeyPatch, ) -> None: install_stream_response(monkeypatch, chunks=(b'provider ', b'evidence')) session = DummySession( headers={'User-Agent': 'shared'}, timeout=core_module.aiohttp.ClientTimeout(total=None), ) result = await AsyncFetcher.fetch_text( 'https://app.example.test', session=session, request_timeout=None, response_byte_limit=32, ) assert result.body == 'provider evidence' assert session.closed is False assert len(DummySession.instances) == 1 @pytest.mark.asyncio async def test_fetch_json_reuses_caller_owned_session_without_closing_it( monkeypatch: pytest.MonkeyPatch, ) -> None: install_stream_response(monkeypatch, chunks=(b'{"provider":"evidence"}',)) session = DummySession( headers={'User-Agent': 'shared'}, timeout=core_module.aiohttp.ClientTimeout(total=None), ) result = await AsyncFetcher.fetch_json( 'https://provider.example/data', session=session, request_timeout=None, ) assert result.body == {'provider': 'evidence'} assert session.closed is False assert len(DummySession.instances) == 1 @pytest.mark.asyncio async def test_fetch_json_accepts_body_at_shared_limit(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=(b'{"a":1}',)) monkeypatch.setattr(core_module, 'MAX_PROVIDER_JSON_BYTES', 7) result = await AsyncFetcher.fetch_json('https://provider.example/data') assert result.body == {'a': 1} @pytest.mark.asyncio @pytest.mark.parametrize('headers', [None, {'Content-Length': '1'}]) async def test_fetch_json_rejects_body_over_shared_limit(monkeypatch, headers: dict[str, str] | None) -> None: install_stream_response(monkeypatch, chunks=(b'{"a":1}', b' '), headers=headers) monkeypatch.setattr(core_module, 'MAX_PROVIDER_JSON_BYTES', 7) with pytest.raises(core_module.ResponseStreamError) as raised: await AsyncFetcher.fetch_json('https://provider.example/data') assert raised.value.reason == 'response-limit' assert DummySession.instances[0].closed is True @pytest.mark.asyncio async def test_fetch_json_rejects_declared_oversized_body(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=(b'{"a":1}',), headers={'Content-Length': '8'}) monkeypatch.setattr(core_module, 'MAX_PROVIDER_JSON_BYTES', 7) with pytest.raises(core_module.ResponseStreamError) as raised: await AsyncFetcher.fetch_json('https://provider.example/data') assert raised.value.reason == 'response-limit' @pytest.mark.asyncio async def test_fetch_json_preserves_non_success_status_without_parsing_body(monkeypatch) -> None: install_stream_response( monkeypatch, chunks=(b'limited',), status=429, headers={'Retry-After': '60'}, ) result = await AsyncFetcher.fetch_json('https://provider.example/data') assert result == FetcherResponse(body=None, status=429, headers={'retry-after': '60'}) @pytest.mark.asyncio async def test_fetch_json_returns_none_for_no_content(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=(), status=204) result = await AsyncFetcher.fetch_json('https://provider.example/data') assert result == FetcherResponse(body=None, status=204, headers={}) @pytest.mark.asyncio @pytest.mark.parametrize('chunks', [(b'\xff',), (b'{',), (b'{"value":NaN}',)]) async def test_fetch_json_rejects_invalid_utf8_or_json(monkeypatch, chunks: tuple[bytes, ...]) -> None: install_stream_response(monkeypatch, chunks=chunks) with pytest.raises(core_module.ResponseStreamError) as raised: await AsyncFetcher.fetch_json('https://provider.example/data') assert raised.value.reason == 'invalid-response' @pytest.mark.asyncio async def test_fetch_json_reports_transport_failure_and_closes_session(monkeypatch) -> None: install_stream_response(monkeypatch, chunks=(), error=TimeoutError()) with pytest.raises(core_module.ResponseStreamError) as raised: await AsyncFetcher.fetch_json('https://provider.example/data') assert raised.value.reason == 'transport-error' assert DummySession.instances[0].closed is True @pytest.mark.asyncio async def test_fetch_json_propagates_cancellation_and_closes_session(monkeypatch) -> None: cancelled = asyncio.CancelledError() install_stream_response(monkeypatch, chunks=(), error=cancelled) with pytest.raises(asyncio.CancelledError) as raised: await AsyncFetcher.fetch_json('https://provider.example/data') assert raised.value is cancelled assert DummySession.instances[0].closed is True @pytest.mark.asyncio async def test_fetch_metadata_distinguishes_transport_failure(monkeypatch) -> None: async def failed_request(*_args: Any, **_kwargs: Any) -> str: raise OSError('network unavailable') monkeypatch.setattr(AsyncFetcher, '_request', failed_request) result = await AsyncFetcher.fetch(session=DummySession(), url='https://example.com', include_metadata=True) assert result is None @pytest.mark.asyncio async def test_fetch_metadata_preserves_non_json_error_body(monkeypatch) -> None: reset_dummy_sessions() monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(core_module.ssl, 'create_default_context', lambda cafile=None: 'ssl-context') monkeypatch.setattr(core_module.certifi, 'where', lambda: '/tmp/cacert.pem') monkeypatch.setattr(core_module.asyncio, 'sleep', fake_sleep) def request_with_invalid_json(self, method: str, url: str, **kwargs): self.requests.append((method, url, kwargs)) return DummyResponse(text_value='upstream error', json_value=ValueError(), status=502) monkeypatch.setattr(DummySession, 'request', request_with_invalid_json) result = await AsyncFetcher.fetch(url='https://example.com', json=True, include_metadata=True) assert result == FetcherResponse(body='upstream error', status=502, headers={}) @pytest.mark.asyncio @pytest.mark.parametrize( ('text_value', 'expected_body'), [('', ''), ('null', None)], ) async def test_fetch_metadata_distinguishes_empty_json_from_null( monkeypatch, text_value: str, expected_body: Any, ) -> None: reset_dummy_sessions() monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(core_module.ssl, 'create_default_context', lambda cafile=None: 'ssl-context') monkeypatch.setattr(core_module.certifi, 'where', lambda: '/tmp/cacert.pem') monkeypatch.setattr(core_module.asyncio, 'sleep', fake_sleep) def request_with_json(self, method: str, url: str, **kwargs): self.requests.append((method, url, kwargs)) return DummyResponse(text_value=text_value, json_value=None) monkeypatch.setattr(DummySession, 'request', request_with_json) result = await AsyncFetcher.fetch(url='https://example.com', json=True, include_metadata=True) assert result == FetcherResponse(body=expected_body, status=200, headers={}) @pytest.mark.asyncio async def test_fetch_all_propagates_metadata_opt_in(monkeypatch) -> None: seen: list[bool] = [] async def fake_fetch(*_args: Any, include_metadata: bool = False, **_kwargs: Any) -> FetcherResponse: seen.append(include_metadata) return FetcherResponse(body='limited', status=429, headers={'retry-after': '60'}) monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(AsyncFetcher, 'fetch', fake_fetch) results = await AsyncFetcher.fetch_all(['https://one.example', 'https://two.example'], include_metadata=True) assert seen == [True, True] assert [result.status for result in results] == [429, 429] @pytest.mark.asyncio async def test_fetch_all_reuses_a_caller_owned_session(monkeypatch: pytest.MonkeyPatch) -> None: session = object() seen_sessions: list[object] = [] async def fake_fetch(*_args: Any, session: object, **_kwargs: Any) -> str: seen_sessions.append(session) return 'ok' monkeypatch.setattr(AsyncFetcher, 'fetch', fake_fetch) results = await AsyncFetcher.fetch_all( ['https://one.example', 'https://two.example'], session=session, ) assert results == ['ok', 'ok'] assert seen_sessions == [session, session] @pytest.mark.asyncio async def test_fetch_all_rejects_proxy_selection_for_a_caller_owned_session() -> None: with pytest.raises(ValueError, match='caller-owned session'): await AsyncFetcher.fetch_all(['https://one.example'], proxy=True, session=object()) @pytest.mark.asyncio async def test_fetch_all_resolves_one_socks_proxy_for_the_owned_session(monkeypatch: pytest.MonkeyPatch) -> None: reset_dummy_sessions() proxy_inputs: list[str | bool | None] = [] connector_inputs: list[tuple[str | None, str | None, object]] = [] monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(AsyncFetcher, '_ssl_context', staticmethod(lambda _verify=True: 'ssl-context')) def resolve_proxy(_cls: type[AsyncFetcher], proxy: str | bool | None) -> tuple[str | None, str | None]: proxy_inputs.append(proxy) return ('socks5://proxy.example:1080', 'socks5') if proxy else (None, None) async def create_connector( proxy_url: str | None, proxy_type: str | None, ssl_context: object, ) -> str: connector_inputs.append((proxy_url, proxy_type, ssl_context)) return 'socks-connector' monkeypatch.setattr(AsyncFetcher, '_resolve_proxy', classmethod(resolve_proxy)) monkeypatch.setattr(AsyncFetcher, '_create_connector', create_connector) results = await AsyncFetcher.fetch_all( ['https://one.example', 'https://two.example'], proxy=True, ) assert results == ['response-text', 'response-text'] assert [proxy for proxy in proxy_inputs if proxy] == [True] assert connector_inputs == [('socks5://proxy.example:1080', 'socks5', 'ssl-context')] assert len(DummySession.instances) == 1 assert DummySession.instances[0].connector == 'socks-connector' assert DummySession.instances[0].closed is True @pytest.mark.asyncio async def test_fetch_uses_http_proxy_when_enabled(monkeypatch) -> None: reset_dummy_sessions() monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(core_module.ssl, 'create_default_context', lambda cafile=None: 'ssl-context') monkeypatch.setattr(core_module.certifi, 'where', lambda: '/tmp/cacert.pem') monkeypatch.setattr(core_module.asyncio, 'sleep', fake_sleep) monkeypatch.setattr(AsyncFetcher, '_get_random_proxy', staticmethod(lambda proxy_dict: ('http://proxy.local:8080', 'http'))) async def fake_create_connector(proxy_url, proxy_type, ssl_context=None): return 'connector' monkeypatch.setattr(AsyncFetcher, '_create_connector', fake_create_connector) result = await AsyncFetcher.fetch(url='https://example.com', proxy=True) assert result == 'response-text' session = DummySession.instances[0] assert session.connector == 'connector' assert session.requests == [ ('GET', 'https://example.com', {'ssl': 'ssl-context', 'proxy': 'http://proxy.local:8080'}) ] @pytest.mark.asyncio async def test_post_fetch_decodes_string_payload_and_posts_params(monkeypatch) -> None: reset_dummy_sessions() monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(core_module.asyncio, 'sleep', fake_sleep) monkeypatch.setattr(core_module.ssl, 'create_default_context', lambda cafile=None: 'ssl-context') monkeypatch.setattr(core_module.certifi, 'where', lambda: '/tmp/cacert.pem') monkeypatch.setattr(Core, 'get_user_agent', staticmethod(lambda: 'test-agent')) async def fake_create_connector( _proxy_url: str | None, _proxy_type: str | None, ssl_context: object, ) -> str: return f'connector:{ssl_context}' monkeypatch.setattr(AsyncFetcher, '_create_connector', fake_create_connector) result = await AsyncFetcher.post_fetch( 'https://example.com/api', data='{"query": "example"}', params={'page': 2}, json=True, ) assert result == {'ok': True} session = DummySession.instances[0] assert session.headers == {'User-Agent': 'test-agent'} assert session.connector == 'connector:ssl-context' assert session.requests == [ ('POST', 'https://example.com/api', {'data': {'query': 'example'}, 'params': {'page': 2}}) ] @pytest.mark.asyncio async def test_post_fetch_sends_json_body(monkeypatch) -> None: reset_dummy_sessions() monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(core_module.asyncio, 'sleep', fake_sleep) monkeypatch.setattr(core_module.ssl, 'create_default_context', lambda cafile=None: 'ssl-context') monkeypatch.setattr(core_module.certifi, 'where', lambda: '/tmp/cacert.pem') async def fake_create_connector( _proxy_url: str | None, _proxy_type: str | None, ssl_context: object, ) -> str: return f'connector:{ssl_context}' monkeypatch.setattr(AsyncFetcher, '_create_connector', fake_create_connector) result = await AsyncFetcher.post_fetch( 'https://example.com/api', json_body={'scan': 'example'}, json=True, ) assert result == {'ok': True} session = DummySession.instances[0] assert session.connector == 'connector:ssl-context' assert session.requests == [ ('POST', 'https://example.com/api', {'json': {'scan': 'example'}}) ] @pytest.mark.asyncio async def test_post_fetch_can_include_response_metadata(monkeypatch) -> None: reset_dummy_sessions() monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(core_module.asyncio, 'sleep', fake_sleep) def request_with_metadata(self, method: str, url: str, **kwargs): self.requests.append((method, url, kwargs)) return DummyResponse(text_value='unavailable', status=503, headers={'Retry-After': '30'}) monkeypatch.setattr(DummySession, 'request', request_with_metadata) result = await AsyncFetcher.post_fetch('https://example.com/api', data='{}', include_metadata=True) assert result == FetcherResponse(body='unavailable', status=503, headers={'retry-after': '30'}) @pytest.mark.asyncio async def test_post_fetch_proxy_branch_posts_body_and_params_with_http_proxy(monkeypatch) -> None: reset_dummy_sessions() created_connectors = [] monkeypatch.setattr(core_module.aiohttp, 'ClientSession', DummySession) monkeypatch.setattr(core_module.asyncio, 'sleep', fake_sleep) monkeypatch.setattr(core_module.ssl, 'create_default_context', lambda cafile=None: 'ssl-context') monkeypatch.setattr(core_module.certifi, 'where', lambda: '/tmp/cacert.pem') monkeypatch.setattr(AsyncFetcher, '_get_random_proxy', staticmethod(lambda proxy_dict: ('http://proxy.local:8080', 'http'))) async def fake_create_connector(proxy_url, proxy_type, ssl_context=None): created_connectors.append((proxy_url, proxy_type, ssl_context)) return 'connector' monkeypatch.setattr(AsyncFetcher, '_create_connector', fake_create_connector) result = await AsyncFetcher.post_fetch( 'https://example.com/resource', json_body={'scan': 'example'}, params={'page': 2}, json=True, proxy=True, ) assert result == {'ok': True} assert created_connectors == [('http://proxy.local:8080', 'http', 'ssl-context')] session = DummySession.instances[0] assert session.connector == 'connector' assert session.proxy == 'http://proxy.local:8080' assert session.requests == [ ( 'POST', 'https://example.com/resource', {'json': {'scan': 'example'}, 'params': {'page': 2}}, ) ]