From 4fc500b1078df0eb4cdfa2307613a953dcf9f3b8 Mon Sep 17 00:00:00 2001 From: mmolenda Date: Mon, 23 Mar 2026 16:51:20 +0100 Subject: [PATCH 01/16] Alerts fetcher --- verity471/__init__.py | 10 +++ verity471/helpers/__init__.py | 4 ++ verity471/helpers/alerts.py | 104 ++++++++++++++++++++++++++++++++ verity471/helpers/url_router.py | 102 +++++++++++++++++++++++++++++++ 4 files changed, 220 insertions(+) create mode 100644 verity471/helpers/__init__.py create mode 100644 verity471/helpers/alerts.py create mode 100644 verity471/helpers/url_router.py diff --git a/verity471/__init__.py b/verity471/__init__.py index a0ffb8a..4163f1d 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", @@ -251,6 +255,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..cc31b0e --- /dev/null +++ b/verity471/helpers/__init__.py @@ -0,0 +1,4 @@ +from verity471.helpers.alerts import AlertTarget, fetch_alert_targets +from verity471.helpers.url_router import call_url, resolve_url + +__all__ = ["AlertTarget", "fetch_alert_targets", "resolve_url", "call_url"] diff --git a/verity471/helpers/alerts.py b/verity471/helpers/alerts.py new file mode 100644 index 0000000..82a41ec --- /dev/null +++ b/verity471/helpers/alerts.py @@ -0,0 +1,104 @@ +from __future__ import annotations + +import logging +from concurrent.futures import ThreadPoolExecutor, as_completed +from dataclasses import dataclass +from typing import Any + +from verity471.api_client import ApiClient +from verity471.models.streaming_alerts_response import StreamingAlertsResponse +from verity471.models.streaming_watcher_alert import StreamingWatcherAlert + +from verity471.helpers.url_router import UnresolvableURL, call_url + +log = logging.getLogger(__name__) + + +@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. + + Convenience properties mirror the most-used fields from ``alert`` so you + rarely need to drill into ``.alert.source_type``. + """ + + alert: StreamingWatcherAlert + target: Any + + @property + def source_type(self) -> str: + return self.alert.source_type + + @property + def source_id(self) -> str: + return self.alert.source_id + + +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.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 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 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 + return AlertTarget(alert=alert, target=target) + + alerts = alerts_response.alerts or [] + results: list[AlertTarget] = [None] * len(alerts) # type: ignore[list-item] + with ThreadPoolExecutor() as executor: + future_to_index = {executor.submit(_fetch, alert): i for i, alert in enumerate(alerts)} + for future in as_completed(future_to_index): + result = future.result() + if result is not None: + results[future_to_index[future]] = result + return [r for r in results if r is not None] 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) From 4eff7151ea1bf5891cebd501e554da430c388f1d Mon Sep 17 00:00:00 2001 From: mmolenda Date: Mon, 23 Mar 2026 17:49:02 +0100 Subject: [PATCH 02/16] Alert summary --- verity471/helpers/alerts.py | 124 +++++++++++++++++++++++++++++++++++- 1 file changed, 123 insertions(+), 1 deletion(-) diff --git a/verity471/helpers/alerts.py b/verity471/helpers/alerts.py index 82a41ec..69d1549 100644 --- a/verity471/helpers/alerts.py +++ b/verity471/helpers/alerts.py @@ -6,13 +6,129 @@ from typing import Any from verity471.api_client import ApiClient +from verity471.models.breach_alert_by_id_response import BreachAlertByIdResponse +from verity471.models.chat_room_message_stream import ChatRoomMessageStream +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.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.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 + + +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([ + target.data.accessed_url, target.data.credential_type, target.last_updated_ts])) + + if isinstance(target, GetCredResponse): + return _prefixed("Credential", _join([ + target.data.credential_login, target.data.credential_domain, 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.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, value, conf])) + + if isinstance(target, IntegrationsEvent): + family = None + if target.threat and target.threat.data and target.threat.data.malware_family: + family = target.threat.data.malware_family.name + label = _type_label(target.type) if target.type else "Event" + return _prefixed(label, _join([ + family, target.data.attack_type if target.data else None])) + + 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: @@ -24,7 +140,9 @@ class AlertTarget: ``target`` is ``None`` when the URL could not be mapped to a known route. Convenience properties mirror the most-used fields from ``alert`` so you - rarely need to drill into ``.alert.source_type``. + rarely need to drill into ``.alert.source_type``. ``target_summary`` + provides a compact, human-readable one-liner for the target (e.g. report + title + date, indicator type + value, credential login + domain). """ alert: StreamingWatcherAlert @@ -38,6 +156,10 @@ def source_type(self) -> str: def source_id(self) -> str: return self.alert.source_id + @property + def target_summary(self) -> str | None: + return _summarize_target(self.target) + def fetch_alert_targets( alerts_response: StreamingAlertsResponse, From 4a402b8ff6489eebd77fdc07d7c622903fbb72b6 Mon Sep 17 00:00:00 2001 From: mmolenda Date: Tue, 24 Mar 2026 17:58:16 +0100 Subject: [PATCH 03/16] Adding tests for alerts helper --- tests/test_helpers.py | 59 +++++++++++++++++++++++++++++++++++++ verity471/helpers/alerts.py | 4 +-- 2 files changed, 61 insertions(+), 2 deletions(-) create mode 100644 tests/test_helpers.py diff --git a/tests/test_helpers.py b/tests/test_helpers.py new file mode 100644 index 0000000..29d4403 --- /dev/null +++ b/tests/test_helpers.py @@ -0,0 +1,59 @@ +import json +from unittest.mock import MagicMock, patch + +import pytest + +from tests.conftest import PREFIX, read_fixture +from verity471 import fetch_alert_targets +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'), + '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.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, filename, query_url, expected_target_summary): + + + + 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 \ No newline at end of file diff --git a/verity471/helpers/alerts.py b/verity471/helpers/alerts.py index 69d1549..d9d4322 100644 --- a/verity471/helpers/alerts.py +++ b/verity471/helpers/alerts.py @@ -108,7 +108,7 @@ def _summarize_target(target: Any) -> str | None: if isinstance(target, IntegrationsIndicator): value = None if target.data: - value = (target.data.domain or target.data.email or target.data.url + value = (target.data.domain or target.data.email or target.data.url or target.data.file.sha256 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, value, conf])) @@ -208,7 +208,7 @@ def _fetch(alert: StreamingWatcherAlert) -> AlertTarget | None: except UnresolvableURL: log.warning("No SDK route for alert %s URL: %s", alert.source_id, url) return AlertTarget(alert=alert, target=None) - except Exception: + except Exception as e: if raise_on_error: raise log.error("Failed to fetch target for alert %s (%s)", alert.source_id, url, exc_info=True) From 52b130808e1ff3a62ef822f05cd413c6d8ae74e4 Mon Sep 17 00:00:00 2001 From: mmolenda Date: Wed, 25 Mar 2026 11:10:43 +0100 Subject: [PATCH 04/16] Added documentation for alerts helper --- README.md | 65 +++++++++++++++++++++++++++++++++++++ verity471/helpers/alerts.py | 23 ++++--------- 2 files changed, 72 insertions(+), 16 deletions(-) diff --git a/README.md b/README.md index c04ebfa..0e8446b 100644 --- a/README.md +++ b/README.md @@ -215,6 +215,71 @@ 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.helpers import fetch_alert_targets, AlertTarget +``` + +or directly from the top-level package: + +```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. | + +### 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: + print(t.alert.source_type, t.alert.source_id, t.alert.status, t.target_summary) +``` + +### Example output + +``` +fintel fintel--abcd1234 read [Fintel] Threat Landscape: Q1 2025 Summary | 2025-03-15T12:00:00Z | Key findings from the first quarter include… +forum_post post--77ef7990-f8f2-5076-9126-1c22b463c515 unread [Forum Post] Selling access to corporate VPN… | 2025-03-14T08:30:00Z +credential_occurrence credential_occurrence--a1b2c3 unread [Credential Occurrence] https://example.com/login | email | 2025-03-13T10:00:00Z +malware_report malware_report--x9y8z7 read [Malware Report] New variant of Lumma Stealer | 2025-03-12T15:45:00Z | A new variant has been observed… +``` + ## Documentation for API Endpoints All URIs are relative to *https://api.intel471.cloud* diff --git a/verity471/helpers/alerts.py b/verity471/helpers/alerts.py index d9d4322..0fcbf6c 100644 --- a/verity471/helpers/alerts.py +++ b/verity471/helpers/alerts.py @@ -1,7 +1,7 @@ from __future__ import annotations import logging -from concurrent.futures import ThreadPoolExecutor, as_completed +import concurrent.futures from dataclasses import dataclass from typing import Any @@ -139,23 +139,14 @@ class AlertTarget: 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. - Convenience properties mirror the most-used fields from ``alert`` so you - rarely need to drill into ``.alert.source_type``. ``target_summary`` - provides a compact, human-readable one-liner for the target (e.g. report - title + date, indicator type + value, credential login + domain). + ``target_summary`` provides a compact, human-readable one-liner for the + target (e.g. report title + date, indicator type + value, credential + login + domain). """ alert: StreamingWatcherAlert target: Any - @property - def source_type(self) -> str: - return self.alert.source_type - - @property - def source_id(self) -> str: - return self.alert.source_id - @property def target_summary(self) -> str | None: return _summarize_target(self.target) @@ -194,7 +185,7 @@ def fetch_alert_targets( alerts = alerts_api.get_alerts_stream(size=10) for r in fetch_alert_targets(alerts, api_client): - print(r.source_type, r.alert.status, r.target) + 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 @@ -217,9 +208,9 @@ def _fetch(alert: StreamingWatcherAlert) -> AlertTarget | None: alerts = alerts_response.alerts or [] results: list[AlertTarget] = [None] * len(alerts) # type: ignore[list-item] - with ThreadPoolExecutor() as executor: + with concurrent.futures.ThreadPoolExecutor() as executor: future_to_index = {executor.submit(_fetch, alert): i for i, alert in enumerate(alerts)} - for future in as_completed(future_to_index): + for future in concurrent.futures.as_completed(future_to_index): result = future.result() if result is not None: results[future_to_index[future]] = result From 3979a566c522d32de564a891222cec1e8a922ed6 Mon Sep 17 00:00:00 2001 From: mmolenda Date: Wed, 25 Mar 2026 11:12:48 +0100 Subject: [PATCH 05/16] Updating readme --- README.md | 6 ------ 1 file changed, 6 deletions(-) diff --git a/README.md b/README.md index 0e8446b..e68f7bb 100644 --- a/README.md +++ b/README.md @@ -222,12 +222,6 @@ The alerts stream endpoint returns lightweight `StreamingWatcherAlert` objects t `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.helpers import fetch_alert_targets, AlertTarget -``` - -or directly from the top-level package: - ```python from verity471 import fetch_alert_targets, AlertTarget ``` From ffe59686adcac27734d5b0ac030d1d7b04413f5c Mon Sep 17 00:00:00 2001 From: mmolenda Date: Wed, 25 Mar 2026 17:35:45 +0100 Subject: [PATCH 06/16] Fetch watcher and group in alerts helper --- tests/test_helpers.py | 50 +++++++++++++++++++++++++++++++++++-- verity471/helpers/alerts.py | 39 ++++++++++++++++++++++++++--- 2 files changed, 84 insertions(+), 5 deletions(-) diff --git a/tests/test_helpers.py b/tests/test_helpers.py index 29d4403..a6b480d 100644 --- a/tests/test_helpers.py +++ b/tests/test_helpers.py @@ -5,6 +5,8 @@ from tests.conftest import PREFIX, read_fixture from verity471 import fetch_alert_targets +from verity471.models.get_watcher_response import GetWatcherResponse +from verity471.models.get_watcher_group_response import GetWatcherGroupResponse import verity471 @@ -30,9 +32,12 @@ '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, filename, query_url, expected_target_summary): +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 = [] @@ -56,4 +61,45 @@ def test_api_responses(rest_client_class_mock, filename, query_url, expected_tar 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 \ No newline at end of file + 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' \ No newline at end of file diff --git a/verity471/helpers/alerts.py b/verity471/helpers/alerts.py index 0fcbf6c..dd8a1af 100644 --- a/verity471/helpers/alerts.py +++ b/verity471/helpers/alerts.py @@ -6,6 +6,7 @@ from typing import Any from verity471.api_client import ApiClient +from verity471.api.watchers_api import WatchersApi from verity471.models.breach_alert_by_id_response import BreachAlertByIdResponse from verity471.models.chat_room_message_stream import ChatRoomMessageStream from verity471.models.fintel_response import FintelResponse @@ -21,6 +22,8 @@ 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 @@ -142,10 +145,16 @@ class AlertTarget: ``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: @@ -187,6 +196,20 @@ def fetch_alert_targets( for r in fetch_alert_targets(alerts, api_client): print(r.alert.source_type, r.alert.status, r.target) """ + watchers_by_id: dict[int, GetWatcherResponse] = {} + groups_by_id: dict[int, GetWatcherGroupResponse] = {} + 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 [])} + except Exception: + log.warning("Failed to fetch watchers/watcher groups; watcher enrichment will be skipped", exc_info=True) + def _fetch(alert: StreamingWatcherAlert) -> AlertTarget | None: url = alert.links.verity_api.href if (alert.links and alert.links.verity_api) else None if not url: @@ -198,13 +221,23 @@ def _fetch(alert: StreamingWatcherAlert) -> AlertTarget | None: 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 Exception as e: + return AlertTarget( + alert=alert, + target=None, + watcher=watchers_by_id.get(alert.watcher_id), + watcher_group=groups_by_id.get(alert.watcher_group_id), + ) + 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 - return AlertTarget(alert=alert, target=target) + return AlertTarget( + alert=alert, + target=target, + watcher=watchers_by_id.get(alert.watcher_id), + watcher_group=groups_by_id.get(alert.watcher_group_id), + ) alerts = alerts_response.alerts or [] results: list[AlertTarget] = [None] * len(alerts) # type: ignore[list-item] From 17429454a17e5916a6aefa8529094f73be2fd15e Mon Sep 17 00:00:00 2001 From: mmolenda Date: Wed, 25 Mar 2026 17:38:44 +0100 Subject: [PATCH 07/16] Updating readme --- README.md | 18 +++++++++++++----- 1 file changed, 13 insertions(+), 5 deletions(-) diff --git a/README.md b/README.md index e68f7bb..33ff27e 100644 --- a/README.md +++ b/README.md @@ -245,6 +245,12 @@ Each `AlertTarget` exposes: | `.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 @@ -262,16 +268,18 @@ with verity471.ApiClient(configuration) as api_client: targets = verity471.fetch_alert_targets(alerts_response, api_client) for t in targets: - print(t.alert.source_type, t.alert.source_id, t.alert.status, t.target_summary) + 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 fintel--abcd1234 read [Fintel] Threat Landscape: Q1 2025 Summary | 2025-03-15T12:00:00Z | Key findings from the first quarter include… -forum_post post--77ef7990-f8f2-5076-9126-1c22b463c515 unread [Forum Post] Selling access to corporate VPN… | 2025-03-14T08:30:00Z -credential_occurrence credential_occurrence--a1b2c3 unread [Credential Occurrence] https://example.com/login | email | 2025-03-13T10:00:00Z -malware_report malware_report--x9y8z7 read [Malware Report] New variant of Lumma Stealer | 2025-03-12T15:45:00Z | A new variant has been observed… +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… ``` ## Documentation for API Endpoints From 9f7b108a260606de5a3a19cfa4761ad8dcbbe8b8 Mon Sep 17 00:00:00 2001 From: mmolenda Date: Fri, 10 Apr 2026 16:16:18 +0200 Subject: [PATCH 08/16] added get_latest helper --- verity471/helpers/__init__.py | 3 +- verity471/helpers/stream_latest.py | 160 +++++++++++++++++++++++++++++ 2 files changed, 162 insertions(+), 1 deletion(-) create mode 100644 verity471/helpers/stream_latest.py diff --git a/verity471/helpers/__init__.py b/verity471/helpers/__init__.py index cc31b0e..f94eff5 100644 --- a/verity471/helpers/__init__.py +++ b/verity471/helpers/__init__.py @@ -1,4 +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", "resolve_url", "call_url"] +__all__ = ["AlertTarget", "fetch_alert_targets", "get_latest", "resolve_url", "call_url"] 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:] From e0ed185bb3fb51925c140d585c08dca841b9c94c Mon Sep 17 00:00:00 2001 From: mmolenda Date: Mon, 13 Apr 2026 12:29:18 +0200 Subject: [PATCH 09/16] Adding test_stream_latest --- tests/test_stream_latest.py | 200 ++++++++++++++++++++++++++++++++++++ 1 file changed, 200 insertions(+) create mode 100644 tests/test_stream_latest.py 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) == [] From 481b61eb68880c84c5941b71683e1c8133aceb9d Mon Sep 17 00:00:00 2001 From: mmolenda Date: Mon, 13 Apr 2026 12:32:20 +0200 Subject: [PATCH 10/16] Adding stream_latest to the Readme --- README.md | 61 +++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 61 insertions(+) diff --git a/README.md b/README.md index 33ff27e..b97608a 100644 --- a/README.md +++ b/README.md @@ -282,6 +282,67 @@ credential_occurrence unread cred_watcher Credential Alerts [Credential Occurren 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* From bb8bfd239d84a24ac39bf315d568d173a5a2cdac Mon Sep 17 00:00:00 2001 From: mmolenda Date: Thu, 23 Apr 2026 14:54:33 +0200 Subject: [PATCH 11/16] Skip marketplace alerts --- verity471/helpers/alerts.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/verity471/helpers/alerts.py b/verity471/helpers/alerts.py index dd8a1af..f8b6361 100644 --- a/verity471/helpers/alerts.py +++ b/verity471/helpers/alerts.py @@ -212,6 +212,8 @@ def fetch_alert_targets( 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) From 04be5aeadf7efb5e4e0ad83517ceb71bab44ba04 Mon Sep 17 00:00:00 2001 From: mmolenda Date: Thu, 30 Apr 2026 17:47:06 +0200 Subject: [PATCH 12/16] Generating missing portal URL for some alerts --- tests/test_helpers.py | 99 ++++++++++++++++++++++++++++++++++++- verity471/helpers/alerts.py | 56 +++++++++++++++++++++ 2 files changed, 154 insertions(+), 1 deletion(-) diff --git a/tests/test_helpers.py b/tests/test_helpers.py index a6b480d..a5d614c 100644 --- a/tests/test_helpers.py +++ b/tests/test_helpers.py @@ -1,12 +1,17 @@ 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.helpers.alerts import _patch_portal_url 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 @@ -102,4 +107,96 @@ def test_alert_target_watcher_enrichment(rest_client_class_mock, watchers_api_mo 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' \ No newline at end of file + 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 \ No newline at end of file diff --git a/verity471/helpers/alerts.py b/verity471/helpers/alerts.py index f8b6361..7db9fe7 100644 --- a/verity471/helpers/alerts.py +++ b/verity471/helpers/alerts.py @@ -4,16 +4,19 @@ import concurrent.futures from dataclasses import dataclass from typing import Any +from urllib.parse import quote from verity471.api_client import ApiClient from verity471.api.watchers_api import WatchersApi 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 @@ -34,6 +37,58 @@ _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. + """ + 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): + 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 _snippet(text: str, limit: int = _SUMMARY_SNIPPET_LEN) -> str: """Truncate *text* to roughly *limit* chars, ending on a word boundary.""" @@ -234,6 +289,7 @@ def _fetch(alert: StreamingWatcherAlert) -> AlertTarget | None: 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, From 56d7982ae70af7dad279aaf920a4b82ecf55a371 Mon Sep 17 00:00:00 2001 From: mmolenda Date: Tue, 5 May 2026 11:58:51 +0200 Subject: [PATCH 13/16] Handling ForbiddenException in fetch_alert_targets --- verity471/helpers/alerts.py | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/verity471/helpers/alerts.py b/verity471/helpers/alerts.py index 7db9fe7..eb2e61b 100644 --- a/verity471/helpers/alerts.py +++ b/verity471/helpers/alerts.py @@ -8,6 +8,7 @@ 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 @@ -284,6 +285,9 @@ def _fetch(alert: StreamingWatcherAlert) -> AlertTarget | None: watcher=watchers_by_id.get(alert.watcher_id), watcher_group=groups_by_id.get(alert.watcher_group_id), ) + 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 From 8bb15458c24bd42fe101cd00024651aee5c453a5 Mon Sep 17 00:00:00 2001 From: mmolenda Date: Wed, 6 May 2026 12:38:43 +0200 Subject: [PATCH 14/16] Optimising fetch_alert_targets helper --- verity471/helpers/alerts.py | 50 +++++++++++++++++-------------------- 1 file changed, 23 insertions(+), 27 deletions(-) diff --git a/verity471/helpers/alerts.py b/verity471/helpers/alerts.py index eb2e61b..8d95641 100644 --- a/verity471/helpers/alerts.py +++ b/verity471/helpers/alerts.py @@ -252,20 +252,6 @@ def fetch_alert_targets( for r in fetch_alert_targets(alerts, api_client): print(r.alert.source_type, r.alert.status, r.target) """ - watchers_by_id: dict[int, GetWatcherResponse] = {} - groups_by_id: dict[int, GetWatcherGroupResponse] = {} - 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 [])} - except Exception: - log.warning("Failed to fetch watchers/watcher groups; watcher enrichment will be skipped", exc_info=True) - 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: @@ -279,12 +265,7 @@ def _fetch(alert: StreamingWatcherAlert) -> AlertTarget | None: 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, - watcher=watchers_by_id.get(alert.watcher_id), - watcher_group=groups_by_id.get(alert.watcher_group_id), - ) + 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 @@ -294,12 +275,7 @@ def _fetch(alert: StreamingWatcherAlert) -> AlertTarget | None: 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, - watcher=watchers_by_id.get(alert.watcher_id), - watcher_group=groups_by_id.get(alert.watcher_group_id), - ) + return AlertTarget(alert=alert, target=target) alerts = alerts_response.alerts or [] results: list[AlertTarget] = [None] * len(alerts) # type: ignore[list-item] @@ -309,4 +285,24 @@ def _fetch(alert: StreamingWatcherAlert) -> AlertTarget | None: result = future.result() if result is not None: results[future_to_index[future]] = result - return [r for r in results if r is not None] + + 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 From 2849048c832a3bcafd2a37aac701474f4f7e8985 Mon Sep 17 00:00:00 2001 From: mmolenda Date: Wed, 6 May 2026 16:52:06 +0200 Subject: [PATCH 15/16] alerts helper: fix geopol URL, safe field access --- verity471/helpers/alerts.py | 20 +++++++++++++++++--- 1 file changed, 17 insertions(+), 3 deletions(-) diff --git a/verity471/helpers/alerts.py b/verity471/helpers/alerts.py index 8d95641..7bef947 100644 --- a/verity471/helpers/alerts.py +++ b/verity471/helpers/alerts.py @@ -49,6 +49,18 @@ def _patch_portal_url(alert: StreamingWatcherAlert, target: Any) -> None: 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 @@ -75,8 +87,9 @@ def _patch_portal_url(alert: StreamingWatcherAlert, target: Any) -> None: url = msg.links.verity_portal.href + "?messageId=" + msg.id elif isinstance(target, GetCredSetResponse): - q = quote("cred_set.name=" + target.data.name, safe="") - url = f"https://verity.intel471.com/search?q={q}&category=creds_cred_set" + 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) @@ -167,7 +180,8 @@ def _summarize_target(target: Any) -> str | None: 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 + 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, value, conf])) From 4c043f1959c0c73347c8b39d093d07a7d13be645 Mon Sep 17 00:00:00 2001 From: mmolenda Date: Thu, 7 May 2026 12:35:43 +0200 Subject: [PATCH 16/16] alerts helper: better formatting of malware event --- tests/test_helpers.py | 119 +++++++++++++++++++++++++++++++++++- verity471/helpers/alerts.py | 53 +++++++++++++--- 2 files changed, 161 insertions(+), 11 deletions(-) diff --git a/tests/test_helpers.py b/tests/test_helpers.py index a5d614c..ef46de7 100644 --- a/tests/test_helpers.py +++ b/tests/test_helpers.py @@ -6,7 +6,9 @@ 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 @@ -20,7 +22,7 @@ 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'), + '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'), @@ -199,4 +201,117 @@ def test_credential_set_patches_url(self): def test_no_op_for_unmatched_target_type(self): alert = _make_alert() _patch_portal_url(alert, None) - assert alert.links.verity_portal is None \ No newline at end of file + 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/verity471/helpers/alerts.py b/verity471/helpers/alerts.py index 7bef947..7f8bea2 100644 --- a/verity471/helpers/alerts.py +++ b/verity471/helpers/alerts.py @@ -1,9 +1,10 @@ from __future__ import annotations import logging +import re import concurrent.futures from dataclasses import dataclass -from typing import Any +from typing import Any, Optional from urllib.parse import quote from verity471.api_client import ApiClient @@ -104,6 +105,11 @@ def _patch_portal_url(alert: StreamingWatcherAlert, target: Any) -> None: # --------------------------------------------------------------------------- +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: @@ -165,11 +171,14 @@ def _summarize_target(target: Any) -> str | None: if isinstance(target, GetCredOccurrenceResponse): return _prefixed("Credential Occurrence", _join([ - target.data.accessed_url, target.data.credential_type, target.last_updated_ts])) + _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, target.data.credential_domain, target.last_updated_ts])) + 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([ @@ -184,15 +193,41 @@ def _summarize_target(target: Any) -> str | None: 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, value, conf])) + return _prefixed("Indicator", _join([target.type, _defang(value) if value else None, conf])) if isinstance(target, IntegrationsEvent): - family = None - if target.threat and target.threat.data and target.threat.data.malware_family: - family = target.threat.data.malware_family.name + 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, target.data.attack_type if target.data else None])) + return _prefixed(label, _join([family_str] + parts)) if isinstance(target, SimplifiedMalwareProfile): aliases = ", ".join(target.aliases[:3]) if target.aliases else None