diff --git a/apps/ai-observability/README.md b/apps/ai-observability/README.md index 701757638..db1c3bef3 100644 --- a/apps/ai-observability/README.md +++ b/apps/ai-observability/README.md @@ -24,6 +24,7 @@ Each app exists to test one thing the others don't: - `openai-agents/python-travel-triage` — tracing processor; the SDK emits the tree - `vercel-ai/nextjs-support-chat` — per-request identity; framework bootstrap - `manual-capture/node-http-chat` — hand-built tree; must reuse the existing client +- `manual-capture/python-multi-route-proxy` — raw proxy routes, streams, failures, and direct calls - `google-adk/node-weather` — framework plugin; identity comes from ADK's own ids - `opentelemetry/go-weather` — Go; no wrapper SDK exists, so the posthog-go OTel bridge @@ -47,8 +48,9 @@ conversation structure. together. A fresh id per call groups nothing and is worse than none. - **Spans come from tool registration** (or explicit `$ai_span` capture on the manual path). Never from hand-authored wrappers around helper functions. -- Static fixtures in the style of [`../mcp-analytics`](../mcp-analytics): not - executed, no lockfiles, no keys. Node type-checks with `tsc --noEmit`. +- Static Wizard fixtures in the style of [`../mcp-analytics`](../mcp-analytics): + no lockfiles or keys. Node type-checks with `tsc --noEmit`. The multi-route + proxy has a local fake-provider contract test; Wizard CI evaluates its diff. - Each README states the tree the app must produce and grades **emitted events, not the diff** — every observed failure mode so far produced plausible code and a broken tree. diff --git a/apps/ai-observability/manual-capture/python-multi-route-proxy/.env.example b/apps/ai-observability/manual-capture/python-multi-route-proxy/.env.example new file mode 100644 index 000000000..25a5a262b --- /dev/null +++ b/apps/ai-observability/manual-capture/python-multi-route-proxy/.env.example @@ -0,0 +1,4 @@ +OPENAI_BASE_URL=http://127.0.0.1:8001 +OLLAMA_BASE_URL=http://127.0.0.1:11434 +PORT=8000 +# OPENAI_API_KEY=your-provider-key diff --git a/apps/ai-observability/manual-capture/python-multi-route-proxy/.gitignore b/apps/ai-observability/manual-capture/python-multi-route-proxy/.gitignore new file mode 100644 index 000000000..77ac75498 --- /dev/null +++ b/apps/ai-observability/manual-capture/python-multi-route-proxy/.gitignore @@ -0,0 +1,3 @@ +.venv/ +__pycache__/ +*.pyc diff --git a/apps/ai-observability/manual-capture/python-multi-route-proxy/README.md b/apps/ai-observability/manual-capture/python-multi-route-proxy/README.md new file mode 100644 index 000000000..f693c84f8 --- /dev/null +++ b/apps/ai-observability/manual-capture/python-multi-route-proxy/README.md @@ -0,0 +1,76 @@ +# wb-aio-manual-python-multi-route-proxy + +An aiohttp chat gateway forwards raw provider HTTP responses. It has no vendor +LLM SDK and no PostHog code. The Wizard must use manual AI capture without +changing the proxy's request or response behavior. + +The gateway supplies `X-User-Id` and `X-Conversation-Id` on every request. A +conversation can make several requests. Each request or WebSocket message is +one trace in that conversation. These headers stand in for identity resolved by +an upstream gateway. They are not provider credentials. + +## Inference paths + +| Entry point | Upstream call | Response shape | Capture target | +| --- | --- | --- | --- | +| `POST /v1/chat/completions` | OpenAI-compatible `/v1/chat/completions` | JSON or SSE | One generation per request | +| `POST /v1/responses` | OpenAI Responses `/v1/responses` | JSON or SSE | One generation per request | +| `POST /api/chat` | Ollama native `/api/chat` | JSON or NDJSON, streaming by default | One generation per request | +| `POST /api/generate` | Ollama native `/api/generate` | JSON or NDJSON, streaming by default | One generation per request | +| `POST /assistant/ask` | Direct OpenAI-compatible calls in `complete()` | Two JSON calls when the order tool runs | Two generations and one tool span | +| `GET /ws/generate` | Direct Ollama `/api/generate` for each WebSocket message | NDJSON relayed as WebSocket frames | One generation per message | + +The two direct paths bypass the four HTTP proxy handlers. Instrumenting only +`forward()` leaves them dark. Instrumenting only `complete()` leaves the public +proxy routes dark. + +## Expected outcome + +- **Coverage.** Every path in the table emits its capture target. A successful + stream emits once when it finishes, not once per chunk. A WebSocket with two + prompt messages emits two generations. +- **Grouping.** Each request or WebSocket message gets a new `$ai_trace_id`. + Calls and the `lookup_order` span within one `/assistant/ask` share that ID. + `$ai_session_id` equals `X-Conversation-Id` across related traces, and the + event `distinct_id` equals `X-User-Id`. +- **Tool flow.** An order question produces `generation → span(lookup_order) → + generation`. The span uses `$ai_parent_id` and records the actual tool + execution, not the mere presence of a tool schema in the model request. +- **Provider formats.** OpenAI-compatible chat usage maps `prompt_tokens` and + `completion_tokens`. OpenAI Responses JSON and `response.completed` SSE map + `input_tokens` and `output_tokens`, with text from `response.output_text.delta` + on streams. Ollama's final JSON or NDJSON object maps top-level + `prompt_eval_count` and `eval_count`. Missing usage stays unknown, not zero. +- **Failures.** A provider HTTP error or connection failure produces an errored + generation with `$ai_is_error` and a safe `$ai_error`. An interrupted SSE, + NDJSON, or WebSocket stream is marked incomplete or errored. No partial stream + is reported as a successful completed generation. +- **Behavior.** HTTP status, body, content type, stream order, and cancellation + still reach the caller. The order tool still runs between its two model calls. + +Fail the evaluation if any listed inference path is missing, token counts use +another provider's field names, a stream is captured before its final outcome, +an error is labeled successful, or the tool span has a separate trace. + +## Run locally + +Install `requirements.txt`, then point `OPENAI_BASE_URL` and `OLLAMA_BASE_URL` +at compatible servers. `OPENAI_API_KEY` is optional for local compatible +servers. The app binds to `127.0.0.1:8000` by default. + +```bash +python -m pip install -r requirements.txt +python app.py +``` + +The contract test uses a fake provider, so it needs no credentials or model: + +```bash +python -m unittest discover -s tests -v +``` + +Run the Wizard evaluation from the workbench root: + +```bash +pnpm wizard-ci --app ai-observability/manual-capture/python-multi-route-proxy --evaluate --local +``` diff --git a/apps/ai-observability/manual-capture/python-multi-route-proxy/app.py b/apps/ai-observability/manual-capture/python-multi-route-proxy/app.py new file mode 100644 index 000000000..d28d3de2c --- /dev/null +++ b/apps/ai-observability/manual-capture/python-multi-route-proxy/app.py @@ -0,0 +1,233 @@ +"""A chat service that forwards several raw LLM protocols without an SDK.""" + +import json +import os +from collections.abc import Sequence +from typing import Any + +from aiohttp import ClientError, ClientSession, ClientTimeout, WSMsgType, web + +from tools import lookup_order + + +CLIENT = web.AppKey("client", ClientSession) +OPENAI_BASE = web.AppKey("openai_base", str) +OLLAMA_BASE = web.AppKey("ollama_base", str) +USER_ID = web.RequestKey("user_id", str) +CONVERSATION_ID = web.RequestKey("conversation_id", str) + +TOOLS = [ + { + "type": "function", + "function": { + "name": "lookup_order", + "description": "Look up the current user's latest order.", + "parameters": {"type": "object", "properties": {}}, + }, + } +] + + +@web.middleware +async def resolve_user(request: web.Request, handler): + """Treat headers as identity already resolved by this fixture's gateway.""" + user_id = request.headers.get("X-User-Id") + conversation_id = request.headers.get("X-Conversation-Id") + if not user_id or not conversation_id: + raise web.HTTPUnauthorized(text="user and conversation required") + request[USER_ID] = user_id + request[CONVERSATION_ID] = conversation_id + return await handler(request) + + +def provider_headers(openai: bool) -> dict[str, str]: + headers = {"Content-Type": "application/json"} + if openai and (api_key := os.environ.get("OPENAI_API_KEY")): + headers["Authorization"] = f"Bearer {api_key}" + return headers + + +async def forward( + request: web.Request, base: str, default_stream: bool, openai: bool +) -> web.StreamResponse: + """Forward one public inference route, including its native stream format.""" + raw = await request.read() + try: + payload = json.loads(raw) + except (ValueError, UnicodeDecodeError) as exc: + raise web.HTTPBadRequest(text="JSON request required") from exc + if not isinstance(payload, dict): + raise web.HTTPBadRequest(text="JSON object required") + + url = f"{base}{request.path}" + try: + async with request.app[CLIENT].post( + url, data=raw, headers=provider_headers(openai) + ) as upstream: + content_type = upstream.headers.get("Content-Type", "application/json") + response_headers = {"Content-Type": content_type} + streamed = payload.get("stream", default_stream) is True + if upstream.status >= 400 or not streamed: + return web.Response( + status=upstream.status, + body=await upstream.read(), + headers=response_headers, + ) + + downstream = web.StreamResponse( + status=upstream.status, headers=response_headers + ) + await downstream.prepare(request) + async for chunk in upstream.content.iter_chunked(4096): + await downstream.write(chunk) + await downstream.write_eof() + return downstream + except ClientError as exc: + raise web.HTTPBadGateway(text="model provider unavailable") from exc + + +async def openai_chat(request: web.Request) -> web.StreamResponse: + return await forward(request, request.app[OPENAI_BASE], False, True) + + +async def openai_responses(request: web.Request) -> web.StreamResponse: + return await forward(request, request.app[OPENAI_BASE], False, True) + + +async def ollama_chat(request: web.Request) -> web.StreamResponse: + return await forward(request, request.app[OLLAMA_BASE], True, False) + + +async def ollama_generate(request: web.Request) -> web.StreamResponse: + return await forward(request, request.app[OLLAMA_BASE], True, False) + + +async def complete( + request: web.Request, messages: Sequence[dict[str, Any]] +) -> dict[str, Any]: + """Make the assistant's direct model call, separate from public proxy routes.""" + payload = { + "model": "local-chat", + "messages": list(messages), + "tools": TOOLS, + "stream": False, + } + try: + async with request.app[CLIENT].post( + f"{request.app[OPENAI_BASE]}/v1/chat/completions", + json=payload, + headers=provider_headers(True), + ) as upstream: + if upstream.status >= 400: + raise web.HTTPBadGateway(text=f"model returned {upstream.status}") + return await upstream.json() + except ClientError as exc: + raise web.HTTPBadGateway(text="model provider unavailable") from exc + + +async def assistant_ask(request: web.Request) -> web.Response: + """Answer one turn, including the registered order tool when requested.""" + body = await request.json() + question = body.get("question") if isinstance(body, dict) else None + if not isinstance(question, str) or not question.strip(): + raise web.HTTPBadRequest(text="question required") + + messages: list[dict[str, Any]] = [ + { + "role": "system", + "content": "Answer briefly. Use lookup_order for order questions.", + }, + {"role": "user", "content": question}, + ] + first = await complete(request, messages) + message = first["choices"][0]["message"] + tool_calls = message.get("tool_calls", []) + if tool_calls: + messages.append(message | {"role": "assistant"}) + for call in tool_calls: + if call["function"]["name"] != "lookup_order": + raise web.HTTPBadGateway(text="unknown model tool") + result = lookup_order(request[USER_ID]) + messages.append( + { + "role": "tool", + "tool_call_id": call["id"], + "content": json.dumps(result), + } + ) + final = await complete(request, messages) + message = final["choices"][0]["message"] + + return web.json_response({"answer": message.get("content") or ""}) + + +async def websocket_generate(request: web.Request) -> web.WebSocketResponse: + """Run native Ollama generation directly for each WebSocket message.""" + socket = web.WebSocketResponse() + await socket.prepare(request) + async for message in socket: + if message.type != WSMsgType.TEXT: + continue + try: + body = json.loads(message.data) + prompt = body.get("prompt") if isinstance(body, dict) else None + except ValueError: + prompt = None + if not isinstance(prompt, str) or not prompt.strip(): + await socket.send_json({"error": "prompt required"}) + continue + + payload = {"model": "local-ollama", "prompt": prompt, "stream": True} + try: + async with request.app[CLIENT].post( + f"{request.app[OLLAMA_BASE]}/api/generate", + json=payload, + headers=provider_headers(False), + ) as upstream: + if upstream.status >= 400: + await socket.send_json( + {"error": "model provider failed", "status": upstream.status} + ) + continue + async for line in upstream.content: + if line.strip(): + await socket.send_str(line.decode().strip()) + except ClientError: + await socket.send_json({"error": "model provider unavailable"}) + return socket + + +async def start_client(app: web.Application) -> None: + app[CLIENT] = ClientSession(timeout=ClientTimeout(total=120)) + + +async def close_client(app: web.Application) -> None: + await app[CLIENT].close() + + +def create_app(openai_base_url: str, ollama_base_url: str) -> web.Application: + """Build the proxy for the two upstream provider bases.""" + app = web.Application(middlewares=[resolve_user]) + app[OPENAI_BASE] = openai_base_url.rstrip("/") + app[OLLAMA_BASE] = ollama_base_url.rstrip("/") + app.on_startup.append(start_client) + app.on_cleanup.append(close_client) + app.router.add_post("/v1/chat/completions", openai_chat) + app.router.add_post("/v1/responses", openai_responses) + app.router.add_post("/api/chat", ollama_chat) + app.router.add_post("/api/generate", ollama_generate) + app.router.add_post("/assistant/ask", assistant_ask) + app.router.add_get("/ws/generate", websocket_generate) + return app + + +if __name__ == "__main__": + web.run_app( + create_app( + os.environ.get("OPENAI_BASE_URL", "http://127.0.0.1:8001"), + os.environ.get("OLLAMA_BASE_URL", "http://127.0.0.1:11434"), + ), + host="127.0.0.1", + port=int(os.environ.get("PORT", "8000")), + handler_cancellation=True, + ) diff --git a/apps/ai-observability/manual-capture/python-multi-route-proxy/requirements.txt b/apps/ai-observability/manual-capture/python-multi-route-proxy/requirements.txt new file mode 100644 index 000000000..e8918badf --- /dev/null +++ b/apps/ai-observability/manual-capture/python-multi-route-proxy/requirements.txt @@ -0,0 +1 @@ +aiohttp>=3.14,<4 diff --git a/apps/ai-observability/manual-capture/python-multi-route-proxy/tests/test_proxy.py b/apps/ai-observability/manual-capture/python-multi-route-proxy/tests/test_proxy.py new file mode 100644 index 000000000..d125ef989 --- /dev/null +++ b/apps/ai-observability/manual-capture/python-multi-route-proxy/tests/test_proxy.py @@ -0,0 +1,238 @@ +"""Contract checks against a fake external model provider.""" + +import asyncio +import json +import unittest + +from aiohttp import ClientSession, web + +from app import create_app + + +async def serve(app: web.Application): + runner = web.AppRunner(app, handler_cancellation=True) + await runner.setup() + site = web.TCPSite(runner, "127.0.0.1", 0) + await site.start() + port = site._server.sockets[0].getsockname()[1] + return runner, f"http://127.0.0.1:{port}" + + +class ProxyContract(unittest.IsolatedAsyncioTestCase): + async def asyncSetUp(self): + self.seen = [] + self.provider_cancelled = asyncio.Event() + provider = web.Application() + provider.router.add_post("/v1/chat/completions", self.chat) + provider.router.add_post("/v1/responses", self.responses) + provider.router.add_post("/api/chat", self.ollama) + provider.router.add_post("/api/generate", self.ollama) + self.provider_runner, provider_url = await serve(provider) + self.proxy_runner, self.proxy_url = await serve( + create_app(provider_url, provider_url) + ) + self.client = ClientSession() + + async def asyncTearDown(self): + await self.client.close() + await self.proxy_runner.cleanup() + await self.provider_runner.cleanup() + + async def chat(self, request): + body = await request.json() + self.seen.append((request.path, body)) + if body.get("fail"): + return web.json_response({"error": "rate limited"}, status=429) + if body.get("stream"): + response = web.StreamResponse(headers={"Content-Type": "text/event-stream"}) + await response.prepare(request) + await response.write(b'data: {"choices":[{"delta":{"content":"Hi"}}]}\n\n') + await response.write(b"data: [DONE]\n\n") + await response.write_eof() + return response + if body.get("tools") and not any( + m.get("role") == "tool" for m in body["messages"] + ): + message = { + "content": None, + "tool_calls": [ + { + "id": "call_1", + "function": {"name": "lookup_order", "arguments": "{}"}, + } + ], + } + else: + message = {"content": "Order is shipped"} + return web.json_response( + { + "choices": [{"message": message}], + "usage": {"prompt_tokens": 11, "completion_tokens": 4}, + } + ) + + async def responses(self, request): + body = await request.json() + self.seen.append((request.path, body)) + if body.get("stream"): + response = web.StreamResponse(headers={"Content-Type": "text/event-stream"}) + await response.prepare(request) + await response.write( + b'event: response.output_text.delta\ndata: {"type":"response.output_text.delta","delta":"Hello"}\n\n' + ) + await response.write( + b'event: response.completed\ndata: {"type":"response.completed","response":{"usage":{"input_tokens":7,"output_tokens":2}}}\n\n' + ) + await response.write_eof() + return response + return web.json_response( + { + "output": [ + { + "type": "message", + "content": [{"type": "output_text", "text": "Hello"}], + } + ], + "usage": {"input_tokens": 7, "output_tokens": 2}, + } + ) + + async def ollama(self, request): + body = await request.json() + self.seen.append((request.path, body)) + chunk = ( + {"response": "Hi"} + if request.path == "/api/generate" + else {"message": {"content": "Hi"}} + ) + if body.get("stream", True): + response = web.StreamResponse( + headers={"Content-Type": "application/x-ndjson"} + ) + await response.prepare(request) + await response.write(json.dumps(chunk | {"done": False}).encode() + b"\n") + if body.get("stall"): + try: + await asyncio.sleep(5) + except asyncio.CancelledError: + self.provider_cancelled.set() + raise + await response.write( + b'{"done":true,"prompt_eval_count":9,"eval_count":3}\n' + ) + await response.write_eof() + return response + return web.json_response( + chunk | {"done": True, "prompt_eval_count": 9, "eval_count": 3} + ) + + def headers(self): + return {"X-User-Id": "user_123", "X-Conversation-Id": "thread_abc"} + + async def test_every_inference_route_preserves_its_response_mode(self): + cases = [ + ( + "/v1/chat/completions", + {"model": "local-chat", "messages": [], "stream": False}, + "application/json", + b'"prompt_tokens": 11', + ), + ( + "/v1/chat/completions", + {"model": "local-chat", "messages": [], "stream": True}, + "text/event-stream", + b"data: [DONE]", + ), + ( + "/v1/responses", + {"model": "local-response", "input": "Hi", "stream": False}, + "application/json", + b'"input_tokens": 7', + ), + ( + "/v1/responses", + {"model": "local-response", "input": "Hi", "stream": True}, + "text/event-stream", + b"response.completed", + ), + ( + "/api/chat", + {"model": "local-ollama", "messages": []}, + "application/x-ndjson", + b'"prompt_eval_count":9', + ), + ( + "/api/generate", + {"model": "local-ollama", "prompt": "Hi"}, + "application/x-ndjson", + b'"eval_count":3', + ), + ( + "/api/chat", + {"model": "local-ollama", "messages": [], "stream": False}, + "application/json", + b'"prompt_eval_count": 9', + ), + ] + for path, body, content_type, expected in cases: + with self.subTest(path=path, stream=body.get("stream")): + async with self.client.post( + self.proxy_url + path, json=body, headers=self.headers() + ) as response: + self.assertEqual(response.status, 200) + self.assertEqual(response.content_type, content_type) + self.assertIn(expected, await response.read()) + self.assertEqual([path for path, _ in self.seen], [path for path, *_ in cases]) + + async def test_provider_failure_is_forwarded(self): + async with self.client.post( + self.proxy_url + "/v1/chat/completions", + json={"fail": True}, + headers=self.headers(), + ) as response: + self.assertEqual(response.status, 429) + self.assertEqual(await response.json(), {"error": "rate limited"}) + + async def test_assistant_executes_tool_between_two_model_calls(self): + async with self.client.post( + self.proxy_url + "/assistant/ask", + json={"question": "Where is my order?"}, + headers=self.headers(), + ) as response: + self.assertEqual(response.status, 200) + self.assertEqual(await response.json(), {"answer": "Order is shipped"}) + calls = [body for path, body in self.seen if path == "/v1/chat/completions"] + self.assertEqual(len(calls), 2) + self.assertEqual(calls[1]["messages"][-1]["role"], "tool") + self.assertIn("shipped", calls[1]["messages"][-1]["content"]) + + async def test_websocket_calls_native_generate_directly(self): + async with self.client.ws_connect( + self.proxy_url + "/ws/generate", headers=self.headers() + ) as socket: + await socket.send_json({"prompt": "Hi"}) + self.assertEqual( + await socket.receive_json(), {"response": "Hi", "done": False} + ) + self.assertEqual( + await socket.receive_json(), + {"done": True, "prompt_eval_count": 9, "eval_count": 3}, + ) + calls = [body for path, body in self.seen if path == "/api/generate"] + self.assertEqual( + calls, [{"model": "local-ollama", "prompt": "Hi", "stream": True}] + ) + + async def test_client_disconnect_cancels_upstream_stream(self): + response = await self.client.post( + self.proxy_url + "/api/generate", + json={"model": "local-ollama", "prompt": "Hi", "stall": True}, + headers=self.headers(), + ) + self.assertIn(b'"done": false', await response.content.readline()) + response.close() + await asyncio.wait_for(self.provider_cancelled.wait(), timeout=1) + + +if __name__ == "__main__": + unittest.main() diff --git a/apps/ai-observability/manual-capture/python-multi-route-proxy/tools.py b/apps/ai-observability/manual-capture/python-multi-route-proxy/tools.py new file mode 100644 index 000000000..693b61afc --- /dev/null +++ b/apps/ai-observability/manual-capture/python-multi-route-proxy/tools.py @@ -0,0 +1,5 @@ +"""A local tool used between two direct model calls.""" + + +def lookup_order(user_id: str) -> dict[str, str]: + return {"user_id": user_id, "order_id": "order_123", "status": "shipped"} diff --git a/services/pr-evaluator/evaluator.test.ts b/services/pr-evaluator/evaluator.test.ts index 4ff3f12ba..9ae16b81b 100644 --- a/services/pr-evaluator/evaluator.test.ts +++ b/services/pr-evaluator/evaluator.test.ts @@ -5,7 +5,7 @@ */ import { describe, it } from "node:test"; import assert from "node:assert/strict"; -import { detectFramework, detectArchType, parseCommandments, parseDocsConfig } from "./prompt-builder.js"; +import { buildSystemPrompt, detectFramework, detectArchType, parseCommandments, parseDocsConfig } from "./prompt-builder.js"; import { repairAndParseJSON, validateAndCorrectScores, @@ -48,6 +48,36 @@ function makeScores(overrides: Partial = {}): EvaluateScores { }; } +it("uses the AI observability criterion keys in the output rubric", async () => { + const prompt = await buildSystemPrompt(undefined, { command: "ai-observability" }); + const block = prompt.match(//); + assert.ok(block, "AI output rubric is present"); + const rubric = JSON.parse(block[1]) as Record>; + + assert.deepEqual(Object.keys(rubric.file_analysis), [ + "fa_changes_relevant", "fa_correct_files", "fa_no_unnecessary_changes", + "fa_code_quality", "fa_imports_valid", "fa_files_complete", + ]); + assert.deepEqual(Object.keys(rubric.app_sanity), [ + "as_builds", "as_preserves_existing", "as_minimal_changes", "as_no_syntax_errors", + "as_correct_imports", "as_env_documented", "as_dependency_version_valid", "as_flushes_before_exit", + "as_manual_capture_ledger_delivered", + ]); + assert.deepEqual(Object.keys(rubric.posthog_implementation), [ + "ph_instrumentation_initialized_once", "ph_generations_captured", "ph_trace_groups_the_request", + "ph_session_id_set", "ph_session_id_correct_key", "ph_session_cardinality_correct", + "ph_identity_scope_correct", "ph_no_hand_authored_spans", "ph_additive_only", + ]); + assert.deepEqual(Object.keys(rubric.event_quality), [ + "eq_events_would_render_as_tree", "eq_output_choices_have_role", "eq_stream_terminal_parsed", + "eq_person_attribution", "eq_span_parenting", + "eq_no_fabricated_structure", "eq_privacy_respected", + ]); + const statedCriteria = [...prompt.matchAll(/^- \*\*((?:fa|as|ph|eq)_[a-z_]+)\*\*/gm)].map((match) => match[1]); + const outputCriteria = Object.values(rubric).flatMap((dimension) => Object.keys(dimension)); + assert.deepEqual(outputCriteria, statedCriteria); +}); + // ── detectFramework ────────────────────────────────────────────────────────── describe("detectFramework", () => { @@ -354,6 +384,37 @@ describe("computeScoresFromRubric", () => { assert.equal(scores.framework, "django"); assert.equal(scores.arch_type, "server-only"); }); + + it("caps AI quality when a generation outcome or output shape is wrong", () => { + const rubric: RubricData = { + file_analysis: { fa_a: "yes" }, + app_sanity: { as_a: "yes" }, + posthog_implementation: { ph_generations_captured: "yes" }, + event_quality: { + eq_events_would_render_as_tree: "yes", + eq_output_choices_have_role: "no", + eq_stream_terminal_parsed: "yes", + }, + }; + const scores = computeScoresFromRubric(rubric, "python", "server-only"); + assert.equal(scores.confidence, 3); + }); + + it("does not give a perfect score when the manual capture ledger is absent", () => { + const rubric: RubricData = { + file_analysis: { fa_a: "yes" }, + app_sanity: { + as_a: "yes", as_b: "yes", as_c: "yes", as_d: "yes", as_e: "yes", + as_f: "yes", as_g: "yes", as_h: "yes", as_i: "yes", + as_manual_capture_ledger_delivered: "no", + }, + posthog_implementation: { ph_generations_captured: "yes" }, + event_quality: { eq_events_would_render_as_tree: "yes" }, + }; + const scores = computeScoresFromRubric(rubric, "python", "server-only"); + assert.equal(scores.app_sanity, 5); + assert.equal(scores.confidence, 4); + }); }); // ── RubricSchema validation ────────────────────────────────────────────────── diff --git a/services/pr-evaluator/evaluator.ts b/services/pr-evaluator/evaluator.ts index 2f1efbc9b..587e1d0f9 100644 --- a/services/pr-evaluator/evaluator.ts +++ b/services/pr-evaluator/evaluator.ts @@ -77,7 +77,18 @@ export function computeScoresFromRubric( const posthog_implementation = computeScoreFromRubric(rubric.posthog_implementation); const event_quality = computeScoreFromRubric(rubric.event_quality); const avg = (file_analysis + app_sanity + posthog_implementation + event_quality) / 4; - const confidence = Math.min(app_sanity, Math.round(avg)); + let confidence = Math.min(app_sanity, Math.round(avg)); + const criticalAiMiss = [ + rubric.posthog_implementation.ph_generations_captured, + rubric.event_quality.eq_events_would_render_as_tree, + rubric.event_quality.eq_output_choices_have_role, + rubric.event_quality.eq_stream_terminal_parsed, + ].some((value) => value === "no"); + if (criticalAiMiss) { + confidence = Math.min(confidence, 3); + } else if (rubric.app_sanity.as_manual_capture_ledger_delivered === "no") { + confidence = Math.min(confidence, 4); + } return { file_analysis, diff --git a/services/pr-evaluator/prompt-builder.ts b/services/pr-evaluator/prompt-builder.ts index 198ffe470..c467aeac9 100644 --- a/services/pr-evaluator/prompt-builder.ts +++ b/services/pr-evaluator/prompt-builder.ts @@ -28,15 +28,17 @@ const OUTPUT_FORMAT_REVENUE = readFileSync( join(__dirname, "prompts/output-format-revenue.md"), "utf-8", ).trim(); +const OUTPUT_FORMAT_AI_OBSERVABILITY = readFileSync( + join(__dirname, "prompts/output-format-ai-observability.md"), + "utf-8", +).trim(); /** Per-command prompt overrides. Extend when a new command gets its own rubric. */ const PROMPTS_BY_COMMAND: Record = { revenue: { rubric: EVALUATION_CRITERIA_REVENUE, outputFormat: OUTPUT_FORMAT_REVENUE }, - // Reuses the default output format — the AIO rubric keeps the same four - // dimensions, only the items differ. "ai-observability": { rubric: EVALUATION_CRITERIA_AI_OBSERVABILITY, - outputFormat: OUTPUT_FORMAT, + outputFormat: OUTPUT_FORMAT_AI_OBSERVABILITY, }, // Reuses the default output format - the replay-vision rubric keeps the same // four dimensions, only the items differ. diff --git a/services/pr-evaluator/prompts/evaluation-ai-observability.md b/services/pr-evaluator/prompts/evaluation-ai-observability.md index 07b048f87..a46d9ac5b 100644 --- a/services/pr-evaluator/prompts/evaluation-ai-observability.md +++ b/services/pr-evaluator/prompts/evaluation-ai-observability.md @@ -31,6 +31,8 @@ Scores are computed server-side from your rubric answers. **Scope of evaluation:** Evaluate ONLY the changes introduced by this PR. This PR is produced by the `ai-observability` wizard command, whose job is to make an **existing** app's LLM calls emit PostHog AI observability events. It is NOT supposed to rewrite the app's LLM logic, restructure its tool loop, or add unrelated PostHog features. If the base app has pre-existing issues, note them separately but do NOT let them affect your YES/NO answers. +For proxy apps, enumerate the inference entry points in the app README before scoring. Check direct model calls behind HTTP and WebSocket handlers as well as shared proxy helpers. A generation on one route does not cover another route. Follow each stream through its final event, provider error, or client disconnect. Check token fields against that provider's response shape. + ## What this wizard is supposed to do The `ai-observability` wizard instruments an app's LLM calls so they land in PostHog as a `session → trace → span → generation` tree. There are three instrumentation paths: @@ -68,11 +70,12 @@ Key structural facts to evaluate against: - **as_env_documented** — Any new environment variables are documented in `.env.example` or README, and **no credentials are hardcoded**. The project token and host must be read from the environment. - **as_dependency_version_valid** — Declared dependency versions can actually resolve to a release containing the APIs the code uses. If the PR uses a module (e.g. `posthog.ai.otel`) that the pinned/declared version range cannot provide, mark NO. If the PR's own report claims a package or module "is not published" or "does not exist", treat that claim as suspect — it is usually a stale resolved version rather than a real gap — and mark NO unless the PR verified it against the registry's latest release. - **as_flushes_before_exit** — For scripts and CLIs, buffered events are flushed before the process exits (`posthog.shutdown()`, `atexit` registration, or `TracerProvider.shutdown()`). Batching means a short-lived process that exits without flushing loses everything. **N/A** for long-running servers. +- **as_manual_capture_ledger_delivered** — For manual capture, the delivered setup report contains a brief inference coverage ledger listing every README inference path, its transport, and capture/verification status. Check that the report file actually exists in the changed files; a wizard outro claiming it was written is not evidence. Missing paths, an absent report, or an unverified path presented as verified is NO. **N/A** for wrapper, OTel, and framework-hook paths. ### 3. PostHog implementation (AI observability-specific) - **ph_instrumentation_initialized_once** — Instrumentation is set up exactly once, at module scope or app startup — not inside a request handler or per-call function. -- **ph_generations_captured** — The app's LLM calls will actually produce `$ai_generation` events: either an instrumentor is attached to the SDK in use, or the client is swapped for the PostHog wrapper, or explicit `$ai_generation` captures exist. An OTel provider configured but never attached to an instrumentor is NO. +- **ph_generations_captured** — Every inference path in the app's README will actually produce `$ai_generation` events: either an instrumentor is attached to the SDK in use, or the client is swapped for the PostHog wrapper, or explicit `$ai_generation` captures exist. One uncovered HTTP or WebSocket path is NO. An OTel provider configured but never attached to an instrumentor is NO. - **ph_trace_groups_the_request** — All LLM calls belonging to one logical request land in **one** trace. On the wrapper path this requires a shared `posthog_trace_id` passed to every call — a call that omits it gets a freshly generated UUID, splitting one request into several single-generation traces. If each call in a multi-call request would get its own trace, mark NO. - **ph_session_id_set** — `$ai_session_id` is set, using the app's own conversation/thread identifier where one exists (see the README). Not set at all is NO. - **ph_session_id_correct_key** — The session is carried by the literal property key **`$ai_session_id`**. Keys like `posthog.session_id`, `session_id`, or `ai_session_id` are silently dropped by ingest — they look plausible and record nothing. Mark NO for any variant spelling. @@ -83,7 +86,9 @@ Key structural facts to evaluate against: ### 4. Event quality (AI observability-specific) -- **eq_events_would_render_as_tree** — Taken together, the emitted events form the hierarchy the README describes: the expected number of sessions, traces per session, and events per trace. Walk the code and count what a single run would emit. If the result is flat generations with no grouping, mark NO regardless of how much instrumentation code is present. +- **eq_events_would_render_as_tree** — Taken together, the emitted events form the hierarchy and outcomes the README describes: the expected number of sessions, traces per session, and events per trace. Walk the code and count what a single run would emit. A partial stream labeled successful, an HTTP provider error without an error outcome, or token counts read from another provider's schema is NO. Flat generations with no grouping are NO regardless of how much instrumentation code is present. +- **eq_output_choices_have_role** — For manual `$ai_output_choices`, inspect the actual provider response fixtures and each normalization path. Every emitted choice must have a `role` such as `assistant`, including direct calls and non-streamed responses. Passing through a provider object without `role` is NO, even when the event itself is captured. **N/A** when an SDK or instrumentor constructs the generation event. +- **eq_stream_terminal_parsed** — For manual streamed capture, completion and error status follow the upstream protocol's parsed terminal event. Check SSE and NDJSON separately, including WebSocket handlers whose public route differs from the upstream route. String checks that reject valid JSON formatting, or a successful stream marked as interrupted, are NO. **N/A** when there are no manually captured streams. - **eq_person_attribution** — Events carry a real user identifier from the app (not a hardcoded placeholder, not a random UUID per run, not left to an anonymous fallback). - **eq_span_parenting** — Any explicitly captured `$ai_span` sets `$ai_trace_id` and links to its parent via `$ai_parent_id`. **N/A** if no manual spans were captured. - **eq_no_fabricated_structure** — The PR does not invent structure the app does not have: no synthetic session id where the app has no conversation concept and none can be inferred, no spans around code that does nothing meaningful. diff --git a/services/pr-evaluator/prompts/output-format-ai-observability.md b/services/pr-evaluator/prompts/output-format-ai-observability.md new file mode 100644 index 000000000..a09ecc64c --- /dev/null +++ b/services/pr-evaluator/prompts/output-format-ai-observability.md @@ -0,0 +1,167 @@ +## Output template + +**Security:** Never include full API keys or secrets in output. Use redacted format like `phc_xxxx...xxxx`. + +Write your review following this Markdown structure: + +--- + +## PR Evaluation Report — AI observability + +### Summary +[1-3 sentence overview of how the PR captures this app's model calls] + +| Files changed | Lines added | Lines removed | +|---------------|-------------|---------------| +| X | +Y | -Z | + +### Confidence score: X/5 🧙 if 5/5 / 👍 if 4/5 / 🤔 if 3/5 / ❌ if 2/5 or 1/5 + +- detailed change or recommendation that's CRITICAL or MEDIUM severity +- detailed change or recommendation that's CRITICAL or MEDIUM severity +- detailed change or recommendation that's CRITICAL or MEDIUM severity + + +--- + +### File changes + +| Filename | Score | Description | +|----------|-------|-------------| +| `path/to/file.ts` | X/5 | Brief description of changes + +--- + +### App sanity check ✅ if all pass / ⚠️ if any NO / ❌ if critical items fail + +| Criteria | Result | Description | +|----------|--------|-------------| +| **App builds and runs** | Yes / No | Description | +| **Preserves existing model behavior** | Yes / No | Description | +| **No syntax or type errors** | Yes / No | Description | +| **Correct imports/exports** | Yes / No | Description | +| **Minimal, focused changes** | Yes / No | Description | +| **Manual capture ledger and report** | Yes / No / N/A | Check the report exists and reconciles every README inference path | +| **Pre-existing issues** | None / List | Issues that exist in the base app, not introduced by this PR | + +#### Issues +- **Issue title**: Description of high severity issue. Description of fix. [CRITICAL] +- **Issue title**: Description of medium severity issue. Description of fix. [MEDIUM] +- **Issue title**: Description of low severity issue. Description of fix. [LOW] + +
+

