Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,7 @@ addopts = "-q --strict-markers --strict-config"
markers = [
"requires_embeddings: needs an installed embeddings bundle (see ./bin/embeddings_manager)",
"requires_retrieval_stack: needs the full langchain/chromadb dependency set installed",
"requires_live_model: calls a real LLM and needs OPENAI_API_KEY plus embeddings",
]

[tool.mypy]
Expand Down
22 changes: 20 additions & 2 deletions specs/010-search-page-answers/spec.md
Original file line number Diff line number Diff line change
Expand Up @@ -200,11 +200,29 @@ classifier's decision matches, and that no LLM answer call happens for the latte
is the goal; a number being met is not, until FR-005a's blocker is removed
- **SC-002**: Zero model calls for requests without a valid token, measured by
counting calls under a load of unauthenticated requests
- **SC-003**: The answer sweep stays green: the endpoint and the chat UI give the
same answer to the same question, because they share a graph
- **SC-003**: The answer sweep stays green, and the two surfaces stay
*configured* the same -- same profile, same shared graph, no difference that can
reach the answer. Textual equality is explicitly **not** the criterion, because
it is not achievable: measured 2026-09-18 on one graph at temperature 0, the
same question asked twice through the *same* surface produced answers 0.331
similar, while endpoint-versus-chat scored 0.356. The surfaces differ from each
other no more than either differs from itself. The sweep is the right mechanism
precisely because it matches patterns rather than literals
- **SC-004**: No search-page request can make the search page itself slower or fail;
verified by taking the service down and confirming the page still renders

## Known: retrieval is not reproducible

Measured 2026-09-18, the same question asked three times returned **12 citations
each time but only 4 pathways common to all three**, a union of 19 and a Jaccard
of 0.26 between two runs. Query expansion is itself a model call, so each run
expands the question differently and retrieves different documents.

This matters beyond wording. A reader who reloads the panel sees different
sources, and FR-007's cache invalidation assumes an answer is a stable artifact
of a release. Deciding what to do about it -- seeding or caching the expansion,
or dropping it for this path -- is open, and is not a blocker for the handover.

## Decisions

