add agregated search
This commit is contained in:
@@ -10,6 +10,7 @@ from .config import ConfigStore
|
||||
from .correlation import correlate_source_ips
|
||||
from .event_context import build_event_context
|
||||
from .feedback import FeedbackStore
|
||||
from .graylog_aggregate import GraylogAggregateSource
|
||||
from .graylog_mcp import GraylogMcpClient
|
||||
from .graylog_source import GraylogStreamSource
|
||||
from .history import HistoryStore, StatusSnapshotStore
|
||||
@@ -87,6 +88,8 @@ def _stream_coverage(runtime_values: dict[str, object], stream_profiles: dict[st
|
||||
"total_fields": total_fields,
|
||||
"readiness": f"{ready_fields}/{total_fields}" if total_fields else "0/0",
|
||||
"events_fetched": int(status.get("events_fetched", 0) or 0),
|
||||
"aggregate_events": int(status.get("aggregate_events", 0) or 0),
|
||||
"aggregate_status": str(status.get("aggregate_status", "")),
|
||||
"latest_event_time": str(status.get("latest_event_time", "")),
|
||||
"truncated": bool(status.get("truncated")),
|
||||
"partial": bool(status.get("partial")),
|
||||
@@ -133,6 +136,10 @@ def build_status(
|
||||
events = []
|
||||
range_seconds = _range_seconds(runtime_values.get("graylog_range_seconds", 300))
|
||||
max_events_per_stream = max(1, int(runtime_values.get("graylog_max_events_per_stream", 5000) or 5000))
|
||||
raw_sample_events = max(1, int(runtime_values.get("graylog_raw_sample_events", 5000) or 5000))
|
||||
fetch_mode = str(runtime_values.get("graylog_fetch_mode", "auto") or "auto")
|
||||
use_aggregate = fetch_mode == "aggregate" or (fetch_mode == "auto" and max_events_per_stream > raw_sample_events)
|
||||
aggregate_events_total = 0
|
||||
for stream_config in stream_configs:
|
||||
stream_id = str(stream_config["id"])
|
||||
profile = stream_profiles.get(stream_id)
|
||||
@@ -144,12 +151,21 @@ def build_status(
|
||||
*tuple(str(field) for field in getattr(profile, "numeric_fields", ())),
|
||||
) if profile else ()
|
||||
stream_name = str(stream_config.get("title", "") or stream_titles.get(stream_id) or stream_id)
|
||||
stream_events, stream_status = GraylogStreamSource(GraylogMcpClient(url, token), stream_id, str(runtime_values.get("graylog_query", "*")), str(runtime_values.get("graylog_field_mapping", "")), stream_name, profile_fields).fetch(max_events=max_events_per_stream, range_seconds=range_seconds)
|
||||
aggregate_status: dict[str, object] = {}
|
||||
if use_aggregate:
|
||||
aggregate_status = GraylogAggregateSource(GraylogMcpClient(url, token), stream_id, str(runtime_values.get("graylog_query", "*"))).fetch_count(range_seconds=range_seconds)
|
||||
aggregate_events_total += int(aggregate_status.get("aggregate_events", 0) or 0)
|
||||
raw_limit = raw_sample_events if use_aggregate else max_events_per_stream
|
||||
stream_events, stream_status = GraylogStreamSource(GraylogMcpClient(url, token), stream_id, str(runtime_values.get("graylog_query", "*")), str(runtime_values.get("graylog_field_mapping", "")), stream_name, profile_fields).fetch(max_events=raw_limit, range_seconds=range_seconds)
|
||||
events.extend(stream_events)
|
||||
stream_statuses.append({"stream_id": stream_id, "stream_name": stream_name, **stream_status})
|
||||
truncated_streams = [item for item in stream_statuses if item.get("truncated")]
|
||||
stream_statuses.append({"stream_id": stream_id, "stream_name": stream_name, **aggregate_status, **stream_status, "raw_sample_limit": raw_limit})
|
||||
sample_limited_streams = [item for item in stream_statuses if item.get("truncated") and use_aggregate]
|
||||
truncated_streams = [item for item in stream_statuses if item.get("truncated") and not use_aggregate]
|
||||
partial_streams = [item for item in stream_statuses if item.get("partial")]
|
||||
aggregate_errors = [item for item in stream_statuses if item.get("aggregate_status") == "error"]
|
||||
warnings = []
|
||||
if aggregate_errors:
|
||||
warnings.append(f"{len(aggregate_errors)} stream(s) returned aggregate MCP errors.")
|
||||
if partial_streams:
|
||||
warnings.append(f"{len(partial_streams)} stream(s) returned a partial MCP fetch; Graylog likely timed out or rejected a large paged query.")
|
||||
if truncated_streams:
|
||||
@@ -158,10 +174,16 @@ def build_status(
|
||||
"status": "partial" if partial_streams else "connected",
|
||||
"streams": stream_statuses,
|
||||
"events_fetched": len(events),
|
||||
"raw_events_fetched": len(events),
|
||||
"aggregate_events": aggregate_events_total,
|
||||
"fetch_mode": "aggregate" if use_aggregate else "raw",
|
||||
"range_seconds": range_seconds,
|
||||
"max_events_per_stream": max_events_per_stream,
|
||||
"raw_sample_events": raw_sample_events,
|
||||
"partial_streams": len(partial_streams),
|
||||
"sample_limited_streams": len(sample_limited_streams),
|
||||
"truncated_streams": len(truncated_streams),
|
||||
"aggregate_error_streams": len(aggregate_errors),
|
||||
"coverage_status": "partial" if partial_streams else "truncated" if truncated_streams else "complete_window",
|
||||
"coverage_warning": " ".join(warnings),
|
||||
}
|
||||
@@ -275,11 +297,16 @@ def build_status(
|
||||
except Exception as exc:
|
||||
policy_error = str(exc)
|
||||
|
||||
summary = summarize_events(events)
|
||||
if isinstance(mcp_status, dict) and int(mcp_status.get("aggregate_events", 0) or 0) > summary.get("total", 0):
|
||||
summary["total"] = int(mcp_status.get("aggregate_events", 0) or 0)
|
||||
summary["raw_sample_total"] = len(events)
|
||||
summary["aggregate_backed"] = True
|
||||
status = {
|
||||
"generated_at": int(time.time()),
|
||||
"log_path": log_path,
|
||||
"policy_path": policy_path,
|
||||
"summary": summarize_events(events),
|
||||
"summary": summary,
|
||||
"anomaly_summary": anomaly_summary(anomalies),
|
||||
"baseline": {"enabled": bool(baseline), "sources_ready": len(profiles), "training_days": baseline_training_days, "new_events_recorded": baseline_events, "profile_fields_recorded": profile_baseline_fields, "maintenance": baseline_maintenance, "size_bytes": baseline_maintenance.get("size_bytes", 0) if isinstance(baseline_maintenance, dict) else 0},
|
||||
"capabilities": {"threat_intel": threat_intel_status, "graylog_mcp": mcp_status, "profile_advisor": profile_advisor_status},
|
||||
|
||||
Reference in New Issue
Block a user