From f99e699da6c99ed399d8c6b83a35328a1e1c32edd393811d49aaac511a1a8f56 Mon Sep 17 00:00:00 2001 From: larssand Date: Sun, 21 Jun 2026 22:31:20 +0200 Subject: [PATCH] add 5000 events --- src/fgai/graylog_source.py | 42 ++++++++++++++++++++---------------- tests/test_graylog_source.py | 1 + 2 files changed, 24 insertions(+), 19 deletions(-) diff --git a/src/fgai/graylog_source.py b/src/fgai/graylog_source.py index 1da8fc8..fed2612 100644 --- a/src/fgai/graylog_source.py +++ b/src/fgai/graylog_source.py @@ -52,34 +52,38 @@ class GraylogStreamSource: if not isinstance(self.mapping, dict): raise RuntimeError("invalid_graylog_field_mapping") - def fetch(self) -> tuple[list[LogEvent], dict[str, object]]: + def fetch(self, *, max_events: int = 5_000) -> tuple[list[LogEvent], dict[str, object]]: status = self.client.probe() mapping_fields = [str(value) for value in self.mapping.values() if isinstance(value, str)] arguments: dict[str, object] = { "query": self.query, - "limit": 1000, "range_seconds": 300, "fields": list(dict.fromkeys([*DEFAULT_FIELDS, *mapping_fields])), } if self.stream: arguments["streams"] = [self.stream] - result = self.client.call_tool("search_messages", arguments) - content = result.get("result", {}).get("content", []) if isinstance(result.get("result"), dict) else [] - if isinstance(result.get("result"), dict) and result["result"].get("isError"): - detail = next( - (str(item.get("text")) for item in content if isinstance(item, dict) and item.get("type") == "text"), - "Graylog search failed", - ) - raise RuntimeError(f"graylog_search_error: {detail}") - records: list[dict[str, object]] = [] - for item in content if isinstance(content, list) else []: - if isinstance(item, dict) and item.get("type") == "text": - try: - records.extend(_records(json.loads(str(item.get("text", ""))))) - except json.JSONDecodeError: - continue - events = [self._event(record) for record in records] - status.update({"source": "graylog_mcp", "events_fetched": len(events)}) + events: list[LogEvent] = [] + page_size = 1_000 + pages = 0 + while len(events) < max_events: + result = self.client.call_tool("search_messages", {**arguments, "limit": page_size, "offset": len(events)}) + content = result.get("result", {}).get("content", []) if isinstance(result.get("result"), dict) else [] + if isinstance(result.get("result"), dict) and result["result"].get("isError"): + detail = next((str(item.get("text")) for item in content if isinstance(item, dict) and item.get("type") == "text"), "Graylog search failed") + raise RuntimeError(f"graylog_search_error: {detail}") + records: list[dict[str, object]] = [] + for item in content if isinstance(content, list) else []: + if isinstance(item, dict) and item.get("type") == "text": + try: + records.extend(_records(json.loads(str(item.get("text", ""))))) + except json.JSONDecodeError: + continue + events.extend(self._event(record) for record in records) + pages += 1 + if len(records) < page_size: + break + latest = max((event.fields.get("eventtime", "") for event in events), default="") + status.update({"source": "graylog_mcp", "events_fetched": len(events), "pages": pages, "truncated": len(events) >= max_events, "latest_event_time": latest}) return events, status def _event(self, record: dict[str, object]) -> LogEvent: diff --git a/tests/test_graylog_source.py b/tests/test_graylog_source.py index e645090..268792e 100644 --- a/tests/test_graylog_source.py +++ b/tests/test_graylog_source.py @@ -28,6 +28,7 @@ class GraylogSourceTests(unittest.TestCase): self.assertEqual(client.arguments["streams"], ["vpn"]) self.assertEqual(client.arguments["range_seconds"], 300) self.assertIn("client", client.arguments["fields"]) + self.assertEqual(client.arguments["offset"], 0) if __name__ == "__main__":