-
-
+
+
-
-
-
-
标准化结果
-
-
模型原始输出
-
+
+
+ | 状态 | 收到 sub2api | 转发上游 | 收到 LLM 回复 | 回复 sub2api | 阶段耗时 | 总耗时 | 判定 / 错误 | |
+
+
+
-
-
+
+
+
+
+
+
+
+ 上游审核模型
保存后立即替换内存配置,后续请求无需读取磁盘。
v0
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/src/sub2api_auditer/web.py b/src/sub2api_auditer/web.py
index 5b8a36c..58ca3a2 100644
--- a/src/sub2api_auditer/web.py
+++ b/src/sub2api_auditer/web.py
@@ -11,13 +11,20 @@
import httpx
from starlette.applications import Starlette
+from starlette.background import BackgroundTask
from starlette.requests import Request
from starlette.responses import HTMLResponse, JSONResponse, PlainTextResponse, Response
from starlette.routing import Route
from . import __version__
from .config import ConfigConflict, ConfigError, ConfigStore
-from .protocol import ProtocolError, extract_audit_text, make_openai_error, make_openai_response
+from .observability import TraceStore
+from .protocol import (
+ ProtocolError,
+ extract_audit_text,
+ make_openai_error,
+ make_openai_response,
+)
from .service import AuditerService, UpstreamError, env_int
LOGGER = logging.getLogger("sub2api_auditer")
@@ -41,7 +48,38 @@ def _error(message: str, code: str, status: int) -> JSONResponse:
return JSONResponse(make_openai_error(message, code), status_code=status)
-async def _json(request: Request) -> Mapping[str, Any]:
+def _traced_response(
+ response: Response,
+ service: AuditerService,
+ trace_id: str,
+) -> Response:
+ response.background = BackgroundTask(
+ service.traces.mark_replied,
+ trace_id,
+ http_status=response.status_code,
+ )
+ response.headers["X-Auditer-Trace-Id"] = trace_id
+ return response
+
+
+def _traced_error(
+ service: AuditerService,
+ trace_id: str,
+ *,
+ message: str,
+ code: str,
+ status: int,
+) -> Response:
+ service.traces.mark_error(
+ trace_id,
+ code=code,
+ message=message,
+ http_status=status,
+ )
+ return _traced_response(_error(message, code, status), service, trace_id)
+
+
+async def _read_json(request: Request) -> tuple[Mapping[str, Any], int]:
limit = env_int("MAX_REQUEST_BODY_BYTES", 2 * 1024 * 1024)
length = request.headers.get("content-length")
if length:
@@ -57,10 +95,15 @@ async def _json(request: Request) -> Mapping[str, Any]:
body.extend(chunk)
try:
value = json.loads(body or b"{}")
- except json.JSONDecodeError as exc:
+ except (UnicodeDecodeError, json.JSONDecodeError) as exc:
raise ProtocolError("请求体不是有效 JSON") from exc
if not isinstance(value, Mapping):
raise ProtocolError("请求体根节点必须是 JSON 对象")
+ return value, len(body)
+
+
+async def _json(request: Request) -> Mapping[str, Any]:
+ value, _ = await _read_json(request)
return value
@@ -71,16 +114,26 @@ def _static(name: str) -> str:
async def home(request: Request) -> Response:
del request
- return HTMLResponse(_static("index.html"))
+ response = HTMLResponse(_static("index.html"))
+ response.headers["Cache-Control"] = "no-cache"
+ response.headers["Content-Security-Policy"] = (
+ "default-src 'self'; script-src 'self'; style-src 'self' 'unsafe-inline'; "
+ "img-src 'self' data:; connect-src 'self'; frame-ancestors *; base-uri 'self'; "
+ "form-action 'self'"
+ )
+ return response
async def asset(request: Request) -> Response:
name = request.path_params["name"]
if name == "app.css":
- return PlainTextResponse(_static(name), media_type="text/css")
- if name == "app.js":
- return PlainTextResponse(_static(name), media_type="application/javascript")
- return Response(status_code=404)
+ response = PlainTextResponse(_static(name), media_type="text/css")
+ elif name == "app.js":
+ response = PlainTextResponse(_static(name), media_type="application/javascript")
+ else:
+ return Response(status_code=404)
+ response.headers["Cache-Control"] = "public, max-age=300"
+ return response
async def healthz(request: Request) -> Response:
@@ -91,25 +144,30 @@ async def healthz(request: Request) -> Response:
async def readyz(request: Request) -> Response:
store: ConfigStore = request.app.state.store
config = store.get()
- status = 200 if config.ready and not store.load_error else 503
- return JSONResponse({
- "status": "ready" if status == 200 else "not_ready",
- "configured": config.ready,
- "config_version": config.version,
- "config_error": store.load_error,
- }, status_code=status)
+ status_code = 200 if config.ready and not store.load_error else 503
+ return JSONResponse(
+ {
+ "status": "ready" if status_code == 200 else "not_ready",
+ "configured": config.ready,
+ "config_version": config.version,
+ "config_error": store.load_error,
+ },
+ status_code=status_code,
+ )
async def get_config(request: Request) -> Response:
if not _authorized(request, request.app.state.admin_token):
return _error("管理员令牌无效", "admin_unauthorized", 401)
store: ConfigStore = request.app.state.store
- return JSONResponse({
- "config": store.get().public_dict(),
- "config_error": store.load_error,
- "admin_auth_enabled": bool(request.app.state.admin_token),
- "proxy_auth_enabled": bool(request.app.state.auditer_token),
- })
+ return JSONResponse(
+ {
+ "config": store.get().public_dict(),
+ "config_error": store.load_error,
+ "admin_auth_enabled": bool(request.app.state.admin_token),
+ "proxy_auth_enabled": bool(request.app.state.auditer_token),
+ }
+ )
async def put_config(request: Request) -> Response:
@@ -132,72 +190,217 @@ async def status(request: Request) -> Response:
return _error("管理员令牌无效", "admin_unauthorized", 401)
store: ConfigStore = request.app.state.store
service: AuditerService = request.app.state.service
- return JSONResponse({
- "version": __version__,
- "ready": store.get().ready and not store.load_error,
- "config_version": store.get().version,
- "config_error": store.load_error,
- "stats": service.stats.public_dict(),
- })
+ return JSONResponse(
+ {
+ "version": __version__,
+ "ready": store.get().ready and not store.load_error,
+ "config_version": store.get().version,
+ "config_error": store.load_error,
+ "stats": service.traces.runtime_stats(),
+ }
+ )
+
+
+async def processing_logs(request: Request) -> Response:
+ if not _authorized(request, request.app.state.admin_token):
+ return _error("管理员令牌无效", "admin_unauthorized", 401)
+ try:
+ limit = int(request.query_params.get("limit", "100"))
+ except ValueError:
+ limit = 100
+ service: AuditerService = request.app.state.service
+ return JSONResponse(
+ {
+ "items": service.traces.list(limit=limit),
+ "capacity": service.traces.capacity,
+ }
+ )
-async def test_audit(request: Request) -> Response:
+async def processing_statistics(request: Request) -> Response:
+ if not _authorized(request, request.app.state.admin_token):
+ return _error("管理员令牌无效", "admin_unauthorized", 401)
+ service: AuditerService = request.app.state.service
+ return JSONResponse(service.traces.statistics())
+
+
+async def clear_processing_logs(request: Request) -> Response:
if not _authorized(request, request.app.state.admin_token):
return _error("管理员令牌无效", "admin_unauthorized", 401)
+ service: AuditerService = request.app.state.service
+ return JSONResponse({"ok": True, "cleared": service.traces.clear()})
+
+
+async def test_audit(request: Request) -> Response:
+ service: AuditerService = request.app.state.service
+ trace_id = service.traces.begin(
+ source="manual_test",
+ client_request_id=request.headers.get("x-request-id", ""),
+ )
+ if not _authorized(request, request.app.state.admin_token):
+ return _traced_error(
+ service,
+ trace_id,
+ message="管理员令牌无效",
+ code="admin_unauthorized",
+ status=401,
+ )
try:
- text = str((await _json(request)).get("text", "")).strip()
+ payload, body_bytes = await _read_json(request)
+ text = str(payload.get("text", "")).strip()
if not text:
raise ProtocolError("测试文本不能为空")
- call = await request.app.state.service.audit(text)
+ call = await service.audit(
+ text,
+ trace_id=trace_id,
+ request_model="manual-test",
+ input_bytes=body_bytes,
+ )
except ProtocolError as exc:
- return _error(str(exc), "invalid_test_input", 400)
+ return _traced_error(
+ service,
+ trace_id,
+ message=str(exc),
+ code="invalid_test_input",
+ status=400,
+ )
+ except ConfigError as exc:
+ return _traced_error(
+ service,
+ trace_id,
+ message=str(exc),
+ code="invalid_upstream_config",
+ status=503,
+ )
except UpstreamError as exc:
- return _error(str(exc), exc.code, exc.status_code)
- return JSONResponse({
- "ok": True,
- "normalized": call.result.as_dict(),
- "raw_model_output": call.raw_output[:8000],
- "latency_ms": call.latency_ms,
- "upstream_request_id": call.upstream_request_id,
- })
+ return _traced_error(
+ service,
+ trace_id,
+ message=str(exc),
+ code=exc.code,
+ status=exc.status_code,
+ )
+ except Exception:
+ LOGGER.exception("unexpected manual audit failure")
+ return _traced_error(
+ service,
+ trace_id,
+ message="审计服务发生内部错误",
+ code="internal_error",
+ status=500,
+ )
+
+ response = JSONResponse(
+ {
+ "ok": True,
+ "trace_id": call.trace_id,
+ "normalized": call.result.as_dict(),
+ "raw_model_output": call.raw_output[:8000],
+ "latency_ms": call.latency_ms,
+ "upstream_request_id": call.upstream_request_id,
+ }
+ )
+ return _traced_response(response, service, trace_id)
async def models(request: Request) -> Response:
if not _authorized(request, request.app.state.auditer_token):
return _error("审计服务访问令牌无效", "unauthorized", 401)
configured = request.app.state.store.get().model
- ids = list(dict.fromkeys(value for value in (configured, "sub2api-auditer") if value))
- return JSONResponse({
- "object": "list",
- "data": [{"id": value, "object": "model", "created": 0, "owned_by": "sub2api-auditer"} for value in ids],
- })
+ ids = list(
+ dict.fromkeys(
+ value for value in (configured, "sub2api-auditer") if value
+ )
+ )
+ return JSONResponse(
+ {
+ "object": "list",
+ "data": [
+ {
+ "id": value,
+ "object": "model",
+ "created": 0,
+ "owned_by": "sub2api-auditer",
+ }
+ for value in ids
+ ],
+ }
+ )
async def completions(request: Request) -> Response:
+ service: AuditerService = request.app.state.service
+ trace_id = service.traces.begin(
+ source="sub2api",
+ client_request_id=(
+ request.headers.get("x-request-id", "")
+ or request.headers.get("x-correlation-id", "")
+ ),
+ )
if not _authorized(request, request.app.state.auditer_token):
- return _error("审计服务访问令牌无效", "unauthorized", 401)
+ return _traced_error(
+ service,
+ trace_id,
+ message="审计服务访问令牌无效",
+ code="unauthorized",
+ status=401,
+ )
+
try:
- payload = await _json(request)
- call = await request.app.state.service.audit(extract_audit_text(payload))
+ payload, body_bytes = await _read_json(request)
+ request_model = str(payload.get("model", "") or "sub2api-auditer")
+ text = extract_audit_text(payload)
+ call = await service.audit(
+ text,
+ trace_id=trace_id,
+ request_model=request_model,
+ input_bytes=body_bytes,
+ )
except ProtocolError as exc:
- return _error(str(exc), "invalid_audit_request", 413 if "过大" in str(exc) else 400)
+ return _traced_error(
+ service,
+ trace_id,
+ message=str(exc),
+ code="invalid_audit_request",
+ status=413 if "过大" in str(exc) else 400,
+ )
except ConfigError as exc:
- return _error(str(exc), "invalid_upstream_config", 503)
+ return _traced_error(
+ service,
+ trace_id,
+ message=str(exc),
+ code="invalid_upstream_config",
+ status=503,
+ )
except UpstreamError as exc:
- return _error(str(exc), exc.code, exc.status_code)
+ return _traced_error(
+ service,
+ trace_id,
+ message=str(exc),
+ code=exc.code,
+ status=exc.status_code,
+ )
except Exception:
LOGGER.exception("unexpected audit failure")
- return _error("审计服务发生内部错误", "internal_error", 500)
+ return _traced_error(
+ service,
+ trace_id,
+ message="审计服务发生内部错误",
+ code="internal_error",
+ status=500,
+ )
- response = JSONResponse(make_openai_response(
- result=call.result,
- request_model=str(payload.get("model", "") or "sub2api-auditer"),
- ))
+ response = JSONResponse(
+ make_openai_response(
+ result=call.result,
+ request_model=request_model,
+ )
+ )
response.headers["X-Auditer-Latency-Ms"] = str(call.latency_ms)
response.headers["X-Auditer-Version"] = __version__
if call.upstream_request_id:
response.headers["X-Upstream-Request-Id"] = call.upstream_request_id[:256]
- return response
+ return _traced_response(response, service, trace_id)
def create_app(
@@ -209,22 +412,39 @@ def create_app(
) -> Starlette:
store = ConfigStore(config_path or os.getenv("CONFIG_PATH", "./data/config.json"))
owns_client = client is None
+ traces = TraceStore(env_int("LOG_CAPACITY", 100))
@asynccontextmanager
async def lifespan(app: Starlette):
await store.load()
app.state.store = store
- app.state.admin_token = os.getenv("ADMIN_TOKEN", "").strip() if admin_token is None else admin_token.strip()
- app.state.auditer_token = os.getenv("AUDITER_TOKEN", "").strip() if auditer_token is None else auditer_token.strip()
+ app.state.admin_token = (
+ os.getenv("ADMIN_TOKEN", "").strip()
+ if admin_token is None
+ else admin_token.strip()
+ )
+ app.state.auditer_token = (
+ os.getenv("AUDITER_TOKEN", "").strip()
+ if auditer_token is None
+ else auditer_token.strip()
+ )
app.state.client = client or httpx.AsyncClient(
- limits=httpx.Limits(max_connections=200, max_keepalive_connections=50, keepalive_expiry=30.0),
+ limits=httpx.Limits(
+ max_connections=env_int("HTTP_MAX_CONNECTIONS", 200),
+ max_keepalive_connections=env_int("HTTP_MAX_KEEPALIVE", 50),
+ keepalive_expiry=30.0,
+ ),
follow_redirects=False,
trust_env=False,
)
- app.state.service = AuditerService(store, app.state.client)
+ app.state.service = AuditerService(store, app.state.client, traces)
LOGGER.info(
- "started version=%s configured=%s admin_auth=%s proxy_auth=%s",
- __version__, store.get().ready, bool(app.state.admin_token), bool(app.state.auditer_token),
+ "started version=%s configured=%s admin_auth=%s proxy_auth=%s log_capacity=%s",
+ __version__,
+ store.get().ready,
+ bool(app.state.admin_token),
+ bool(app.state.auditer_token),
+ traces.capacity,
)
try:
yield
@@ -240,6 +460,9 @@ async def lifespan(app: Starlette):
Route("/api/config", get_config, methods=["GET"]),
Route("/api/config", put_config, methods=["PUT"]),
Route("/api/status", status, methods=["GET"]),
+ Route("/api/logs", processing_logs, methods=["GET"]),
+ Route("/api/logs", clear_processing_logs, methods=["DELETE"]),
+ Route("/api/statistics", processing_statistics, methods=["GET"]),
Route("/api/test", test_audit, methods=["POST"]),
Route("/v1/models", models, methods=["GET"]),
Route("/models", models, methods=["GET"]),
diff --git a/tests/test_app.py b/tests/test_app.py
index 1ab9ce6..e9c347f 100644
--- a/tests/test_app.py
+++ b/tests/test_app.py
@@ -1,19 +1,34 @@
+from __future__ import annotations
+
+import itertools
+
import httpx
from sub2api_auditer.app import create_app
-async def _configured_app(tmp_path, handler, *, auditer_token=""):
+async def _configured_app(tmp_path, handler, *, admin_token="", auditer_token=""):
upstream = httpx.AsyncClient(transport=httpx.MockTransport(handler))
app = create_app(
config_path=str(tmp_path / "config.json"),
client=upstream,
- admin_token="",
+ admin_token=admin_token,
auditer_token=auditer_token,
)
return app, upstream
+async def _configure(app):
+ await app.state.store.update(
+ {
+ "base_url": "https://gateway.example.com/v1",
+ "api_key": "sk-upstream",
+ "model": "audit-model",
+ "prompt": "自定义审核策略",
+ }
+ )
+
+
async def test_chat_completions_returns_sub2api_format(tmp_path):
observed = {}
@@ -37,18 +52,12 @@ def handler(request: httpx.Request) -> httpx.Response:
app, upstream = await _configured_app(tmp_path, handler)
async with app.router.lifespan_context(app):
- await app.state.store.update(
- {
- "base_url": "https://gateway.example.com/v1",
- "api_key": "sk-upstream",
- "model": "audit-model",
- "prompt": "自定义审核策略",
- }
- )
+ await _configure(app)
transport = httpx.ASGITransport(app=app)
async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
response = await client.post(
"/v1/chat/completions",
+ headers={"X-Request-ID": "sub2api-123"},
json={"model": "sub2api-model", "messages": [{"role": "user", "content": "test"}]},
)
@@ -60,6 +69,7 @@ def handler(request: httpx.Request) -> httpx.Response:
assert observed["authorization"] == "Bearer sk-upstream"
assert '"model":"audit-model"' in observed["payload"].replace(" ", "")
assert "自定义审核策略" in observed["payload"]
+ assert response.headers["x-auditer-trace-id"].startswith("aud-")
async def test_models_endpoint_supports_sub2api_probe(tmp_path):
@@ -68,13 +78,7 @@ def handler(request: httpx.Request) -> httpx.Response:
app, upstream = await _configured_app(tmp_path, handler, auditer_token="secret")
async with app.router.lifespan_context(app):
- await app.state.store.update(
- {
- "base_url": "https://gateway.example.com",
- "model": "audit-model",
- "prompt": "审核策略",
- }
- )
+ await _configure(app)
transport = httpx.ASGITransport(app=app)
async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
unauthorized = await client.get("/v1/models")
@@ -97,13 +101,7 @@ def handler(request: httpx.Request) -> httpx.Response:
app, upstream = await _configured_app(tmp_path, handler)
async with app.router.lifespan_context(app):
- await app.state.store.update(
- {
- "base_url": "https://gateway.example.com",
- "model": "audit-model",
- "prompt": "审核策略",
- }
- )
+ await _configure(app)
transport = httpx.ASGITransport(app=app)
async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
response = await client.post(
@@ -114,3 +112,110 @@ def handler(request: httpx.Request) -> httpx.Response:
await upstream.aclose()
assert response.status_code == 502
assert response.json()["error"]["code"] == "audit_model_invalid_response"
+
+
+async def test_processing_log_captures_four_timestamps_and_phases(tmp_path):
+ def handler(request: httpx.Request) -> httpx.Response:
+ return httpx.Response(
+ 200,
+ headers={"x-request-id": "upstream-timing"},
+ json={"choices": [{"message": {"content": '{"safety":"Safe","categories":[]}'}}]},
+ )
+
+ app, upstream = await _configured_app(tmp_path, handler)
+ async with app.router.lifespan_context(app):
+ await _configure(app)
+ transport = httpx.ASGITransport(app=app)
+ async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
+ response = await client.post(
+ "/v1/chat/completions",
+ headers={"X-Request-ID": "sub2api-timing"},
+ json={"model": "sub2api-auditer", "messages": [{"role": "user", "content": "hello"}]},
+ )
+ logs_response = await client.get("/api/logs")
+
+ await upstream.aclose()
+ assert response.status_code == 200
+ assert logs_response.status_code == 200
+ logs = logs_response.json()["items"]
+ assert len(logs) == 1
+ trace = logs[0]
+ assert trace["source"] == "sub2api"
+ assert trace["client_request_id"] == "sub2api-timing"
+ assert trace["upstream_request_id"] == "upstream-timing"
+ assert trace["status"] == "success"
+ assert trace["safety"] == "Safe"
+ assert trace["received_at"]
+ assert trace["forwarded_at"]
+ assert trace["llm_replied_at"]
+ assert trace["sub2api_replied_at"]
+ assert trace["preprocess_ms"] is not None and trace["preprocess_ms"] >= 0
+ assert trace["upstream_ms"] is not None and trace["upstream_ms"] >= 0
+ assert trace["response_ms"] is not None and trace["response_ms"] >= 0
+ assert trace["total_ms"] is not None and trace["total_ms"] >= 0
+ assert trace["input_chars"] == 5
+ assert trace["upstream_response_bytes"] > 0
+
+
+async def test_statistics_are_derived_from_current_log_window_and_clearable(tmp_path):
+ counter = itertools.count()
+
+ def handler(request: httpx.Request) -> httpx.Response:
+ if next(counter) == 0:
+ content = '{"safety":"Controversial","categories":["Copyright Violation"]}'
+ else:
+ content = "unparseable output"
+ return httpx.Response(200, json={"choices": [{"message": {"content": content}}]})
+
+ app, upstream = await _configured_app(tmp_path, handler)
+ async with app.router.lifespan_context(app):
+ await _configure(app)
+ transport = httpx.ASGITransport(app=app)
+ async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
+ first = await client.post(
+ "/v1/chat/completions",
+ json={"messages": [{"role": "user", "content": "first"}]},
+ )
+ second = await client.post(
+ "/v1/chat/completions",
+ json={"messages": [{"role": "user", "content": "second"}]},
+ )
+ statistics = await client.get("/api/statistics")
+ cleared = await client.delete("/api/logs")
+ after_clear = await client.get("/api/statistics")
+
+ await upstream.aclose()
+ assert first.status_code == 200
+ assert second.status_code == 502
+ data = statistics.json()
+ assert data["window_size"] == 2
+ assert data["completed"] == 2
+ assert data["success"] == 1
+ assert data["failed"] == 1
+ assert data["success_rate"] == 50.0
+ assert data["decisions"]["Controversial"] == 1
+ assert data["decisions"]["Unclassified"] == 1
+ assert data["errors"] == [{"code": "audit_model_invalid_response", "count": 1}]
+ assert len(data["series"]) == 2
+ assert cleared.json() == {"ok": True, "cleared": 2}
+ assert after_clear.json()["window_size"] == 0
+
+
+async def test_logs_and_statistics_require_admin_token(tmp_path):
+ def handler(request: httpx.Request) -> httpx.Response:
+ return httpx.Response(200, json={"choices": [{"message": {"content": "Safety: Safe\nCategories: None"}}]})
+
+ app, upstream = await _configured_app(tmp_path, handler, admin_token="admin-secret")
+ async with app.router.lifespan_context(app):
+ transport = httpx.ASGITransport(app=app)
+ async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
+ logs = await client.get("/api/logs")
+ stats = await client.get("/api/statistics")
+ authorized = await client.get(
+ "/api/logs", headers={"Authorization": "Bearer admin-secret"}
+ )
+
+ await upstream.aclose()
+ assert logs.status_code == 401
+ assert stats.status_code == 401
+ assert authorized.status_code == 200
diff --git a/tests/test_observability.py b/tests/test_observability.py
new file mode 100644
index 0000000..bce985d
--- /dev/null
+++ b/tests/test_observability.py
@@ -0,0 +1,54 @@
+from __future__ import annotations
+
+import time
+
+from sub2api_auditer.observability import TraceStore
+
+
+def test_trace_store_keeps_latest_capacity_items():
+ store = TraceStore(capacity=100)
+ ids = [store.begin(source="sub2api", client_request_id=str(index)) for index in range(105)]
+ items = store.list(limit=100)
+
+ assert len(items) == 100
+ assert items[0]["id"] == ids[-1]
+ assert items[-1]["id"] == ids[5]
+
+
+def test_trace_store_records_monotonic_phase_durations_and_statistics():
+ store = TraceStore(capacity=100)
+ trace_id = store.begin(source="sub2api", request_model="requested")
+ store.update_request(
+ trace_id,
+ request_model="requested",
+ upstream_model="actual",
+ input_chars=4,
+ input_bytes=4,
+ )
+ time.sleep(0.001)
+ store.mark_forwarded(trace_id)
+ time.sleep(0.001)
+ store.mark_llm_replied(
+ trace_id,
+ upstream_http_status=200,
+ upstream_request_id="upstream-id",
+ response_bytes=128,
+ )
+ store.mark_result(trace_id, safety="Unsafe", categories=("Jailbreak",))
+ time.sleep(0.001)
+ store.mark_replied(trace_id, http_status=200)
+
+ trace = store.list()[0]
+ assert trace["preprocess_ms"] > 0
+ assert trace["upstream_ms"] > 0
+ assert trace["response_ms"] > 0
+ assert trace["total_ms"] >= trace["preprocess_ms"] + trace["upstream_ms"]
+ assert trace["upstream_model"] == "actual"
+ assert trace["upstream_request_id"] == "upstream-id"
+
+ stats = store.statistics()
+ assert stats["window_size"] == 1
+ assert stats["success"] == 1
+ assert stats["decisions"]["Unsafe"] == 1
+ assert stats["latency"]["p95_ms"] == trace["total_ms"]
+ assert stats["phases"]["upstream_average_ms"] == trace["upstream_ms"]