Skip to content

Commit bfbec1b

Browse files
committed
feat(28-02): emit structured JSON failure telemetry in provider fallback and cap CLI subprocess size
- Add maf_starter/telemetry.py with emit_failure_telemetry(event, **fields), a stdlib-only, non-raising helper that writes one line of parseable JSON to stderr so failures are observable to log-based monitoring, not just in-memory response metadata. - Wire emit_failure_telemetry into provider_fallback.py's fallback_middleware and _wrap_stream_with_fallback at four points: primary-provider failure, each fallback step failure, fallback success/recovery, and full-chain exhaustion, in both the non-streaming and streaming code paths. - Add MAX_CLI_OUTPUT_BYTES and MAX_CLI_PROMPT_BYTES (1MB, mirroring loop_worker_cli.py's existing MAX_REQUEST_BYTES convention) and enforce them in _run_subprocess and _messages_to_prompt, closing the remaining C6 gap on the CLI-fallback subprocess path. - Add tests/test_provider_fallback_telemetry.py: parseable single-line JSON contract, non-raising stderr-write-failure behavior, ordered telemetry events for chain exhaustion vs. fallback recovery, and size-limit rejection tests for both new constants.
1 parent e52e6aa commit bfbec1b

3 files changed

Lines changed: 258 additions & 1 deletion

File tree

maf_starter/provider_fallback.py

Lines changed: 60 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -68,9 +68,12 @@ async def get_final_response(self):
6868
from maf_starter.execution_profile import LOCAL_PROFILE, ExecutionProfile
6969
from maf_starter.routing_policy import RoutingPlan, build_routing_plan
7070
from maf_starter.routing_types import CapabilityChange, ChainStep, RouteAttempt
71+
from maf_starter.telemetry import emit_failure_telemetry
7172

7273

7374
FALLBACK_NOTICE = "[Fallback provider used because the earlier model/provider in the chain failed.]\n"
75+
MAX_CLI_OUTPUT_BYTES = 1_000_000
76+
MAX_CLI_PROMPT_BYTES = 1_000_000
7477
FALLBACK_ERROR_MARKERS = (
7578
"resource_exhausted",
7679
"quota",
@@ -132,6 +135,12 @@ async def fallback_middleware(context: ChatContext, call_next):
132135
error=last_error,
133136
)
134137
)
138+
emit_failure_telemetry(
139+
"provider_failed",
140+
provider=route.primary_provider,
141+
model=route.primary_model,
142+
error=str(exc),
143+
)
135144
for step in route.fallback_steps:
136145
try:
137146
context.result = await _execute_chain_step(
@@ -143,6 +152,12 @@ async def fallback_middleware(context: ChatContext, call_next):
143152
attempt_log=attempt_log,
144153
fallback_index=len(attempt_log),
145154
)
155+
emit_failure_telemetry(
156+
"fallback_succeeded",
157+
provider=step.provider,
158+
model=step.model or step.label,
159+
primary_error=str(last_error),
160+
)
146161
return
147162
except Exception as fallback_exc:
148163
last_error = fallback_exc
@@ -156,8 +171,21 @@ async def fallback_middleware(context: ChatContext, call_next):
156171
error=fallback_exc,
157172
)
158173
)
174+
emit_failure_telemetry(
175+
"fallback_step_failed",
176+
provider=step.provider,
177+
model=step.model or step.label,
178+
error=str(fallback_exc),
179+
fallback_index=len(attempt_log) - 1,
180+
)
159181
continue
160182

183+
emit_failure_telemetry(
184+
"fallback_exhausted",
185+
primary_provider=route.primary_provider,
186+
attempted_providers=[attempt.provider for attempt in attempt_log],
187+
final_error=str(last_error),
188+
)
161189
raise last_error
162190
finally:
163191
reset_run_scope(scope_tokens)
@@ -233,6 +261,12 @@ async def _stream():
233261
error=last_error,
234262
)
235263
)
264+
emit_failure_telemetry(
265+
"provider_failed",
266+
provider=route.primary_provider,
267+
model=route.primary_model,
268+
error=str(exc),
269+
)
236270
for step in route.fallback_steps:
237271
try:
238272
fallback_result = await _execute_chain_step(
@@ -244,6 +278,12 @@ async def _stream():
244278
attempt_log=attempt_log,
245279
fallback_index=len(attempt_log),
246280
)
281+
emit_failure_telemetry(
282+
"fallback_succeeded",
283+
provider=step.provider,
284+
model=step.model or step.label,
285+
primary_error=str(last_error),
286+
)
247287
if isinstance(fallback_result, ResponseStream):
248288
async for update in fallback_result:
249289
yield update
@@ -274,6 +314,19 @@ async def _stream():
274314
error=fallback_exc,
275315
)
276316
)
317+
emit_failure_telemetry(
318+
"fallback_step_failed",
319+
provider=step.provider,
320+
model=step.model or step.label,
321+
error=str(fallback_exc),
322+
fallback_index=len(attempt_log) - 1,
323+
)
324+
emit_failure_telemetry(
325+
"fallback_exhausted",
326+
primary_provider=route.primary_provider,
327+
attempted_providers=[attempt.provider for attempt in attempt_log],
328+
final_error=str(last_error),
329+
)
277330
raise last_error
278331

