"""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