Skip to content
10 changes: 10 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -386,6 +386,16 @@ exist on the post.

## CI gates

Do not classify an upstream TimeoutError as proof that a local deadline expired.
Use the owning timer's expiration state and keep shutdown cancellation distinct
from failure settlement. Test upstream failure, actual timer expiry, and shutdown
independently; exception names alone do not identify the terminating boundary.

For optional model timeouts, verify omission, explicit null, and explicit seconds
through factories and transport. Dropping a null keyword can silently restore a
downstream default; keep null intact and distinguish remaining worker limits
from a client's default. Null transport timeouts do not prove socket cancellation.

`.github/workflows/tests.yml` runs the full suite on every PR to `main`.
Do not weaken, skip, or `continue-on-error` a failing check -- fix the
underlying cause or, for a genuine false positive in a third-party scanner,
Expand Down
15 changes: 9 additions & 6 deletions backend/app/global_ask_queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -537,6 +537,7 @@ async def process_global_ask_job(
)
if row is None:
return
answer_timeout: asyncio.Timeout | None = None
try:
async with pool.acquire() as conn:
(
Expand All @@ -554,8 +555,8 @@ async def process_global_ask_job(
raise _SafeJobError(
"Ask Agent is unavailable: set ORCHESTRATOR_BASE_URL / ORCHESTRATOR_API_KEY"
)
payload = await asyncio.wait_for(
compute_global_ask_answer(
async with asyncio.timeout(JOB_DEADLINE_SECONDS) as answer_timeout:
payload = await compute_global_ask_answer(
pool,
question_text=str(row["question_text"]),
corporate_entity_ids=entity_ids,
Expand All @@ -567,9 +568,7 @@ async def process_global_ask_job(
verify_external=bool(row["verify_external_requested"]),
claim_verification_client=claim_verification_factory(),
knowledge_cutoff=row["knowledge_cutoff"],
),
timeout=JOB_DEADLINE_SECONDS,
)
)
except asyncio.CancelledError:
# Shutdown: leave the row `running`; the recovery sweep re-queues
# it after the orphan window on the next process start.
Expand All @@ -587,7 +586,11 @@ async def process_global_ask_job(
# Raised locally with a pre-authored, safe message (permission
# state / missing config) — never a provider-boundary leak.
detail = str(exc)
elif isinstance(exc, asyncio.TimeoutError):
elif (
isinstance(exc, asyncio.TimeoutError)
and answer_timeout is not None
and answer_timeout.expired()
):
detail = f"job exceeded the {JOB_DEADLINE_SECONDS}s deadline"
else:
# Provider responses/exceptions can carry credentials, gateway
Expand Down
13 changes: 5 additions & 8 deletions backend/app/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -323,8 +323,8 @@ async def lifespan(app: FastAPI):
app.state.post_content_worker = content_worker
# Late-bound lambda so tests that monkeypatch _post_chat_client reach
# the worker too (the name resolves in module globals at call time).
# Only this worker gets the long answer timeout; the per-post chat
# endpoint keeps the client's interactive default.
# This worker still has an explicit answer socket limit; per-post
# chat retains the default null transport timeout.
global_ask_worker = asyncio.create_task(
run_global_ask_worker(
valkey,
Expand Down Expand Up @@ -508,19 +508,16 @@ def _post_structure_client():
def _post_chat_client(timeout: float | None = None):
"""Live orchestrator client when configured; otherwise the unavailable null.

``timeout`` overrides the client's socket timeout. Only the Ask worker
passes the long answer timeout — the synchronous per-post chat endpoint
keeps the client default so an interactive request never hangs a reader
for the worker's full budget.
Preserve an omitted or explicit null timeout through the transport boundary.
The Ask worker still supplies its separately configured answer socket limit.
"""
settings = load_settings()
if not (settings.orchestrator_base_url and settings.orchestrator_api_key):
return NullPostChatClient()
kwargs = {} if timeout is None else {"timeout": timeout}
return ContextualOrchestratorPostChatClient(
base_url=settings.orchestrator_base_url,
api_key=settings.orchestrator_api_key,
**kwargs,
timeout=timeout,
)


Expand Down
22 changes: 22 additions & 0 deletions docs/adr/0083-orchestrator-runtime-commit-pin.md
Original file line number Diff line number Diff line change
Expand Up @@ -47,3 +47,25 @@ The runtime contract is:
- Rebuilding the image is required after the upstream pin changes.
- Protected-branch review and merge remain external gates; this pin does not
bypass upstream review.

## Proposed amendment: post-chat transport timeout (2026-09-07)

Status: Proposed; the Accepted runtime pin decision above is unchanged.

The post-chat client silently supplies 180 seconds when a caller omits a limit.
The factory also drops an explicit null, restoring that limit. This contradicts
the requested default-null model lifetime even when the upstream owner has no
implicit limit. Increasing the constant merely postpones the same failure; a
second per-model policy store would duplicate contextual-orchestrator.

Use null as the post-chat transport default and pass it unchanged through the
factory and shared HTTP transport. Preserve explicit caller limits while their
separate migration is pending. This avoids client abandonment by default but can
leave a synchronous chat waiting until transport/provider termination. Do not
claim cancellation of a blocking socket merely because an async task is cancelled.

Confirm omitted/null/explicit values at the client and factory boundaries. The
Ask worker's explicit 570-second setting, 600-second execution deadline, and
age-based recovery remain unresolved. Other model clients and upstream model
administration require separate owner-aligned verification. No runtime pin is
changed and no open upstream PR becomes a released contract.
37 changes: 37 additions & 0 deletions docs/product-technical-gap-baseline.md
Original file line number Diff line number Diff line change
Expand Up @@ -932,3 +932,40 @@ The ONET rows stacked into base branches (#743/#745/#746/#740/#732) reached
`main` together through the #759 promotion; their per-base merge records are
historical evidence only. The job-architecture artifact ship originally via
#749 is now re-verified on `main` from the promotion.

### Ask timeout attribution repair (2026-09-07)

The queue worker treated every `asyncio.TimeoutError` as proof its 600-second
execution deadline expired. A provider that terminated immediately therefore
produced an incorrect durable deadline explanation. A paired synthetic regression
reproduced that mismatch (provider case failed; actual zero-duration worker timer
passed). The worker now uses the standard asyncio timeout context's expiration
state to attribute only its own expiry. Other failures retain ADR 0123's existing
bounded unavailable message; provider exception content is not persisted.

The focused queue/service suite passed 22 tests in 1.59 s, including provider
failure, actual timer expiry, and shutdown cancellation with no failed settlement.
Compilation and diff checks passed. This does not remove the execution deadline,
change model policy, or repair age-based orphan recovery: default-null execution
still requires worker liveness and claim fencing. No database migration, new
container, provider call, deployment, or protected merge was performed.

Python Software Foundation. (2026). *Coroutines and tasks: Timeouts*.
https://docs.python.org/3/library/asyncio-task.html#timeouts

### Post-chat null timeout propagation (2026-09-07; proposed ADR 0083 amendment)

Three synthetic assertions reproduced an implicit 180-second limit: direct
construction with no timeout, factory construction with no timeout, and factory
construction with explicit null. The post-chat client now defaults to null and
the factory passes the value unchanged. The shared HTTP request and JSON POST
annotations accept the native transport's null timeout without a new adapter.
Explicit caller seconds remain intact. The runtime pin is unchanged.

Focused post-chat, HTTP, queue, and service tests passed 69 cases with two
real-provider cases skipped in 13.84 s. The existing local HTTP server test covers
both null and numeric limits. Compilation and diff checks passed. These results
do not prove blocking-socket cancellation, upstream policy enforcement, or
unlimited Ask execution: the explicit 570-second Ask socket setting, 600-second
worker deadline, and age-based recovery remain unresolved. The policy amendment
is Proposed, not a protected acceptance or release claim.
16 changes: 12 additions & 4 deletions lineageweave/http_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,9 +23,18 @@
from .llm_context import current_llm_metadata
from .observability import current_session_id, inject_trace_context, traced


def _build_ssl_context() -> ssl.SSLContext:
"""Build the owned client context with its protocol floor before it escapes."""

context = ssl.create_default_context(cafile=certifi.where())
context.minimum_version = ssl.TLSVersion.TLSv1_2
return context


# Some interpreter distributions don't reliably inherit the OS trust store.
# Pointing at certifi keeps full chain validation without weakening TLS.
_SSL_CONTEXT = ssl.create_default_context(cafile=certifi.where())
_SSL_CONTEXT = _build_ssl_context()
_ALLOWED_SCHEMES = frozenset({"http", "https"})
_SESSION_HEADER_PEERS = frozenset({"contextual-orchestrator", "tepp"})

Expand Down Expand Up @@ -148,7 +157,7 @@ def _request(
*,
body: bytes | None,
headers: dict[str, str],
timeout: float,
timeout: float | None,
maximum_response_bytes: int | None = None,
expected_response_media_type: str | None = None,
) -> tuple[int, bytes]:
Expand Down Expand Up @@ -218,7 +227,6 @@ def _request(

def _decode_json(raw: bytes, hostname: str) -> object:
"""Decode UTF-8 JSON without exposing response content in errors."""

try:
return json.loads(raw.decode("utf-8"))
except (UnicodeDecodeError, json.JSONDecodeError) as exc:
Expand Down Expand Up @@ -248,7 +256,7 @@ def post_json(
payload: dict,
*,
headers: dict[str, str],
timeout: float,
timeout: float | None,
service_peer_name: str = "contextual-orchestrator",
) -> dict:
"""POST ``payload`` as JSON to ``url`` and return the decoded object.
Expand Down
6 changes: 2 additions & 4 deletions lineageweave/post_chat.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,6 @@
CANONICAL_INVOLVED_QUESTION = "Who is involved?"
CANONICAL_COMMITMENT_QUESTION = "What is the next commitment?"

_TRAILING_PUNCT = re.compile(r"[?.!\s]+$")
_CANONICAL_QUESTION_NORM = "what happened between these events"
_INVOLVED_QUESTION_NORM = "who is involved"
_COMMITMENT_QUESTION_NORM = "what is the next commitment"
Expand All @@ -45,7 +44,7 @@ def normalize_chat_question(question: str) -> str:
the second seeded chip that names Keymen. ``What's the next
commitment?`` aliases the third chip that names the Calendar ticket.
"""
folded = _TRAILING_PUNCT.sub("", " ".join(question.strip().lower().split()))
folded = " ".join(question.strip().lower().split()).rstrip("?.!")
if folded == "what happened":
return _CANONICAL_QUESTION_NORM
if folded in {"who's involved", "who is involved here"}:
Expand Down Expand Up @@ -181,7 +180,6 @@ def _buyer_evidence_kind(fact: str) -> str:
return "source_field"



def _buyer_evidence_text(fact: str) -> str:
cleaned = re.sub(r"\s*\|\s*(?:ontology_iri|extraction_method|confidence):\s*[^|\[]+", "", fact)
cleaned = re.sub(r"\s*\[provenance=[^]]+\]", "", cleaned)
Expand Down Expand Up @@ -377,7 +375,7 @@ class ContextualOrchestratorPostChatClient:
available = True

def __init__(
self, base_url: str, api_key: str, *, reasoning_effort: str = "auto", timeout: float = 180.0
self, base_url: str, api_key: str, *, reasoning_effort: str = "auto", timeout: float | None = None
) -> None:
self._base_url = base_url.rstrip("/")
self._api_key = api_key
Expand Down
36 changes: 29 additions & 7 deletions tests/test_global_ask_queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@
from contextlib import asynccontextmanager
from datetime import UTC, datetime

import pytest

from backend.app import global_ask_queue
from backend.app.global_ask_queue import load_job_visibility
from lineageweave import claim_verification as cv
Expand Down Expand Up @@ -435,25 +437,39 @@ async def _fake_compute_global_ask_answer(*_args, **_kwargs):
assert settle_args[-1] == "account lacks the post_read permission"


def test_job_deadline_timeout_settles_with_a_specific_but_still_generic_detail(
monkeypatch,
@pytest.mark.parametrize("timeout_source", ["provider", "worker", "shutdown"])
def test_timeout_detail_identifies_only_an_expired_worker_deadline(
monkeypatch, timeout_source,
) -> None:
"""A bare `asyncio.TimeoutError` (no message) still gets a useful,
non-empty detail rather than an empty string."""
"""Provider timeout is not proof that the worker deadline expired."""
connection = _Connection(_queued_row())
pool = _Pool(connection)

async def _fake_load_job_visibility(_conn, _job_id, _account_id):
return {"corp-1"}, set(), False, True

async def _fake_compute_global_ask_answer(*_args, **_kwargs):
raise asyncio.TimeoutError()

if timeout_source == "shutdown":
raise asyncio.CancelledError()
if timeout_source == "worker":
await asyncio.Event().wait()
raise asyncio.TimeoutError("synthetic private upstream detail")

if timeout_source == "worker":
monkeypatch.setattr(global_ask_queue, "JOB_DEADLINE_SECONDS", 0)
monkeypatch.setattr(global_ask_queue, "load_job_visibility", _fake_load_job_visibility)
monkeypatch.setattr(
global_ask_queue, "compute_global_ask_answer", _fake_compute_global_ask_answer
)

if timeout_source == "shutdown":
with pytest.raises(asyncio.CancelledError):
asyncio.run(global_ask_queue.process_global_ask_job(
pool, job_id="job-1", chat_factory=_AvailableClient,
))
assert connection.executed == []
return

asyncio.run(
global_ask_queue.process_global_ask_job(
pool,
Expand All @@ -463,7 +479,13 @@ async def _fake_compute_global_ask_answer(*_args, **_kwargs):
)

_settle_query, settle_args = connection.executed[-1]
assert settle_args[-1] == f"job exceeded the {global_ask_queue.JOB_DEADLINE_SECONDS}s deadline"
if timeout_source == "worker":
assert settle_args[-1] == "job exceeded the 0s deadline"
else:
assert settle_args[-1] == (
"Ask Agent is unavailable: contextual-orchestrator returned "
"no complete evidence object"
)


def test_job_visibility_never_expands_past_queued_scope() -> None:
Expand Down
5 changes: 3 additions & 2 deletions tests/test_http_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -130,15 +130,16 @@ def test_post_json_refuses_missing_hostname() -> None:
post_json("https:///v1/embeddings", {}, headers={}, timeout=1.0)


def test_post_json_posts_json_to_http_endpoint() -> None:
@pytest.mark.parametrize("request_timeout", [None, 2.0])
def test_post_json_posts_json_to_http_endpoint(request_timeout) -> None:
_JsonHandler.received = {}
server, base = _serve(_JsonHandler)
try:
body = post_json(
f"{base}/v1/embeddings",
{"model": "demo", "input": "hello"},
headers={"authorization": "Bearer test-token"},
timeout=2.0,
timeout=request_timeout,
)
finally:
server.shutdown()
Expand Down
37 changes: 37 additions & 0 deletions tests/test_http_client_tls_floor_contract.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
from __future__ import annotations

import ssl
from pathlib import Path
from types import SimpleNamespace
from unittest.mock import patch

from lineageweave import http_client
from lineageweave.post_chat import CANONICAL_CHAT_QUESTION, normalize_chat_question


def test_http_client_declares_tls_1_2_floor_in_owned_transport_boundary() -> None:
"""The owned HTTPS boundary must set its protocol floor before escape."""

context = SimpleNamespace(minimum_version=ssl.TLSVersion.MINIMUM_SUPPORTED)
with patch.object(http_client.ssl, "create_default_context", return_value=context):
built_context = http_client._build_ssl_context()

assert built_context is context
assert context.minimum_version == ssl.TLSVersion.TLSv1_2
assert http_client._SSL_CONTEXT.minimum_version >= ssl.TLSVersion.TLSv1_2


def test_chat_question_trailing_strip_avoids_end_anchored_regex_backtracking() -> None:
"""User punctuation normalization must use the linear string boundary."""

source = (
Path(__file__).resolve().parents[1] / "lineageweave" / "post_chat.py"
).read_text(encoding="utf-8")

assert "_TRAILING_PUNCT = re.compile" not in source
assert '.rstrip("?.!")' in source
canonical = normalize_chat_question(CANONICAL_CHAT_QUESTION)
assert normalize_chat_question("What happened" + "?" * 10_000) == canonical
assert normalize_chat_question("When was the bid sent" + "!" * 10_000) == (
"when was the bid sent"
)
25 changes: 25 additions & 0 deletions tests/test_post_chat.py
Original file line number Diff line number Diff line change
Expand Up @@ -469,3 +469,28 @@ def fake_post_json(url, payload, *, headers, timeout):
assert observed["payload"]["reasoning_effort"] == "auto"
assert observed["payload"]["mode"] == "auto"
assert "CITED SOURCES" in observed["payload"]["messages"][0]["content"]


@pytest.mark.parametrize("creation_path", ["direct", "factory"])
@pytest.mark.parametrize("timeout_options, expected", [({}, None), ({"timeout": None}, None), ({"timeout": 7.5}, 7.5)])
def test_post_chat_preserves_optional_transport_timeout(monkeypatch, timeout_options, expected, creation_path) -> None:
"""Omitted or null limits stay null; explicit limits reach transport unchanged."""
observed = {}

def fake_post_json(_url, _payload, *, headers, timeout):
observed["timeout"] = timeout
return {"choices": [{"message": {"content": "Answer\nCITED SOURCES: 1"}}]}

monkeypatch.setattr("lineageweave.post_chat.post_json", fake_post_json)
if creation_path == "factory":
from types import SimpleNamespace
from backend.app import main

monkeypatch.setattr(main, "load_settings", lambda: SimpleNamespace(
orchestrator_base_url="https://orchestrator.test", orchestrator_api_key="synthetic-token",
))
client = main._post_chat_client(**timeout_options)
else:
client = ContextualOrchestratorPostChatClient("https://orchestrator.test", "synthetic-token", **timeout_options)
client.answer("Question", _SOURCES)
assert observed["timeout"] == expected
Loading