diff --git a/mcp-server/src/streamops_mcp/agent/monitor.py b/mcp-server/src/streamops_mcp/agent/monitor.py index f8c9394..6dbbd4c 100644 --- a/mcp-server/src/streamops_mcp/agent/monitor.py +++ b/mcp-server/src/streamops_mcp/agent/monitor.py @@ -25,6 +25,7 @@ from streamops_mcp.agent.executor import execute_tool from streamops_mcp.agent.schemas import ( Confidence, + DetectedAnomaly, DiagnosisReport, DiagnosticToReportHandoff, IncidentReport, @@ -92,11 +93,15 @@ async def _retry_subagent(self, name: str, coro_factory): retryable = self._is_retryable(exc) logger.error( "%s failed (attempt %d/%d, retryable=%s): %s", - name, attempt + 1, 1 + max_retries, retryable, exc, + name, + attempt + 1, + 1 + max_retries, + retryable, + exc, ) if not retryable or attempt >= max_retries: raise - delay = base_delay * (2 ** attempt) + delay = base_delay * (2**attempt) logger.info("Retrying %s in %.1fs", name, delay) await asyncio.sleep(delay) @@ -117,7 +122,8 @@ async def run_cycle(self) -> IncidentReport | None: try: logger.info( "Starting monitoring cycle (cycle_id=%s, multi_agent=%s)", - cycle_id, self.multi_agent, + cycle_id, + self.multi_agent, ) detection = await self._detect_anomalies() @@ -133,7 +139,8 @@ async def run_cycle(self) -> IncidentReport | None: ) except Exception as exc: logger.error( - "Diagnostic Agent failed after retries, cycle aborted: %s", exc, + "Diagnostic Agent failed after retries, cycle aborted: %s", + exc, ) return None @@ -156,7 +163,8 @@ async def run_cycle(self) -> IncidentReport | None: ) except Exception as exc: logger.warning( - "Report Agent failed after retries, producing fallback report: %s", exc, + "Report Agent failed after retries, producing fallback report: %s", + exc, ) report = self._fallback_report(diagnosis) @@ -174,8 +182,10 @@ async def run_cycle(self) -> IncidentReport | None: for conflict in diagnosis.conflicts: logger.warning( "Conflict %s [%s]: claims %s vs %s", - conflict.conflict_id, conflict.topic, - conflict.claim_a_id, conflict.claim_b_id, + conflict.conflict_id, + conflict.topic, + conflict.claim_a_id, + conflict.claim_b_id, ) await escalate(report, diagnosis=diagnosis) @@ -222,12 +232,26 @@ def _extract_low_confidence_claims(diagnosis: DiagnosisReport) -> list[str]: low_levels = {Confidence.LOW, Confidence.UNSOURCED} return [c.text for c in diagnosis.claims if c.confidence in low_levels] - async def _detect_anomalies(self) -> str | None: + async def _detect_anomalies(self) -> DetectedAnomaly | None: """Run the agentic loop to poll infrastructure and detect anomalies. - Returns the assistant's final text if anomalies were found, None if healthy. + Returns a structured DetectedAnomaly if an anomaly was found, None if + healthy. The agent explores with tools, then, on concluding an anomaly + exists, emits a DetectedAnomaly JSON so the handoff to the Diagnostic + agent carries typed context instead of prose. """ - messages: list[dict[str, Any]] = [{"role": "user", "content": "Run a health check on the streaming infrastructure. Check Flink jobs, consumer lag, and recent events. Report any anomalies you find."}] + detection_schema = DetectedAnomaly.model_json_schema() + messages: list[dict[str, Any]] = [ + { + "role": "user", + "content": ( + "Run a health check on the streaming infrastructure. Check Flink jobs, " + "consumer lag, and recent events. If everything is healthy, say so briefly. " + "If you detect an anomaly, respond with ONLY a JSON object matching this " + f"DetectedAnomaly schema:\n{detection_schema}" + ), + } + ] for round_num in range(self.max_tool_rounds): logger.debug("Detection loop round %d", round_num + 1) @@ -256,9 +280,9 @@ async def _detect_anomalies(self) -> str | None: if response.stop_reason == "end_turn": logger.info("Detection complete after %d rounds", round_num + 1) if self._mentions_anomaly(assistant_text): - summary = assistant_text[:300].replace("\n", " ") - logger.info("Anomaly detected: %s", summary) - return assistant_text + anomaly = self._parse_detection(assistant_text) + logger.info("Anomaly detected: %s", anomaly.summary[:300]) + return anomaly logger.info("No anomalies found in detection response") return None @@ -266,11 +290,13 @@ async def _detect_anomalies(self) -> str | None: tool_results = [] for tool_call in tool_calls: result = await execute_tool(tool_call.name, tool_call.input) - tool_results.append({ - "type": "tool_result", - "tool_use_id": tool_call.id, - "content": result, - }) + tool_results.append( + { + "type": "tool_result", + "tool_use_id": tool_call.id, + "content": result, + } + ) messages.append({"role": "user", "content": tool_results}) else: @@ -278,41 +304,50 @@ async def _detect_anomalies(self) -> str | None: break logger.warning("Detection loop hit max rounds (%d)", self.max_tool_rounds) - return assistant_text if self._mentions_anomaly(assistant_text) else None + return ( + self._parse_detection(assistant_text) + if self._mentions_anomaly(assistant_text) + else None + ) - async def _spawn_diagnostic_agent(self, anomaly_context: str) -> DiagnosisReport: + async def _spawn_diagnostic_agent(self, anomaly: DetectedAnomaly) -> DiagnosisReport: """Spawn a Diagnostic sub-agent with scoped context and tools. The sub-agent starts with a blank context. All relevant information must be injected explicitly via the prompt (not inherited from the - coordinator's conversation history). + coordinator's conversation history), as a typed DetectedAnomaly rather + than a prose string. """ logger.info("Spawning Diagnostic Agent") schema_hint = DiagnosisReport.model_json_schema() handoff = MonitorToDiagnosticHandoff( - anomaly_context=anomaly_context, + anomaly=anomaly, schema_hint=schema_hint, ) + anomaly_json = handoff.anomaly.model_dump_json(indent=2) logger.info( - "Monitor->Diagnostic handoff validated (%d chars)", - len(handoff.anomaly_context), + "Monitor->Diagnostic handoff validated (type=%s, %d chars)", + handoff.anomaly.anomaly_type, + len(anomaly_json), ) system_prompt = DIAGNOSTIC_SYSTEM_PROMPT - runbook_section = self._resolve_runbooks(handoff.anomaly_context) + runbook_section = self._resolve_runbooks(handoff.anomaly.summary) if runbook_section: system_prompt = system_prompt + "\n\n" + runbook_section - messages: list[dict[str, Any]] = [{ - "role": "user", - "content": f"""Investigate the following anomaly detected by the monitoring system: + messages: list[dict[str, Any]] = [ + { + "role": "user", + "content": f"""Investigate the following anomaly detected by the monitoring system: -{handoff.anomaly_context} +{anomaly_json} Use the available tools to determine the root cause. Respond with a JSON object matching the DiagnosisReport schema: {handoff.schema_hint}""", - }] + } + ] for round_num in range(self.max_tool_rounds): if not messages or messages[-1]["role"] != "user": @@ -347,11 +382,13 @@ async def _spawn_diagnostic_agent(self, anomaly_context: str) -> DiagnosisReport tool_results = [] for tool_call in tool_calls: result = await execute_tool(tool_call.name, tool_call.input) - tool_results.append({ - "type": "tool_result", - "tool_use_id": tool_call.id, - "content": result, - }) + tool_results.append( + { + "type": "tool_result", + "tool_use_id": tool_call.id, + "content": result, + } + ) messages.append({"role": "user", "content": tool_results}) logger.warning("Diagnostic Agent hit max rounds") @@ -381,15 +418,17 @@ async def _spawn_report_agent(self, diagnosis: DiagnosisReport) -> IncidentRepor model=self.model, max_tokens=config.agent_max_tokens, system=REPORT_SYSTEM_PROMPT, - messages=[{ - "role": "user", - "content": f"""Produce an incident report from this diagnosis: + messages=[ + { + "role": "user", + "content": f"""Produce an incident report from this diagnosis: {handoff.diagnosis_json} Respond with a JSON object matching the IncidentReport schema: {handoff.schema_hint}""", - }], + } + ], ) text = "".join(b.text for b in response.content if b.type == "text") @@ -423,28 +462,57 @@ def _resolve_runbooks(anomaly_context: str) -> str: def _mentions_anomaly(self, text: str) -> bool: """Simple heuristic: does the text suggest an anomaly was found?""" anomaly_keywords = [ - "anomaly", "spike", "degraded", "failing", "exceeded", "threshold", - "timeout", "error", "critical", "backpressure", "lag", "pressure", + "anomaly", + "spike", + "degraded", + "failing", + "exceeded", + "threshold", + "timeout", + "error", + "critical", + "backpressure", + "lag", + "pressure", ] text_lower = text.lower() return any(kw in text_lower for kw in anomaly_keywords) - def _extract_diagnosis_from_detection(self, detection_text: str) -> DiagnosisReport: - """In single-agent mode, build a diagnosis from the detection text.""" + def _extract_diagnosis_from_detection(self, anomaly: DetectedAnomaly) -> DiagnosisReport: + """In single-agent mode, build a diagnosis from the detected anomaly.""" return DiagnosisReport( - anomaly_type="unknown", - detected_at=datetime.now(UTC).isoformat(), + anomaly_type=anomaly.anomaly_type, + detected_at=anomaly.detected_at, affected_components=[], root_cause=RootCause( - summary="See detection notes below", + summary=anomaly.summary, confidence="medium", - reasoning=detection_text[:1000], - supporting_metrics=[], + reasoning=anomaly.summary, + supporting_metrics=[ + m for m in (anomaly.metric, anomaly.observed_value) if m is not None + ], ), tools_used=[], - raw_evidence=[detection_text[:500]], + raw_evidence=[anomaly.summary], ) + def _parse_detection(self, text: str) -> DetectedAnomaly: + """Extract a DetectedAnomaly from the monitor's response. + + Falls back to wrapping the raw text in ``summary`` when the model + emitted prose instead of JSON, so detection always yields typed context. + """ + try: + json_str = self._extract_json(text) + return DetectedAnomaly.model_validate_json(json_str) + except Exception as e: + logger.warning("Failed to parse DetectedAnomaly: %s, using prose fallback", e) + return DetectedAnomaly( + anomaly_type="unknown", + summary=text[:1000], + detected_at=datetime.now(UTC).isoformat(), + ) + def _parse_diagnosis(self, text: str) -> DiagnosisReport: """Extract a DiagnosisReport from the agent's text response.""" try: diff --git a/mcp-server/src/streamops_mcp/agent/schemas/__init__.py b/mcp-server/src/streamops_mcp/agent/schemas/__init__.py index fa2fe3e..7ef847b 100644 --- a/mcp-server/src/streamops_mcp/agent/schemas/__init__.py +++ b/mcp-server/src/streamops_mcp/agent/schemas/__init__.py @@ -7,6 +7,7 @@ SourceRecord, ) from streamops_mcp.agent.schemas.handoff import ( + DetectedAnomaly, DiagnosticToReportHandoff, MonitorToDiagnosticHandoff, ) @@ -16,6 +17,7 @@ "ClaimRecord", "Confidence", "ConflictRecord", + "DetectedAnomaly", "DiagnosisReport", "DiagnosticToReportHandoff", "IncidentReport", diff --git a/mcp-server/src/streamops_mcp/agent/schemas/handoff.py b/mcp-server/src/streamops_mcp/agent/schemas/handoff.py index 9c5a680..8531748 100644 --- a/mcp-server/src/streamops_mcp/agent/schemas/handoff.py +++ b/mcp-server/src/streamops_mcp/agent/schemas/handoff.py @@ -3,6 +3,11 @@ Each handoff between agents is validated through a typed Pydantic model. This catches malformed or oversized payloads before they reach the LLM, preventing silent failures or token overflow. + +Context between agents is passed as structured data, never prose summaries a +downstream agent would have to re-parse. The Monitor hands the Diagnostic agent +a typed DetectedAnomaly (what breached, by how much, where); the Diagnostic +hands the Report agent the full serialized DiagnosisReport. """ from pydantic import BaseModel, Field, model_validator @@ -10,11 +15,58 @@ from streamops_mcp.config import config +class DetectedAnomaly(BaseModel): + """Typed description of an anomaly the Monitor detected. + + Replaces the old free-form ``anomaly_context`` string so the Diagnostic + sub-agent receives unambiguous, typed context (metric, observed vs baseline, + breach direction, component) instead of prose it has to re-parse. ``summary`` + keeps a one-line human-readable form for logging and runbook matching; the + optional fields are best-effort, since not every signal exposes all of them. + """ + + anomaly_type: str = Field( + description=( + "Category: latency_spike, throughput_drop, backpressure, " + "checkpoint_failure, memory_pressure, error_burst, or unknown" + ) + ) + summary: str = Field( + description="One-line human-readable description citing the metric and value" + ) + detected_at: str = Field(description="ISO-8601 timestamp when the anomaly was detected") + metric: str | None = Field( + default=None, + description="Name of the breached metric (e.g. 'consumer_lag', 'processing_latency_ms')", + ) + observed_value: str | None = Field( + default=None, + description="Observed value, as a string to preserve units (e.g. '2,340ms', '85000')", + ) + baseline: str | None = Field( + default=None, description="Normal/expected value for context (e.g. '200ms')" + ) + threshold: str | None = Field( + default=None, description="The threshold that was breached, if applicable" + ) + breach_direction: str | None = Field( + default=None, description="Direction of the breach: 'above' or 'below'" + ) + affected_component: str | None = Field( + default=None, + description="Component the anomaly centers on (e.g. 'sink-kafka', 'streamops-processor')", + ) + source_signal_ids: list[str] = Field( + default_factory=list, + description="Identifiers of the tool signals/metrics that triggered detection (audit trail)", + ) + + class MonitorToDiagnosticHandoff(BaseModel): """Payload passed from the Monitor agent to the Diagnostic sub-agent.""" - anomaly_context: str = Field( - description="Raw anomaly description from the monitoring loop", + anomaly: DetectedAnomaly = Field( + description="Structured description of the detected anomaly", ) schema_hint: dict = Field( description="JSON schema the diagnostic agent must conform to", @@ -22,9 +74,11 @@ class MonitorToDiagnosticHandoff(BaseModel): @model_validator(mode="after") def validate_context_size(self): + # summary is the only unbounded field; truncate it (not raise) so an + # over-long detection narrative can never overflow the diagnostic prompt. max_chars = config.agent_handoff_max_context_chars - if len(self.anomaly_context) > max_chars: - self.anomaly_context = self.anomaly_context[:max_chars] + if len(self.anomaly.summary) > max_chars: + self.anomaly.summary = self.anomaly.summary[:max_chars] return self diff --git a/mcp-server/tests/test_monitor.py b/mcp-server/tests/test_monitor.py index 2de65c8..62831b1 100644 --- a/mcp-server/tests/test_monitor.py +++ b/mcp-server/tests/test_monitor.py @@ -36,7 +36,6 @@ def agent(): class TestAnomalyDetection: - def test_detects_anomaly_keywords(self, agent): # Arrange text = "Consumer lag spike detected on partition 3, lag is 450,000 records" @@ -68,8 +67,84 @@ def test_case_insensitive(self, agent): assert result is True -class TestJsonExtraction: +class TestDetectionParsing: + """_detect_anomalies must hand the Diagnostic agent a typed DetectedAnomaly.""" + + def test_parses_structured_detection(self, agent): + # Arrange: the monitor emitted a well-formed DetectedAnomaly JSON + text = json.dumps( + { + "anomaly_type": "latency_spike", + "summary": "Processing latency 2,340ms (threshold 200ms)", + "detected_at": "2026-07-09T12:00:00Z", + "metric": "processing_latency_ms", + "observed_value": "2340", + "threshold": "200", + "breach_direction": "above", + "affected_component": "streamops-processor", + "source_signal_ids": ["query_flink_metrics:latency"], + } + ) + + # Act + anomaly = agent._parse_detection(text) + + # Assert: typed fields, not a string to re-parse + assert anomaly.anomaly_type == "latency_spike" + assert anomaly.observed_value == "2340" + assert anomaly.source_signal_ids == ["query_flink_metrics:latency"] + + def test_prose_falls_back_to_summary(self, agent): + # Arrange: the monitor wrote prose instead of JSON + text = "Consumer lag spike detected on partition 3, lag is 450,000 records" + + # Act + anomaly = agent._parse_detection(text) + + # Assert: still a typed object, prose preserved in summary, never crashes + assert anomaly.anomaly_type == "unknown" + assert "450,000" in anomaly.summary + assert anomaly.detected_at # stamped + + @pytest.mark.asyncio + async def test_detect_anomalies_returns_typed_anomaly(self, agent): + # Arrange: detection concludes with a structured anomaly + agent.client = _mock_client( + _api_response( + json.dumps( + { + "anomaly_type": "backpressure", + "summary": "Backpressure ratio 0.87 on sink-kafka", + "detected_at": "2026-07-09T12:00:00Z", + "affected_component": "sink-kafka", + } + ) + ) + ) + + # Act + result = await agent._detect_anomalies() + + # Assert + assert result is not None + assert result.anomaly_type == "backpressure" + assert result.affected_component == "sink-kafka" + + @pytest.mark.asyncio + async def test_detect_anomalies_returns_none_when_healthy(self, agent): + # Arrange: no anomaly keywords in the response + agent.client = _mock_client( + _api_response("All systems healthy. Flink job running, everything nominal.") + ) + + # Act + result = await agent._detect_anomalies() + # Assert + assert result is None + + +class TestJsonExtraction: def test_extracts_from_code_block(self, agent): # Arrange text = 'Here is the report:\n```json\n{"key": "value"}\n```\nDone.' @@ -102,10 +177,9 @@ def test_extracts_from_generic_code_block(self, agent): class TestDiagnosisParsing: - def test_valid_json_parses(self, agent): # Arrange - text = '''```json + text = """```json { "anomaly_type": "latency_spike", "detected_at": "2026-06-18T15:00:00Z", @@ -117,7 +191,7 @@ def test_valid_json_parses(self, agent): }, "tools_used": ["query_flink_jobs"] } -```''' +```""" # Act result = agent._parse_diagnosis(text) @@ -141,7 +215,6 @@ def test_invalid_json_returns_fallback(self, agent): class TestIncidentParsing: - def test_valid_json_parses(self, agent): # Arrange diagnosis = DiagnosisReport( @@ -151,7 +224,7 @@ def test_valid_json_parses(self, agent): root_cause={"summary": "test", "confidence": "low", "reasoning": "test"}, tools_used=[], ) - text = '''{ + text = """{ "incident_id": "inc-001", "title": "Test incident", "severity": "HIGH", @@ -162,7 +235,7 @@ def test_valid_json_parses(self, agent): "timeline": ["event 1"], "recommended_actions": [{"action": "fix", "rationale": "why", "risk": "low", "requires_downtime": false}], "monitoring_notes": "watch it" -}''' +}""" # Act result = agent._parse_incident(text, diagnosis) @@ -192,7 +265,6 @@ def test_invalid_json_returns_fallback(self, agent): class TestIsRetryable: - def test_timeout_is_retryable(self, agent): # Arrange exc = anthropic.APITimeoutError(request=None) @@ -253,7 +325,6 @@ def test_value_error_is_not_retryable(self, agent): class TestRetrySubagent: - @pytest.mark.asyncio async def test_succeeds_on_first_try(self, agent): # Arrange @@ -323,16 +394,24 @@ async def factory(): class TestFallbackReport: - def test_produces_valid_incident_report(self): # Arrange diagnosis = DiagnosisReport( anomaly_type="latency_spike", detected_at="2026-06-23T12:00:00Z", affected_components=[ - {"name": "flink-job", "role": "processor", "status": "degraded", "evidence": "p99 > 5s"}, + { + "name": "flink-job", + "role": "processor", + "status": "degraded", + "evidence": "p99 > 5s", + }, ], - root_cause={"summary": "GC pressure on TaskManager", "confidence": "high", "reasoning": "heap at 95%"}, + root_cause={ + "summary": "GC pressure on TaskManager", + "confidence": "high", + "reasoning": "heap at 95%", + }, tools_used=["query_flink_jobs"], ) @@ -353,7 +432,11 @@ def test_handles_empty_components(self): anomaly_type="unknown", detected_at="2026-06-23T12:00:00Z", affected_components=[], - root_cause={"summary": "unclear", "confidence": "low", "reasoning": "insufficient data"}, + root_cause={ + "summary": "unclear", + "confidence": "low", + "reasoning": "insufficient data", + }, tools_used=[], ) @@ -397,12 +480,16 @@ def _make_diagnosis_with_claims(claim_confidences: list[Confidence]) -> Diagnosi class TestConfidenceDistribution: - def test_logs_distribution(self, agent, caplog): # Arrange - diagnosis = _make_diagnosis_with_claims([ - Confidence.HIGH, Confidence.HIGH, Confidence.MEDIUM, Confidence.LOW, - ]) + diagnosis = _make_diagnosis_with_claims( + [ + Confidence.HIGH, + Confidence.HIGH, + Confidence.MEDIUM, + Confidence.LOW, + ] + ) # Act with caplog.at_level("INFO", logger="streamops-mcp.monitor"): @@ -416,7 +503,6 @@ def test_logs_distribution(self, agent, caplog): class TestAllClaimsLowConfidence: - def test_all_low_returns_true(self): # Arrange diagnosis = _make_diagnosis_with_claims([Confidence.LOW, Confidence.LOW]) @@ -440,9 +526,13 @@ def test_mixed_low_unsourced_returns_true(self): def test_one_medium_returns_false(self): # Arrange - diagnosis = _make_diagnosis_with_claims([ - Confidence.LOW, Confidence.MEDIUM, Confidence.UNSOURCED, - ]) + diagnosis = _make_diagnosis_with_claims( + [ + Confidence.LOW, + Confidence.MEDIUM, + Confidence.UNSOURCED, + ] + ) # Act / Assert assert MonitorAgent._all_claims_low_confidence(diagnosis) is False @@ -456,12 +546,16 @@ def test_no_claims_returns_false(self): class TestExtractLowConfidenceClaims: - def test_extracts_low_and_unsourced(self): # Arrange - diagnosis = _make_diagnosis_with_claims([ - Confidence.HIGH, Confidence.LOW, Confidence.MEDIUM, Confidence.UNSOURCED, - ]) + diagnosis = _make_diagnosis_with_claims( + [ + Confidence.HIGH, + Confidence.LOW, + Confidence.MEDIUM, + Confidence.UNSOURCED, + ] + ) # Act result = MonitorAgent._extract_low_confidence_claims(diagnosis) @@ -503,63 +597,95 @@ def _api_response(text: str, stop_reason: str = "end_turn") -> SimpleNamespace: def _mock_client(*responses) -> SimpleNamespace: """A stand-in Anthropic async client whose messages.create yields the given responses (or raises, if a response is an Exception) in order.""" - return SimpleNamespace( - messages=SimpleNamespace(create=AsyncMock(side_effect=list(responses))) - ) - + return SimpleNamespace(messages=SimpleNamespace(create=AsyncMock(side_effect=list(responses)))) + + +_DIAGNOSIS_JSON = json.dumps( + { + "anomaly_type": "throughput_drop", + "detected_at": "2026-07-03T12:00:00Z", + "sources": [ + { + "source_id": "src-001", + "tool_name": "query_flink_jobs", + "retrieved_at": "2026-07-03T12:00:00Z", + "raw_output": "{}", + }, + ], + "claims": [ + { + "claim_id": "C01", + "text": "Consumer lag is 45000 on partition 2", + "source_id": "src-001", + "confidence": "HIGH", + }, + ], + "affected_components": [ + { + "name": "kafka-consumer", + "role": "consumer", + "status": "degraded", + "evidence": "lag 45000", + }, + ], + "root_cause": { + "summary": "slow downstream sink", + "confidence": "high", + "reasoning": "lag climbing steadily", + "supporting_metrics": [], + }, + "tools_used": ["query_flink_jobs"], + "raw_evidence": [], + } +) -_DIAGNOSIS_JSON = json.dumps({ - "anomaly_type": "throughput_drop", - "detected_at": "2026-07-03T12:00:00Z", - "sources": [ - {"source_id": "src-001", "tool_name": "query_flink_jobs", - "retrieved_at": "2026-07-03T12:00:00Z", "raw_output": "{}"}, - ], - "claims": [ - {"claim_id": "C01", "text": "Consumer lag is 45000 on partition 2", - "source_id": "src-001", "confidence": "HIGH"}, - ], - "affected_components": [ - {"name": "kafka-consumer", "role": "consumer", "status": "degraded", - "evidence": "lag 45000"}, - ], - "root_cause": {"summary": "slow downstream sink", "confidence": "high", - "reasoning": "lag climbing steadily", "supporting_metrics": []}, - "tools_used": ["query_flink_jobs"], - "raw_evidence": [], -}) - -_LOW_DIAGNOSIS_JSON = json.dumps({ - "anomaly_type": "latency_spike", - "detected_at": "2026-07-03T12:00:00Z", - "sources": [ - {"source_id": "src-001", "tool_name": "query_prometheus", - "retrieved_at": "2026-07-03T12:00:00Z", "raw_output": "{}"}, - ], - "claims": [ - {"claim_id": "C01", "text": "Possibly elevated latency", - "source_id": "src-001", "confidence": "LOW"}, - ], - "affected_components": [], - "root_cause": {"summary": "unclear", "confidence": "low", - "reasoning": "single weak signal", "supporting_metrics": []}, - "tools_used": ["query_prometheus"], - "raw_evidence": [], -}) +_LOW_DIAGNOSIS_JSON = json.dumps( + { + "anomaly_type": "latency_spike", + "detected_at": "2026-07-03T12:00:00Z", + "sources": [ + { + "source_id": "src-001", + "tool_name": "query_prometheus", + "retrieved_at": "2026-07-03T12:00:00Z", + "raw_output": "{}", + }, + ], + "claims": [ + { + "claim_id": "C01", + "text": "Possibly elevated latency", + "source_id": "src-001", + "confidence": "LOW", + }, + ], + "affected_components": [], + "root_cause": { + "summary": "unclear", + "confidence": "low", + "reasoning": "single weak signal", + "supporting_metrics": [], + }, + "tools_used": ["query_prometheus"], + "raw_evidence": [], + } +) -_REPORT_JSON = json.dumps({ - "incident_id": "inc-001", - "title": "Consumer lag spike on partition 2", - "severity": "LOW", - "summary": "Consumer lag elevated but within recoverable range.", - "anomaly_type": "throughput_drop", - "root_cause": "slow downstream sink", - "affected_components": ["kafka-consumer"], - "timeline": ["Lag began climbing at 12:00"], - "recommended_actions": [], - "monitoring_notes": "Watch consumer lag over the next 15 minutes.", - "requires_human_approval": False, -}) +_REPORT_JSON = json.dumps( + { + "incident_id": "inc-001", + "title": "Consumer lag spike on partition 2", + "severity": "LOW", + "summary": "Consumer lag elevated but within recoverable range.", + "anomaly_type": "throughput_drop", + "root_cause": "slow downstream sink", + "affected_components": ["kafka-consumer"], + "timeline": ["Lag began climbing at 12:00"], + "recommended_actions": [], + "monitoring_notes": "Watch consumer lag over the next 15 minutes.", + "requires_human_approval": False, + } +) class TestRunCycleOrchestration: @@ -580,7 +706,8 @@ async def test_full_handoff_produces_incident_and_escalates(self, agent): # Act with patch( - "streamops_mcp.agent.monitor.escalate", new_callable=AsyncMock, + "streamops_mcp.agent.monitor.escalate", + new_callable=AsyncMock, ) as mock_escalate: result = await agent.run_cycle() @@ -600,7 +727,8 @@ async def test_no_anomaly_returns_none_and_skips_escalation(self, agent): # Act with patch( - "streamops_mcp.agent.monitor.escalate", new_callable=AsyncMock, + "streamops_mcp.agent.monitor.escalate", + new_callable=AsyncMock, ) as mock_escalate: result = await agent.run_cycle() @@ -619,7 +747,8 @@ async def test_all_low_confidence_skips_report_agent(self, agent): # Act with patch( - "streamops_mcp.agent.monitor.escalate", new_callable=AsyncMock, + "streamops_mcp.agent.monitor.escalate", + new_callable=AsyncMock, ) as mock_escalate: result = await agent.run_cycle() @@ -641,7 +770,8 @@ async def test_report_agent_failure_uses_fallback(self, agent, monkeypatch): # Act with patch( - "streamops_mcp.agent.monitor.escalate", new_callable=AsyncMock, + "streamops_mcp.agent.monitor.escalate", + new_callable=AsyncMock, ) as mock_escalate: result = await agent.run_cycle() @@ -662,7 +792,8 @@ async def test_diagnostic_failure_aborts_cycle(self, agent, monkeypatch): # Act with patch( - "streamops_mcp.agent.monitor.escalate", new_callable=AsyncMock, + "streamops_mcp.agent.monitor.escalate", + new_callable=AsyncMock, ) as mock_escalate: result = await agent.run_cycle() @@ -687,7 +818,8 @@ async def _capture(*args, **kwargs): # Act with patch( - "streamops_mcp.agent.monitor.escalate", new=_capture, + "streamops_mcp.agent.monitor.escalate", + new=_capture, ): await agent.run_cycle() diff --git a/mcp-server/tests/test_schemas.py b/mcp-server/tests/test_schemas.py index cf9f94f..4c0bca7 100644 --- a/mcp-server/tests/test_schemas.py +++ b/mcp-server/tests/test_schemas.py @@ -9,6 +9,7 @@ ClaimRecord, Confidence, ConflictRecord, + DetectedAnomaly, DiagnosisReport, DiagnosticToReportHandoff, IncidentReport, @@ -784,19 +785,66 @@ def test_requires_human_approval_in_json_schema(self): assert "requires_human_approval" in str(schema) +def _make_anomaly(summary: str = "Latency spike: 2,340ms (threshold 200ms)") -> DetectedAnomaly: + return DetectedAnomaly( + anomaly_type="latency_spike", + summary=summary, + detected_at="2026-07-09T12:00:00Z", + metric="processing_latency_ms", + observed_value="2340", + baseline="180", + threshold="200", + breach_direction="above", + affected_component="streamops-processor", + source_signal_ids=["query_flink_metrics:latency"], + ) + + +class TestDetectedAnomaly: + def test_minimal_required_fields(self): + # Arrange + Act: only the three required fields + anomaly = DetectedAnomaly( + anomaly_type="backpressure", + summary="Backpressure ratio 0.87 on sink-kafka", + detected_at="2026-07-09T12:00:00Z", + ) + + # Assert: optional typed fields default cleanly, ids default to empty + assert anomaly.metric is None + assert anomaly.source_signal_ids == [] + + def test_missing_required_field_rejected(self): + # Arrange + Act + Assert: summary is required + with pytest.raises(ValidationError): + DetectedAnomaly(anomaly_type="latency_spike", detected_at="2026-07-09T12:00:00Z") + + def test_typed_fields_round_trip(self): + # Arrange + anomaly = _make_anomaly() + + # Act + restored = DetectedAnomaly.model_validate_json(anomaly.model_dump_json()) + + # Assert + assert restored.observed_value == "2340" + assert restored.breach_direction == "above" + assert restored.source_signal_ids == ["query_flink_metrics:latency"] + + class TestMonitorToDiagnosticHandoff: def test_valid_handoff(self): # Arrange + Act handoff = MonitorToDiagnosticHandoff( - anomaly_context="Latency spike detected on partition 2", + anomaly=_make_anomaly(), schema_hint=DiagnosisReport.model_json_schema(), ) - # Assert - assert handoff.anomaly_context == "Latency spike detected on partition 2" + # Assert: typed context, not a string to re-parse + assert handoff.anomaly.anomaly_type == "latency_spike" + assert handoff.anomaly.observed_value == "2340" assert "anomaly_type" in str(handoff.schema_hint) - def test_truncates_oversized_context(self, monkeypatch): + def test_truncates_oversized_summary(self, monkeypatch): # Arrange monkeypatch.setenv("STREAMOPS_AGENT_HANDOFF_MAX_CONTEXT_CHARS", "100") from streamops_mcp.config import StreamOpsConfig @@ -804,21 +852,19 @@ def test_truncates_oversized_context(self, monkeypatch): test_config = StreamOpsConfig() monkeypatch.setattr("streamops_mcp.agent.schemas.handoff.config", test_config) - oversized = "x" * 200 - - # Act + # Act: an over-long narrative summary must be truncated, not overflow the prompt handoff = MonitorToDiagnosticHandoff( - anomaly_context=oversized, + anomaly=_make_anomaly(summary="x" * 200), schema_hint={}, ) # Assert - assert len(handoff.anomaly_context) == 100 + assert len(handoff.anomaly.summary) == 100 def test_round_trip_serialization(self): # Arrange handoff = MonitorToDiagnosticHandoff( - anomaly_context="Backpressure ratio exceeded threshold", + anomaly=_make_anomaly(summary="Backpressure ratio exceeded threshold"), schema_hint=DiagnosisReport.model_json_schema(), ) @@ -827,7 +873,8 @@ def test_round_trip_serialization(self): restored = MonitorToDiagnosticHandoff.model_validate_json(json_str) # Assert - assert restored.anomaly_context == handoff.anomaly_context + assert restored.anomaly.summary == handoff.anomaly.summary + assert restored.anomaly.anomaly_type == handoff.anomaly.anomaly_type class TestDiagnosticToReportHandoff: