diff --git a/README.md b/README.md index 115a4af..8d22357 100644 --- a/README.md +++ b/README.md @@ -215,6 +215,134 @@ Client's class/method | API endpoint | Produced outcome *Empty cells inherit the value from the previous row.* +## Helper: `fetch_alert_targets` + +The alerts stream endpoint returns lightweight `StreamingWatcherAlert` objects that carry metadata +(status, watcher IDs, timestamps, highlights) but not the actual content the alert refers to. +`fetch_alert_targets` resolves each alert's API link in parallel and pairs it with the fully +fetched target object — a report, forum post, credential, indicator, or any other supported type. + +```python +from verity471 import fetch_alert_targets, AlertTarget +``` + +### Parameters + +| Parameter | Type | Default | Description | +|---|---|---|---| +| `alerts_response` | `StreamingAlertsResponse` | *(required)* | The page returned by `AlertsApi.get_alerts_stream()`. | +| `api_client` | `ApiClient` | *(required)* | An active `ApiClient` instance (must share credentials with the alerts call). | +| `raise_on_error` | `bool` | `False` | When `True`, re-raise exceptions instead of logging and skipping the alert. | + +### Returns + +A list of `AlertTarget` objects in the same order as `alerts_response.alerts`. + +Each `AlertTarget` exposes: + +| Attribute | Type | Description | +|---|---|---| +| `.alert` | `StreamingWatcherAlert` | The original alert object (status, watcher IDs, timestamps, highlights, etc.). | +| `.target` | model instance or `None` | The resolved API object (report, post, credential, …). `None` when the URL could not be mapped to a known SDK route. | +| `.target_summary` | `str \| None` | A compact, human-readable one-liner describing the target. | +| `.watcher` | `GetWatcherResponse \| None` | The full watcher object that triggered this alert (name, DSL query, mute status, etc.). `None` if the watcher ID was not found in the user's watcher list. | +| `.watcher_group` | `GetWatcherGroupResponse \| None` | The full watcher group object the watcher belongs to (name, description, etc.). `None` if not found. | + +Watcher and group data is fetched once per `fetch_alert_targets` call (two parallel API requests) and +shared by reference across all `AlertTarget` objects — there is no duplication even when many alerts +share the same watcher. + +### Example usage + +```python +import verity471 + +configuration = verity471.Configuration( + username="your_username", + password="your_password", +) + +with verity471.ApiClient(configuration) as api_client: + alerts_api = verity471.AlertsApi(api_client) + alerts_response = alerts_api.get_alerts_stream(size=10) + + targets = verity471.fetch_alert_targets(alerts_response, api_client) + for t in targets: + watcher_name = t.watcher.name if t.watcher else None + group_name = t.watcher_group.name if t.watcher_group else None + print(t.alert.source_type, t.alert.status, watcher_name, group_name, t.target_summary) +``` + +### Example output + +``` +fintel read threat_actor Ransomware actors [Fintel] Threat Landscape: Q1 2025 Summary | 2025-03-15T12:00:00Z | Key findings from the first quarter include… +forum_post unread ddos_monitor My Watchers [Forum Post] Selling access to corporate VPN… | 2025-03-14T08:30:00Z +credential_occurrence unread cred_watcher Credential Alerts [Credential Occurrence] https://example.com/login | email | 2025-03-13T10:00:00Z +malware_report read malware_tracker My Watchers [Malware Report] New variant of Lumma Stealer | 2025-03-12T15:45:00Z | A new variant has been observed… +``` + +## Helper: `get_latest` + +Stream endpoints return data in ascending (oldest-first) order only — there is no way to sort +descending or jump directly to the most recent page. `get_latest` works around this by probing the +API with a small time window, expanding it until enough items exist, and then fetching and slicing +the tail. + +```python +from verity471.helpers import get_latest +``` + +### Parameters + +| Parameter | Type | Default | Description | +|---|---|---|---| +| `stream_method` | callable | *(required)* | A bound stream API method, e.g. `reports_api.get_reports_spot_stream`. Must end with `_stream`. | +| `n` | `int` | *(required)* | Number of most-recent items to return. | +| `**kwargs` | | | Additional filter parameters forwarded verbatim to the stream method (e.g. `malware_family_name`, `threat_type`, `girs`). Do not pass `var_from`, `until`, `size`, or `cursor` — these are managed internally. | + +### Returns + +A list of at most `n` items in ascending order (oldest first), so the last element is always the +single most-recent item. When fewer than `n` items exist in the full history the function returns +however many are available. + +### How probing works + +The helper issues lightweight `size=1` probe calls to count matching items in progressively wider +time windows (doubling each round), starting from an endpoint-specific seed tuned to typical data +density. Once `count >= n`, a single fetch retrieves all items in that window and the tail is +sliced. For very high-density endpoints where the window contains more than 1 000 items the helper +paginates automatically and keeps only the last `n`. + +### Example usage + +```python +import verity471 +from verity471.helpers import get_latest + +configuration = verity471.Configuration( + username="your_username", + password="your_password", +) + +with verity471.ApiClient(configuration) as api_client: + reports_api = verity471.ReportsApi(api_client) + + # 20 most recent spot reports + reports = get_latest(reports_api.get_reports_spot_stream, n=20) + for r in reports: + print(r.title, r.released_ts) + + # 50 most recent malware events for a specific family + events_api = verity471.EventsApi(api_client) + events = get_latest( + events_api.get_events_stream, + n=50, + malware_family_name="Cobalt Strike", + ) +``` + ## Documentation for API Endpoints All URIs are relative to *https://api.intel471.cloud* diff --git a/tests/test_helpers.py b/tests/test_helpers.py new file mode 100644 index 0000000..ef46de7 --- /dev/null +++ b/tests/test_helpers.py @@ -0,0 +1,317 @@ +import json +from datetime import datetime, timezone +from unittest.mock import MagicMock, patch + +import pytest + +from tests.conftest import PREFIX, read_fixture +from verity471 import fetch_alert_targets +from verity471.exceptions import ForbiddenException +from verity471.helpers.alerts import _patch_portal_url +from verity471.helpers.url_router import UnresolvableURL +from verity471.models.get_watcher_response import GetWatcherResponse +from verity471.models.get_watcher_group_response import GetWatcherGroupResponse +from verity471.models.href import Href +from verity471.models.links import Links +from verity471.models.streaming_watcher_alert import StreamingWatcherAlert +import verity471 + + +configuration = verity471.Configuration() + + +test_params = { + 'IndicatorsApi:get_indicator_by_id': ('IntegrationsIndicator', 'https://api.intel471.cloud/integrations/indicators/v1/indicators/malware-indicator--00000000-0000-0000-0000-000000000000', '[Indicator] file | 0000000000000000000000000000000000000000000000000000000000000000 | confidence: 50'), + 'EventsApi:get_event_by_id': ('IntegrationsEvent', 'https://api.intel471.cloud/integrations/malware-intel/v1/events/malware-event--00000000-0000-0000-0000-000000000000', '[Artifact Extraction] dummy | C2: hxxps://example[.]com'), + 'MalwareApi:get_malware_family_by_id': ('SimplifiedMalwareProfile', 'https://api.intel471.cloud/integrations/malware-intel/v1/malware/malware-family--00000000-0000-0000-0000-000000000000', '[Malware] dummy | dummy'), + 'CredentialsApi:get_credential_sets_id': ('GetCredSetResponse', 'https://api.intel471.cloud/integrations/creds/v1/credential-sets/cred-set--84c92b87-ed31-5103-8101-97b87c03a47a', '[Credential Set] dummy | 913706 records | 2023-01-16 00:00:00+00:00'), + 'CredentialsApi:get_credentials_id': ('GetCredResponse', 'https://api.intel471.cloud/integrations/creds/v1/credentials/cred--3f2abe55-8469-59db-b25a-f8268eb31f34', '[Credential] user@example.com | dummy | 2023-01-18T08:08:19.994Z'), + 'CredentialsApi:get_credentials_occurrences_id': ('GetCredOccurrenceResponse', 'https://api.intel471.cloud/integrations/creds/v1/credentials/occurrences/cred-occurrence--1edd10b4-e75d-5aa2-9b43-5e08c6a682cb', '[Credential Occurrence] dummy | 2023-01-18T08:08:19.994Z'), + 'ReportsApi:get_reports_breach_alert_id': ('BreachAlertByIdResponse', 'https://api.intel471.cloud/integrations/intel-report/v1/reports/breach-alert/report--fbbb23d6-713f-5f41-9ee4-45b3ff027017', '[Breach Alert] dummy | 2021-07-01T09:17:33Z'), + 'ReportsApi:get_reports_fintel_id': ('FintelResponse', 'https://api.intel471.cloud/integrations/intel-report/v1/reports/fintel/report--e71d387f-325e-5bfc-a43d-143876c6cfc0', '[Fintel] dummy | 2020-01-03T20:41:55Z |
Actor summaryThe actor AD0 is a long-standing member of the Russian-speaking ...
'), + 'ReportsApi:get_reports_geopol_id': ('GeopolReportDetailsResponse', 'https://api.intel471.cloud/integrations/intel-report/v1/reports/geopol/report--464ad694-6e92-5983-93f2-f0a7f4d84d7e', '[Geopol Report] dummy | 2024-04-16T16:41:10Z |Event backgroundFollowing the Oct
'), + 'ReportsApi:get_reports_info_id': ('InfoReportResponse', 'https://api.intel471.cloud/integrations/intel-report/v1/reports/info/report--1d4f77cb-ee3b-5ec2-9291-8cf9356bdfb8', '[Info Report] dummy | 2014-06-25T23:49:04Z |Within the last few days the online service Indexeus http://indexeus
'), + 'ReportsApi:get_reports_malware_id': ('MalwareReportResponse', 'https://api.intel471.cloud/integrations/intel-report/v1/reports/malware/report--8d11b63b-f7d6-5061-bb17-290ee5af9464', '[Malware Report] dummy | 2019-01-31T14:16:22Z |Malware Analysis Report # Summary # Pony loader, aka Fareit, is a credential ...
'), + 'ReportsApi:get_reports_spot_id': ('SpotReportResponse', 'https://api.intel471.cloud/integrations/intel-report/v1/reports/spot/report--cb89fbf0-4a56-5f0c-8bd4-166b2115362f', '[Spot Report] dummy | 2019-01-17T16:59:57Z | dummy'), + 'ReportsApi:get_reports_vulnerability_id': ('VulnerabilitiesReportDetailsResponse', 'https://api.intel471.cloud/integrations/intel-report/v1/reports/vulnerability/vulnerability--451a1d7b-e555-5c25-bb21-f544d2ce6997', '[Vulnerability Report] dummy | dummy | dummy | RiskLevel.HIGH | VulnerabilityStatus.HISTORICAL'), + 'SourcesApi:get_forums_posts_post_id': ('PostDetails1', 'https://api.intel471.cloud/integrations/sources/v1/forums/posts/post--44a97352-e0bf-537a-8b51-13e16992586b', '[Forum Post] 2022-10-13T16:05:37Z'), + 'SourcesApi:get_forums_private_messages_private_message_id': ('PrivateMessageDetails1', 'https://api.intel471.cloud/integrations/sources/v1/forums/private-messages/private-message--d2c24f11-ed5a-5d6c-b4cc-83a5ff96c0b4', '[Forum PM] dummy | dummy | 2010-08-24T23:25:34Z'), + 'SourcesApi:get_messaging_services_messages_message_id': ('ChatRoomMessageStream', 'https://api.intel471.cloud/integrations/sources/v1/messaging-services/messages/message--da2a22d1-4d3e-5b79-b557-c275453f31f9', '[Message] dummy | 2017-02-27T02:37:48Z'), +} + +@patch('verity471.helpers.alerts.WatchersApi') +@patch('verity471.rest.RESTClientObject') +@pytest.mark.parametrize('filename, query_url, expected_target_summary', test_params.values(), ids=test_params.keys()) +def test_api_responses(rest_client_class_mock, watchers_api_mock, filename, query_url, expected_target_summary): + watchers_api_mock.return_value.get_watchers.return_value.watchers = [] + watchers_api_mock.return_value.get_watcher_groups.return_value.watchers_groups = [] + + + + rest_client_response = MagicMock(name='rest_client_response') + rest_client_response.status = 200 + rest_client_response.reason = 'OK' + rest_client_response.headers = {'content-type': 'application/json; charset=utf-8'} + response = read_fixture(f'{PREFIX}/fixtures/api_responses/{filename}.json') + rest_client_response.data = json.dumps(response).encode('utf-8') + + rest_client_instance_mock = MagicMock(name='rest_client_instance') + rest_client_instance_mock.request.return_value = rest_client_response + + rest_client_class_mock.side_effect = [rest_client_instance_mock] + + with verity471.ApiClient(configuration) as api_client: + + mock_alerts_response = MagicMock(name='alerts_response') + mock_alert = MagicMock(name='alert') + mock_alert.links.verity_api.href = query_url + mock_alerts_response.alerts = [mock_alert] + response = fetch_alert_targets(mock_alerts_response, api_client) + assert response[0].target is not None + assert response[0].target_summary == expected_target_summary + + +@patch('verity471.helpers.alerts.WatchersApi') +@patch('verity471.rest.RESTClientObject') +def test_alert_target_watcher_enrichment(rest_client_class_mock, watchers_api_mock): + filename, query_url, _ = list(test_params.values())[0] + + rest_client_response = MagicMock(name='rest_client_response') + rest_client_response.status = 200 + rest_client_response.reason = 'OK' + rest_client_response.headers = {'content-type': 'application/json; charset=utf-8'} + rest_client_response.data = json.dumps(read_fixture(f'{PREFIX}/fixtures/api_responses/{filename}.json')).encode('utf-8') + rest_client_instance_mock = MagicMock(name='rest_client_instance') + rest_client_instance_mock.request.return_value = rest_client_response + rest_client_class_mock.side_effect = [rest_client_instance_mock] + + mock_watcher = MagicMock(spec=GetWatcherResponse) + mock_watcher.id = 42 + mock_watcher.name = 'my_watcher' + + mock_group = MagicMock(spec=GetWatcherGroupResponse) + mock_group.id = 7 + mock_group.name = 'my_group' + + watchers_api_mock.return_value.get_watchers.return_value.watchers = [mock_watcher] + watchers_api_mock.return_value.get_watcher_groups.return_value.watchers_groups = [mock_group] + + with verity471.ApiClient(verity471.Configuration()) as api_client: + mock_alerts_response = MagicMock(name='alerts_response') + mock_alert = MagicMock(name='alert') + mock_alert.links.verity_api.href = query_url + mock_alert.watcher_id = 42 + mock_alert.watcher_group_id = 7 + mock_alerts_response.alerts = [mock_alert] + + response = fetch_alert_targets(mock_alerts_response, api_client) + + assert response[0].watcher is mock_watcher + assert response[0].watcher.name == 'my_watcher' + assert response[0].watcher_group is mock_group + assert response[0].watcher_group.name == 'my_group' + + +# --------------------------------------------------------------------------- +# Tests for the temporary portal URL workaround (_patch_portal_url). +# Remove this class when the API bug is fixed and _patch_portal_url is deleted. +# --------------------------------------------------------------------------- + +def _make_alert(portal_href=None): + return StreamingWatcherAlert( + id=1, + watcher_group_id=1, + watcher_id=1, + status="pending", + source_type="test", + source_id="test--id", + links=Links(verity_portal=Href(href=portal_href) if portal_href else None), + creation_ts=datetime.now(timezone.utc), + is_trashed=False, + ) + + +class TestPatchPortalUrl: + def test_forum_post_patches_url(self): + alert = _make_alert() + target = verity471.PostDetails1.from_dict( + read_fixture(f"{PREFIX}/fixtures/api_responses/PostDetails1.json") + ) + _patch_portal_url(alert, target) + assert alert.links.verity_portal.href == ( + "https://example.com?postId=post--00000000-0000-0000-0000-000000000000" + ) + + def test_forum_post_skips_when_portal_already_set(self): + alert = _make_alert(portal_href="https://existing.example.com") + target = verity471.PostDetails1.from_dict( + read_fixture(f"{PREFIX}/fixtures/api_responses/PostDetails1.json") + ) + _patch_portal_url(alert, target) + assert alert.links.verity_portal.href == "https://existing.example.com" + + def test_forum_post_skips_when_thread_has_no_portal_url(self): + alert = _make_alert() + data = read_fixture(f"{PREFIX}/fixtures/api_responses/PostDetails1.json") + data["thread"]["links"].pop("verity_portal", None) + target = verity471.PostDetails1.from_dict(data) + _patch_portal_url(alert, target) + assert alert.links.verity_portal is None + + def test_data_leak_site_patches_url(self): + alert = _make_alert() + data = read_fixture(f"{PREFIX}/fixtures/api_responses/DataLeakSitePostsStreamingPage.json") + target = verity471.DataLeakSitePostItem.from_dict(data["posts"][0]) + _patch_portal_url(alert, target) + assert alert.links.verity_portal.href == ( + "https://verity.intel471.com/sources/data-leak-sites/" + "website--00000000-0000-0000-0000-000000000000" + "/threads/thread--00000000-0000-0000-0000-000000000000" + ) + + def test_messaging_service_patches_url(self): + alert = _make_alert() + target = verity471.ChatRoomMessageStream.from_dict( + read_fixture(f"{PREFIX}/fixtures/api_responses/ChatRoomMessageStream.json") + ) + _patch_portal_url(alert, target) + assert alert.links.verity_portal.href == ( + "https://example.com?messageId=message--00000000-0000-0000-0000-000000000000" + ) + + def test_messaging_service_skips_when_message_has_no_portal_url(self): + alert = _make_alert() + data = read_fixture(f"{PREFIX}/fixtures/api_responses/ChatRoomMessageStream.json") + data["message"]["links"].pop("verity_portal", None) + target = verity471.ChatRoomMessageStream.from_dict(data) + _patch_portal_url(alert, target) + assert alert.links.verity_portal is None + + def test_credential_set_patches_url(self): + alert = _make_alert() + target = verity471.GetCredSetResponse.from_dict( + read_fixture(f"{PREFIX}/fixtures/api_responses/GetCredSetResponse.json") + ) + _patch_portal_url(alert, target) + assert alert.links.verity_portal.href == ( + "https://verity.intel471.com/search" + "?q=cred_set.name%3Ddummy&category=creds_cred_set" + ) + + def test_no_op_for_unmatched_target_type(self): + alert = _make_alert() + _patch_portal_url(alert, None) + assert alert.links.verity_portal is None + + +# --------------------------------------------------------------------------- +# Tests for fetch_alert_targets control flow and error handling. +# --------------------------------------------------------------------------- + +_SPOT_URL = "https://api.intel471.cloud/integrations/intel-report/v1/reports/spot/report--test" +_MARKETPLACE_URL = "https://api.intel471.cloud/integrations/marketplaces/v1/product/marketplace-product--test" + + +def _mock_alert(url=_SPOT_URL, source_id="src--test", watcher_id=1, watcher_group_id=1): + alert = MagicMock() + alert.links.verity_api.href = url + alert.source_id = source_id + alert.watcher_id = watcher_id + alert.watcher_group_id = watcher_group_id + return alert + + +def _alerts_response(*alerts): + resp = MagicMock() + resp.alerts = list(alerts) + return resp + + +def _no_watchers(watchers_mock): + watchers_mock.return_value.get_watchers.return_value.watchers = [] + watchers_mock.return_value.get_watcher_groups.return_value.watchers_groups = [] + + +@patch('verity471.helpers.alerts.WatchersApi') +@patch('verity471.helpers.alerts.call_url') +class TestFetchAlertTargets: + + def test_empty_alerts_returns_empty_list(self, call_url_mock, _watchers_mock): + with verity471.ApiClient(configuration) as api_client: + result = fetch_alert_targets(_alerts_response(), api_client) + assert result == [] + call_url_mock.assert_not_called() + + def test_marketplace_url_is_skipped(self, call_url_mock, _watchers_mock): + alert = _mock_alert(url=_MARKETPLACE_URL) + with verity471.ApiClient(configuration) as api_client: + result = fetch_alert_targets(_alerts_response(alert), api_client) + assert result == [] + call_url_mock.assert_not_called() + + def test_missing_link_is_skipped(self, _call_url_mock, _watchers_mock): + alert = _mock_alert() + alert.links = None + with verity471.ApiClient(configuration) as api_client: + result = fetch_alert_targets(_alerts_response(alert), api_client) + assert result == [] + + def test_missing_link_raises_when_requested(self, _call_url_mock, _watchers_mock): + alert = _mock_alert() + alert.links = None + with verity471.ApiClient(configuration) as api_client: + with pytest.raises(ValueError): + fetch_alert_targets(_alerts_response(alert), api_client, raise_on_error=True) + + def test_unresolvable_url_yields_none_target(self, call_url_mock, watchers_mock): + call_url_mock.side_effect = UnresolvableURL("no route") + _no_watchers(watchers_mock) + alert = _mock_alert() + with verity471.ApiClient(configuration) as api_client: + result = fetch_alert_targets(_alerts_response(alert), api_client) + assert len(result) == 1 + assert result[0].alert is alert + assert result[0].target is None + + def test_forbidden_is_skipped(self, call_url_mock, _watchers_mock): + call_url_mock.side_effect = ForbiddenException() + with verity471.ApiClient(configuration) as api_client: + result = fetch_alert_targets(_alerts_response(_mock_alert()), api_client) + assert result == [] + + def test_api_error_skipped_by_default(self, call_url_mock, _watchers_mock): + call_url_mock.side_effect = RuntimeError("boom") + with verity471.ApiClient(configuration) as api_client: + result = fetch_alert_targets(_alerts_response(_mock_alert()), api_client) + assert result == [] + + def test_api_error_raises_when_requested(self, call_url_mock, _watchers_mock): + call_url_mock.side_effect = RuntimeError("boom") + with verity471.ApiClient(configuration) as api_client: + with pytest.raises(RuntimeError): + fetch_alert_targets(_alerts_response(_mock_alert()), api_client, raise_on_error=True) + + def test_watcher_enrichment_failure_still_returns_results(self, call_url_mock, watchers_mock): + target = MagicMock() + call_url_mock.return_value = target + watchers_mock.return_value.get_watchers.side_effect = RuntimeError("watcher API down") + alert = _mock_alert() + with verity471.ApiClient(configuration) as api_client: + result = fetch_alert_targets(_alerts_response(alert), api_client) + assert len(result) == 1 + assert result[0].target is target + assert result[0].watcher is None + assert result[0].watcher_group is None + + def test_result_order_preserved(self, call_url_mock, watchers_mock): + _no_watchers(watchers_mock) + targets = {i: MagicMock(name=f"target_{i}") for i in range(5)} + urls = {f"{_SPOT_URL}--{i}": targets[i] for i in range(5)} + call_url_mock.side_effect = lambda _client, url: urls[url] + alerts = [_mock_alert(url=f"{_SPOT_URL}--{i}", source_id=f"src--{i}") for i in range(5)] + with verity471.ApiClient(configuration) as api_client: + result = fetch_alert_targets(_alerts_response(*alerts), api_client) + assert len(result) == 5 + for i, r in enumerate(result): + assert r.alert is alerts[i] + assert r.target is targets[i] \ No newline at end of file diff --git a/tests/test_stream_latest.py b/tests/test_stream_latest.py new file mode 100644 index 0000000..2745f34 --- /dev/null +++ b/tests/test_stream_latest.py @@ -0,0 +1,200 @@ +from unittest.mock import MagicMock, patch + +import pytest + +from verity471.helpers.stream_latest import _extract_items, get_latest + + +def _make_response(count, items, cursor_next=None, items_field="reports"): + resp = MagicMock() + resp.count = count + resp.cursor_next = cursor_next + resp.model_fields = {"count": None, "cursor_next": None, items_field: None} + setattr(resp, items_field, items) + return resp + + +def _method(name="get_reports_spot_stream"): + m = MagicMock() + m.__name__ = name + return m + + +# --------------------------------------------------------------------------- +# Validation +# --------------------------------------------------------------------------- + +def test_non_stream_method_raises(): + with pytest.raises(ValueError, match="stream"): + get_latest(_method("get_reports_spot"), n=5) + + +@pytest.mark.parametrize("kwarg", ["var_from", "until", "size", "cursor"]) +def test_reserved_kwarg_raises(kwarg): + with pytest.raises(ValueError, match="internally"): + get_latest(_method(), n=5, **{kwarg: "x"}) + + +# --------------------------------------------------------------------------- +# Empty / sparse data +# --------------------------------------------------------------------------- + +def test_empty_history_returns_empty_list(): + m = _method() + m.return_value = _make_response(count=0, items=[]) + assert get_latest(m, n=5) == [] + + +def test_fewer_items_than_n_returns_all_available(): + # count=3 < n=10 even at the max window; returns whatever is there + items = ["a", "b", "c"] + m = _method() + m.return_value = _make_response(count=3, items=items) + assert get_latest(m, n=10) == items + + +# --------------------------------------------------------------------------- +# Single-fetch path (count <= 1000) +# --------------------------------------------------------------------------- + +@patch("verity471.helpers.stream_latest.time.time", return_value=1_000_000.0) +def test_single_fetch_returns_last_n(_): + items = list(range(33)) + m = _method() + m.side_effect = [ + _make_response(count=33, items=items[:1]), # probe + _make_response(count=33, items=items), # fetch + ] + assert get_latest(m, n=20) == items[-20:] + + +@patch("verity471.helpers.stream_latest.time.time", return_value=1_000_000.0) +def test_until_is_pinned_across_all_calls(_): + now_ms = 1_000_000_000 + items = ["x"] + m = _method() + m.side_effect = [ + _make_response(count=1, items=items), # probe + _make_response(count=1, items=items), # fetch + ] + get_latest(m, n=1) + for c in m.call_args_list: + assert c.kwargs["until"] == now_ms + + +@patch("verity471.helpers.stream_latest.time.time", return_value=1_000_000.0) +def test_extra_kwargs_forwarded_to_every_call(_): + items = list(range(5)) + m = _method() + m.side_effect = [ + _make_response(count=5, items=items[:1]), # probe + _make_response(count=5, items=items), # fetch + ] + get_latest(m, n=3, threat_type="apt") + for c in m.call_args_list: + assert c.kwargs.get("threat_type") == "apt" + + +# --------------------------------------------------------------------------- +# Window expansion +# --------------------------------------------------------------------------- + +@patch("verity471.helpers.stream_latest.time.time", return_value=1_000_000.0) +def test_window_expands_until_count_sufficient(_): + # Unknown method → default 86 400 s initial window. + # Probe 1: count=2 < n=10 → expand. Probe 2: count=15 >= n=10 → fetch. + items = list(range(15)) + m = _method("get_unknown_endpoint_stream") + m.side_effect = [ + _make_response(count=2, items=items[:1]), # probe 1 + _make_response(count=15, items=items[:1]), # probe 2 + _make_response(count=15, items=items), # fetch + ] + result = get_latest(m, n=10) + assert result == items[-10:] + assert m.call_count == 3 + + +@patch("verity471.helpers.stream_latest.time.time", return_value=1_000_000.0) +def test_window_doubles_between_probes(_): + now_ms = 1_000_000_000 + initial_s = 86_400 # default for unknown method + items = list(range(5)) + m = _method("get_unknown_endpoint_stream") + m.side_effect = [ + _make_response(count=0, items=[]), # probe 1: window=86 400 + _make_response(count=5, items=items[:1]), # probe 2: window=172 800 + _make_response(count=5, items=items), # fetch + ] + get_latest(m, n=5) + probe1, probe2, _ = m.call_args_list + assert probe1.kwargs["var_from"] == now_ms - initial_s * 1000 + assert probe2.kwargs["var_from"] == now_ms - initial_s * 2 * 1000 + + +# --------------------------------------------------------------------------- +# Pagination path (count > 1000) +# --------------------------------------------------------------------------- + +@patch("verity471.helpers.stream_latest.time.time", return_value=1_000_000.0) +def test_pagination_collects_across_multiple_pages(_): + page1 = list(range(0, 1000)) + page2 = list(range(1000, 2000)) + page3 = list(range(2000, 2500)) + m = _method() + m.side_effect = [ + _make_response(count=2500, items=page1[:1]), # probe + _make_response(count=2500, items=page1, cursor_next="c1"), # page 1 + _make_response(count=2500, items=page2, cursor_next="c2"), # page 2 + _make_response(count=2500, items=page3, cursor_next="c3"), # page 3 + ] + assert get_latest(m, n=100) == (page1 + page2 + page3)[-100:] + + +@patch("verity471.helpers.stream_latest.time.time", return_value=1_000_000.0) +def test_pagination_stops_at_count_not_cursor(_): + # count=1500; cursor is always present but loop must stop after 2 pages + page1 = list(range(1000)) + page2 = list(range(1000, 1500)) + m = _method() + m.side_effect = [ + _make_response(count=1500, items=page1[:1]), # probe + _make_response(count=1500, items=page1, cursor_next="c1"), # page 1 + _make_response(count=1500, items=page2, cursor_next="c2"), # page 2 → stop + _make_response(count=1500, items=[], cursor_next="c3"), # must NOT be called + ] + get_latest(m, n=50) + assert m.call_count == 3 # probe + 2 pages + + +# --------------------------------------------------------------------------- +# _extract_items +# --------------------------------------------------------------------------- + +def test_extract_items_returns_list_field(): + resp = MagicMock() + resp.model_fields = {"count": None, "cursor_next": None, "reports": None} + resp.reports = ["item1", "item2"] + assert _extract_items(resp) == ["item1", "item2"] + + +def test_extract_items_skips_integer_count_fields(): + # ReportResponseStream carries multiple *_count integer sub-totals + resp = MagicMock() + resp.model_fields = { + "count": None, + "info_report_count": None, + "fintel_report_count": None, + "cursor_next": None, + "reports": None, + } + resp.info_report_count = 14363 + resp.fintel_report_count = 1824 + resp.reports = ["r1", "r2"] + assert _extract_items(resp) == ["r1", "r2"] + + +def test_extract_items_returns_empty_when_no_list_field(): + resp = MagicMock() + resp.model_fields = {"count": None, "cursor_next": None} + assert _extract_items(resp) == [] diff --git a/verity471/__init__.py b/verity471/__init__.py index 245cd39..3c817cf 100644 --- a/verity471/__init__.py +++ b/verity471/__init__.py @@ -40,6 +40,10 @@ "WatchersApi", "EntitiesApi", "ObservablesApi", + "AlertTarget", + "fetch_alert_targets", + "resolve_url", + "call_url", "ApiResponse", "ApiClient", "Configuration", @@ -249,6 +253,12 @@ from verity471.api.entities_api import EntitiesApi as EntitiesApi from verity471.api.observables_api import ObservablesApi as ObservablesApi +# import helpers +from verity471.helpers import AlertTarget as AlertTarget +from verity471.helpers import fetch_alert_targets as fetch_alert_targets +from verity471.helpers import resolve_url as resolve_url +from verity471.helpers import call_url as call_url + # import ApiClient from verity471.api_response import ApiResponse as ApiResponse from verity471.api_client import ApiClient as ApiClient diff --git a/verity471/helpers/__init__.py b/verity471/helpers/__init__.py new file mode 100644 index 0000000..f94eff5 --- /dev/null +++ b/verity471/helpers/__init__.py @@ -0,0 +1,5 @@ +from verity471.helpers.alerts import AlertTarget, fetch_alert_targets +from verity471.helpers.stream_latest import get_latest +from verity471.helpers.url_router import call_url, resolve_url + +__all__ = ["AlertTarget", "fetch_alert_targets", "get_latest", "resolve_url", "call_url"] diff --git a/verity471/helpers/alerts.py b/verity471/helpers/alerts.py new file mode 100644 index 0000000..7f8bea2 --- /dev/null +++ b/verity471/helpers/alerts.py @@ -0,0 +1,357 @@ +from __future__ import annotations + +import logging +import re +import concurrent.futures +from dataclasses import dataclass +from typing import Any, Optional +from urllib.parse import quote + +from verity471.api_client import ApiClient +from verity471.api.watchers_api import WatchersApi +from verity471.exceptions import ForbiddenException +from verity471.models.breach_alert_by_id_response import BreachAlertByIdResponse +from verity471.models.chat_room_message_stream import ChatRoomMessageStream +from verity471.models.data_leak_site_post_item import DataLeakSitePostItem +from verity471.models.fintel_response import FintelResponse +from verity471.models.geopol_report_details_response import GeopolReportDetailsResponse +from verity471.models.get_cred_occurrence_response import GetCredOccurrenceResponse +from verity471.models.get_cred_response import GetCredResponse +from verity471.models.get_cred_set_response import GetCredSetResponse +from verity471.models.href import Href +from verity471.models.info_report_response import InfoReportResponse +from verity471.models.integrations_event import IntegrationsEvent +from verity471.models.integrations_indicator import IntegrationsIndicator +from verity471.models.malware_report_response import MalwareReportResponse +from verity471.models.post_details1 import PostDetails1 +from verity471.models.private_message_details1 import PrivateMessageDetails1 +from verity471.models.simplified_malware_profile import SimplifiedMalwareProfile +from verity471.models.spot_report_response import SpotReportResponse +from verity471.models.get_watcher_group_response import GetWatcherGroupResponse +from verity471.models.get_watcher_response import GetWatcherResponse +from verity471.models.streaming_alerts_response import StreamingAlertsResponse +from verity471.models.streaming_watcher_alert import StreamingWatcherAlert +from verity471.models.vulnerabilities_report_details_response import VulnerabilitiesReportDetailsResponse + +from verity471.helpers.url_router import UnresolvableURL, call_url + +log = logging.getLogger(__name__) + +_SUMMARY_SNIPPET_LEN = 256 # soft char limit for text snippets; expands to end of current word + +# --------------------------------------------------------------------------- +# TEMPORARY WORKAROUND — remove this block (and its call site in _fetch) +# once the API populates links.verity_portal on all alert types. +# --------------------------------------------------------------------------- + +def _patch_portal_url(alert: StreamingWatcherAlert, target: Any) -> None: + """Backfill alert.links.verity_portal when the API omits it. + + Temporary workaround for an API bug where certain source types do not + include a verity_portal link. Remove once the API is fixed. + """ + # Geopol reports return /intelligence/geopolReportView/{id} but the correct + # path prefix is /geopol/. + if isinstance(target, GeopolReportDetailsResponse): + if (alert.links and alert.links.verity_portal + and alert.links.verity_portal.href): + alert.links.verity_portal.href = alert.links.verity_portal.href.replace( + "intel471.com/intelligence/geopolReportView/", + "intel471.com/geopol/geopolReportView/", + 1, + ) + return + + if alert.links and alert.links.verity_portal: + return + + url: str | None = None + + if isinstance(target, PostDetails1): + thread = target.thread + if (thread and thread.links and thread.links.verity_portal + and thread.links.verity_portal.href): + url = thread.links.verity_portal.href + "?postId=" + target.post.id + + elif isinstance(target, DataLeakSitePostItem): + url = ( + "https://verity.intel471.com/sources/data-leak-sites/" + + target.website.id + + "/threads/" + + target.thread.id + ) + + elif isinstance(target, ChatRoomMessageStream): + msg = target.message + if (msg.links and msg.links.verity_portal + and msg.links.verity_portal.href): + url = msg.links.verity_portal.href + "?messageId=" + msg.id + + elif isinstance(target, GetCredSetResponse): + if target.data.name: + q = quote("cred_set.name=" + target.data.name, safe="") + url = f"https://verity.intel471.com/search?q={q}&category=creds_cred_set" + + if url: + alert.links.verity_portal = Href(href=url) + else: + log.debug( + "Could not build portal URL for alert %s (source_type=%s)", + alert.source_id, alert.source_type, + ) + +# --------------------------------------------------------------------------- +# END TEMPORARY WORKAROUND +# --------------------------------------------------------------------------- + + +def _defang(text: Optional[str]) -> str: + text = str(text).replace("http://", "hxxp://").replace("https://", "hxxps://") + return re.sub(r"(\w)\.(\w)", r"\1[.]\2", text) + + +def _snippet(text: str, limit: int = _SUMMARY_SNIPPET_LEN) -> str: + """Truncate *text* to roughly *limit* chars, ending on a word boundary.""" + if len(text) <= limit: + return text + end = text.find(" ", limit) + return text[:end] + "\u2026" if end != -1 else text[:limit] + "\u2026" + + +def _join(parts: list) -> str | None: + joined = " | ".join(str(p) for p in parts if p) + return joined or None + + +def _type_label(snake: str) -> str: + return snake.replace("_", " ").title() + + +def _prefixed(label: str, rest: str | None) -> str | None: + prefix = f"[{label}]" + return f"{prefix} {rest}" if rest else prefix + + +def _summarize_target(target: Any) -> str | None: + if target is None or isinstance(target, (bytes, bytearray)): + return None + + if isinstance(target, PostDetails1): + p = target.post + return _prefixed("Forum Post", _join([ + _snippet(p.message) if p.message else None, p.creation_ts])) + + if isinstance(target, PrivateMessageDetails1): + pm = target.private_message + return _prefixed("Forum PM", _join([ + pm.subject, _snippet(pm.message) if pm.message else None, pm.creation_ts])) + + if isinstance(target, ChatRoomMessageStream): + m = target.message + return _prefixed("Message", _join([ + _snippet(m.text) if m.text else None, m.creation_ts])) + + if isinstance(target, (BreachAlertByIdResponse, FintelResponse, + GeopolReportDetailsResponse, MalwareReportResponse, + SpotReportResponse)): + return _prefixed(_type_label(target.type), _join([ + target.title, target.released_ts, + _snippet(target.body) if target.body else None])) + + if isinstance(target, InfoReportResponse): + summary = target.executive_summary or target.body + return _prefixed(_type_label(target.type), _join([ + target.title, target.released_ts, + _snippet(summary) if summary else None])) + + if isinstance(target, VulnerabilitiesReportDetailsResponse): + return _prefixed(_type_label(target.type), _join([ + target.name, target.vendor_name, target.product_name, + str(target.risk_level), str(target.status)])) + + if isinstance(target, GetCredOccurrenceResponse): + return _prefixed("Credential Occurrence", _join([ + _defang(target.data.accessed_url) if target.data.accessed_url else None, + target.data.credential_type, target.last_updated_ts])) + + if isinstance(target, GetCredResponse): + return _prefixed("Credential", _join([ + target.data.credential_login, + _defang(target.data.credential_domain) if target.data.credential_domain else None, + target.last_updated_ts])) + + if isinstance(target, GetCredSetResponse): + return _prefixed("Credential Set", _join([ + target.data.name, + f"{target.data.record_count} records" if target.data.record_count else None, + target.data.breach_ts])) + + if isinstance(target, IntegrationsIndicator): + value = None + if target.data: + value = (target.data.domain or target.data.email or target.data.url + or (target.data.file.sha256 if target.data.file else None) + or (target.data.ipv4.ip_address if target.data.ipv4 else None)) + conf = f"confidence: {target.confidence}" if target.confidence is not None else None + return _prefixed("Indicator", _join([target.type, _defang(value) if value else None, conf])) + + if isinstance(target, IntegrationsEvent): + family_str = None + if target.threat and target.threat.data: + td = target.threat.data + name = td.malware_family.name if td.malware_family else None + version = td.malware.version if td.malware else None + if name: + family_str = f"{name} v{version}" if version else name + + parts: list = [] + d = target.data + if d: + if d.attack_type: + parts.append(d.attack_type) + if d.inject_type: + parts.append(d.inject_type) + if d.plugin_type: + parts.append(d.plugin_type) + elif d.plugin_name: + parts.append(d.plugin_name) + if d.component_type: + parts.append(d.component_type) + if d.target_type: + parts.append(f"target: {d.target_type}") + if d.exfil_location: + parts.append(f"exfil: {_defang(d.exfil_location)}") + elif d.controllers and d.controllers[0].url: + parts.append(f"C2: {_defang(d.controllers[0].url)}") + elif d.controller and d.controller.url: + parts.append(f"C2: {_defang(d.controller.url)}") + + label = _type_label(target.type) if target.type else "Event" + return _prefixed(label, _join([family_str] + parts)) + + if isinstance(target, SimplifiedMalwareProfile): + aliases = ", ".join(target.aliases[:3]) if target.aliases else None + return _prefixed("Malware", _join([ + target.name, aliases, + _snippet(target.description) if target.description else None])) + + return None + + +@dataclass +class AlertTarget: + """An alert paired with its fully fetched target object. + + ``alert`` is the original :class:`StreamingWatcherAlert` (carries status, + watcher IDs, timestamps, highlights, etc.). ``target`` is the resolved API + object — a report, forum post, credential, or whatever the alert refers to. + ``target`` is ``None`` when the URL could not be mapped to a known route. + + ``target_summary`` provides a compact, human-readable one-liner for the + target (e.g. report title + date, indicator type + value, credential + login + domain). + + ``watcher`` is the full :class:`GetWatcherResponse` for the watcher that + triggered this alert, or ``None`` if not found in the fetched list. + ``watcher_group`` is the corresponding :class:`GetWatcherGroupResponse`. + """ + + alert: StreamingWatcherAlert + target: Any + watcher: GetWatcherResponse | None = None + watcher_group: GetWatcherGroupResponse | None = None + + @property + def target_summary(self) -> str | None: + return _summarize_target(self.target) + + +def fetch_alert_targets( + alerts_response: StreamingAlertsResponse, + api_client: ApiClient, + raise_on_error: bool = False, +) -> list[AlertTarget]: + """Fetch the full target object for every alert in *alerts_response*. + + Each :class:`StreamingWatcherAlert` only carries ``source_type``, + ``source_id``, and a ``links.verity_api.href``. This helper resolves that + URL and returns :class:`AlertTarget` pairs so you can work with the actual + content (report body, forum post text, etc.) alongside the alert metadata. + + URLs that cannot be mapped to a known SDK route always produce an + :class:`AlertTarget` with ``target=None`` (and emit a warning). Other + errors (missing link, API call failure) follow *raise_on_error*: when + ``True`` the exception propagates; when ``False`` an error is logged and + the alert is omitted from the result. + + Args: + alerts_response: The page returned by :meth:`AlertsApi.get_alerts_stream`. + api_client: An active :class:`ApiClient` (must share credentials with + the alerts call). + raise_on_error: When ``True``, re-raise unexpected errors instead of + logging and skipping the alert. Defaults to ``False``. + + Returns: + A list of :class:`AlertTarget` objects in the same order as + ``alerts_response.alerts``. + + Example:: + + alerts = alerts_api.get_alerts_stream(size=10) + for r in fetch_alert_targets(alerts, api_client): + print(r.alert.source_type, r.alert.status, r.target) + """ + def _fetch(alert: StreamingWatcherAlert) -> AlertTarget | None: + url = alert.links.verity_api.href if (alert.links and alert.links.verity_api) else None + if url and "/integrations/marketplaces/" in url: + return None + if not url: + if raise_on_error: + raise ValueError("Alert %s has no verity_api link" % alert.source_id) + log.error("Alert %s has no verity_api link", alert.source_id) + return None + try: + target = call_url(api_client, url) + except UnresolvableURL: + log.warning("No SDK route for alert %s URL: %s", alert.source_id, url) + return AlertTarget(alert=alert, target=None) + except ForbiddenException: + log.debug("Failed to fetch target for alert %s (%s) - Forbidden", alert.source_id, url) + return None + except Exception: + if raise_on_error: + raise + log.error("Failed to fetch target for alert %s (%s)", alert.source_id, url, exc_info=True) + return None + _patch_portal_url(alert, target) # TEMPORARY WORKAROUND — remove once API is fixed + return AlertTarget(alert=alert, target=target) + + alerts = alerts_response.alerts or [] + results: list[AlertTarget] = [None] * len(alerts) # type: ignore[list-item] + with concurrent.futures.ThreadPoolExecutor() as executor: + future_to_index = {executor.submit(_fetch, alert): i for i, alert in enumerate(alerts)} + for future in concurrent.futures.as_completed(future_to_index): + result = future.result() + if result is not None: + results[future_to_index[future]] = result + + enriched = [r for r in results if r is not None] + if not enriched: + return [] + + try: + watchers_api = WatchersApi(api_client) + with concurrent.futures.ThreadPoolExecutor(max_workers=2) as watcher_executor: + future_watchers = watcher_executor.submit(watchers_api.get_watchers) + future_groups = watcher_executor.submit(watchers_api.get_watcher_groups) + watchers_resp = future_watchers.result() + groups_resp = future_groups.result() + watchers_by_id = {w.id: w for w in (watchers_resp.watchers or [])} + groups_by_id = {g.id: g for g in (groups_resp.watchers_groups or [])} + for r in enriched: + r.watcher = watchers_by_id.get(r.alert.watcher_id) + r.watcher_group = groups_by_id.get(r.alert.watcher_group_id) + except Exception: + log.warning("Failed to fetch watchers/watcher groups; watcher enrichment will be skipped", exc_info=True) + + return enriched diff --git a/verity471/helpers/stream_latest.py b/verity471/helpers/stream_latest.py new file mode 100644 index 0000000..19d3c6c --- /dev/null +++ b/verity471/helpers/stream_latest.py @@ -0,0 +1,160 @@ +# coding: utf-8 + +"""Helper for fetching the most recent N items from any stream endpoint. + +Stream endpoints only return data in ascending (oldest-first) order. To get +the *latest* N results the helper probes the API with progressively wider time +windows until ``count >= n``, then performs a single fetch (or a short +pagination run for very high-density endpoints) and slices the tail. + +Usage:: + + import verity471 + from verity471.helpers import get_latest + + cfg = verity471.Configuration(username=..., password=...) + with verity471.ApiClient(cfg) as client: + api = verity471.ReportsApi(client) + reports = get_latest(api.get_reports_spot_stream, n=20) + + # Domain-specific filters are forwarded as keyword arguments: + events = get_latest( + verity471.EventsApi(client).get_events_stream, + n=50, + malware_family_name="Cobalt Strike", + ) +""" + +from __future__ import annotations + +import logging +import time +from typing import Any, Callable + +log = logging.getLogger(__name__) + +# Initial probe window in seconds per stream method, derived from observed +# data density (counts at 1 h / 24 h / 7 d / 30 d). The algorithm doubles +# this window until count >= n, so a good starting point reduces round-trips. +_INITIAL_WINDOW_SECONDS: dict[str, int] = { + "get_messaging_services_messages_stream": 60, # ~23 700/h + "get_events_stream": 300, # ~3 690/h (malware events) + "get_forums_posts_stream": 300, # ~2 897/h + "get_alerts_stream": 3_600, # ~19/h (1 h) + "get_reports_breach_alert_stream": 43_200, # ~1.4/h (12 h) + "get_reports_stream": 86_400, # ~2/h (24 h) + "get_reports_fintel_stream": 604_800, # ~0.5/day (7 d) + "get_reports_geopol_stream": 604_800, # ~0.5/day (7 d) + "get_reports_info_stream": 604_800, # ~0.25/day (7 d) + "get_reports_spot_stream": 604_800, # ~0.5/day (7 d) + "get_reports_vulnerability_stream": 604_800, # ~5/day (7 d) + "get_reports_malware_stream": 1_209_600, # ~0.1/day (14 d) + "get_credentials_stream": 1_209_600, # ~0.1/day (14 d) + "get_credentials_occurrences_stream": 1_209_600, + "get_credential_sets_stream": 1_209_600, + "get_credential_sets_accessed_urls_stream": 1_209_600, + "get_data_leak_sites_posts_stream": 5_184_000, # ~1/month (60 d) + "get_forums_private_messages_stream": 5_184_000, + "get_indicators_stream": 300, # assumed similar to events + "get_actors_stream": 86_400, # requires search term + "get_observables_stream": 86_400, + "get_entities_stream": 86_400, +} +_DEFAULT_INITIAL_WINDOW_SECONDS = 86_400 # 24 h fallback for unknown methods +_MAX_WINDOW_SECONDS = 365 * 24 * 3_600 # 1-year hard cap + + +def _extract_items(response: Any) -> list: + """Return the items list from a stream response object. + + Every stream response has exactly three fields: ``count``, ``cursor_next``, + and one items field (e.g. ``reports``, ``events``, ``credentials``). This + function returns that third field without needing to know its name. + """ + for field_name in response.model_fields: + if field_name not in ("count", "cursor_next"): + value = getattr(response, field_name) + if isinstance(value, list): + return value + return [] + + +def get_latest(stream_method: Callable, n: int, **kwargs: Any) -> list: + """Return the *n* most recent items from a stream endpoint. + + Args: + stream_method: A bound stream API method, e.g. + ``reports_api.get_reports_spot_stream``. + n: Number of most-recent items to return. When fewer than *n* items + exist across the entire history the function returns however many + are available (possibly an empty list). + **kwargs: Additional filter parameters forwarded verbatim to the stream + method (e.g. ``malware_family_name``, ``threat_type``, ``girs``). + Do **not** pass ``var_from``, ``until``, ``size``, or ``cursor`` — + these are managed internally. + + Returns: + A list of at most *n* items in ascending order (oldest first), i.e. the + last element is always the single most-recent item. + + Raises: + ValueError: If any of the internally-managed parameters (``var_from``, + ``until``, ``size``, ``cursor``) appear in *kwargs*. + """ + method_name = getattr(stream_method, "__name__", "") + + if not method_name.endswith("_stream"): + raise ValueError( + f"{method_name!r} does not appear to be a stream method; " + "method name must end with '_stream'." + ) + + reserved = {"var_from", "until", "size", "cursor"} + if conflicts := reserved & kwargs.keys(): + raise ValueError( + f"get_latest manages {sorted(conflicts)} internally; " + "do not pass them as keyword arguments." + ) + + now_ms = int(time.time() * 1000) + window_s = _INITIAL_WINDOW_SECONDS.get(method_name, _DEFAULT_INITIAL_WINDOW_SECONDS) + + # Phase 1: exponential expansion until count >= n or the 1-year cap is hit. + from_ms = now_ms - window_s * 1000 + count = 0 + while True: + from_ms = now_ms - window_s * 1000 + resp = stream_method(var_from=from_ms, until=now_ms, size=1, **kwargs) + count = resp.count + log.debug( + "%s probe: window=%ds from_ms=%d count=%d (target n=%d)", + method_name, window_s, from_ms, count, n, + ) + if count >= n or window_s >= _MAX_WINDOW_SECONDS: + break + window_s = min(window_s * 2, _MAX_WINDOW_SECONDS) + + if count == 0: + return [] + + # Phase 2a: everything fits in one page — single fetch, slice the tail. + if count <= 1000: + resp = stream_method(var_from=from_ms, until=now_ms, size=count, **kwargs) + return _extract_items(resp)[-n:] + + # Phase 2b: high-density burst — paginate and keep only the last n items. + log.debug( + "%s: count=%d > 1000, paginating to collect last %d items", method_name, count, n + ) + buffer: list = [] + cursor: str | None = None + while len(buffer) < count: + if cursor: + resp = stream_method(cursor=cursor, size=1000, **kwargs) + else: + resp = stream_method(var_from=from_ms, until=now_ms, size=1000, **kwargs) + buffer.extend(_extract_items(resp)) + cursor = resp.cursor_next + if not cursor: + break + return buffer[-n:] diff --git a/verity471/helpers/url_router.py b/verity471/helpers/url_router.py new file mode 100644 index 0000000..0bf4831 --- /dev/null +++ b/verity471/helpers/url_router.py @@ -0,0 +1,102 @@ +from __future__ import annotations + +import re +from typing import Any +from urllib.parse import urlparse + +from verity471.api_client import ApiClient +from verity471.api.credentials_api import CredentialsApi + + +class UnresolvableURL(Exception): + """Raised when a URL cannot be mapped to a known SDK route.""" +from verity471.api.events_api import EventsApi +from verity471.api.indicators_api import IndicatorsApi +from verity471.api.malware_api import MalwareApi +from verity471.api.reports_api import ReportsApi +from verity471.api.sources_api import SourcesApi + +# Ordered list of (path_template, ApiClass, method_name). +# More-specific paths must come before shorter prefix matches +# (e.g. /credentials/occurrences/{id} before /credentials/{id}). +_RAW_ROUTES: list[tuple[str, type, str]] = [ + # Sources + ("/integrations/sources/v1/forums/posts/{post_id}", SourcesApi, "get_forums_posts_post_id"), + ("/integrations/sources/v1/forums/private-messages/{private_message_id}", SourcesApi, "get_forums_private_messages_private_message_id"), + ("/integrations/sources/v1/data-leak-sites/file-listings/{id}", SourcesApi, "get_data_leak_sites_file_listings_id"), + ("/integrations/sources/v1/messaging-services/messages/{message_id}", SourcesApi, "get_messaging_services_messages_message_id"), + # Reports + ("/integrations/intel-report/v1/reports/breach-alert/{id}", ReportsApi, "get_reports_breach_alert_id"), + ("/integrations/intel-report/v1/reports/fintel/{id}", ReportsApi, "get_reports_fintel_id"), + ("/integrations/intel-report/v1/reports/geopol/{id}", ReportsApi, "get_reports_geopol_id"), + ("/integrations/intel-report/v1/reports/info/{id}", ReportsApi, "get_reports_info_id"), + ("/integrations/intel-report/v1/reports/malware/{id}", ReportsApi, "get_reports_malware_id"), + ("/integrations/intel-report/v1/reports/spot/{id}", ReportsApi, "get_reports_spot_id"), + ("/integrations/intel-report/v1/reports/vulnerability/{id}", ReportsApi, "get_reports_vulnerability_id"), + # Credentials + ("/integrations/creds/v1/credentials/occurrences/{id}", CredentialsApi, "get_credentials_occurrences_id"), + ("/integrations/creds/v1/credentials/{id}", CredentialsApi, "get_credentials_id"), + ("/integrations/creds/v1/credential-sets/{id}", CredentialsApi, "get_credential_sets_id"), + # Indicators + ("/integrations/indicators/v1/indicators/{id}", IndicatorsApi, "get_indicator_by_id"), + # Events and Malware + ("/integrations/malware-intel/v1/events/{id}", EventsApi, "get_event_by_id"), + ("/integrations/malware-intel/v1/malware/{id}", MalwareApi, "get_malware_family_by_id"), +] + + +def _template_to_regex(template: str) -> re.Pattern[str]: + parts = re.split(r'\{(\w+)\}', template) + segments = [] + for i, part in enumerate(parts): + if i % 2 == 0: + segments.append(re.escape(part)) + else: + segments.append(f'(?P<{part}>[^/]+)') + return re.compile(''.join(segments) + '$') + + +_COMPILED_ROUTES: list[tuple[re.Pattern[str], type, str]] = [ + (_template_to_regex(template), api_class, method_name) + for template, api_class, method_name in _RAW_ROUTES +] + + +def resolve_url(url: str) -> tuple[type, str, dict[str, str]] | None: + """Parse a Verity API URL and return (ApiClass, method_name, path_params). + + Returns None if no route matches. + + Example:: + + api_class, method, params = resolve_url( + "https://api.intel471.cloud/integrations/sources/v1/forums/posts/post--abc" + ) + # -> (SourcesApi, 'get_forums_posts_post_id', {'post_id': 'post--abc'}) + """ + path = urlparse(url).path + for pattern, api_class, method_name in _COMPILED_ROUTES: + m = pattern.match(path) + if m: + return api_class, method_name, m.groupdict() + return None + + +def call_url(api_client: ApiClient, url: str) -> Any: + """Resolve a Verity API URL and call the corresponding API method. + + Raises ValueError if the URL does not match any known route. + + Example:: + + obj = call_url(api_client, + "https://api.intel471.cloud/integrations/sources/v1/forums/posts/post--abc" + ) + # -> ForumsPost object + """ + resolved = resolve_url(url) + if resolved is None: + raise UnresolvableURL(f"No API route found for URL: {url}") + api_class, method_name, path_params = resolved + instance = api_class(api_client) + return getattr(instance, method_name)(**path_params)