Files
6krrt/tests/test_router_wall_clock.py
adlee-was-taken 0111da9bbf feat(telemetry): router-observed wall-clock and TTFT columns, with plan docs
Adds router_wall_seconds and router_ttft_seconds to energy_observations:

- router_wall_seconds: time.monotonic() from just before the accepted
  candidate's POST/connection-open to the complete response body (buffered)
  or the last byte forwarded (streaming). Re-marked per candidate so
  failover time is excluded — a dead model's 30s stall is not charged to
  the healthy one that replaced it.
- router_ttft_seconds: streaming-only. First delta carrying non-empty
  content or a tool_calls fragment, excluding the role-only opening delta.
  NULL on buffered rows (not applicable) and on streams that produced no
  output token (a broken upstream — not the same as a literal 0).

Separate from duration_seconds (provider's reported serving time) on
purpose: OpenRouter reports duration_seconds on 0 of 3,850 rows, while
the router can always measure its own clock. The two quantities are stored
independently and never written into each other.

log_observation() defaults both new params to None so seed_energy.py and
eval_proficiency.py pass unchanged — a reprise of 6e729ad's bug where new
keyword-only arguments killed the seed timer.

Also:
- docs/data-model.md: document the new columns, span definitions, and
  the distinction from duration_seconds
- config/schema.sql: full CREATE TABLE declaration
- plans/token-waste-waves.md: status update for Wave 1
- plans/ten-thousand-foot-review.md: companion diagnosis
2026-09-13 12:25:39 -04:00

553 lines
19 KiB
Python

"""Router-observed latency: `router_wall_seconds` and `router_ttft_seconds`.
`energy_observations.duration_seconds` is the PROVIDER's reported serving time.
On the live database it is present on 31,243 of 31,309 NeuralWatt rows and on
0 of 3,850 OpenRouter rows, and no request-body opt-in will change that --
OpenRouter serves generation timing only from its separate
`/api/v1/generation?id=` endpoint. OpenRouter is ~73% of routed decisions, so a
latency term in the objective needs a number that exists for every provider.
These pin the router's own measurement: that it is written on both dispatch
paths, that it is the ACCEPTED candidate's span rather than the whole request's
history, that TTFT excludes the role-only opening delta every provider sends,
and that the columns arrive by an idempotent ALTER on a database that predates
them. Nothing here calls a real provider or reads a real clock except the one
plausibility test that deliberately does.
"""
from __future__ import annotations
import json
import sqlite3
from pathlib import Path
import pytest
from starlette.testclient import TestClient
import dispatcher
import session_cache
ROOT = Path(__file__).resolve().parent.parent
SCHEMA_SQL = (ROOT / "config" / "schema.sql").read_text()
MODEL = "cheap-model"
class _Clock:
"""A monotonic stand-in the test drives by hand.
Only ever moves forward, like the real one. Advanced explicitly by the
fakes below so every asserted duration is an exact figure rather than a
band -- a timing test that asserts a range can pass while measuring the
wrong span.
"""
def __init__(self, start: float = 1000.0) -> None:
self.now = start
def __call__(self) -> float:
return self.now
def advance(self, seconds: float) -> None:
self.now += seconds
class _FakeResponse:
"""Minimal requests.Response stand-in; advances the clock as it streams."""
def __init__(
self,
payload=None,
*,
status_code=200,
lines=None,
clock=None,
step=0.0,
):
self.status_code = status_code
self._payload = payload or {}
# (delay_before_yield, line) pairs, so a test can put the clock
# exactly where it wants it when a given SSE line is parsed.
self._lines = lines or []
self._clock = clock
self._step = step
self.text = json.dumps(self._payload)
self.headers = {"Content-Type": "text/event-stream; charset=utf-8"}
self.reason = ""
self.request = None
def json(self):
return self._payload
@property
def encoding(self):
return "utf-8"
def iter_lines(self, decode_unicode=False):
for entry in self._lines:
delay, line = entry if isinstance(entry, tuple) else (self._step, entry)
if self._clock is not None and delay:
self._clock.advance(delay)
wire = line.encode("utf-8") if isinstance(line, str) else line
yield wire.decode("utf-8") if decode_unicode else wire
def close(self):
pass
def _completion(model=MODEL):
return {
"id": "chatcmpl-buffered-1",
"model": model,
"choices": [
{
"message": {"role": "assistant", "content": "hello there"},
"finish_reason": "stop",
}
],
"usage": {"prompt_tokens": 31, "completion_tokens": 12},
}
# A realistic SSE shape: providers open with a role-only delta carrying no
# content. That is protocol, not answer, and must NOT count as the first token.
STREAM_LINES = [
(0.5, 'data: {"id":"chatcmpl-s-1","choices":[{"delta":{"role":"assistant","content":""}}]}'),
(0.0, ""),
(0.5, 'data: {"id":"chatcmpl-s-1","choices":[{"delta":{"content":"hel"}}]}'),
(0.0, ""),
(1.0, 'data: {"id":"chatcmpl-s-1","choices":[{"delta":{"content":"lo"},'
'"finish_reason":"stop"}],"usage":{"prompt_tokens":5,"completion_tokens":2}}'),
(0.0, ""),
(1.0, "data: [DONE]"),
(0.0, ""),
]
@pytest.fixture(autouse=True)
def _clean_state():
session_cache.clear()
dispatcher._provider_refusal_since.clear()
yield
dispatcher._provider_refusal_since.clear()
session_cache.clear()
@pytest.fixture
def router(tmp_path, monkeypatch):
"""A chat-completions client over a one-model catalog, with a fake clock."""
db_path = tmp_path / "test.db"
conn = sqlite3.connect(db_path)
conn.executescript(SCHEMA_SQL)
conn.execute(
"""
INSERT INTO models (
model_id, provider, base_model_id, tier, context_window,
effective_context_window, max_output_tokens,
cost_per_1m_prompt, cost_per_1m_completion,
supports_vision, supports_json_mode,
latency_class, reasoning_mode, context_variant,
access_level, availability, last_updated
) VALUES (?, 'neuralwatt', ?, 2, 262128, 192500, 16384, 0.10, 0.30,
1, 1, 'standard', 'default', 'full', 'public', 'active',
'2026-08-22T00:00:00+00:00')
""",
(MODEL, MODEL),
)
conn.commit()
conn.close()
monkeypatch.setattr(dispatcher.cfg.database, "path", str(db_path))
monkeypatch.setattr(dispatcher.cfg.verification, "local_llm_enabled", False)
monkeypatch.setattr(dispatcher.cfg.local_vision, "enabled", False)
monkeypatch.setattr(dispatcher.cfg.session_cache, "enabled", False)
monkeypatch.setattr(dispatcher.cfg.exploration, "epsilon", 0.0)
monkeypatch.setenv("NEURALWATT_API_KEY", "test-key")
clock = _Clock()
monkeypatch.setattr(dispatcher.time, "monotonic", clock)
state = {"buffered_delay": 2.5, "responses": None, "post": None}
def fake_post(url, headers=None, json=None, stream=False, timeout=None):
if state["post"] is not None:
return state["post"](
url, headers=headers, json=json, stream=stream, timeout=timeout
)
if state["responses"] is not None:
return state["responses"].pop(0)
if stream:
return _FakeResponse(lines=STREAM_LINES, clock=clock)
clock.advance(state["buffered_delay"])
return _FakeResponse(_completion(json["model"]))
monkeypatch.setattr(dispatcher.requests, "post", fake_post)
yield TestClient(dispatcher.app), clock, db_path, state
def _latency_row(db_path):
conn = sqlite3.connect(db_path)
conn.row_factory = sqlite3.Row
row = conn.execute(
"SELECT model_id, duration_seconds, router_wall_seconds, "
"router_ttft_seconds FROM energy_observations ORDER BY id DESC LIMIT 1"
).fetchone()
conn.close()
return row
# --- the two dispatch paths ------------------------------------------------
def test_buffered_dispatch_records_router_wall_clock(router):
"""The non-streaming path times its own POST, start to complete body."""
client, clock, db_path, _state = router
resp = client.post(
"/v1/chat/completions",
json={"model": MODEL, "messages": [{"role": "user", "content": "hi"}]},
)
assert resp.status_code == 200
row = _latency_row(db_path)
assert row is not None
assert row["router_wall_seconds"] == 2.5
# Not applicable, deliberately NULL rather than 0: on a buffered answer the
# first token and the last arrive together, so there is no TTFT to report.
assert row["router_ttft_seconds"] is None
def test_streaming_dispatch_records_wall_clock_and_ttft(router):
"""The streamed path times connection-open to last byte, plus first token.
3.0 s total is every advance in STREAM_LINES; 1.0 s TTFT is the first two,
i.e. the first delta that carried real content. The role-only opener at
0.5 s is excluded on purpose -- it is protocol, not an answer, and counting
it would make TTFT read better than anything the operator ever saw.
"""
client, clock, db_path, _state = router
resp = client.post(
"/v1/chat/completions",
json={
"model": MODEL,
"messages": [{"role": "user", "content": "hi"}],
"stream": True,
},
)
assert resp.status_code == 200
assert "data: [DONE]" in resp.text
row = _latency_row(db_path)
assert row is not None
assert row["router_wall_seconds"] == 3.0
assert row["router_ttft_seconds"] == 1.0
# TTFT is a prefix of the total, never the other way round.
assert row["router_ttft_seconds"] < row["router_wall_seconds"]
def test_a_tool_call_turn_still_gets_a_ttft(router):
"""An agent turn emits no content at all, and is the traffic that matters.
Keying TTFT on content alone would leave every tool-calling turn NULL --
which is most of this router's real traffic, since opencode sends tools on
essentially every request.
"""
client, clock, db_path, _state = router
tool_stream = [
(0.5, 'data: {"id":"chatcmpl-t-1","choices":[{"delta":{"role":"assistant"}}]}'),
(0.0, ""),
(0.7, 'data: {"id":"chatcmpl-t-1","choices":[{"delta":{"tool_calls":'
'[{"index":0,"function":{"name":"read"}}]}}]}'),
(0.0, ""),
(0.3, 'data: {"id":"chatcmpl-t-1","choices":[{"delta":{},'
'"finish_reason":"tool_calls"}],"usage":{"prompt_tokens":9,'
'"completion_tokens":4}}'),
(0.0, ""),
(0.0, "data: [DONE]"),
(0.0, ""),
]
_state["responses"] = [_FakeResponse(lines=tool_stream, clock=clock)]
resp = client.post(
"/v1/chat/completions",
json={
"model": MODEL,
"messages": [{"role": "user", "content": "read the config"}],
"stream": True,
},
)
assert resp.status_code == 200
row = _latency_row(db_path)
assert row["router_ttft_seconds"] == 1.2
assert row["router_wall_seconds"] == 1.5
def test_a_stream_with_no_output_delta_records_wall_but_null_ttft(router):
"""No first token ever arrived. That is a real outcome, not a zero."""
client, clock, db_path, _state = router
empty_stream = [
(0.4, 'data: {"id":"chatcmpl-e-1","choices":[{"delta":{"role":"assistant"}}]}'),
(0.0, ""),
(0.6, 'data: {"id":"chatcmpl-e-1","choices":[{"delta":{},'
'"finish_reason":"stop"}],"usage":{"prompt_tokens":3,'
'"completion_tokens":0}}'),
(0.0, ""),
(0.0, "data: [DONE]"),
(0.0, ""),
]
_state["responses"] = [_FakeResponse(lines=empty_stream, clock=clock)]
resp = client.post(
"/v1/chat/completions",
json={
"model": MODEL,
"messages": [{"role": "user", "content": "hi"}],
"stream": True,
},
)
assert resp.status_code == 200
row = _latency_row(db_path)
assert row["router_wall_seconds"] == 1.0
assert row["router_ttft_seconds"] is None
# --- what the span excludes ------------------------------------------------
def test_a_failed_attempt_is_not_charged_to_the_model_that_answered(
router, monkeypatch
):
"""The mark is re-taken per candidate, so failover time is excluded.
The question a latency term asks is "how long did this model take", not
"how long did the request take to satisfy". A 30 s stall on a dead
candidate must not be recorded against the healthy one that picked it up --
which is exactly backwards, since the healthy one is what the ranking would
then be pushed away from.
"""
client, clock, db_path, _state = router
monkeypatch.setattr(dispatcher.cfg.circuit_breaker, "enabled", False)
conn = sqlite3.connect(db_path)
conn.execute(
"""
INSERT INTO models (
model_id, provider, base_model_id, tier, context_window,
effective_context_window, max_output_tokens,
cost_per_1m_prompt, cost_per_1m_completion,
supports_vision, supports_json_mode,
latency_class, reasoning_mode, context_variant,
access_level, availability, last_updated
) VALUES ('backup-model', 'neuralwatt', 'backup-model', 2, 262128,
192500, 16384, 0.20, 0.60, 1, 1, 'standard', 'default',
'full', 'public', 'active', '2026-08-22T00:00:00+00:00')
"""
)
conn.commit()
conn.close()
attempts = []
def failing_then_ok(url, headers=None, json=None, stream=False, timeout=None):
if not stream:
return _FakeResponse(_completion(json["model"]))
attempts.append(json["model"])
if len(attempts) == 1:
clock.advance(30.0)
return _FakeResponse({"error": "boom"}, status_code=503)
return _FakeResponse(lines=STREAM_LINES, clock=clock)
_state["post"] = failing_then_ok
resp = client.post(
"/v1/chat/completions",
json={
"model": "auto",
"messages": [{"role": "user", "content": "refactor this"}],
"stream": True,
},
)
assert resp.status_code == 200
assert len(attempts) == 2, "expected a failover onto the runner-up"
row = _latency_row(db_path)
# 3.0, not 33.0: the 30 s the refusing candidate burned belongs to it.
assert row["router_wall_seconds"] == 3.0
assert row["router_ttft_seconds"] == 1.0
# --- plausibility, against the real clock ----------------------------------
def test_against_the_real_clock_the_value_is_bounded(tmp_path, monkeypatch):
"""No fake clock: a positive, small, finite number reaches the column.
A mocked clock proves the arithmetic and nothing about the wiring. This one
proves a real `time.monotonic()` difference survives the round trip through
log_observation and SQLite as a sane REAL.
"""
db_path = tmp_path / "real.db"
conn = sqlite3.connect(db_path)
conn.executescript(SCHEMA_SQL)
conn.execute(
"""
INSERT INTO models (
model_id, provider, base_model_id, tier, context_window,
effective_context_window, max_output_tokens,
cost_per_1m_prompt, cost_per_1m_completion,
supports_vision, supports_json_mode,
latency_class, reasoning_mode, context_variant,
access_level, availability, last_updated
) VALUES (?, 'neuralwatt', ?, 2, 262128, 192500, 16384, 0.10, 0.30,
1, 1, 'standard', 'default', 'full', 'public', 'active',
'2026-08-22T00:00:00+00:00')
""",
(MODEL, MODEL),
)
conn.commit()
conn.close()
monkeypatch.setattr(dispatcher.cfg.database, "path", str(db_path))
monkeypatch.setattr(dispatcher.cfg.verification, "local_llm_enabled", False)
monkeypatch.setattr(dispatcher.cfg.local_vision, "enabled", False)
monkeypatch.setattr(dispatcher.cfg.session_cache, "enabled", False)
monkeypatch.setattr(dispatcher.cfg.exploration, "epsilon", 0.0)
monkeypatch.setenv("NEURALWATT_API_KEY", "test-key")
def fake_post(url, headers=None, json=None, stream=False, timeout=None):
return _FakeResponse(_completion(json["model"]))
monkeypatch.setattr(dispatcher.requests, "post", fake_post)
client = TestClient(dispatcher.app)
resp = client.post(
"/v1/chat/completions",
json={"model": MODEL, "messages": [{"role": "user", "content": "hi"}]},
)
assert resp.status_code == 200
row = _latency_row(db_path)
wall = row["router_wall_seconds"]
assert wall is not None
# Not zero (the mark must not be taken after the call) and not absurd
# (an in-process fake cannot take a minute).
assert 0.0 < wall < 60.0
# --- the migration ---------------------------------------------------------
def test_alter_adds_both_columns_to_a_database_that_predates_them(tmp_path):
"""The live router.db is exactly this shape until its next restart."""
db_path = tmp_path / "old.db"
conn = sqlite3.connect(db_path)
conn.execute(
"""
CREATE TABLE energy_observations (
id INTEGER PRIMARY KEY AUTOINCREMENT,
model_id TEXT NOT NULL,
provider TEXT NOT NULL,
observed_at TEXT NOT NULL
)
"""
)
conn.commit()
dispatcher._ensure_energy_observations_table(conn)
cols = {r[1] for r in conn.execute("PRAGMA table_info(energy_observations)")}
assert "router_wall_seconds" in cols
assert "router_ttft_seconds" in cols
# Idempotent: a second start-up must not raise "duplicate column name".
dispatcher._ensure_energy_observations_table(conn)
again = {r[1] for r in conn.execute("PRAGMA table_info(energy_observations)")}
assert again == cols
conn.close()
def test_existing_rows_are_left_null_not_backfilled(tmp_path):
"""Nothing in an old row records how long the router waited.
A substituted 0 would read as an instantaneous response, which is a
measurement nobody made -- the same reasoning af18009 applied to
cached_tokens_source.
"""
db_path = tmp_path / "old.db"
conn = sqlite3.connect(db_path)
conn.execute(
"""
CREATE TABLE energy_observations (
id INTEGER PRIMARY KEY AUTOINCREMENT,
model_id TEXT NOT NULL,
provider TEXT NOT NULL,
duration_seconds REAL,
observed_at TEXT NOT NULL
)
"""
)
conn.execute(
"INSERT INTO energy_observations (model_id, provider, duration_seconds, "
"observed_at) VALUES ('m', 'neuralwatt', 4.2, '2026-09-01T00:00:00+00:00')"
)
conn.commit()
dispatcher._ensure_energy_observations_table(conn)
row = conn.execute(
"SELECT duration_seconds, router_wall_seconds, router_ttft_seconds "
"FROM energy_observations"
).fetchone()
conn.close()
# The provider's own figure is untouched, and the router's is NULL.
assert row == (4.2, None, None)
def test_log_observation_defaults_leave_the_columns_null(tmp_path, monkeypatch):
"""The harness call sites (seed_energy, eval_proficiency) pass no timings.
They must keep working unchanged: 6e729ad added trailing keyword-only
arguments to this function without updating seed_energy.py and killed the
seed timer for weeks. Both new parameters default to None for that reason.
"""
db_path = tmp_path / "defaults.db"
conn = sqlite3.connect(db_path)
conn.executescript(SCHEMA_SQL)
conn.commit()
conn.close()
monkeypatch.setattr(dispatcher.cfg.database, "path", str(db_path))
dispatcher.log_observation(
MODEL, "neuralwatt", "coding_general", "chatcmpl-x",
prompt_tokens=10, completion_tokens=5,
telemetry=dispatcher.Telemetry(),
)
row = _latency_row(db_path)
assert row["router_wall_seconds"] is None
assert row["router_ttft_seconds"] is None
def test_the_two_quantities_are_stored_separately(tmp_path, monkeypatch):
"""A router-side figure must never land in the provider's column.
duration_seconds is the one clean serving-time measurement this project
has; overwriting it with wall clock would destroy it silently, on the
provider that actually reports it.
"""
db_path = tmp_path / "separate.db"
conn = sqlite3.connect(db_path)
conn.executescript(SCHEMA_SQL)
conn.commit()
conn.close()
monkeypatch.setattr(dispatcher.cfg.database, "path", str(db_path))
dispatcher.log_observation(
MODEL, "neuralwatt", "coding_general", "chatcmpl-y",
prompt_tokens=10, completion_tokens=5,
telemetry=dispatcher.Telemetry(duration_seconds=1.25),
router_wall_seconds=3.75,
router_ttft_seconds=0.5,
)
row = _latency_row(db_path)
assert row["duration_seconds"] == 1.25
assert row["router_wall_seconds"] == 3.75
assert row["router_ttft_seconds"] == 0.5