Compare commits
10 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 85143b866c | |||
| 00854ce618 | |||
| 78bfcadb9e | |||
| 44138590ad | |||
| d516537b08 | |||
| 92aa05c688 | |||
| 209427ab23 | |||
| 139771c8d8 | |||
| 489cfee1f0 | |||
| 11ef6830ab |
+13
-5
@@ -7,11 +7,19 @@ documents the pin, the vendored artifacts, and the bump procedure.
|
||||
|
||||
| Field | Value |
|
||||
|---|---|
|
||||
| Worldtree git SHA | `55101e909abcd2219833266b6f905c5bc956e0f0` |
|
||||
| Worldtree HEAD message | `memory: snapshot — #177 Vili v1 + persona async-decouple shipped as v0.19.0` |
|
||||
| Pinned on | 2026-05-20 |
|
||||
| Pinned by | brokkr-smithy-dev (initial scaffold) |
|
||||
| Worldtree version at pin | `v0.19.0` |
|
||||
| Worldtree git SHA | `562001af28d752c3a60d449c7ddd09f44fa9dc9a` |
|
||||
| Worldtree HEAD message | `feat(#201): v0.29.0 — awaiting_llm_first_token SSE heartbeat` |
|
||||
| Pinned on | 2026-05-26 |
|
||||
| Pinned by | ratatoskr-dev (bump for #201 awaiting_llm_first_token SSE) |
|
||||
| Worldtree version at pin | `v0.29.0` |
|
||||
|
||||
## Pin history
|
||||
|
||||
| Date | SHA | Version | Notable deltas consumed |
|
||||
|---|---|---|---|
|
||||
| 2026-05-26 | `562001a` | v0.29.0 | #201 — new SSE event `awaiting_llm_first_token` (heartbeat during BuildingPrompt → CallingLLM gap, default 5s interval) |
|
||||
| 2026-05-25 | `da93ca7` | v0.28.0 | #204 — new SSE event `affect_update` (current/scheduled), new endpoint `GET /agents/{id}/persona_state`, auth-model doc edits |
|
||||
| 2026-05-20 | `55101e9` | v0.19.0 | initial scaffold pin |
|
||||
|
||||
## Vendored artifacts
|
||||
|
||||
|
||||
@@ -55,6 +55,69 @@ conversation_api:
|
||||
|
||||
---
|
||||
|
||||
## Authorization model — agent invocation
|
||||
|
||||
When you call `POST /sessions` against an agent, the authorization check that fires depends on **which kind of agent** you target. There are two distinct scope namespaces — the spelling differs by one character (`agent` vs `agents`) and the granting mechanism differs entirely. Confusing the two is a common source of bug reports.
|
||||
|
||||
### Tier 1 — foundational agents (no `:` in agent_id)
|
||||
|
||||
Agents bundled with Worldtree: `mimir`, `lofn`, `soong`, `forseti`, `domari`, `vili`, `actor`, `saga`, `bragi`, `leif`, `troi`, `cara`, `glados`, and any future Asgardian. The agent_id is a simple slug like `mimir` — no colon.
|
||||
|
||||
> **About tiers:** Your `tier` is set on the `users` table row your API key resolves to, assigned at key-mint time (see `POST /admin/keys`). Tiers are `anonymous` (dev-mode unauthenticated), `user` (default for newly-issued keys), `free`/`pro` (subscription-shaped, not actively differentiated), and `admin`. The tier you have is visible via `GET /me`'s `tier` field. Tier-derived scopes come from `config/policies.yaml > tiers.<tier>.scopes` — there is no per-key scope override.
|
||||
|
||||
**Authorization rule (singular `agent`):**
|
||||
|
||||
```yaml
|
||||
- id: agent-call-baseline-allow
|
||||
principal:
|
||||
tiers: ["anonymous", "user", "free", "pro", "admin"]
|
||||
action: "agent.call:*"
|
||||
resource: "*"
|
||||
effect: allow
|
||||
```
|
||||
|
||||
This baseline rule lives at `config/policies.yaml`. Every authenticated tier — including the `user` tier that newly-issued keys default to — already passes this check for every Tier 1 agent. **There is no per-agent scope you can add to "grant" Tier 1 access; it's covered by tier.**
|
||||
|
||||
If you get a 422 calling a Tier 1 agent (e.g., `lofn` rejecting with `end_user_id_required`), that's a **request-body validation**, not a scope denial. Check the `error_code` in the response detail — `END_USER_ID_REQUIRED` means the agent requires an `end_user_id` field in the request body; `AUTH_SCOPE_DENIED` (403) would be the actual scope problem. They're not interchangeable.
|
||||
|
||||
### Tier 3 — consumer-defined agents (`:` in agent_id)
|
||||
|
||||
Agents created at runtime via `POST /agents/define`. The agent_id is `<owner_user_id>:<agent_name>`, e.g., `acme:support-bot`. The `:` in the path is the trigger that switches the auth model.
|
||||
|
||||
**Authorization is DB-backed per-resource, NOT policy-driven (plural `agents`):**
|
||||
|
||||
```
|
||||
scope action checked: agents.call:<owner_user_id>:<agent_name>
|
||||
^^^^^^
|
||||
PLURAL — different namespace from Tier 1
|
||||
```
|
||||
|
||||
There is **no blanket allow rule** for `agents.call:*` in policy. The grant comes from the live `consumer_agents` table:
|
||||
|
||||
- A non-soft-deleted row in `consumer_agents` owned by `ctx.user_id` IS the grant.
|
||||
- Cascade soft-delete and owner-initiated `DELETE` revoke it.
|
||||
- Missing row → policy defaults to deny (403 `auth_scope_denied`).
|
||||
|
||||
To "add the scope" for a Tier 3 agent, you don't amend any config or call an admin endpoint — you `POST /agents/define` to register it under your `user_id`. Owning the row IS the grant. You cannot call another user's Tier 3 agent; ownership is checked at session-create (`row.user_id == ctx.user_id`).
|
||||
|
||||
### Common pitfalls
|
||||
|
||||
- **Singular vs plural.** Tier 1 uses `agent.call:*` (singular `agent`). Tier 3 uses `agents.call:<owner>:<name>` (plural `agents`). One character difference, two completely different mechanisms. There is no Tier 1 scope named `agent.call:mimir` or `agents.call:mimir` — Tier 1 is granted by baseline rule, not per-agent name.
|
||||
- **No scope-mutation API.** `POST /admin/keys` accepts `{user_id, label, tier}` only. There is no per-key scope override mechanism in the storage schema. To change a user's effective scopes, change their `tier`, not their key. Per-resource Tier 3 grants flow through `POST /agents/define` (and its DELETE counterpart), not through admin endpoints.
|
||||
- **422 vs 403.** A 422 is body-validation (e.g., `end_user_id_required`); a 403 is auth-policy denial (`auth_scope_denied`). Different fix paths. Read the `error_code` in `detail`.
|
||||
|
||||
### Quick decision table for consumers
|
||||
|
||||
| Target | Auth requirement |
|
||||
|---|---|
|
||||
| Tier 1 agent (e.g., `mimir`) | Authenticated tier ≥ `user`. No additional body requirements |
|
||||
| Tier 1 agent `lofn` (the default welcoming intermediary) | Authenticated tier ≥ `user` + `end_user_id` field required in request body. 422 `END_USER_ID_REQUIRED` if absent |
|
||||
| Tier 3 agent (any agent_id containing `:`) | `end_user_id` field required in body. AND the row must be owner-matched: `POST /agents/define` first to create a row under your `user_id`, then session-create works against your existing key. Cross-user Tier 3 invocation is rejected with 403 |
|
||||
|
||||
> **Programmatic discovery of `end_user_id` requirements:** as of v0.22.x there is no field on `GET /agents` indicating which agents require `end_user_id` — the spec line above (lofn + Tier 3) is the authoritative list, and 422 `END_USER_ID_REQUIRED` is the fallback signal at request time. Adding a discoverable `requires_end_user_id` field on `AgentInfoResponse` is on the table as a small future capability; ping if you want to drive it.
|
||||
|
||||
---
|
||||
|
||||
## GET /me
|
||||
|
||||
Returns the authenticated principal's identity and key metadata. Lets a client verify its key on boot without triggering agent-config-loading side effects.
|
||||
@@ -1878,6 +1941,65 @@ Tool-using turns cycle through `CallingLLM → ProcessingTools → CallingLLM
|
||||
|
||||
Clients that don't need phase events can filter on `event["type"] != "worker_phase"` client-side. Existing SSE consumers that switch on `event["type"]` ignore this event type without code changes.
|
||||
|
||||
### affect_update
|
||||
|
||||
Persona-state observability event (issue #204). Fires twice per turn for agents with `persona.enabled: true` on non-ephemeral sessions; suppressed entirely for persona-disabled agents (e.g. `domari`, `muninn`), Tier 3 consumer-defined agents (Phase 2.0), and ephemeral sessions.
|
||||
|
||||
**Start-of-turn — `status: "current"`:**
|
||||
|
||||
Emitted immediately at the start of each qualifying turn, before any `worker_phase` event. Carries the agent's current persona snapshot reflecting all prior turns' completed appraisals.
|
||||
|
||||
```json
|
||||
{
|
||||
"type": "affect_update",
|
||||
"status": "current",
|
||||
"turn_id": 42,
|
||||
"snapshot": {
|
||||
"agent_id": "mimir",
|
||||
"pad": {"pleasure": 0.52, "arousal": 0.47, "dominance": 0.50},
|
||||
"dominant_emotion": "curiosity",
|
||||
"emotions_active": [
|
||||
{"type": "curiosity", "intensity": 0.6, "decay_remaining_s": 202.7}
|
||||
],
|
||||
"baseline_pad": {"pleasure": 0.50, "arousal": 0.40, "dominance": 0.50},
|
||||
"mood_drift": {"valence_delta": 0.02, "arousal_delta": 0.07},
|
||||
"last_updated_at": "2026-05-25T22:30:18+00:00"
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
**End-of-turn — `status: "scheduled"`:**
|
||||
|
||||
Emitted after the post-turn appraisal task has been scheduled (per #177 Phase A's fire-and-forget discipline) and before `done`. Lightweight notification — no PAD numbers, since the appraisal is still running asynchronously. The result lands in the NEXT turn's `status: "current"` snapshot.
|
||||
|
||||
```json
|
||||
{"type": "affect_update", "status": "scheduled", "turn_id": 42}
|
||||
```
|
||||
|
||||
`scheduled` is skipped on turn failure/cancel paths (the appraisal was never reached); `current` still fires unconditionally for qualifying turns.
|
||||
|
||||
Bootstrap reads available via `GET /agents/{agent_id}/persona_state` (same `snapshot` shape, requires `persona.read` scope).
|
||||
|
||||
### awaiting_llm_first_token
|
||||
|
||||
Periodic heartbeat event (issue #201) emitted at a configurable interval during the gap between `worker_phase: phase="BuildingPrompt"` and `worker_phase: phase="CallingLLM"`. Solves the legitimate-slow first-token visibility gap: consumer TUIs can render a "thinking for Ns…" timer rather than a frozen line during heavy-CoT prompt warmup.
|
||||
|
||||
```json
|
||||
{
|
||||
"type": "awaiting_llm_first_token",
|
||||
"turn_id": 42,
|
||||
"elapsed_ms_since_building_prompt": 5012.3
|
||||
}
|
||||
```
|
||||
|
||||
`elapsed_ms_since_building_prompt` is the server-authoritative wall-clock milliseconds since `BuildingPrompt` was emitted. Independent of network latency or clock skew.
|
||||
|
||||
Heartbeats stop the moment the engine produces its first event (the `CallingLLM` marker). They do NOT re-fire during tool-roundtrip `CallingLLM` re-entries — the heartbeat is scoped to the FIRST `BuildingPrompt → CallingLLM` gap only.
|
||||
|
||||
**Configuration:** `conversation_api.awaiting_llm_first_token_heartbeat_s` (default `5.0`). Per-agent override via `agent.conversation.awaiting_llm_first_token_heartbeat_s`. Value `0.0` disables emission entirely.
|
||||
|
||||
Cancellation paths (stall watchdog, user-cancel) also stop the heartbeat — no `awaiting_llm_first_token` event appears after the terminal `cancelled` event.
|
||||
|
||||
### thinking
|
||||
|
||||
Incremental reasoning/thinking content (from thinking-enabled models).
|
||||
|
||||
@@ -1851,6 +1851,103 @@ SQLite `consumer_agents` table.
|
||||
before any other processing; non-slug user_ids return 403
|
||||
`tier3_user_id_unsupported`.
|
||||
|
||||
### Persona-state observability (issue #204)
|
||||
|
||||
- **INV-204-1 (affect_update event type)**: `affect_update` is a
|
||||
top-level SSE event `type` discriminator, sibling to `worker_phase`
|
||||
/ `tool_*` / `text` / `thinking` / `done`. Not a `worker_phase` sub-
|
||||
phase. INV-061's "BuildingPrompt is the FIRST event" property is
|
||||
scoped to `worker_phase` events only — `affect_update status="current"`
|
||||
may precede BuildingPrompt for persona-enabled agents.
|
||||
- **INV-204-2 (per-turn emission)**: For agents with persona enabled
|
||||
on non-ephemeral sessions, `stream_turn` emits `status="current"`
|
||||
before any other SSE event on a successful or failed turn, and
|
||||
`status="scheduled"` after `update_after_turn` schedules the
|
||||
appraisal task (success path only — skipped on cancel / error
|
||||
before update_after_turn was reached). See contract
|
||||
`docs/contracts/issues/204.contract.md`.
|
||||
- **INV-204-3 (emission suppression)**: Persona-disabled agents and
|
||||
ephemeral sessions emit ZERO `affect_update` events.
|
||||
- **INV-204-6 / INV-204-7 (persona_state endpoint)**: New
|
||||
`GET /agents/{agent_id}/persona_state` gated on Heimdall scope
|
||||
`persona.read`. Route ordering: auth → Tier 3 short-circuit (404
|
||||
`persona_not_configured`) → Tier 1/2 existence (404
|
||||
`agent_not_available`) → persona-enabled check (404
|
||||
`persona_not_configured`) → snapshot (200).
|
||||
- **INV-204-9 (read-only registry primitive)**: `PersonaRegistry.get_state`
|
||||
is mutex-free and never mutates `persona.emotions`. Eventual
|
||||
consistency under concurrent `_appraisal_wrapper` mutations.
|
||||
- **INV-204-14 (replay participation)**: `affect_update` events flow
|
||||
through `_publish`, so SSE resume / replay handles them with no
|
||||
special case.
|
||||
|
||||
## Amendment — AwaitingLLMFirstToken heartbeat (issue #201, INV-201-1..7)
|
||||
|
||||
Adds a periodic SSE heartbeat event during the gap between
|
||||
`BuildingPrompt` and `CallingLLM` so consumers can distinguish
|
||||
"engine is thinking" from "engine is wedged" without out-of-band
|
||||
server inspection. Filed by ratatoskr-dev; ships in v0.29.0.
|
||||
|
||||
- **INV-201-1 (new top-level event type)**: `awaiting_llm_first_token`
|
||||
is a new top-level SSE event type, sibling to `worker_phase` /
|
||||
`tool_*` / `text` / `thinking` / `debug` / `done` / `affect_update`.
|
||||
`_WORKER_PHASE_VOCAB` is NOT extended; INV-053 / INV-054 unchanged.
|
||||
Same precedent as #204's `affect_update`.
|
||||
|
||||
- **INV-201-2 (config-gated emission)**: Heartbeat emission requires
|
||||
`awaiting_llm_first_token_heartbeat_s > 0.0`. When the resolved
|
||||
value is `0.0`, the heartbeat task is never started and zero
|
||||
`awaiting_llm_first_token` events emit for the turn. When > 0.0,
|
||||
the task starts immediately after `_publish_phase("BuildingPrompt")`
|
||||
and emits an event every `interval` seconds until cancelled.
|
||||
|
||||
- **INV-201-3 (defense-in-depth cancellation)**: The heartbeat task
|
||||
is cancelled at three sites (idempotent via the `_cancel_heartbeat`
|
||||
helper): (a) immediately before `_publish_phase("CallingLLM")` on
|
||||
the engine-first-event path; (b) inside the `cancelled`/`error`
|
||||
handling that wraps `_handle_cancel` (covers stall + user-cancel
|
||||
paths); (c) in the outer `finally` block alongside
|
||||
`_clear_stall_timer`. After cancellation, no further
|
||||
`awaiting_llm_first_token` events emit.
|
||||
|
||||
- **INV-201-4 (wire shape)**: Payload is exactly `{type:
|
||||
"awaiting_llm_first_token", turn_id: <int>,
|
||||
elapsed_ms_since_building_prompt: <float>}` plus the composite `id:
|
||||
"<turn_id>:<seq>"` stamped by `_publish`. No additional fields.
|
||||
`elapsed_ms_since_building_prompt` is `(time.monotonic() -
|
||||
building_prompt_t) * 1000.0` where `building_prompt_t` is captured
|
||||
immediately before `BuildingPrompt` is published.
|
||||
|
||||
- **INV-201-5 (first-gap-only scope)**: Heartbeat is scoped to the
|
||||
FIRST `BuildingPrompt → CallingLLM` gap of the turn. Tool round-trip
|
||||
`CallingLLM` re-entries (INV-058) emit ZERO
|
||||
`awaiting_llm_first_token` events. Out-of-scope sub-phases
|
||||
(`AwaitingToolResult`, `AwaitingNextLLMCall`) would be separate
|
||||
follow-up features.
|
||||
|
||||
- **INV-201-6 (replay participation)**: Heartbeat events flow through
|
||||
`_publish → _replay_buffer + queue` per INV-060 — same replay
|
||||
semantics as worker_phase events. On `Last-Event-ID` reconnect,
|
||||
prior heartbeats replay identically.
|
||||
|
||||
- **INV-201-7 (config resolution precedence)**: Per-agent
|
||||
`agent.conversation.awaiting_llm_first_token_heartbeat_s` →
|
||||
`api_cfg.awaiting_llm_first_token_heartbeat_s` → built-in `5.0`.
|
||||
Negative values raise `ConfigurationError` at agent load; `0.0`
|
||||
is valid and means "disabled." Mirrors the `_resolve_stall_timeout_s`
|
||||
precedence pattern (INV-038).
|
||||
|
||||
### Mechanism note
|
||||
|
||||
The heartbeat task is a separate `asyncio.Task` (NOT `loop.call_later`,
|
||||
because heartbeats repeat at an interval rather than fire once at a
|
||||
timeout). An `asyncio.Queue` shared between the heartbeat task and the
|
||||
generator carries events; the generator uses
|
||||
`asyncio.wait(return_when=FIRST_COMPLETED)` to race the engine's
|
||||
`__anext__` against the heartbeat queue's `get` ONLY during the first
|
||||
iteration. After `CallingLLM` fires, the heartbeat task is cancelled
|
||||
and subsequent iterations use the original non-race pattern.
|
||||
|
||||
### Storage extension
|
||||
|
||||
The `consumer_agents` table lives in `core/heimdall/storage/sqlite.py`
|
||||
|
||||
@@ -32,9 +32,9 @@ separate dev team rather than an in-tree Worldtree tool.
|
||||
|
||||
## Current state / in-flight
|
||||
|
||||
_As of 2026-05-25 (post-v0.8.0 local tier-3 agent index in picker):_
|
||||
_As of 2026-05-25 (post-v0.8.2 drop double-print; v0.9.0 live-md next):_
|
||||
|
||||
**Status: v0.8.0 shipped.** Eleven core features complete (`sse_client`
|
||||
**Status: v0.8.2 shipped.** Eleven core features complete (`sse_client`
|
||||
#1, `sessions` #2, `cli` #3, `tui` #4, `--end-user-id` #5, TUI
|
||||
startup error visibility #6, presenter contract semantics amendment
|
||||
#12, startup agent picker #8, §5 layout reshape + Tools pane #13)
|
||||
@@ -51,7 +51,9 @@ Static in the footer (static "Tools" v1; dynamic when more tabs
|
||||
land). CLI mode (--send) unaffected by design — INV-018.
|
||||
|
||||
Last commits on `main`:
|
||||
- v0.8.0 feat(local_agents): JSON-backed local tier-3 index + picker merge
|
||||
- v0.8.2 fix(tui): drop post-Done Markdown body re-render (no double-print)
|
||||
- `11ef683` fix(tui,sse): inline Text streaming + empty-id keepalive skip (v0.8.1)
|
||||
- `9fade55` feat(local_agents): JSON-backed local tier-3 index + picker merge (v0.8.0)
|
||||
- `9918c10` fix(tui): coalesce thinking deltas on `\n` (v0.7.1)
|
||||
- `c086ae2` feat(tier3): ratatoskr.tier3 module + CLI (v0.7.0)
|
||||
- `d356990` refactor(tui): thinking streams into thinking-log (v0.6.5)
|
||||
|
||||
+4
-4
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
|
||||
|
||||
[project]
|
||||
name = "ratatoskr"
|
||||
version = "0.8.0"
|
||||
version = "0.14.2"
|
||||
description = "Worldtree Conversation API debug TUI — multi-pane observability dashboard"
|
||||
readme = "README.md"
|
||||
requires-python = ">=3.12"
|
||||
@@ -42,9 +42,9 @@ Repository = "https://gitea.phasefinal.com/vh/ratatoskr"
|
||||
# Ratatoskr is built against Worldtree at this commit; the vendored
|
||||
# spec snapshot in docs/ reflects that SHA.
|
||||
[tool.ratatoskr.spec-pin]
|
||||
worldtree-spec-rev = "55101e909abcd2219833266b6f905c5bc956e0f0"
|
||||
worldtree-version = "v0.19.0"
|
||||
pinned-on = "2026-05-20"
|
||||
worldtree-spec-rev = "562001af28d752c3a60d449c7ddd09f44fa9dc9a"
|
||||
worldtree-version = "v0.29.0"
|
||||
pinned-on = "2026-05-26"
|
||||
|
||||
[tool.hatch.build.targets.wheel]
|
||||
packages = ["src/ratatoskr"]
|
||||
|
||||
@@ -18,6 +18,8 @@ import httpx
|
||||
|
||||
from ratatoskr.sessions import AgentNotFound, SessionApiFailed, create_session
|
||||
from ratatoskr.sse_client import (
|
||||
AffectUpdate,
|
||||
AwaitingLlmFirstToken,
|
||||
CancelAlreadyCompleted,
|
||||
CancelFailed,
|
||||
Cancelled,
|
||||
@@ -205,6 +207,7 @@ class CliPresenterState:
|
||||
(
|
||||
WorkerPhase, Thinking, Text, TextBoundary,
|
||||
ToolStart, ToolResult, Done, Error, Cancelled,
|
||||
AffectUpdate, AwaitingLlmFirstToken,
|
||||
),
|
||||
)
|
||||
# Thinking events accumulate into the open run.
|
||||
@@ -275,6 +278,28 @@ class CliPresenterState:
|
||||
f". text_boundary: kind={event.kind} char_offset={event.char_offset}\n"
|
||||
)
|
||||
return
|
||||
if isinstance(event, AffectUpdate):
|
||||
# Worldtree #204 / v0.28.0. CLI surface is debug telemetry —
|
||||
# one line to stderr with status + (for current) dominant_emotion.
|
||||
if event.snapshot is not None:
|
||||
dom = event.snapshot.get("dominant_emotion")
|
||||
stderr.write(
|
||||
f". affect_update: status={event.status} turn_id={event.turn_id} "
|
||||
f"dominant_emotion={dom!r}\n"
|
||||
)
|
||||
else:
|
||||
stderr.write(
|
||||
f". affect_update: status={event.status} turn_id={event.turn_id}\n"
|
||||
)
|
||||
return
|
||||
if isinstance(event, AwaitingLlmFirstToken):
|
||||
# Worldtree #201 / v0.29.0. Heartbeat during BuildingPrompt →
|
||||
# CallingLLM gap. Stderr surface, one line per heartbeat.
|
||||
secs = event.elapsed_ms_since_building_prompt / 1000.0
|
||||
stderr.write(
|
||||
f". awaiting_llm_first_token: turn_id={event.turn_id} elapsed={secs:.1f}s\n"
|
||||
)
|
||||
return
|
||||
|
||||
|
||||
async def _cancel_and_log(
|
||||
|
||||
@@ -89,6 +89,46 @@ class SessionApiFailed(Exception):
|
||||
self.body = body
|
||||
|
||||
|
||||
# Worldtree #204 / v0.28.0 — persona_state endpoint failure modes.
|
||||
class PersonaNotConfigured(Exception):
|
||||
"""Raised on HTTP 404 `persona_not_configured` from GET persona_state.
|
||||
|
||||
Agent exists but has no persona surface: persona-disabled Tier 1/2
|
||||
agents (e.g. `domari`, `muninn`) and all Tier 3 consumer-defined
|
||||
agents (Phase 2.0). Distinct from `AgentNotAvailable` which means the
|
||||
agent_id is unknown entirely.
|
||||
"""
|
||||
|
||||
def __init__(self, *, agent_id: str) -> None:
|
||||
super().__init__(f"persona not configured for agent_id: {agent_id!r}")
|
||||
self.agent_id = agent_id
|
||||
|
||||
|
||||
class AgentNotAvailable(Exception):
|
||||
"""Raised on HTTP 404 `agent_not_available` from GET persona_state.
|
||||
|
||||
The agent_id is unknown to the server. Distinct from
|
||||
`PersonaNotConfigured` (agent exists but has no persona).
|
||||
"""
|
||||
|
||||
def __init__(self, *, agent_id: str) -> None:
|
||||
super().__init__(f"agent not available: {agent_id!r}")
|
||||
self.agent_id = agent_id
|
||||
|
||||
|
||||
class AuthScopeDenied(Exception):
|
||||
"""Raised on HTTP 403 `auth_scope_denied` from a Heimdall-scoped endpoint.
|
||||
|
||||
The API key lacks the required scope (e.g. `persona.read` for
|
||||
GET /agents/{id}/persona_state). User-tier keys carry `persona.read`
|
||||
by default; this surfaces when a narrower key is in use.
|
||||
"""
|
||||
|
||||
def __init__(self, *, scope: str) -> None:
|
||||
super().__init__(f"auth scope denied: required={scope!r}")
|
||||
self.scope = scope
|
||||
|
||||
|
||||
async def list_sessions(
|
||||
client: httpx.AsyncClient,
|
||||
*,
|
||||
@@ -200,3 +240,45 @@ async def list_agents(client: httpx.AsyncClient) -> list[AgentInfo]:
|
||||
)
|
||||
for item in body
|
||||
]
|
||||
|
||||
|
||||
async def get_persona_state(
|
||||
client: httpx.AsyncClient, agent_id: str
|
||||
) -> dict[str, Any]:
|
||||
"""GET /agents/{agent_id}/persona_state — fetch current persona snapshot.
|
||||
|
||||
Worldtree #204 / v0.28.0. Returns the same `snapshot` dict shape as the
|
||||
`affect_update` SSE event's `status="current"` emission: pad,
|
||||
dominant_emotion, emotions_active, baseline_pad, mood_drift,
|
||||
last_updated_at. Bootstrap read for clients that want to populate a
|
||||
persona pane on session-open without waiting for turn-1's `affect_update`.
|
||||
|
||||
Auth: requires Heimdall `persona.read` scope (user-tier default).
|
||||
|
||||
Failure modes (mapped to typed exceptions per the spec error_codes):
|
||||
- 404 `persona_not_configured` → PersonaNotConfigured (persona-disabled
|
||||
agents: domari / muninn, and all Tier 3 in Phase 2.0)
|
||||
- 404 `agent_not_available` → AgentNotAvailable (unknown agent_id)
|
||||
- 403 `auth_scope_denied` → AuthScopeDenied (key lacks persona.read)
|
||||
- any other non-2xx → SessionApiFailed (preserves the broader-error
|
||||
precedent from list_agents / list_sessions / create_session)
|
||||
"""
|
||||
assert client is not None
|
||||
assert agent_id and isinstance(agent_id, str)
|
||||
|
||||
resp = await client.get(f"/agents/{agent_id}/persona_state")
|
||||
if resp.status_code == 200:
|
||||
return resp.json()
|
||||
# Discriminate the 4xx error_code sub-codes; everything else falls through.
|
||||
try:
|
||||
err = resp.json()
|
||||
error_code = err.get("error_code") if isinstance(err, dict) else None
|
||||
except ValueError:
|
||||
error_code = None
|
||||
if resp.status_code == 404 and error_code == "persona_not_configured":
|
||||
raise PersonaNotConfigured(agent_id=agent_id)
|
||||
if resp.status_code == 404 and error_code == "agent_not_available":
|
||||
raise AgentNotAvailable(agent_id=agent_id)
|
||||
if resp.status_code == 403 and error_code == "auth_scope_denied":
|
||||
raise AuthScopeDenied(scope="persona.read")
|
||||
raise SessionApiFailed(status=resp.status_code, body=resp.content)
|
||||
|
||||
@@ -111,6 +111,55 @@ class Cancelled:
|
||||
partial_message_id: int | None
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class AwaitingLlmFirstToken:
|
||||
"""SSE event `awaiting_llm_first_token`: heartbeat during slow first-token.
|
||||
|
||||
Fires at the configured interval (default 5s) during the gap between
|
||||
`worker_phase` phase=BuildingPrompt and phase=CallingLLM. Lets clients
|
||||
render a live "thinking for Ns…" indicator instead of a frozen line
|
||||
during legitimate-slow first-token latency. Stops the moment CallingLLM
|
||||
fires (defense-in-depth at three sites); no heartbeat after Cancelled
|
||||
or stalled terminal events. Tool round-trip re-entries do NOT re-fire
|
||||
heartbeats — INV-201-5 scopes the mechanism to the FIRST gap only.
|
||||
|
||||
`elapsed_ms_since_building_prompt` is server-authoritative
|
||||
`time.monotonic()`-based — independent of network latency or clock
|
||||
skew, monotonically increasing across the heartbeat sequence.
|
||||
|
||||
See docs/conversation-api-spec.md § awaiting_llm_first_token
|
||||
(Worldtree #201, v0.29.0).
|
||||
"""
|
||||
|
||||
sse_id: SseId
|
||||
turn_id: int
|
||||
elapsed_ms_since_building_prompt: float
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class AffectUpdate:
|
||||
"""SSE event `affect_update`: persona-state observability snapshot.
|
||||
|
||||
Two emissions per qualifying turn (persona-enabled agent on non-
|
||||
ephemeral session): `status="current"` at turn start carrying the full
|
||||
snapshot, `status="scheduled"` after post-turn appraisal kicks off
|
||||
(lightweight — `snapshot` is None). Suppressed entirely for persona-
|
||||
disabled agents (e.g. `domari`, `muninn`), Tier 3 consumer-defined
|
||||
agents (Phase 2.0), and ephemeral sessions.
|
||||
|
||||
Bootstrap reads available via `GET /agents/{agent_id}/persona_state`
|
||||
(same `snapshot` shape, requires `persona.read` scope).
|
||||
|
||||
See docs/conversation-api-spec.md § affect_update (Worldtree #204,
|
||||
v0.28.0).
|
||||
"""
|
||||
|
||||
sse_id: SseId
|
||||
status: str # "current" | "scheduled"
|
||||
turn_id: int
|
||||
snapshot: dict[str, Any] | None # None when status="scheduled"
|
||||
|
||||
|
||||
Event = (
|
||||
WorkerPhase
|
||||
| Thinking
|
||||
@@ -121,6 +170,8 @@ Event = (
|
||||
| Done
|
||||
| Error
|
||||
| Cancelled
|
||||
| AffectUpdate
|
||||
| AwaitingLlmFirstToken
|
||||
)
|
||||
|
||||
|
||||
@@ -284,6 +335,26 @@ def _envelope_for_type(body: dict[str, Any], sse_id: SseId) -> Event:
|
||||
reason=body.get("reason"),
|
||||
partial_message_id=body.get("partial_message_id"),
|
||||
)
|
||||
if t == "awaiting_llm_first_token":
|
||||
# Worldtree #201 / v0.29.0: top-level heartbeat during BuildingPrompt
|
||||
# → CallingLLM gap. Lets clients render live elapsed-time indicators
|
||||
# instead of frozen lines on legitimate-slow first-token latency.
|
||||
return AwaitingLlmFirstToken(
|
||||
sse_id=sse_id,
|
||||
turn_id=body["turn_id"],
|
||||
elapsed_ms_since_building_prompt=body["elapsed_ms_since_building_prompt"],
|
||||
)
|
||||
if t == "affect_update":
|
||||
# Worldtree #204 / v0.28.0: persona-state observability event.
|
||||
# status="current" carries full snapshot at turn start;
|
||||
# status="scheduled" omits snapshot (lightweight post-appraisal-
|
||||
# kickoff notification).
|
||||
return AffectUpdate(
|
||||
sse_id=sse_id,
|
||||
status=body["status"],
|
||||
turn_id=body["turn_id"],
|
||||
snapshot=body.get("snapshot"),
|
||||
)
|
||||
raise ValueError(f"unknown SSE event type: {t!r}")
|
||||
|
||||
|
||||
@@ -309,6 +380,14 @@ async def _iter_events(
|
||||
# with a bad id is still a keepalive). Don't reorder.
|
||||
if sse.data == "":
|
||||
continue
|
||||
# v0.8.1: empty-id frames are also treated as keepalives. Worldtree
|
||||
# SOMETIMES emits events without an `id:` line (observed mid-stream
|
||||
# on the qwen3.6-35-a3b-heretic provider, 2026-05-25). Per the SSE
|
||||
# RFC, events without ids are legitimate (they just don't update
|
||||
# Last-Event-ID); the previous strict behavior crashed every turn
|
||||
# on the offending agent. Treat same as empty-data: skip silently.
|
||||
if sse.id == "":
|
||||
continue
|
||||
try:
|
||||
sse_id = _parse_sse_id(sse.id)
|
||||
except ValueError as exc:
|
||||
|
||||
+661
-112
File diff suppressed because it is too large
Load Diff
@@ -6,11 +6,15 @@ import respx
|
||||
|
||||
from ratatoskr.sessions import (
|
||||
AgentInfo,
|
||||
AgentNotAvailable,
|
||||
AgentNotFound,
|
||||
AuthScopeDenied,
|
||||
InvalidCursor,
|
||||
PersonaNotConfigured,
|
||||
SessionApiFailed,
|
||||
SessionPage,
|
||||
create_session,
|
||||
get_persona_state,
|
||||
list_agents,
|
||||
list_sessions,
|
||||
)
|
||||
@@ -556,3 +560,108 @@ class TestListAgents:
|
||||
with pytest.raises(SessionApiFailed) as excinfo:
|
||||
await list_agents(client)
|
||||
assert excinfo.value.status == 401
|
||||
|
||||
|
||||
class TestGetPersonaState:
|
||||
"""Worldtree #204 / v0.28.0 — GET /agents/{agent_id}/persona_state.
|
||||
|
||||
Bootstrap read for the persona snapshot — same shape as `affect_update`'s
|
||||
`current` snapshot. Auth via `persona.read` scope (user-tier default).
|
||||
"""
|
||||
|
||||
@respx.mock
|
||||
async def test_happy_full_snapshot(self) -> None:
|
||||
"""happy_full_snapshot [happy,tracer]: 200 → snapshot dict with pad +
|
||||
dominant_emotion + emotions_active + baseline_pad + mood_drift.
|
||||
"""
|
||||
snapshot = {
|
||||
"agent_id": "mimir",
|
||||
"pad": {"pleasure": 0.52, "arousal": 0.47, "dominance": 0.50},
|
||||
"dominant_emotion": "curiosity",
|
||||
"emotions_active": [
|
||||
{"type": "curiosity", "intensity": 0.6, "decay_remaining_s": 202.7}
|
||||
],
|
||||
"baseline_pad": {"pleasure": 0.50, "arousal": 0.40, "dominance": 0.50},
|
||||
"mood_drift": {"valence_delta": 0.02, "arousal_delta": 0.07},
|
||||
"last_updated_at": "2026-05-25T22:30:18+00:00",
|
||||
}
|
||||
respx.get("https://w.example/agents/mimir/persona_state").mock(
|
||||
return_value=httpx.Response(200, json=snapshot)
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
result = await get_persona_state(client, "mimir")
|
||||
assert result == snapshot
|
||||
|
||||
@respx.mock
|
||||
async def test_persona_not_configured_404(self) -> None:
|
||||
"""persona_not_configured_404 [error]: 404 with error_code
|
||||
persona_not_configured → PersonaNotConfigured. Agent exists but has
|
||||
no persona surface (e.g. domari, muninn, Tier 3).
|
||||
"""
|
||||
respx.get("https://w.example/agents/domari/persona_state").mock(
|
||||
return_value=httpx.Response(
|
||||
404, json={"error_code": "persona_not_configured", "message": "no persona"}
|
||||
)
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
with pytest.raises(PersonaNotConfigured) as exc_info:
|
||||
await get_persona_state(client, "domari")
|
||||
assert exc_info.value.agent_id == "domari"
|
||||
|
||||
@respx.mock
|
||||
async def test_agent_not_available_404(self) -> None:
|
||||
"""agent_not_available_404 [error]: 404 with error_code
|
||||
agent_not_available → AgentNotAvailable. Distinct from
|
||||
persona_not_configured — the agent_id itself is unknown.
|
||||
"""
|
||||
respx.get("https://w.example/agents/bogus/persona_state").mock(
|
||||
return_value=httpx.Response(
|
||||
404, json={"error_code": "agent_not_available", "message": "unknown agent"}
|
||||
)
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
with pytest.raises(AgentNotAvailable) as exc_info:
|
||||
await get_persona_state(client, "bogus")
|
||||
assert exc_info.value.agent_id == "bogus"
|
||||
|
||||
@respx.mock
|
||||
async def test_auth_scope_denied_403(self) -> None:
|
||||
"""auth_scope_denied_403 [error]: 403 with error_code auth_scope_denied
|
||||
→ AuthScopeDenied. Key lacks `persona.read` scope.
|
||||
"""
|
||||
respx.get("https://w.example/agents/mimir/persona_state").mock(
|
||||
return_value=httpx.Response(
|
||||
403,
|
||||
json={"error_code": "auth_scope_denied", "message": "missing persona.read"},
|
||||
)
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
with pytest.raises(AuthScopeDenied) as exc_info:
|
||||
await get_persona_state(client, "mimir")
|
||||
assert exc_info.value.scope == "persona.read"
|
||||
|
||||
@respx.mock
|
||||
async def test_404_unknown_error_code_falls_through(self) -> None:
|
||||
"""404_unknown_error_code_falls_through [adversarial]: 404 without the
|
||||
two known error codes → SessionApiFailed (don't swallow novel failure
|
||||
modes as something more specific than they are).
|
||||
"""
|
||||
respx.get("https://w.example/agents/mimir/persona_state").mock(
|
||||
return_value=httpx.Response(404, json={"error_code": "novel_404"})
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
with pytest.raises(SessionApiFailed) as exc_info:
|
||||
await get_persona_state(client, "mimir")
|
||||
assert exc_info.value.status == 404
|
||||
|
||||
@respx.mock
|
||||
async def test_500_unexpected_status(self) -> None:
|
||||
"""500_unexpected_status [error]: 5xx → SessionApiFailed (matches the
|
||||
list_agents / list_sessions / create_session precedent)."""
|
||||
respx.get("https://w.example/agents/mimir/persona_state").mock(
|
||||
return_value=httpx.Response(500, content=b"boom")
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
with pytest.raises(SessionApiFailed) as exc_info:
|
||||
await get_persona_state(client, "mimir")
|
||||
assert exc_info.value.status == 500
|
||||
|
||||
@@ -5,6 +5,8 @@ import pytest
|
||||
import respx
|
||||
|
||||
from ratatoskr.sse_client import (
|
||||
AffectUpdate,
|
||||
AwaitingLlmFirstToken,
|
||||
CancelAlreadyCompleted,
|
||||
Cancelled,
|
||||
CancelResult,
|
||||
@@ -711,6 +713,47 @@ def _sse_raw_chunk(sse_id: str, raw_data: str) -> bytes:
|
||||
return f"id: {sse_id}\ndata: {raw_data}\n\n".encode()
|
||||
|
||||
|
||||
def _sse_no_id_chunk(data: str) -> bytes:
|
||||
"""SSE frame with NO id line + arbitrary data (v0.8.1: keepalive shape)."""
|
||||
return f"data: {data}\n\n".encode()
|
||||
|
||||
|
||||
class TestEmptyIdSkipped:
|
||||
@respx.mock
|
||||
async def test_empty_id_on_first_event_skipped(self) -> None:
|
||||
"""empty_id_on_first_event_skipped [v0.8.1]: stream starts with an
|
||||
event carrying NO `id:` line → httpx_sse exposes sse.id == ''
|
||||
(no prior id to inherit). Pre-v0.8.1: MalformedSseId raw='' crashed
|
||||
the turn. v0.8.1: treat same as empty-data keepalive — skip silently.
|
||||
|
||||
Observed 2026-05-25 on Worldtree's qwen3.6-35-a3b-heretic provider:
|
||||
the first stream frame had no id line, every turn died with
|
||||
`[malformed_sse_id] raw=''`.
|
||||
"""
|
||||
from ratatoskr.sse_client import Done as _Done
|
||||
from ratatoskr.sse_client import Text as _Text
|
||||
|
||||
# First frame: no id line (httpx_sse → sse.id = ""). Skip it.
|
||||
# Subsequent frames have ids; normal processing resumes.
|
||||
stream = (
|
||||
_sse_no_id_chunk('{"type":"keepalive"}') # ← skipped (sse.id == "")
|
||||
+ _sse_chunk("42:1", {"type": "text", "content": "first"})
|
||||
+ _sse_chunk("42:2", _DONE_42_6)
|
||||
)
|
||||
respx.post("https://w.example/sessions/s1/messages").mock(
|
||||
return_value=httpx.Response(
|
||||
200, headers={"content-type": "text/event-stream"}, content=stream
|
||||
)
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
events = [e async for e in stream_turn(client, "s1", "hi")]
|
||||
# 2 events — the no-id frame is invisible (no MalformedSseId crash).
|
||||
assert len(events) == 2
|
||||
assert isinstance(events[0], _Text)
|
||||
assert events[0].content == "first"
|
||||
assert isinstance(events[1], _Done)
|
||||
|
||||
|
||||
class TestEmptyDataSkipped:
|
||||
@respx.mock
|
||||
async def test_empty_data_skipped(self) -> None:
|
||||
@@ -836,3 +879,164 @@ class TestEmptyDataSkipped:
|
||||
assert exc_info.value.raw == "x" * 200
|
||||
# Exception message also only contains the truncated form
|
||||
assert "x" * 5000 not in str(exc_info.value)
|
||||
|
||||
|
||||
class TestAffectUpdate:
|
||||
"""Worldtree #204 / v0.28.0 — persona-state observability SSE event.
|
||||
|
||||
Two emissions per qualifying turn (persona-enabled agent, non-ephemeral
|
||||
session): `status: "current"` at turn start with full snapshot, then
|
||||
`status: "scheduled"` near turn end (lightweight, no snapshot).
|
||||
|
||||
See docs/conversation-api-spec.md § affect_update.
|
||||
"""
|
||||
|
||||
@respx.mock
|
||||
async def test_current_status_parsed_with_snapshot(self) -> None:
|
||||
"""current_status_parsed_with_snapshot [tracer]: status=current carries
|
||||
the full snapshot dict; AffectUpdate.snapshot is populated with the
|
||||
nested PAD / dominant_emotion / emotions_active fields.
|
||||
"""
|
||||
snapshot = {
|
||||
"agent_id": "mimir",
|
||||
"pad": {"pleasure": 0.52, "arousal": 0.47, "dominance": 0.50},
|
||||
"dominant_emotion": "curiosity",
|
||||
"emotions_active": [
|
||||
{"type": "curiosity", "intensity": 0.6, "decay_remaining_s": 202.7}
|
||||
],
|
||||
"baseline_pad": {"pleasure": 0.50, "arousal": 0.40, "dominance": 0.50},
|
||||
"mood_drift": {"valence_delta": 0.02, "arousal_delta": 0.07},
|
||||
"last_updated_at": "2026-05-25T22:30:18+00:00",
|
||||
}
|
||||
stream = _sse_chunk(
|
||||
"42:1",
|
||||
{
|
||||
"type": "affect_update",
|
||||
"status": "current",
|
||||
"turn_id": 42,
|
||||
"snapshot": snapshot,
|
||||
},
|
||||
) + _sse_chunk("42:2", _DONE_42_6)
|
||||
respx.post("https://w.example/sessions/s1/messages").mock(
|
||||
return_value=httpx.Response(
|
||||
200, headers={"content-type": "text/event-stream"}, content=stream
|
||||
)
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
events = [e async for e in stream_turn(client, "s1", "hi")]
|
||||
affect = events[0]
|
||||
assert isinstance(affect, AffectUpdate)
|
||||
assert affect.status == "current"
|
||||
assert affect.turn_id == 42
|
||||
assert affect.snapshot == snapshot
|
||||
assert affect.sse_id == SseId(42, 1)
|
||||
|
||||
@respx.mock
|
||||
async def test_scheduled_status_parsed_no_snapshot(self) -> None:
|
||||
"""scheduled_status_parsed_no_snapshot [trace]: status=scheduled carries
|
||||
no snapshot field; AffectUpdate.snapshot is None.
|
||||
"""
|
||||
stream = (
|
||||
_sse_chunk("42:1", {"type": "text", "content": "x"})
|
||||
+ _sse_chunk(
|
||||
"42:2",
|
||||
{"type": "affect_update", "status": "scheduled", "turn_id": 42},
|
||||
)
|
||||
+ _sse_chunk("42:3", _DONE_42_6)
|
||||
)
|
||||
respx.post("https://w.example/sessions/s1/messages").mock(
|
||||
return_value=httpx.Response(
|
||||
200, headers={"content-type": "text/event-stream"}, content=stream
|
||||
)
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
events = [e async for e in stream_turn(client, "s1", "hi")]
|
||||
affect = next(e for e in events if isinstance(e, AffectUpdate))
|
||||
assert affect.status == "scheduled"
|
||||
assert affect.turn_id == 42
|
||||
assert affect.snapshot is None
|
||||
assert affect.sse_id == SseId(42, 2)
|
||||
|
||||
|
||||
class TestAwaitingLlmFirstToken:
|
||||
"""Worldtree #201 / v0.29.0 — `awaiting_llm_first_token` SSE heartbeat.
|
||||
|
||||
Top-level event (not a worker_phase extension) fired during the
|
||||
BuildingPrompt → CallingLLM gap at the configured interval (default
|
||||
5s). Server-authoritative elapsed_ms is time.monotonic()-based and
|
||||
monotonically increasing across the heartbeat sequence.
|
||||
|
||||
See docs/conversation-api-spec.md § awaiting_llm_first_token.
|
||||
"""
|
||||
|
||||
@respx.mock
|
||||
async def test_single_heartbeat_parsed(self) -> None:
|
||||
"""single_heartbeat_parsed [tracer]: type=awaiting_llm_first_token →
|
||||
AwaitingLlmFirstToken(turn_id, elapsed_ms_since_building_prompt).
|
||||
"""
|
||||
stream = _sse_chunk(
|
||||
"42:1",
|
||||
{
|
||||
"type": "awaiting_llm_first_token",
|
||||
"turn_id": 42,
|
||||
"elapsed_ms_since_building_prompt": 5012.3,
|
||||
},
|
||||
) + _sse_chunk("42:2", _DONE_42_6)
|
||||
respx.post("https://w.example/sessions/s1/messages").mock(
|
||||
return_value=httpx.Response(
|
||||
200, headers={"content-type": "text/event-stream"}, content=stream
|
||||
)
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
events = [e async for e in stream_turn(client, "s1", "hi")]
|
||||
beat = events[0]
|
||||
assert isinstance(beat, AwaitingLlmFirstToken)
|
||||
assert beat.turn_id == 42
|
||||
assert beat.elapsed_ms_since_building_prompt == 5012.3
|
||||
assert beat.sse_id == SseId(42, 1)
|
||||
|
||||
@respx.mock
|
||||
async def test_heartbeat_sequence_monotonic(self) -> None:
|
||||
"""heartbeat_sequence_monotonic [scenario]: three consecutive heartbeats
|
||||
in one turn — elapsed_ms_since_building_prompt monotonically increases,
|
||||
all carry the same turn_id.
|
||||
"""
|
||||
stream = (
|
||||
_sse_chunk(
|
||||
"42:1",
|
||||
{
|
||||
"type": "awaiting_llm_first_token",
|
||||
"turn_id": 42,
|
||||
"elapsed_ms_since_building_prompt": 5000.0,
|
||||
},
|
||||
)
|
||||
+ _sse_chunk(
|
||||
"42:2",
|
||||
{
|
||||
"type": "awaiting_llm_first_token",
|
||||
"turn_id": 42,
|
||||
"elapsed_ms_since_building_prompt": 10005.4,
|
||||
},
|
||||
)
|
||||
+ _sse_chunk(
|
||||
"42:3",
|
||||
{
|
||||
"type": "awaiting_llm_first_token",
|
||||
"turn_id": 42,
|
||||
"elapsed_ms_since_building_prompt": 15011.8,
|
||||
},
|
||||
)
|
||||
+ _sse_chunk("42:4", _DONE_42_6)
|
||||
)
|
||||
respx.post("https://w.example/sessions/s1/messages").mock(
|
||||
return_value=httpx.Response(
|
||||
200, headers={"content-type": "text/event-stream"}, content=stream
|
||||
)
|
||||
)
|
||||
async with httpx.AsyncClient(base_url="https://w.example") as client:
|
||||
events = [e async for e in stream_turn(client, "s1", "hi")]
|
||||
beats = [e for e in events if isinstance(e, AwaitingLlmFirstToken)]
|
||||
assert len(beats) == 3
|
||||
elapsed = [b.elapsed_ms_since_building_prompt for b in beats]
|
||||
assert elapsed == sorted(elapsed) # monotonically increasing
|
||||
assert all(b.turn_id == 42 for b in beats)
|
||||
|
||||
+757
-183
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user