diff --git a/tests/discovery/test_gitlabsearch.py b/tests/discovery/test_gitlabsearch.py index e2ae7a73..7f4e1504 100644 --- a/tests/discovery/test_gitlabsearch.py +++ b/tests/discovery/test_gitlabsearch.py @@ -22,7 +22,17 @@ if TYPE_CHECKING: async def test_public_discovery_normalizes_evidence_and_uses_bounded_requests( monkeypatch: pytest.MonkeyPatch, ) -> None: + import contextlib + requests: list[dict[str, object]] = [] + session_proxy: list[object] = [] + session = object() + + @contextlib.asynccontextmanager + async def fake_open_session(**kwargs: Any): + session_proxy.append(kwargs.get('proxy')) + yield session + projects = [ { 'id': 'group/project', @@ -55,11 +65,11 @@ async def test_public_discovery_normalizes_evidence_and_uses_bounded_requests( async def fake_fetch_all( urls: list[str] | set[str], headers: dict[str, str] | None = None, - proxy: bool = False, + session: object | None = None, **_kwargs: Any, ) -> list[FetcherResponse]: url = next(iter(urls)) - requests.append({'url': url, 'headers': headers, 'proxy': proxy}) + requests.append({'url': url, 'headers': headers, 'session': session}) responses = { 'https://gitlab.com/api/v4/projects?search=example.test&per_page=100&page=1': json.dumps(projects), 'https://gitlab.com/api/v4/projects/group%2Fproject/repository/files/README.md/raw?ref=feature%2Freadme': ( @@ -73,13 +83,15 @@ async def test_public_discovery_normalizes_evidence_and_uses_bounded_requests( return [FetcherResponse(responses[url], 200, {})] monkeypatch.setattr(gitlabsearch.Core, 'get_user_agent', staticmethod(lambda: 'UA')) + monkeypatch.setattr(gitlabsearch.AsyncFetcher, 'open_session', fake_open_session) monkeypatch.setattr(gitlabsearch.AsyncFetcher, 'fetch_all', fake_fetch_all) search = gitlabsearch.SearchGitlab('example.test') await search.process(proxy=True) + assert session_proxy == [True] assert requests == [ - {'url': url, 'headers': {'User-agent': 'UA'}, 'proxy': True} + {'url': url, 'headers': {'User-agent': 'UA'}, 'session': session} for url in ( 'https://gitlab.com/api/v4/projects?search=example.test&per_page=100&page=1', 'https://gitlab.com/api/v4/projects/group%2Fproject/repository/files/README.md/raw?ref=feature%2Freadme', diff --git a/tests/discovery/test_huntersearch.py b/tests/discovery/test_huntersearch.py index bf86d611..1a273c81 100644 --- a/tests/discovery/test_huntersearch.py +++ b/tests/discovery/test_huntersearch.py @@ -87,7 +87,17 @@ async def test_hunter_empty_or_malformed_response_returns_no_results( @pytest.mark.asyncio async def test_paid_hunter_search_honors_limit_and_offset(monkeypatch) -> None: - requests: list[tuple[str, bool]] = [] + import contextlib + + requests: list[tuple[str, object]] = [] + session_proxy: list[object] = [] + session = object() + + @contextlib.asynccontextmanager + async def fake_open_session(**kwargs: Any): + session_proxy.append(kwargs.get('proxy')) + yield session + responses = iter( [ {'data': {'plan_name': 'Growth', 'requests': {'searches': {'available': 10, 'used': 0}}}}, @@ -109,8 +119,8 @@ async def test_paid_hunter_search_honors_limit_and_offset(monkeypatch) -> None: ] ) - async def fake_fetch_all(urls, *, proxy=False, **_kwargs): - requests.append((urls[0], proxy)) + async def fake_fetch_all(urls, *, session=None, **_kwargs): + requests.append((urls[0], session)) return [FetcherResponse(body=next(responses), status=200, headers={})] async def no_sleep(_seconds: float) -> None: @@ -118,23 +128,20 @@ async def test_paid_hunter_search_honors_limit_and_offset(monkeypatch) -> None: monkeypatch.setattr(huntersearch.Core, 'hunter_key', lambda: 'test-key') monkeypatch.setattr(huntersearch.Core, 'get_user_agent', lambda: 'test-agent') + monkeypatch.setattr(huntersearch.AsyncFetcher, 'open_session', fake_open_session) monkeypatch.setattr(huntersearch.AsyncFetcher, 'fetch_all', fake_fetch_all) monkeypatch.setattr(huntersearch.asyncio, 'sleep', no_sleep) search = huntersearch.SearchHunter('example.test', 150, 25) await search.process(proxy=True) - assert requests == [ - ('https://api.hunter.io/v2/account?api_key=test-key', True), - ('https://api.hunter.io/v2/email-count?domain=example.test', True), - ( - 'https://api.hunter.io/v2/domain-search?domain=example.test&api_key=test-key&limit=100&offset=25', - True, - ), - ( - 'https://api.hunter.io/v2/domain-search?domain=example.test&api_key=test-key&limit=50&offset=125', - True, - ), + assert session_proxy == [True] + assert all(entry[1] is session for entry in requests) + assert [entry[0] for entry in requests] == [ + 'https://api.hunter.io/v2/account?api_key=test-key', + 'https://api.hunter.io/v2/email-count?domain=example.test', + 'https://api.hunter.io/v2/domain-search?domain=example.test&api_key=test-key&limit=100&offset=25', + 'https://api.hunter.io/v2/domain-search?domain=example.test&api_key=test-key&limit=50&offset=125', ] assert await search.get_emails() == ['alice@example.test', 'bob@example.test'] assert await search.get_hostnames() == ['api.example.test', 'www.example.test'] diff --git a/tests/discovery/test_pentesttools.py b/tests/discovery/test_pentesttools.py index af86ed95..0427a9c3 100644 --- a/tests/discovery/test_pentesttools.py +++ b/tests/discovery/test_pentesttools.py @@ -1,4 +1,5 @@ import asyncio +import contextlib import logging import pytest @@ -52,6 +53,13 @@ async def test_successful_scan_returns_normalized_hostnames_and_ips(monkeypatch) async def no_sleep(*_args, **_kwargs): return None + session = object() + + @contextlib.asynccontextmanager + async def fake_open_session(**_kwargs): + yield session + + monkeypatch.setattr(pentesttools.AsyncFetcher, 'open_session', fake_open_session) monkeypatch.setattr(pentesttools.AsyncFetcher, 'post_fetch', fake_post_fetch) monkeypatch.setattr(pentesttools.AsyncFetcher, 'fetch', fake_fetch) monkeypatch.setattr(pentesttools.asyncio, 'sleep', no_sleep) @@ -76,7 +84,7 @@ async def test_successful_scan_returns_normalized_hostnames_and_ips(monkeypatch) 'tool_params': {'scan_type': 'light', 'web_details': False, 'unresolved_results': True}, }, 'json': True, - 'proxy': False, + 'session': session, } ] assert get_requests == [ @@ -84,13 +92,13 @@ async def test_successful_scan_returns_normalized_hostnames_and_ips(monkeypatch) 'url': 'https://app.pentest-tools.com/api/v2/scans/420323', 'headers': headers, 'json': True, - 'proxy': False, + 'session': session, }, { 'url': 'https://app.pentest-tools.com/api/v2/scans/420323/output', 'headers': headers, 'json': True, - 'proxy': False, + 'session': session, }, ] diff --git a/tests/discovery/test_tombasearch.py b/tests/discovery/test_tombasearch.py index a277bdc0..7224a5e6 100644 --- a/tests/discovery/test_tombasearch.py +++ b/tests/discovery/test_tombasearch.py @@ -109,7 +109,17 @@ async def test_tomba_empty_or_malformed_response_returns_no_results( @pytest.mark.asyncio async def test_paid_tomba_search_uses_documented_pages_and_page_size(monkeypatch) -> None: - requests: list[tuple[str, bool]] = [] + import contextlib + + requests: list[tuple[str, object]] = [] + session_proxy: list[object] = [] + session = object() + + @contextlib.asynccontextmanager + async def fake_open_session(**kwargs: Any): + session_proxy.append(kwargs.get('proxy')) + yield session + responses = iter( [ { @@ -125,8 +135,8 @@ async def test_paid_tomba_search_uses_documented_pages_and_page_size(monkeypatch ] ) - async def fake_fetch_all(urls, *, proxy=False, **_kwargs): - requests.append((urls[0], proxy)) + async def fake_fetch_all(urls, *, session=None, **_kwargs): + requests.append((urls[0], session)) return [FetcherResponse(body=next(responses), status=200, headers={})] async def no_sleep(_seconds: float) -> None: @@ -134,18 +144,21 @@ async def test_paid_tomba_search_uses_documented_pages_and_page_size(monkeypatch monkeypatch.setattr(tombasearch.Core, 'tomba_key', lambda: ('test-key', 'test-secret')) monkeypatch.setattr(tombasearch.Core, 'get_user_agent', lambda: 'test-agent') + monkeypatch.setattr(tombasearch.AsyncFetcher, 'open_session', fake_open_session) monkeypatch.setattr(tombasearch.AsyncFetcher, 'fetch_all', fake_fetch_all) monkeypatch.setattr(tombasearch.asyncio, 'sleep', no_sleep) search = tombasearch.SearchTomba('example.test', 120, 0) await search.process(proxy=True) - assert requests == [ - ('https://api.tomba.io/v1/me', True), - ('https://api.tomba.io/v1/email-count?domain=example.test', True), - ('https://api.tomba.io/v1/domain-search?domain=example.test&limit=50&page=1', True), - ('https://api.tomba.io/v1/domain-search?domain=example.test&limit=50&page=2', True), - ('https://api.tomba.io/v1/domain-search?domain=example.test&limit=50&page=3', True), + assert session_proxy == [True] + assert all(entry[1] is session for entry in requests) + assert [entry[0] for entry in requests] == [ + 'https://api.tomba.io/v1/me', + 'https://api.tomba.io/v1/email-count?domain=example.test', + 'https://api.tomba.io/v1/domain-search?domain=example.test&limit=50&page=1', + 'https://api.tomba.io/v1/domain-search?domain=example.test&limit=50&page=2', + 'https://api.tomba.io/v1/domain-search?domain=example.test&limit=50&page=3', ] emails = await search.get_emails() hostnames = await search.get_hostnames() diff --git a/tests/discovery/test_windvane.py b/tests/discovery/test_windvane.py index c2722ac2..a394a76f 100644 --- a/tests/discovery/test_windvane.py +++ b/tests/discovery/test_windvane.py @@ -42,7 +42,7 @@ async def test_authenticated_results_are_normalized_and_scoped(monkeypatch) -> N }, } - async def fake_post_fetch(url, headers=None, data=None, proxy=False): + async def fake_post_fetch(url, headers=None, data=None, session=None): endpoint = url.rsplit('/', 1)[-1] page = json.loads(data)['page_request']['page'] assert headers['X-Api-Key'] == 'test-key' @@ -87,7 +87,7 @@ async def test_authenticated_unlimited_search_follows_all_endpoint_pagination(mo monkeypatch.setattr(windvane.Core, 'windvane_key', lambda: 'test-key') requests: list[tuple[str, int, int]] = [] - async def fake_post_fetch(url, headers=None, data=None, proxy=False): + async def fake_post_fetch(url, headers=None, data=None, session=None): endpoint = url.rsplit('/', 1)[-1] page_request = json.loads(data)['page_request'] page = page_request['page'] @@ -172,7 +172,7 @@ async def test_finite_limit_stops_each_windvane_endpoint(monkeypatch) -> None: monkeypatch.setattr(windvane.Core, 'windvane_key', lambda: 'test-key') requests: list[tuple[str, int, int]] = [] - async def fake_post_fetch(url, headers=None, data=None, proxy=False): + async def fake_post_fetch(url, headers=None, data=None, session=None): endpoint = url.rsplit('/', 1)[-1] page_request = json.loads(data)['page_request'] requests.append((endpoint, page_request['page'], page_request['count'])) diff --git a/tests/discovery/test_yahoosearch.py b/tests/discovery/test_yahoosearch.py index f62d17d7..61b45ced 100644 --- a/tests/discovery/test_yahoosearch.py +++ b/tests/discovery/test_yahoosearch.py @@ -1,4 +1,5 @@ import asyncio +import contextlib from typing import Any import pytest @@ -10,26 +11,37 @@ from theHarvester.lib.source_execution import SourceExecutionReport @pytest.mark.asyncio async def test_yahoo_uses_exact_pages_and_normalizes_evidence(monkeypatch: pytest.MonkeyPatch) -> None: + import contextlib + requests: list[dict[str, Any]] = [] + session_proxy: list[object] = [] + session = object() + + @contextlib.asynccontextmanager + async def fake_open_session(**kwargs: Any): + session_proxy.append(kwargs.get('proxy')) + yield session async def fake_fetch_all( urls: list[str] | set[str], headers: dict[str, str] | None = None, - proxy: bool = False, + session: object | None = None, **_kwargs: Any, ) -> list[str]: - requests.append({'urls': list(urls), 'headers': headers, 'proxy': proxy}) + requests.append({'urls': list(urls), 'headers': headers, 'session': session}) return [ 'Contact Admin@Example.COM. at Blog.Example.COM.', 'Ignore outsider@example.net and api.example.net', ] monkeypatch.setattr(yahoosearch.Core, 'get_browser_user_agent', staticmethod(lambda: 'UA')) + monkeypatch.setattr(yahoosearch.AsyncFetcher, 'open_session', fake_open_session) monkeypatch.setattr(yahoosearch.AsyncFetcher, 'fetch_all', fake_fetch_all) search = yahoosearch.SearchYahoo('example.com', 20) await search.process(proxy=True) + assert session_proxy == [True] assert requests == [ { 'urls': [ @@ -37,7 +49,7 @@ async def test_yahoo_uses_exact_pages_and_normalizes_evidence(monkeypatch: pytes 'https://search.yahoo.com/search?p=%40example.com&b=10&pz=10', ], 'headers': {'Host': 'search.yahoo.com', 'User-Agent': 'UA'}, - 'proxy': True, + 'session': session, } ] assert set(await search.get_emails()) == {'admin@example.com'} diff --git a/tests/test_mojeek.py b/tests/test_mojeek.py index 3d68f87c..8ecd2c34 100644 --- a/tests/test_mojeek.py +++ b/tests/test_mojeek.py @@ -1,4 +1,5 @@ import asyncio +import contextlib from typing import Any import pytest @@ -7,6 +8,20 @@ from theHarvester.discovery import mojeek from theHarvester.lib.core import FetcherResponse +def install_session(monkeypatch: pytest.MonkeyPatch) -> tuple[object, list[object]]: + """Patch AsyncFetcher.open_session to yield one sentinel session per process call.""" + session = object() + proxies: list[object] = [] + + @contextlib.asynccontextmanager + async def fake_open_session(**kwargs: Any): + proxies.append(kwargs.get('proxy')) + yield session + + monkeypatch.setattr(mojeek.AsyncFetcher, 'open_session', fake_open_session) + return session, proxies + + class TestMojeekSearch: @pytest.mark.asyncio async def test_unlimited_keyless_stops_when_provider_repeats_a_page(self, monkeypatch: pytest.MonkeyPatch) -> None: @@ -127,6 +142,7 @@ class TestMojeekSearch: self, monkeypatch: pytest.MonkeyPatch, ) -> None: + session, session_proxies = install_session(monkeypatch) calls: list[dict[str, Any]] = [] delays: list[float] = [] responses = iter( @@ -160,14 +176,15 @@ class TestMojeekSearch: report = await search.process(proxy=True) + assert session_proxies == [True] assert [call['url'] for call in calls] == [ 'https://www.mojeek.com/search?q=example.com&s=0', 'https://www.mojeek.com/search?q=example.com&s=10', ] + assert all(call['session'] is session for call in calls) assert all(call['include_metadata'] is True for call in calls) assert all(call['headers'] == {'User-Agent': 'UA'} for call in calls) assert all(call['follow_redirects'] is False for call in calls) - assert all(call['proxy'] is True for call in calls) assert delays == [1.0] assert await search.get_hostnames() == ['docs.example.com', 'example.com'] assert await search.get_emails() == {'admin@example.com'} @@ -240,6 +257,7 @@ class TestMojeekSearch: @pytest.mark.asyncio async def test_failed_keyed_api_does_not_fall_back_to_scraping(self, monkeypatch: pytest.MonkeyPatch) -> None: + session, session_proxies = install_session(monkeypatch) calls: list[dict[str, Any]] = [] async def fake_fetch_all(urls: list[str], **kwargs: Any) -> list[FetcherResponse]: @@ -256,9 +274,10 @@ class TestMojeekSearch: report = await search.process(proxy=True) + assert session_proxies == [True] assert len(calls) == 1 + assert calls[0]['session'] is session assert calls[0]['include_metadata'] is True - assert calls[0]['proxy'] is True assert report.status == 'failed' assert report.stop_reason == 'access-denied' assert await search.get_hostnames() == [] @@ -268,6 +287,7 @@ class TestMojeekSearch: self, monkeypatch: pytest.MonkeyPatch, ) -> None: + session, _session_proxies = install_session(monkeypatch) requests: list[dict[str, Any]] = [] responses = iter( [ @@ -308,15 +328,15 @@ class TestMojeekSearch: assert requests == [ { 'urls': ['https://api.mojeek.com/search?api_key=test-key&q=example.com&fmt=json&s=1'], + 'session': session, 'headers': {'User-Agent': 'UA'}, - 'proxy': True, 'json': True, 'include_metadata': True, }, { 'urls': ['https://api.mojeek.com/search?api_key=test-key&q=example.com&fmt=json&s=11'], + 'session': session, 'headers': {'User-Agent': 'UA'}, - 'proxy': True, 'json': True, 'include_metadata': True, }, diff --git a/theHarvester/discovery/arquivo.py b/theHarvester/discovery/arquivo.py index 3cfeb648..c3f9596f 100644 --- a/theHarvester/discovery/arquivo.py +++ b/theHarvester/discovery/arquivo.py @@ -21,6 +21,11 @@ class SearchArquivo: async def process(self, proxy: bool = False) -> SourceExecutionReport | None: self.proxy = proxy + headers = {'User-agent': Core.get_user_agent()} + async with AsyncFetcher.open_session(headers=headers, proxy=self.proxy) as session: + return await self._search(session, headers) + + async def _search(self, session, headers: dict[str, str]) -> SourceExecutionReport | None: offset = 0 previous_page = None report = None @@ -39,8 +44,8 @@ class SearchArquivo: try: responses: list[FetcherResponse | None] = await AsyncFetcher.fetch_all( [f'https://arquivo.pt/wayback/cdx?{query}'], - headers={'User-agent': Core.get_user_agent()}, - proxy=self.proxy, + headers=headers, + session=session, include_metadata=True, ) except asyncio.CancelledError: diff --git a/theHarvester/discovery/certspottersearch.py b/theHarvester/discovery/certspottersearch.py index 0f280d28..6c7c6e5f 100644 --- a/theHarvester/discovery/certspottersearch.py +++ b/theHarvester/discovery/certspottersearch.py @@ -24,7 +24,7 @@ class SearchCertspoter: status: SourceReportStatus = 'rate-limited' if rate_limited else 'partial' self._report = SourceExecutionReport(status, reason) - async def do_search(self) -> None: + async def do_search(self, session) -> None: base_url = 'https://api.certspotter.com/v1/issuances' cursor = None seen_cursors: set[str] = set() @@ -39,7 +39,7 @@ class SearchCertspoter: params['after'] = cursor responses = await AsyncFetcher.fetch_all( - [f'{base_url}?{urlencode(params)}'], json=True, proxy=self.proxy, include_metadata=True + [f'{base_url}?{urlencode(params)}'], json=True, session=session, include_metadata=True ) if not responses: self._mark_incomplete('no-response') @@ -140,6 +140,7 @@ class SearchCertspoter: async def process(self, proxy: bool = False) -> SourceExecutionReport | None: self.proxy = proxy self._report = None - await self.do_search() + async with AsyncFetcher.open_session(proxy=self.proxy) as session: + await self.do_search(session) logger.info('\tSearching results.') return self._report diff --git a/theHarvester/discovery/commoncrawl.py b/theHarvester/discovery/commoncrawl.py index 4a25c38b..04bb0aea 100644 --- a/theHarvester/discovery/commoncrawl.py +++ b/theHarvester/discovery/commoncrawl.py @@ -106,14 +106,14 @@ class SearchCommoncrawl: selected.append(entry) return selected - async def do_search(self) -> SourceExecutionReport | None: + async def do_search(self, session) -> SourceExecutionReport | None: try: if self.limit == 0: return None headers = {'User-agent': Core.get_user_agent()} catalog_response = await AsyncFetcher.fetch_all( - [f'{self.hostname}/collinfo.json'], headers=headers, proxy=self.proxy, json=True + [f'{self.hostname}/collinfo.json'], headers=headers, session=session, json=True ) if not catalog_response or not isinstance(catalog_response[0], list) or not catalog_response[0]: logger.error('Common Crawl API error: invalid index catalog') @@ -149,7 +149,7 @@ class SearchCommoncrawl: try: query_had_errors = False count_url = f'{endpoint}?{urlencode({"url": query, "output": "json", "pageSize": self.PAGE_SIZE, "showNumPages": "true"})}' - count_response = await AsyncFetcher.fetch_all([count_url], headers=headers, proxy=self.proxy) + count_response = await AsyncFetcher.fetch_all([count_url], headers=headers, session=session) count_payload = json.loads(count_response[0]) page_count = count_payload.get('pages') if isinstance(count_payload, dict) else None if isinstance(page_count, bool) or not isinstance(page_count, int) or page_count < 0: @@ -165,7 +165,7 @@ class SearchCommoncrawl: return None page_url = f'{endpoint}?{urlencode({"url": query, "output": "json", "pageSize": self.PAGE_SIZE, "page": first_page, "limit": min(remaining, self.MAX_RECORDS_PER_REQUEST)})}' first_page += 1 - responses = await AsyncFetcher.fetch_all([page_url], headers=headers, proxy=self.proxy) + responses = await AsyncFetcher.fetch_all([page_url], headers=headers, session=session) if not isinstance(responses, list) or not responses: raise ValueError('invalid page response') try: @@ -216,8 +216,11 @@ class SearchCommoncrawl: async def process(self, proxy: bool = False) -> SourceExecutionReport | None: self.proxy = proxy try: - async with asyncio.timeout(self.RUNTIME_SECONDS): - return await self.do_search() + async with ( + AsyncFetcher.open_session(headers={'User-agent': Core.get_user_agent()}, proxy=self.proxy) as session, + asyncio.timeout(self.RUNTIME_SECONDS), + ): + return await self.do_search(session) except TimeoutError: logger.info( f'Common Crawl runtime limit reached after {self.RUNTIME_SECONDS:g}s; preserved {len(self.totalhosts)} hosts' diff --git a/theHarvester/discovery/gitlabsearch.py b/theHarvester/discovery/gitlabsearch.py index 4a867de2..02ec3e14 100644 --- a/theHarvester/discovery/gitlabsearch.py +++ b/theHarvester/discovery/gitlabsearch.py @@ -80,6 +80,7 @@ class SearchGitlab: async def _fetch_page( self, + session, endpoint: str, term: str, page: int, @@ -89,7 +90,7 @@ class SearchGitlab: response = await AsyncFetcher.fetch_all( [url], headers={'User-agent': Core.get_user_agent()}, - proxy=self.proxy, + session=session, json=True, include_metadata=True, ) @@ -110,7 +111,7 @@ class SearchGitlab: next_page = str(page + 1) if len(records) >= per_page else None return records, next_page, None - async def search_projects(self) -> SourceExecutionReport | None: + async def search_projects(self, session) -> SourceExecutionReport | None: """Search GitLab projects for references to the target domain.""" try: headers = {'User-agent': Core.get_user_agent()} @@ -123,7 +124,7 @@ class SearchGitlab: seen_cursors: set[str] = set() while self.limit is None or records_seen < self.limit: per_page = min(100, self.limit - records_seen) if self.limit is not None else 100 - projects, next_page, page_report = await self._fetch_page('projects', term, page, per_page) + projects, next_page, page_report = await self._fetch_page(session, 'projects', term, page, per_page) if page_report is not None: report = self._combine_reports(report, page_report) break @@ -151,7 +152,7 @@ class SearchGitlab: f'/repository/files/README.md/raw?ref={quote(default_branch, safe="")}' ) try: - readme_response = await AsyncFetcher.fetch_all([readme_url], headers=headers, proxy=self.proxy) + readme_response = await AsyncFetcher.fetch_all([readme_url], headers=headers, session=session) if readme_response and readme_response[0]: readme_text = ( readme_response[0] if isinstance(readme_response[0], str) else str(readme_response[0]) @@ -181,7 +182,7 @@ class SearchGitlab: logger.info(f'GitLab API projects search error: {e}') return SourceExecutionReport('failed', 'transport-error') - async def search_users(self) -> SourceExecutionReport | None: + async def search_users(self, session) -> SourceExecutionReport | None: """Search GitLab users for references to the target domain.""" try: page = 1 @@ -190,7 +191,7 @@ class SearchGitlab: seen_cursors: set[str] = set() while self.limit is None or records_seen < self.limit: per_page = min(100, self.limit - records_seen) if self.limit is not None else 100 - users, next_page, report = await self._fetch_page('users', self.word, page, per_page) + users, next_page, report = await self._fetch_page(session, 'users', self.word, page, per_page) if report is not None: return report signature = json.dumps(users, sort_keys=True, default=str) @@ -240,9 +241,9 @@ class SearchGitlab: logger.info(f'GitLab API users search error: {e}') return SourceExecutionReport('failed', 'transport-error') - async def do_search(self) -> SourceExecutionReport | None: - project_report = await self.search_projects() - user_report = await self.search_users() + async def do_search(self, session) -> SourceExecutionReport | None: + project_report = await self.search_projects(session) + user_report = await self.search_users(session) return self._combine_reports(project_report, user_report) async def get_hostnames(self) -> set: @@ -256,4 +257,5 @@ class SearchGitlab: async def process(self, proxy: bool = False) -> SourceExecutionReport | None: self.proxy = proxy - return await self.do_search() + async with AsyncFetcher.open_session(proxy=self.proxy) as session: + return await self.do_search(session) diff --git a/theHarvester/discovery/huntersearch.py b/theHarvester/discovery/huntersearch.py index 318758ce..b4928742 100644 --- a/theHarvester/discovery/huntersearch.py +++ b/theHarvester/discovery/huntersearch.py @@ -29,11 +29,11 @@ class SearchHunter: self.hostnames: list = [] self.emails: list = [] - async def _fetch_json(self, url: str, headers: dict[str, str]) -> dict | SourceExecutionReport: + async def _fetch_json(self, url: str, headers: dict[str, str], session) -> dict | SourceExecutionReport: response = await AsyncFetcher.fetch_all( [url], headers=headers, - proxy=self.proxy, + session=session, json=True, include_metadata=True, ) @@ -49,11 +49,11 @@ class SearchHunter: return SourceExecutionReport('failed', 'invalid-response') return metadata.body - async def do_search(self) -> SourceExecutionReport | None: + async def do_search(self, session) -> SourceExecutionReport | None: # First determine if a user account is not a free account, this call is free headers = {'User-Agent': Core.get_user_agent()} acc_info_url = f'https://api.hunter.io/v2/account?api_key={self.key}' - response = await self._fetch_json(acc_info_url, headers) + response = await self._fetch_json(acc_info_url, headers, session) if isinstance(response, SourceExecutionReport): return response is_free = 'plan_name' in response['data'].keys() and response['data']['plan_name'].lower() == 'free' @@ -63,7 +63,7 @@ class SearchHunter: response['data']['requests']['searches']['available'] - response['data']['requests']['searches']['used'] ) if is_free: - response = await self._fetch_json(self.database, headers) + response = await self._fetch_json(self.database, headers, session) if isinstance(response, SourceExecutionReport): return response self.emails, self.hostnames = await self.parse_resp(json_resp=response) @@ -79,7 +79,7 @@ class SearchHunter: # As the most emails you can get within one query are 100 # This is only done where paid accounts are in play hunter_dinfo_url = f'https://api.hunter.io/v2/email-count?domain={self.word}' - response = await self._fetch_json(hunter_dinfo_url, headers) + response = await self._fetch_json(hunter_dinfo_url, headers, session) if isinstance(response, SourceExecutionReport): return response available_results = max(0, response['data']['total'] - self.start) @@ -101,7 +101,7 @@ class SearchHunter: for offset in range(self.start, result_end, 100): page_limit = min(100, result_end - offset) req_url = f'https://api.hunter.io/v2/domain-search?domain={self.word}&api_key={self.key}&limit={page_limit}&offset={offset}' - response = await self._fetch_json(req_url, headers) + response = await self._fetch_json(req_url, headers, session) if isinstance(response, SourceExecutionReport): return response temp_emails, temp_hostnames = await self.parse_resp(response) @@ -129,7 +129,8 @@ class SearchHunter: async def process(self, proxy: bool = False) -> SourceExecutionReport | None: self.proxy = proxy try: - return await self.do_search() # Only need to do it once. + async with AsyncFetcher.open_session(proxy=self.proxy) as session: + return await self.do_search(session) # Only need to do it once. except AttributeError, KeyError, TypeError: logger.info('Hunter returned malformed data') return SourceExecutionReport('failed', 'invalid-response') diff --git a/theHarvester/discovery/mojeek.py b/theHarvester/discovery/mojeek.py index b7d08a86..c81863bf 100644 --- a/theHarvester/discovery/mojeek.py +++ b/theHarvester/discovery/mojeek.py @@ -75,14 +75,14 @@ class SearchMojeek: parsed_results.append(f'{url} {title} {description}') return parsed_results - async def _search_api(self, headers: dict[str, str]) -> None: + async def _search_api(self, headers: dict[str, str], session) -> None: if self.limit is None: seen_pages: set[tuple[str, ...]] = set() offset = 1 while True: url = f'https://{self.api_server}/search?api_key={self.api_key}&q={self.word}&fmt=json&s={offset}' responses = await AsyncFetcher.fetch_all( - [url], headers=headers, proxy=self.proxy, json=True, include_metadata=True + [url], headers=headers, session=session, json=True, include_metadata=True ) if len(responses) != 1 or not isinstance(responses[0], FetcherResponse): self._stop('failed', 'transport-error') @@ -106,7 +106,7 @@ class SearchMojeek: offset = 1 while offset <= result_limit: url = f'https://{self.api_server}/search?api_key={self.api_key}&q={self.word}&fmt=json&s={offset}' - responses = await AsyncFetcher.fetch_all([url], headers=headers, proxy=self.proxy, json=True, include_metadata=True) + responses = await AsyncFetcher.fetch_all([url], headers=headers, session=session, json=True, include_metadata=True) if len(responses) != 1 or not isinstance(responses[0], FetcherResponse): self._stop('failed', 'transport-error') return @@ -124,7 +124,7 @@ class SearchMojeek: offset += 10 logger.info('[*] Mojeek: API search completed successfully.') - async def _search_keyless(self, headers: dict[str, str]) -> None: + async def _search_keyless(self, headers: dict[str, str], session) -> None: seen_bodies: set[str] = set() offset = 0 page = 0 @@ -133,9 +133,9 @@ class SearchMojeek: if page: await asyncio.sleep(self.REQUEST_DELAY_SECONDS) response = await AsyncFetcher.fetch( + session=session, url=url, headers=headers, - proxy=self.proxy, request_timeout=60, follow_redirects=False, include_metadata=True, @@ -175,19 +175,20 @@ class SearchMojeek: offset += 10 page += 1 - async def do_search(self) -> SourceExecutionReport | None: + async def do_search(self, session) -> SourceExecutionReport | None: self._report = None user_agent = Core.get_user_agent() if self.api_key else Core.get_browser_user_agent() headers = {'User-Agent': user_agent} if self.api_key: - await self._search_api(headers) + await self._search_api(headers, session) else: - await self._search_keyless(headers) + await self._search_keyless(headers, session) return self._report async def process(self, proxy: bool = False) -> SourceExecutionReport | None: self.proxy = proxy - return await self.do_search() + async with AsyncFetcher.open_session(proxy=self.proxy) as session: + return await self.do_search(session) async def get_emails(self): rawres = myparser.Parser(self.total_results, self.word) diff --git a/theHarvester/discovery/pentesttools.py b/theHarvester/discovery/pentesttools.py index 85862428..f456816e 100644 --- a/theHarvester/discovery/pentesttools.py +++ b/theHarvester/discovery/pentesttools.py @@ -33,15 +33,15 @@ class SearchPentestTools: return None return data - async def poll(self, scan_id: int) -> SourceExecutionReport | None: + async def poll(self, scan_id: int, session) -> SourceExecutionReport | None: for _attempt in range(10): await asyncio.sleep(3) status = self._response_data( await AsyncFetcher.fetch( + session=session, url=f'{self.api}/scans/{scan_id}', headers=self.headers, json=True, - proxy=self.proxy, ) ) if status is None: @@ -55,10 +55,10 @@ class SearchPentestTools: if status_name == 'finished': output = self._response_data( await AsyncFetcher.fetch( + session=session, url=f'{self.api}/scans/{scan_id}/output', headers=self.headers, json=True, - proxy=self.proxy, ) ) if output is not None: @@ -103,7 +103,7 @@ class SearchPentestTools: async def get_ips(self) -> set[str]: return self.totalips - async def do_search(self) -> SourceExecutionReport | None: + async def do_search(self, session) -> SourceExecutionReport | None: # Pentest-Tools documents Subdomain Finder as tool 20: # https://pentest-tools.com/docs/api-reference/scans/start-a-scan subdomain_payload = { @@ -121,7 +121,7 @@ class SearchPentestTools: headers=self.headers, json_body=subdomain_payload, json=True, - proxy=self.proxy, + session=session, ) ) if response is None: @@ -130,12 +130,13 @@ class SearchPentestTools: if not isinstance(scan_id, int) or isinstance(scan_id, bool): logger.info('Pentest-Tools returned a malformed start response') return SourceExecutionReport('failed', 'invalid-response') - return await self.poll(scan_id) + return await self.poll(scan_id, session) async def process(self, proxy: bool = False) -> SourceExecutionReport | None: self.proxy = proxy try: - return await self.do_search() # Only need to do it once. + async with AsyncFetcher.open_session(headers=self.headers, proxy=self.proxy) as session: + return await self.do_search(session) # Only need to do it once. except asyncio.CancelledError: raise except Exception: diff --git a/theHarvester/discovery/tombasearch.py b/theHarvester/discovery/tombasearch.py index 92aa6218..0bc7489c 100644 --- a/theHarvester/discovery/tombasearch.py +++ b/theHarvester/discovery/tombasearch.py @@ -27,11 +27,11 @@ class SearchTomba: self.hostnames: list = [] self.emails: list = [] - async def _fetch_json(self, url: str, headers: dict[str, str]) -> dict | SourceExecutionReport: + async def _fetch_json(self, url: str, headers: dict[str, str], session) -> dict | SourceExecutionReport: response = await AsyncFetcher.fetch_all( [url], headers=headers, - proxy=self.proxy, + session=session, json=True, include_metadata=True, ) @@ -47,7 +47,7 @@ class SearchTomba: return SourceExecutionReport('failed', 'invalid-response') return metadata.body - async def do_search(self) -> SourceExecutionReport | None: + async def do_search(self, session) -> SourceExecutionReport | None: # First determine if a user account is not a free account, this call is free headers = { 'User-Agent': Core.get_user_agent(), @@ -55,7 +55,7 @@ class SearchTomba: 'X-Tomba-Secret': self.key[1], } acc_info_url = 'https://api.tomba.io/v1/me' - response = await self._fetch_json(acc_info_url, headers) + response = await self._fetch_json(acc_info_url, headers, session) if isinstance(response, SourceExecutionReport): return response is_free = 'name' in response['data']['pricing'].keys() and response['data']['pricing']['name'].lower() == 'free' @@ -70,7 +70,7 @@ class SearchTomba: total_results = self.limit else: tomba_counter = f'https://api.tomba.io/v1/email-count?domain={self.word}' - response = await self._fetch_json(tomba_counter, headers) + response = await self._fetch_json(tomba_counter, headers, session) if isinstance(response, SourceExecutionReport): return response available_results = max(0, response['data']['total'] - self.start) @@ -91,7 +91,7 @@ class SearchTomba: pages_to_fetch = min(total_number_reqs, max(total_requests_avail, 0)) for page in range(first_page, first_page + pages_to_fetch): req_url = f'https://api.tomba.io/v1/domain-search?domain={self.word}&limit={page_size}&page={page}' - response = await self._fetch_json(req_url, headers) + response = await self._fetch_json(req_url, headers, session) if isinstance(response, SourceExecutionReport): return response skip = first_page_skip if page == first_page else 0 @@ -139,7 +139,8 @@ class SearchTomba: async def process(self, proxy: bool = False) -> SourceExecutionReport | None: self.proxy = proxy try: - return await self.do_search() # Only need to do it once. + async with AsyncFetcher.open_session(proxy=self.proxy) as session: + return await self.do_search(session) # Only need to do it once. except AttributeError, KeyError, TypeError: logger.info('Tomba returned malformed data') return SourceExecutionReport('failed', 'invalid-response') diff --git a/theHarvester/discovery/waybackarchive.py b/theHarvester/discovery/waybackarchive.py index 812db123..84490910 100644 --- a/theHarvester/discovery/waybackarchive.py +++ b/theHarvester/discovery/waybackarchive.py @@ -3,7 +3,7 @@ import logging from urllib.parse import unquote_plus, urlencode, urlsplit from theHarvester.discovery.provider_response import provider_http_error -from theHarvester.lib.core import AsyncFetcher, Core, FetcherResponse +from theHarvester.lib.core import AsyncFetcher, Core, FetcherResponse, ProxyUnavailableError from theHarvester.lib.source_execution import SourceExecutionReport logger = logging.getLogger(__name__) @@ -57,7 +57,7 @@ class SearchWaybackarchive: continuation_lines = [line.strip() for line in continuation.splitlines() if line.strip()] return body.splitlines(), continuation_lines[0] if len(continuation_lines) == 1 else None - async def _search_pattern(self, pattern: str, headers: dict[str, str]) -> SourceExecutionReport | None: + async def _search_pattern(self, pattern: str, headers: dict[str, str], session) -> SourceExecutionReport | None: resume_key: str | None = None seen_resume_keys: set[str] = set() page_number = 0 @@ -74,7 +74,7 @@ class SearchWaybackarchive: query['resumeKey'] = unquote_plus(resume_key) url = f'{self.hostname}/cdx/search/cdx?{urlencode(query)}' - response = await AsyncFetcher.fetch_all([url], headers=headers, proxy=self.proxy, include_metadata=True) + response = await AsyncFetcher.fetch_all([url], headers=headers, session=session, include_metadata=True) if not response or not isinstance(response, list): logger.info(f'Wayback Archive returned an invalid response container for pattern {pattern}') return SourceExecutionReport('failed', 'invalid-response') @@ -119,29 +119,32 @@ class SearchWaybackarchive: try: headers = {'User-agent': Core.get_user_agent()} degraded: SourceExecutionReport | None = None - try: - async with asyncio.timeout(self.RUNTIME_SECONDS): - for pattern in (f'*.{self.word}', f'{self.word}/*'): - try: - outcome = await self._search_pattern(pattern, headers) - except Exception as e: - degraded = degraded or SourceExecutionReport('failed', 'request-error') - logger.info(f'Wayback Archive API error for pattern {pattern}: {e}') - continue - if outcome is not None and outcome.status == 'completed': - if degraded is None: - return outcome - break - if outcome is not None: - degraded = degraded or outcome - except TimeoutError: - logger.info( - f'Wayback Archive runtime limit reached after {self.RUNTIME_SECONDS:g}s; ' - f'preserved {len(self.totalhosts)} hosts' - ) - return SourceExecutionReport('partial' if self.totalhosts else 'failed', 'runtime-limit') + async with AsyncFetcher.open_session(headers=headers, proxy=self.proxy) as session: + try: + async with asyncio.timeout(self.RUNTIME_SECONDS): + for pattern in (f'*.{self.word}', f'{self.word}/*'): + try: + outcome = await self._search_pattern(pattern, headers, session) + except Exception as e: + degraded = degraded or SourceExecutionReport('failed', 'request-error') + logger.info(f'Wayback Archive API error for pattern {pattern}: {e}') + continue + if outcome is not None and outcome.status == 'completed': + if degraded is None: + return outcome + break + if outcome is not None: + degraded = degraded or outcome + except TimeoutError: + logger.info( + f'Wayback Archive runtime limit reached after {self.RUNTIME_SECONDS:g}s; ' + f'preserved {len(self.totalhosts)} hosts' + ) + return SourceExecutionReport('partial' if self.totalhosts else 'failed', 'runtime-limit') if degraded is not None: return degraded + except ProxyUnavailableError: + return SourceExecutionReport('failed', 'proxy-unavailable') except Exception as e: logger.info(f'Wayback Archive API error: {e}') return SourceExecutionReport('failed', 'unexpected-error') diff --git a/theHarvester/discovery/windvane.py b/theHarvester/discovery/windvane.py index ae89c535..929ca6f4 100644 --- a/theHarvester/discovery/windvane.py +++ b/theHarvester/discovery/windvane.py @@ -98,6 +98,7 @@ class SearchWindvane: async def _paginate( self, headers: dict[str, str], + session, endpoint: str, query: dict[str, str], page_size: int, @@ -115,7 +116,7 @@ class SearchWindvane: url, headers=headers, data=json.dumps(request_data, separators=(',', ':')), - proxy=self.proxy, + session=session, ) if not response: return SourceExecutionReport('failed', 'transport-error') @@ -154,7 +155,7 @@ class SearchWindvane: page = next_page return None - async def do_search(self) -> SourceExecutionReport | None: + async def do_search(self, session) -> SourceExecutionReport | None: """Query the Windvane endpoints used by this source.""" try: headers = {'User-agent': Core.get_user_agent(), 'Content-Type': 'application/json', 'Accept': 'application/json'} @@ -165,14 +166,14 @@ class SearchWindvane: # With API key, use full API endpoints reports = [ - await self._search_subdomains(headers), - await self._search_dns_history(headers), - await self._search_emails(headers), + await self._search_subdomains(headers, session), + await self._search_dns_history(headers, session), + await self._search_emails(headers, session), ] else: # Without API key, use the provider's limited endpoint only. logger.info('[*] Windvane API key not found. Using limited unauthenticated access.') - reports = [await self._search_subdomains_limited(headers)] + reports = [await self._search_subdomains_limited(headers, session)] retained = [report for report in reports if report is not None] return next((report for report in retained if report.status != 'completed'), retained[0] if retained else None) @@ -181,17 +182,18 @@ class SearchWindvane: logger.info(f'Windvane API error: {e}') return SourceExecutionReport('failed', 'transport-error') - async def _search_subdomains(self, headers: dict[str, str]) -> SourceExecutionReport | None: + async def _search_subdomains(self, headers: dict[str, str], session) -> SourceExecutionReport | None: """Search for subdomains with ``/ListSubDomain``.""" return await self._paginate( headers, + session, 'ListSubDomain', {'domain': self.word}, 30, lambda item: self._add_host(item.get('domain')) if isinstance(item, dict) else None, ) - async def _search_dns_history(self, headers: dict[str, str]) -> SourceExecutionReport | None: + async def _search_dns_history(self, headers: dict[str, str], session) -> SourceExecutionReport | None: """Collect subdomains and IP addresses from ``/ListDNS`` history.""" def consume(record: object) -> None: @@ -206,9 +208,9 @@ class SearchWindvane: ): self.totalips.add(answer) - return await self._paginate(headers, 'ListDNS', {'domain': self.word}, 30, consume) + return await self._paginate(headers, session, 'ListDNS', {'domain': self.word}, 30, consume) - async def _search_emails(self, headers: dict[str, str]) -> SourceExecutionReport | None: + async def _search_emails(self, headers: dict[str, str], session) -> SourceExecutionReport | None: """Search for email addresses with ``/ListEmail``.""" def consume(item: object) -> None: @@ -216,12 +218,13 @@ class SearchWindvane: self._add_email(item.get('email')) self._add_host(item.get('domain')) - return await self._paginate(headers, 'ListEmail', {'email': self.word}, 50, consume) + return await self._paginate(headers, session, 'ListEmail', {'email': self.word}, 50, consume) - async def _search_subdomains_limited(self, headers: dict[str, str]) -> SourceExecutionReport | None: + async def _search_subdomains_limited(self, headers: dict[str, str], session) -> SourceExecutionReport | None: """Search the unauthenticated subdomain endpoints.""" report = await self._paginate( headers, + session, 'ListSubDomain', {'domain': self.word}, 10, @@ -266,5 +269,5 @@ class SearchWindvane: self.proxy = proxy # API key is already set via _get_api_key() method - - return await self.do_search() + async with AsyncFetcher.open_session(proxy=self.proxy) as session: + return await self.do_search(session) diff --git a/theHarvester/discovery/yahoosearch.py b/theHarvester/discovery/yahoosearch.py index 1c7ba72d..4e2b8edc 100644 --- a/theHarvester/discovery/yahoosearch.py +++ b/theHarvester/discovery/yahoosearch.py @@ -9,71 +9,79 @@ class SearchYahoo: def __init__(self, word, limit: int | None) -> None: self.word = word self.total_results = '' + self._pages: list[str] = [] self.server = 'search.yahoo.com' self.limit = limit self.proxy = False def _page(self, response: object) -> tuple[str | None, SourceExecutionReport | None]: if response is None: - return None, SourceExecutionReport('partial' if self.total_results else 'failed', 'transport-error') + return None, SourceExecutionReport('partial' if self._pages else 'failed', 'transport-error') if isinstance(response, FetcherResponse): if not 200 <= response.status < 300: - return None, SourceExecutionReport('partial' if self.total_results else 'failed', f'http-{response.status}') + return None, SourceExecutionReport('partial' if self._pages else 'failed', f'http-{response.status}') response = response.body if not isinstance(response, str): - return None, SourceExecutionReport('partial' if self.total_results else 'failed', 'invalid-response') + return None, SourceExecutionReport('partial' if self._pages else 'failed', 'invalid-response') normalized = response.casefold() if not response.strip() or 'no results' in normalized or 'no-results' in normalized: return '', None if 'captcha' in normalized or 'verify you are human' in normalized: - return None, SourceExecutionReport('partial' if self.total_results else 'failed', 'security-verification') + return None, SourceExecutionReport('partial' if self._pages else 'failed', 'security-verification') if 'access denied' in normalized or 'temporarily blocked' in normalized: - return None, SourceExecutionReport('partial' if self.total_results else 'failed', 'access-denied') + return None, SourceExecutionReport('partial' if self._pages else 'failed', 'access-denied') return response, None - async def do_search(self) -> SourceExecutionReport | None: + async def do_search(self, session) -> SourceExecutionReport | None: base_url = f'https://{self.server}/search?p=%40{self.word}&b=xx&pz=10' headers = {'Host': self.server, 'User-Agent': Core.get_browser_user_agent()} + self._pages = [] + report: SourceExecutionReport | None = None if self.limit is None: seen_pages: set[str] = set() offset = 0 while True: response = await AsyncFetcher.fetch( + session=session, url=base_url.replace('xx', str(offset)), headers=headers, - proxy=self.proxy, include_metadata=True, ) - body, report = self._page(response) - if report is not None: - return report + body, page_report = self._page(response) + if page_report is not None: + report = page_report + break if not body: - return None + break if body in seen_pages: - return SourceExecutionReport('partial', 'repeated-page') + report = SourceExecutionReport('partial', 'repeated-page') + break seen_pages.add(body) - self.total_results += f' {body}' + self._pages.append(body) offset += 10 - urls = [base_url.replace('xx', str(num)) for num in range(0, self.limit, 10)] - responses = await AsyncFetcher.fetch_all(urls, headers=headers, proxy=self.proxy, include_metadata=True) - if urls and not responses: - return SourceExecutionReport('failed', 'transport-error') - seen_finite_pages: set[str] = set() - for response in responses: - body, report = self._page(response) - if report is not None: - return report - if not body: - break - if body in seen_finite_pages: - return SourceExecutionReport('partial', 'repeated-page') - seen_finite_pages.add(body) - self.total_results += f' {body}' - return None + else: + urls = [base_url.replace('xx', str(num)) for num in range(0, self.limit, 10)] + responses = await AsyncFetcher.fetch_all(urls, headers=headers, session=session, include_metadata=True) + seen_finite_pages: set[str] = set() + for response in responses: + body, page_report = self._page(response) + if page_report is not None: + report = page_report + break + if not body: + break + if body in seen_finite_pages: + report = SourceExecutionReport('partial', 'repeated-page') + break + seen_finite_pages.add(body) + self._pages.append(body) + self.total_results = ' '.join(self._pages) + return report async def process(self, proxy: bool = False) -> SourceExecutionReport | None: self.proxy = proxy - return await self.do_search() + async with AsyncFetcher.open_session(headers={'Host': self.server}, proxy=self.proxy) as session: + return await self.do_search(session) async def get_emails(self): rawres = myparser.Parser(self.total_results, self.word) @@ -83,8 +91,7 @@ class SearchYahoo: for email in toparse_emails: email = str(email) if '-' in email and email[0].isdigit() and email.index('-') <= 9: - while email[0] == '-' or email[0].isdigit(): - email = email[1:] + email = email.lstrip('-0123456789') emails.add(email) return list(emails)