279332
if isinstance(original_stream, ResponseStream):
@@ -791,6 +844,9 @@ def _run_subprocess(
791844
raise RuntimeError(f"{provider_name} failed: {detail}")
792845

793846
output = completed.stdout.strip()
847+
output_bytes = len(output.encode("utf-8", errors="replace"))
848+
if output_bytes > MAX_CLI_OUTPUT_BYTES:
849+
raise RuntimeError(f"{provider_name} output exceeds {MAX_CLI_OUTPUT_BYTES} bytes")
794850
if not output:
795851
raise RuntimeError(f"{provider_name} returned empty output")
796852

@@ -818,7 +874,10 @@ def _messages_to_prompt(messages: list[Message] | tuple[Message, ...]) -> str:
818874
for message in messages:
819875
rendered.append(f"{str(message.role).upper()}: {_message_text(message)}")
820876
rendered.append("ASSISTANT:")
821-
return "\n\n".join(rendered)
877+
prompt = "\n\n".join(rendered)
878+
if len(prompt.encode("utf-8", errors="replace")) > MAX_CLI_PROMPT_BYTES:
879+
raise ValueError(f"Rendered CLI prompt exceeds {MAX_CLI_PROMPT_BYTES} bytes")
880+
return prompt
822881

823882

824883
def _message_text(message: Message) -> str:

maf_starter/telemetry.py

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
from __future__ import annotations
2+
3+
import json
4+
import sys
5+
from datetime import datetime, timezone
6+
from typing import Any
7+
8+
9+
def emit_failure_telemetry(event: str, **fields: Any) -> None:
10+
payload = {
11+
"event": event,
12+
"timestamp": datetime.now(timezone.utc).isoformat(),
13+
**fields,
14+
}
15+
try:
16+
sys.stderr.write(json.dumps(payload, separators=(",", ":")) + "\n")
17+
sys.stderr.flush()
18+
except Exception: # pragma: no cover - telemetry must never mask the original failure
19+
return
Lines changed: 179 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,179 @@
1+
from __future__ import annotations
2+
3+
import contextlib
4+
import io
5+
import json
6+
import subprocess
7+
import unittest
8+
import uuid
9+
from pathlib import Path
10+
from types import SimpleNamespace
11+
from unittest.mock import patch
12+
13+
from agent_framework import ChatResponse, Message
14+
15+
from maf_starter.config import Settings, load_settings
16+
from maf_starter.provider_fallback import (
17+
_messages_to_prompt,
18+
_run_subprocess,
19+
build_fallback_middleware,
20+
)
21+
from maf_starter.routing_policy import RoutingPlan
22+
from maf_starter.routing_types import ChainStep
23+
from maf_starter.telemetry import emit_failure_telemetry
24+
25+
26+
SCRATCH_ROOT = Path(__file__).resolve().parents[1] / ".tmp-tests"
27+
28+
29+
class RepoScratchTestCase(unittest.TestCase):
30+
def make_scratch_dir(self) -> Path:
31+
path = SCRATCH_ROOT / uuid.uuid4().hex
32+
path.mkdir(parents=True, exist_ok=False)
33+
self.addCleanup(lambda: path.exists() and __import__("shutil").rmtree(path, ignore_errors=True))
34+
return path
35+
36+
37+
class ProviderFallbackTelemetryTests(RepoScratchTestCase):
38+
def _make_settings(self, root: Path) -> Settings:
39+
entities = root / "entities"
40+
repo = root / "repo"
41+
entities.mkdir()
42+
repo.mkdir()
43+
env = {
44+
"MAF_API_KEY": "test-key",
45+
"MAF_REPO_ROOT": str(repo),
46+
"MAF_ENTITIES_DIR": str(entities),
47+
"MAF_FALLBACK_CHAIN": "gemini-cli:gemini-2.5-pro,claude-cli:claude-sonnet-4-6",
48+
}
49+
with patch.dict("os.environ", env, clear=False):
50+
return load_settings(project_root=root, env_path=root / ".missing-env")
51+
52+
@staticmethod
53+
def _make_route_plan() -> RoutingPlan:
54+
return RoutingPlan(
55+
mode="auto",
56+
route_lane="auto",
57+
tier="standard",
58+
rationale="test route",
59+
primary_provider="gemini",
60+
primary_model="gemini-2.5-pro",
61+
requested_provider="gemini",
62+
requested_model="gemini-2.5-pro",
63+
fallback_steps=(
64+
ChainStep("gemini-cli", "gemini-2.5-pro"),
65+
ChainStep("claude-cli", "claude-sonnet-4-6"),
66+
),
67+
)
68+
69+
@staticmethod
70+
def _make_context() -> SimpleNamespace:
71+
return SimpleNamespace(
72+
messages=[Message(role="user", text="hello fallback")],
73+
options={},
74+
kwargs={},
75+
metadata={},
76+
result=None,
77+
)
78+
79+
def test_emit_failure_telemetry_writes_single_parseable_json_line(self) -> None:
80+
stream = io.StringIO()
81+
with contextlib.redirect_stderr(stream):
82+
emit_failure_telemetry("provider_failed", provider="gemini", error="quota exceeded")
83+
84+
lines = stream.getvalue().strip().splitlines()
85+
self.assertEqual(len(lines), 1)
86+
payload = json.loads(lines[0])
87+
self.assertEqual(payload["event"], "provider_failed")
88+
self.assertEqual(payload["provider"], "gemini")
89+
self.assertEqual(payload["error"], "quota exceeded")
90+
self.assertIn("timestamp", payload)
91+
92+
def test_emit_failure_telemetry_swallows_stderr_write_failures(self) -> None:
93+
fake_stderr = SimpleNamespace(
94+
write=lambda _: (_ for _ in ()).throw(RuntimeError("stderr failed")),
95+
flush=lambda: None,
96+
)
97+
with patch("sys.stderr", fake_stderr):
98+
self.assertIsNone(emit_failure_telemetry("provider_failed", provider="gemini", error="boom"))
99+
100+
def test_fallback_middleware_emits_primary_step_and_exhausted_events(self) -> None:
101+
root = self.make_scratch_dir()
102+
settings = self._make_settings(root)
103+
middleware = build_fallback_middleware(
104+
settings,
105+
primary_provider="gemini",
106+
primary_model="gemini-2.5-pro",
107+
routing_mode="auto",
108+
)
109+
context = self._make_context()
110+
stderr = io.StringIO()
111+
112+
async def call_next():
113+
raise RuntimeError("rate limit from primary")
114+
115+
attempts = [
116+
RuntimeError("rate limit in gemini-cli fallback"),
117+
RuntimeError("quota exhausted in claude fallback"),
118+
]
119+
120+
async def fake_execute_chain_step(**_kwargs):
121+
raise attempts.pop(0)
122+
123+
with (
124+
patch("maf_starter.provider_fallback.build_routing_plan", return_value=self._make_route_plan()),
125+
patch("maf_starter.provider_fallback._execute_chain_step", side_effect=fake_execute_chain_step),
126+
contextlib.redirect_stderr(stderr),
127+
):
128+
with self.assertRaisesRegex(RuntimeError, "quota exhausted in claude fallback"):
129+
__import__("asyncio").run(middleware(context, call_next))
130+
131+
events = [json.loads(line)["event"] for line in stderr.getvalue().strip().splitlines()]
132+
self.assertEqual(events, ["provider_failed", "fallback_step_failed", "fallback_step_failed", "fallback_exhausted"])
133+
134+
def test_fallback_middleware_emits_recovery_event_when_fallback_succeeds(self) -> None:
135+
root = self.make_scratch_dir()
136+
settings = self._make_settings(root)
137+
middleware = build_fallback_middleware(
138+
settings,
139+
primary_provider="gemini",
140+
primary_model="gemini-2.5-pro",
141+
routing_mode="auto",
142+
)
143+
context = self._make_context()
144+
stderr = io.StringIO()
145+
146+
async def call_next():
147+
raise RuntimeError("rate limit from primary")
148+
149+
async def fake_execute_chain_step(**_kwargs):
150+
return ChatResponse(messages=[Message(role="assistant", text="fallback ok")])
151+
152+
with (
153+
patch("maf_starter.provider_fallback.build_routing_plan", return_value=self._make_route_plan()),
154+
patch("maf_starter.provider_fallback._execute_chain_step", side_effect=fake_execute_chain_step),
155+
contextlib.redirect_stderr(stderr),
156+
):
157+
__import__("asyncio").run(middleware(context, call_next))
158+
159+
payloads = [json.loads(line) for line in stderr.getvalue().strip().splitlines()]
160+
self.assertEqual(payloads[0]["event"], "provider_failed")
161+
self.assertEqual(payloads[1]["event"], "fallback_succeeded")
162+
self.assertNotIn("fallback_exhausted", [payload["event"] for payload in payloads])
163+
self.assertIsInstance(context.result, ChatResponse)
164+
165+
def test_run_subprocess_rejects_oversized_output(self) -> None:
166+
completed = subprocess.CompletedProcess(args=["codex.cmd"], returncode=0, stdout="x" * 32, stderr="")
167+
with (
168+
patch("maf_starter.provider_fallback.subprocess.run", return_value=completed),
169+
patch("maf_starter.provider_fallback.shutil.which", return_value="codex.cmd"),
170+
patch("maf_starter.provider_fallback.MAX_CLI_OUTPUT_BYTES", 8),
171+
):
172+
with self.assertRaisesRegex(RuntimeError, "output exceeds 8 bytes"):
173+
_run_subprocess("codex-cli", ["codex.cmd", "exec"], Path.cwd(), None)
174+
175+
def test_messages_to_prompt_rejects_oversized_prompt(self) -> None:
176+
messages = [Message(role="user", text="x" * 32)]
177+
with patch("maf_starter.provider_fallback.MAX_CLI_PROMPT_BYTES", 8):
178+
with self.assertRaisesRegex(ValueError, "Rendered CLI prompt exceeds 8 bytes"):
179+
_messages_to_prompt(messages)

0 commit comments

Comments
 (0)