### D1 -- what proves a person is human?
Expand Down
6 changes: 4 additions & 2 deletions specs/010-search-page-answers/tasks.md
Original file line number Diff line number Diff line change
Expand Up @@ -34,15 +34,17 @@ state.
- [x] T011 [US1] Emit citations from retrieved documents' `st_id` metadata, deduplicated — never by parsing anchors out of the model's prose
- [x] T012 [US1] Mount the router in bin/chat-fastapi.py and let the captcha middleware pass /chat/api/ through, since the endpoint verifies its own caller
- [x] T013 [US1] Test over HTTP with a real client in tests/api/test_answer_endpoint.py, not by calling the handler — mounting order and middleware only interact in the served path (Principle I)
- [ ] T014 [US1] Assert the answer matches the chat UI's for the same question (SC-003); two surfaces that can disagree is a defect
- [x] T014 [US1] SC-003, restated against the measurement: pin that the two surfaces are *configured* the same in tests/api/test_answer_matches_chat.py. Answer equality is not assertable -- the same surface asked twice scores 0.331 similarity, endpoint-vs-chat 0.356 -- so the spec's criterion was corrected rather than the test weakened (PR #237)
- [x] T014a [US1] Give each request its own checkpointer thread in src/api/answer.py; `id(body)` put 192 of 200 requests on a shared thread, and `chat_history` is checkpointed state the rephraser reads (PR #236)
- [x] T014b [US1] Send `release` on start and `seconds` on done per contracts/answer_endpoint.md; both were promised to the website and neither was implemented (PR #236)

## Phase 4: User Story 2 — no answer without a verified person (P1)

- [x] T015 [US2] Refuse missing, expired, malformed and wrongly-signed tokens before any model call, in src/api/answer.py
- [x] T016 [US2] Test that no model call happens for a refused request in tests/api/test_answer_endpoint.py, by asserting on a patched graph rather than on timing (SC-002)
- [ ] T017 [P] [US2] Rate limit per token as a backstop; the budget is the website's, enforced before the call reaches here (FR-008)
- [x] T017 [P] [US2] Rate limit per token as a backstop in src/util/rate_limit.py; the budget is the website's, enforced before the call reaches here (FR-008). 30 per 10 minutes, keyed on `sub`/`jti` when D1 provides one and a token hash until then (PR #237)
- [x] T017a [US2] Stop paying for a discarded web search: the endpoint took `enable_postprocess` at its default, so every answer ran a Tavily search that `astream_answer` has no event to return (PR #237)
- [ ] T020a Decide what to do about non-reproducible retrieval: three runs of one question shared only 4 of 19 citations (Jaccard 0.26) because query expansion is itself a model call. Affects what FR-007 can cache
- [x] T018 [US2] Return `state: failed` with no partial answer on any internal error, so the page renders no panel (FR-006)
- [x] T011a [US1] Strip inline HTML anchors from the token stream in src/util/anchor_strip.py; the contract promises prose without them and the chat prompt emits them, split across ~20 fragments (PR #236)
- [x] T013a [US1] Run the endpoint end to end against a real graph: release 97, answered in 19.5-44.2s, 12 citations, anchors 0 (PR #236)
Expand Down
24 changes: 22 additions & 2 deletions src/api/answer.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,13 +32,18 @@
from util.anchor_strip import AnchorStripper
from util.human_token import TokenRejectedError, verify
from util.logging import logging
from util.rate_limit import identity_of, limiter_from_env

logger = logging.getLogger(__name__)

router = APIRouter()

PROFILE = "react-to-me"

# One limiter for the process, built at import so the window is not reset by a
# request. FR-008: a backstop behind the website's own budget.
_limiter = limiter_from_env()

# FR-006 names timeout alongside error, and nothing here implemented it. The only
# bound was the LLM client's `request_timeout=360.0` -- six minutes per model call,
# and six calls run around one answer, so a pathological request could hold a
Expand Down Expand Up @@ -85,10 +90,16 @@ async def answer(request: Request, body: AnswerRequest) -> StreamingResponse:
return _refusal("no verifying key on the app")

try:
verify(body.human_token, verifying_key)
claims = verify(body.human_token, verifying_key)
except TokenRejectedError as rejected:
return _refusal(rejected.reason)

# After verification, so an unsigned token cannot consume someone else's
# budget by claiming their `sub`, and before the graph, so a caller over the
# limit costs nothing.
if not _limiter.allow(identity_of(claims, body.human_token)):
return _refusal("rate limited")

graph = get_graph()

# A fresh thread per request, never `id(body)`. `chat_history` is checkpointed
Expand All @@ -115,7 +126,16 @@ async def stream() -> AsyncIterator[str]:
try:
async with asyncio.timeout(ANSWER_TIMEOUT_SECONDS):
async for event in graph.astream_answer(
body.question, PROFILE, thread_id=f"search-{uuid.uuid4()}"
body.question,
PROFILE,
thread_id=f"search-{uuid.uuid4()}",
# The postprocess node runs a Tavily web search after the
# answer, and `astream_answer` has no event to carry the
# result -- so on this path it was paid for and discarded,
# delaying `done` by the length of a web search. The chat UI
# renders those results; the search page has no place for
# them.
enable_postprocess=False,
):
if event.kind == "token":
text = stripper.feed(event.text)
Expand Down
99 changes: 99 additions & 0 deletions src/util/rate_limit.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,99 @@
"""A per-caller request limit for the answer endpoint (FR-008).

A backstop, not the budget. The website enforces the real one before a call
reaches here; this exists so a leaked or shared token cannot run up an unbounded
bill against a service whose every answer costs six model calls.

In process and in memory, because there is one process serving this. If the
service is ever scaled out, a shared store has to replace this, and the limit
becomes per instance until it is.
"""

import hashlib
import os
import time
from collections import deque


def _positive_int(name: str, default: int) -> int:
"""Configuration that is absent, empty or nonsense falls back to the default."""
raw = os.getenv(name, "")
if not raw.strip():
return default
try:
value = int(raw)
except ValueError:
return default
return value if value > 0 else default


def identity_of(claims: dict[str, object], token: str) -> str:
"""Who to count against.

The token's claims are D1 and not yet settled with the website, so `sub` may
never arrive. `sub` then `jti` are used when present, so this starts keying on
a real person the moment D1 lands; until then a hash of the token itself is
the best available proxy -- one issuance, short lived, one person.

Hashed, never raw: this lands in a dict that lives as long as the process, and
a bearer token is a credential.
"""
for claim in ("sub", "jti"):
value = claims.get(claim)
if isinstance(value, str) and value:
return f"{claim}:{value}"
return "token:" + hashlib.sha256(token.encode()).hexdigest()[:32]


class SlidingWindowLimiter:
"""Allow `limit` requests per `window` seconds, per key.

No lock. Every mutation happens between awaits on one event loop, so a
request cannot be interleaved mid-update. Adding an await inside `allow`
would break that, which is why it does no I/O.
"""

def __init__(self, limit: int, window: float) -> None:
self.limit = limit
self.window = window
self._hits: dict[str, deque[float]] = {}
self._last_sweep = 0.0

def allow(self, key: str) -> bool:
now = time.monotonic()
self._sweep(now)
hits = self._hits.setdefault(key, deque())
cutoff = now - self.window
while hits and hits[0] <= cutoff:
hits.popleft()
if len(hits) >= self.limit:
return False
hits.append(now)
return True

def _sweep(self, now: float) -> None:
"""Drop keys with nothing left in the window.

Without this the dict grows with every distinct token forever, which on a
search page is every visitor. Swept once per window rather than per
request, so the cost is amortised.
"""
if now - self._last_sweep < self.window:
return
self._last_sweep = now
cutoff = now - self.window
self._hits = {
key: hits for key, hits in self._hits.items() if hits and hits[-1] > cutoff
}


def limiter_from_env() -> SlidingWindowLimiter:
"""30 requests per 10 minutes by default.

Generous for a person -- an answer takes 20-40 seconds, so thirty is far more
than anyone reads -- and it still caps a leaked token at 180 an hour.
"""
return SlidingWindowLimiter(
limit=_positive_int("ANSWER_RATE_LIMIT", 30),
window=float(_positive_int("ANSWER_RATE_WINDOW_SECONDS", 600)),
)
96 changes: 96 additions & 0 deletions tests/api/test_answer_endpoint.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@

from agent.graph import AnswerEvent
from api.answer import router
from util.rate_limit import SlidingWindowLimiter

PREFIX = "/chat/guest/api"

Expand Down Expand Up @@ -72,6 +73,19 @@ def keys() -> tuple[str, str]:
)


@pytest.fixture(autouse=True)
def _fresh_limiter(monkeypatch: pytest.MonkeyPatch) -> None:
"""A private limiter per test.

`_limiter` is module state shared by the whole process, so without this the
fifty requests below would exhaust the real limit and refuse later tests --
and which tests failed would depend on the order they ran in.
"""
monkeypatch.setattr(
"api.answer._limiter", SlidingWindowLimiter(limit=10_000, window=600.0)
)


@pytest.fixture
def stub(monkeypatch: pytest.MonkeyPatch) -> _StubGraph:
graph = _StubGraph()
Expand Down Expand Up @@ -308,3 +322,85 @@ def test_a_hanging_upstream_still_ends_the_stream(
events = dict(_events(response.text))
assert json.loads(events["done"])["state"] == "failed"
assert elapsed < 2, f"stream ran {elapsed:.1f}s; the timeout did not fire"


def test_a_caller_over_the_limit_is_refused_before_the_model(
keys: tuple[str, str], stub: _StubGraph, monkeypatch: pytest.MonkeyPatch
) -> None:
"""FR-008. A backstop: the website enforces the real budget upstream.

Refused like any other refusal -- one `done` shape, no HTTP error -- so the
page renders no panel rather than a broken one.
"""
private, public = keys
monkeypatch.setattr(
"api.answer._limiter", SlidingWindowLimiter(limit=2, window=600.0)
)
client = _client(public)
token = _token(private)

states = []
for _ in range(4):
response = client.post(
f"{PREFIX}/answer", json={"question": "what is CDK5", "human_token": token}
)
assert response.status_code == 200
states.append(json.loads(dict(_events(response.text))["done"])["state"])

assert states == ["answered", "answered", "refused", "refused"]
assert stub.calls == 2, "a refused request still reached the model"


def test_the_limit_is_per_caller(
keys: tuple[str, str], stub: _StubGraph, monkeypatch: pytest.MonkeyPatch
) -> None:
"""One visitor exhausting their budget must not silence the page for others."""
private, public = keys
monkeypatch.setattr(
"api.answer._limiter", SlidingWindowLimiter(limit=1, window=600.0)
)
client = _client(public)

first = _token(private)
second = _token(private, seconds=301) # A different token, so a different key.

def ask(token: str) -> str:
response = client.post(
f"{PREFIX}/answer", json={"question": "what is CDK5", "human_token": token}
)
return str(json.loads(dict(_events(response.text))["done"])["state"])

assert ask(first) == "answered"
assert ask(first) == "refused"
assert ask(second) == "answered"


def test_the_search_page_does_not_pay_for_a_web_search(
keys: tuple[str, str], monkeypatch: pytest.MonkeyPatch
) -> None:
"""The postprocess node runs a Tavily search whose result this path drops.

`astream_answer` has no event carrying `additional_content`, so with the
default the endpoint paid for a web search, discarded it, and delayed `done`
by its duration. The chat UI renders those results; the search page has no
place for them.
"""
private, public = keys
seen: dict[str, object] = {}

class _RecordingGraph(_StubGraph):
async def astream_answer(
self, *_a: Any, **kwargs: Any
) -> AsyncIterator[AnswerEvent]:
seen.update(kwargs)
for event in self._events:
yield event

monkeypatch.setattr("api.answer.get_graph", lambda: _RecordingGraph())
client = _client(public)
client.post(
f"{PREFIX}/answer",
json={"question": "what is CDK5", "human_token": _token(private)},
)

assert seen["enable_postprocess"] is False
Loading
Loading