Other completed criteria

+ +- Other criterion met +- Other criterion met +
+ +--- + +### AI observability implementation ✅ if all pass / ⚠️ if any NO / ❌ if critical items fail + +| Criteria | Result | Description | +|----------|--------|-------------| +| **Every model path captured** | Yes / No | Check each path in the app README, including direct calls | +| **Trace and session grouping** | Yes / No | State the actual request and conversation IDs | +| **Tool spans** | Yes / No / N/A | Check executed tools and parent links | +| **Stream and error outcomes** | Yes / No | Check completion, provider errors, and cancellation | +| **Provider token fields** | Yes / No | Trace source fields to emitted token properties | +| **Output choice roles** | Yes / No / N/A | Inspect actual response shapes for a `role` on every manual output choice | +| **Protocol terminal events** | Yes / No / N/A | Parse each stream's terminal signal, including valid SSE formatting and WebSocket upstream paths | + +#### Issues +- **Issue title**: Description of high severity issue. Description of fix. [CRITICAL] +- **Issue title**: Description of medium severity issue. Description of fix. [MEDIUM] +- **Issue title**: Description of low severity issue. Description of fix. [LOW] + +
+

Other completed criteria

+ +- Other criterion met +- Other criterion met +
+ +--- + +### AI events ✅ if all pass / ⚠️ if any NO / ❌ if critical items fail + +| Filename | AI events | Description | +|----------|-----------------|-------------| +| `filename` | `$ai_generation`, `$ai_span`, or `$ai_embedding` | Which real model or tool actions emit them | + +#### Issues +- **Issue title**: Description of high severity issue. Description of fix. [CRITICAL] +- **Issue title**: Description of medium severity issue. Description of fix. [MEDIUM] +- **Issue title**: Description of low severity issue. Description of fix. [LOW] + +
+

Other completed criteria

+ +- Other criterion met +- Other criterion met +
+ +--- + + + +IMPORTANT: In the RUBRIC block above, replace each placeholder with exactly "yes", "no", or "n/a" (lowercase). These keys are the AI observability criteria. Use "n/a" only where the rubric permits it: `as_flushes_before_exit` for long-running servers, `as_manual_capture_ledger_delivered` outside manual capture, `ph_no_hand_authored_spans` when the app has no tools or spans, `eq_span_parenting` when no manual spans are captured, `eq_output_choices_have_role` when the event is built by an SDK or instrumentor, and `eq_stream_terminal_parsed` when no streams are manually captured. The manual capture path also permits "n/a" for `fa_imports_valid`. + + + +IMPORTANT: Leave all score values as 0 — they are computed server-side from the rubric. Only fill in "framework" and "arch_type". + +Reviewed by wizard workbench PR evaluator (AI observability)