Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 92aa05c688 | |||
| 209427ab23 | |||
| 139771c8d8 | |||
| 489cfee1f0 | |||
| 11ef6830ab | |||
| 9fade55901 |
+12
-5
@@ -7,11 +7,18 @@ 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 | `da93ca7cf613f1dc229a7a44d07fc1d7efc78e25` |
|
||||
| Worldtree HEAD message | `feat(#204): v0.28.0 — persona-state observability surface` |
|
||||
| Pinned on | 2026-05-25 |
|
||||
| Pinned by | ratatoskr-dev (bump for #204 affect_update SSE) |
|
||||
| Worldtree version at pin | `v0.28.0` |
|
||||
|
||||
## Pin history
|
||||
|
||||
| Date | SHA | Version | Notable deltas consumed |
|
||||
|---|---|---|---|
|
||||
| 2026-05-25 | `da93ca7` | v0.28.0 | #204 — new SSE event `affect_update` (current/scheduled), new endpoint `GET /agents/{id}/persona_state` (not yet consumed), 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,45 @@ 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).
|
||||
|
||||
### thinking
|
||||
|
||||
Incremental reasoning/thinking content (from thinking-enabled models).
|
||||
|
||||
@@ -1851,6 +1851,36 @@ 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.
|
||||
|
||||
### 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.7.1 thinking coalesce-by-newline):_
|
||||
_As of 2026-05-25 (post-v0.8.2 drop double-print; v0.9.0 live-md next):_
|
||||
|
||||
**Status: v0.7.1 shipped.** Ten 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,10 @@ 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.7.1 fix(tui): coalesce thinking deltas on `\n` — no more per-token newlines
|
||||
- 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)
|
||||
- `82437bd` style(tui): picker highlighted item → Aurora blue (v0.6.4)
|
||||
|
||||
+4
-4
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
|
||||
|
||||
[project]
|
||||
name = "ratatoskr"
|
||||
version = "0.7.1"
|
||||
version = "0.11.0"
|
||||
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 = "da93ca7cf613f1dc229a7a44d07fc1d7efc78e25"
|
||||
worldtree-version = "v0.28.0"
|
||||
pinned-on = "2026-05-25"
|
||||
|
||||
[tool.hatch.build.targets.wheel]
|
||||
packages = ["src/ratatoskr"]
|
||||
|
||||
@@ -0,0 +1,134 @@
|
||||
"""Local index of tier-3 agents defined via `python -m ratatoskr.tier3`.
|
||||
|
||||
Workaround for Worldtree's ``GET /agents`` not returning consumer-defined
|
||||
agents (the public list excludes tier-3 per-spec; see issue #15 smoke
|
||||
findings). Local file maintains a list of agent_ids + display metadata so
|
||||
the picker can show them alongside foundational agents.
|
||||
|
||||
Storage shape: JSON at ``$XDG_CONFIG_HOME/ratatoskr/local_agents.json``
|
||||
(default ``~/.config/ratatoskr/local_agents.json``). Override via
|
||||
``$RATATOSKR_LOCAL_AGENTS`` env var for tests / per-machine isolation.
|
||||
|
||||
If Worldtree later starts returning tier-3 agents in ``GET /agents``, this
|
||||
module's role narrows to redundant local cache; can be removed cleanly
|
||||
since the picker's dedup-by-agent-id keeps remote-wins behavior.
|
||||
|
||||
Failure modes are lenient: missing file → empty index; corrupt JSON or
|
||||
schema mismatch → empty index (no crash). The picker continues to show
|
||||
foundational agents either way; the local-tier-3 surface degrades to
|
||||
"operator passes --agent ratatoskr:<name> explicitly" — the
|
||||
pre-v0.8.0 workflow.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
from dataclasses import asdict, dataclass
|
||||
from pathlib import Path
|
||||
|
||||
_SCHEMA_VERSION = 1
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class LocalAgentEntry:
|
||||
"""One row in the local tier-3 agent index.
|
||||
|
||||
Schema:
|
||||
- ``agent_id``: full "user_id:agent_name" string (Worldtree-owned).
|
||||
- ``agent_name``: slug from define (display name).
|
||||
- ``model``: provider model ID at last define/patch.
|
||||
- ``description``: synthetic display string (typically derived from
|
||||
the system_prompt's first line + a "(tier 3)" prefix; the picker
|
||||
uses this in its ``{id} · {name} — {description}`` rendering).
|
||||
- ``defined_at``: ISO-8601 timestamp from the Tier3AgentInfo response.
|
||||
"""
|
||||
|
||||
agent_id: str
|
||||
agent_name: str
|
||||
model: str
|
||||
description: str
|
||||
defined_at: str
|
||||
|
||||
|
||||
def _local_agents_path() -> Path:
|
||||
"""Resolve the local index file path with XDG + env-var override."""
|
||||
override = os.environ.get("RATATOSKR_LOCAL_AGENTS")
|
||||
if override:
|
||||
return Path(override)
|
||||
xdg = os.environ.get("XDG_CONFIG_HOME")
|
||||
base = Path(xdg) if xdg else (Path.home() / ".config")
|
||||
return base / "ratatoskr" / "local_agents.json"
|
||||
|
||||
|
||||
def load_local_agents() -> list[LocalAgentEntry]:
|
||||
"""Read the local index. Returns ``[]`` on missing file, corrupt JSON,
|
||||
schema mismatch, or any read error — never raises.
|
||||
"""
|
||||
path = _local_agents_path()
|
||||
if not path.exists():
|
||||
return []
|
||||
try:
|
||||
raw = json.loads(path.read_text())
|
||||
except (json.JSONDecodeError, OSError):
|
||||
return []
|
||||
if not isinstance(raw, dict) or raw.get("version") != _SCHEMA_VERSION:
|
||||
return []
|
||||
agents = raw.get("agents", [])
|
||||
if not isinstance(agents, list):
|
||||
return []
|
||||
out: list[LocalAgentEntry] = []
|
||||
for item in agents:
|
||||
if not isinstance(item, dict):
|
||||
continue
|
||||
try:
|
||||
out.append(LocalAgentEntry(**item))
|
||||
except TypeError:
|
||||
# Malformed row (missing/extra fields) — skip silently.
|
||||
continue
|
||||
return out
|
||||
|
||||
|
||||
def _save_local_agents(agents: list[LocalAgentEntry]) -> None:
|
||||
"""Persist the index. Creates parent dir as needed."""
|
||||
path = _local_agents_path()
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
payload = {"version": _SCHEMA_VERSION, "agents": [asdict(a) for a in agents]}
|
||||
path.write_text(json.dumps(payload, indent=2))
|
||||
|
||||
|
||||
def add_local_agent(entry: LocalAgentEntry) -> None:
|
||||
"""Add (or replace) an agent in the local index. agent_id is the key."""
|
||||
agents = [a for a in load_local_agents() if a.agent_id != entry.agent_id]
|
||||
agents.append(entry)
|
||||
_save_local_agents(agents)
|
||||
|
||||
|
||||
def update_local_agent(entry: LocalAgentEntry) -> None:
|
||||
"""Update an existing entry. Identical semantics to ``add_local_agent``
|
||||
(agent_id is the dedup key), exposed separately so callers can
|
||||
self-document intent.
|
||||
"""
|
||||
add_local_agent(entry)
|
||||
|
||||
|
||||
def remove_local_agent(agent_id: str) -> None:
|
||||
"""Remove an entry by agent_id. No-op if absent (idempotent)."""
|
||||
agents = [a for a in load_local_agents() if a.agent_id != agent_id]
|
||||
_save_local_agents(agents)
|
||||
|
||||
|
||||
def make_description(system_prompt: str) -> str:
|
||||
"""Synthesize a one-line description for the picker from a system prompt.
|
||||
|
||||
Strategy: first non-empty line, stripped of leading markdown heading
|
||||
markers and whitespace, prefixed with "(tier 3) ", truncated to 80
|
||||
chars. Falls back to "(tier 3) custom system prompt" if the prompt is
|
||||
empty (defensive — define rejects empty prompts at PRE-002).
|
||||
"""
|
||||
for line in system_prompt.splitlines():
|
||||
stripped = line.lstrip("# ").strip()
|
||||
if stripped:
|
||||
label = f"(tier 3) {stripped}"
|
||||
return label[:80] + ("…" if len(label) > 80 else "")
|
||||
return "(tier 3) custom system prompt"
|
||||
@@ -111,6 +111,30 @@ class Cancelled:
|
||||
partial_message_id: int | None
|
||||
|
||||
|
||||
@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 +145,7 @@ Event = (
|
||||
| Done
|
||||
| Error
|
||||
| Cancelled
|
||||
| AffectUpdate
|
||||
)
|
||||
|
||||
|
||||
@@ -284,6 +309,17 @@ 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 == "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 +345,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:
|
||||
|
||||
@@ -309,6 +309,11 @@ def _resolve_auth(ns: argparse.Namespace) -> tuple[str, str]:
|
||||
async def _run_define(ns: argparse.Namespace) -> int:
|
||||
api_key, server_url = _resolve_auth(ns)
|
||||
from ratatoskr.cli import USER_AGENT
|
||||
from ratatoskr.local_agents import (
|
||||
LocalAgentEntry,
|
||||
add_local_agent,
|
||||
make_description,
|
||||
)
|
||||
|
||||
async with httpx.AsyncClient(
|
||||
base_url=server_url,
|
||||
@@ -324,6 +329,16 @@ async def _run_define(ns: argparse.Namespace) -> int:
|
||||
system_prompt=ns.system_prompt,
|
||||
model=ns.model,
|
||||
)
|
||||
# v0.8.0: persist to local index so the picker can show it.
|
||||
add_local_agent(
|
||||
LocalAgentEntry(
|
||||
agent_id=info.agent_id,
|
||||
agent_name=info.agent_name,
|
||||
model=info.model,
|
||||
description=make_description(info.system_prompt),
|
||||
defined_at=info.created_at,
|
||||
)
|
||||
)
|
||||
print(f"defined {info.agent_id} ({info.model})")
|
||||
return 0
|
||||
|
||||
@@ -331,6 +346,11 @@ async def _run_define(ns: argparse.Namespace) -> int:
|
||||
async def _run_patch(ns: argparse.Namespace) -> int:
|
||||
api_key, server_url = _resolve_auth(ns)
|
||||
from ratatoskr.cli import USER_AGENT
|
||||
from ratatoskr.local_agents import (
|
||||
LocalAgentEntry,
|
||||
make_description,
|
||||
update_local_agent,
|
||||
)
|
||||
|
||||
if ns.system_prompt is None and ns.model is None:
|
||||
raise _Tier3UsageError(
|
||||
@@ -350,6 +370,16 @@ async def _run_patch(ns: argparse.Namespace) -> int:
|
||||
system_prompt=ns.system_prompt,
|
||||
model=ns.model,
|
||||
)
|
||||
# v0.8.0: refresh local index with the post-patch state.
|
||||
update_local_agent(
|
||||
LocalAgentEntry(
|
||||
agent_id=info.agent_id,
|
||||
agent_name=info.agent_name,
|
||||
model=info.model,
|
||||
description=make_description(info.system_prompt),
|
||||
defined_at=info.updated_at,
|
||||
)
|
||||
)
|
||||
print(f"patched {info.agent_id}")
|
||||
return 0
|
||||
|
||||
@@ -357,6 +387,7 @@ async def _run_patch(ns: argparse.Namespace) -> int:
|
||||
async def _run_delete(ns: argparse.Namespace) -> int:
|
||||
api_key, server_url = _resolve_auth(ns)
|
||||
from ratatoskr.cli import USER_AGENT
|
||||
from ratatoskr.local_agents import remove_local_agent
|
||||
|
||||
async with httpx.AsyncClient(
|
||||
base_url=server_url,
|
||||
@@ -367,6 +398,8 @@ async def _run_delete(ns: argparse.Namespace) -> int:
|
||||
timeout=httpx.Timeout(connect=10.0, read=30.0, write=10.0, pool=10.0),
|
||||
) as client:
|
||||
await delete_agent(client, ns.agent_id)
|
||||
# v0.8.0: drop from local index so the picker stops listing it.
|
||||
remove_local_agent(ns.agent_id)
|
||||
print(f"deleted {ns.agent_id}")
|
||||
return 0
|
||||
|
||||
|
||||
+416
-109
@@ -10,13 +10,15 @@ from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import sys
|
||||
from dataclasses import dataclass, field
|
||||
import time as _time
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime as _datetime
|
||||
from typing import ClassVar, Literal
|
||||
|
||||
import httpx
|
||||
from textual.app import App, ComposeResult
|
||||
from textual.binding import Binding
|
||||
from textual.containers import Horizontal, Vertical
|
||||
from textual.containers import Horizontal, Vertical, VerticalScroll
|
||||
from textual.theme import Theme
|
||||
from textual.widgets import (
|
||||
Footer,
|
||||
@@ -39,6 +41,7 @@ from ratatoskr.sessions import (
|
||||
list_agents,
|
||||
)
|
||||
from ratatoskr.sse_client import (
|
||||
AffectUpdate,
|
||||
CancelAlreadyCompleted,
|
||||
CancelFailed,
|
||||
Cancelled,
|
||||
@@ -178,6 +181,62 @@ def _plain_label(event: Event) -> str:
|
||||
return f"[unknown_event] {type(event).__name__}"
|
||||
|
||||
|
||||
def _ts() -> str:
|
||||
"""HH:MM:SS.fff wall-clock timestamp for debug-pane log lines."""
|
||||
now = _datetime.now()
|
||||
return now.strftime("%H:%M:%S") + f".{now.microsecond // 1000:03d}"
|
||||
|
||||
|
||||
def _audit_line(event: Event) -> str:
|
||||
"""One-line wire-level audit summary for the debug pane.
|
||||
|
||||
v0.10.0: every SSE event arrival lands as one of these in the debug
|
||||
pane (Text and Thinking deltas are aggregated into the turn summary
|
||||
instead — token-rate per-delta lines would drown the pane). Shape:
|
||||
`[HH:MM:SS.fff] event_type sse_id=T:S key=val …`.
|
||||
"""
|
||||
sid = getattr(event, "sse_id", None)
|
||||
sid_str = f"{sid.turn_id}:{sid.seq}" if sid is not None else "-"
|
||||
kind = type(event).__name__.lower()
|
||||
if isinstance(event, WorkerPhase):
|
||||
detail = f"phase={event.phase} turn_id={event.turn_id}"
|
||||
elif isinstance(event, ToolStart):
|
||||
detail = f"name={event.name} args={event.arguments!r:.80}"
|
||||
elif isinstance(event, ToolResult):
|
||||
detail = f"name={event.name} duration_ms={event.duration_ms}"
|
||||
elif isinstance(event, TextBoundary):
|
||||
detail = f"kind={event.kind} char_offset={event.char_offset}"
|
||||
elif isinstance(event, Done):
|
||||
detail = (
|
||||
f"turn_id={event.sse_id.turn_id} model={event.model} "
|
||||
f"duration_ms={event.duration_ms}"
|
||||
)
|
||||
elif isinstance(event, Error):
|
||||
detail = (
|
||||
f"turn_id={event.sse_id.turn_id} code={event.error_code} "
|
||||
f"message={event.message!r:.80}"
|
||||
)
|
||||
elif isinstance(event, Cancelled):
|
||||
detail = f"turn_id={event.turn_id} reason={event.reason!r}"
|
||||
elif isinstance(event, AffectUpdate):
|
||||
# Worldtree #204 / v0.28.0. status="current" carries the full
|
||||
# snapshot; surface dominant_emotion + PAD inline so the operator
|
||||
# sees persona drift at a glance. status="scheduled" is
|
||||
# lightweight — no PAD, just the appraisal-kickoff marker.
|
||||
if event.snapshot is not None:
|
||||
pad = event.snapshot.get("pad") or {}
|
||||
detail = (
|
||||
f"status={event.status} turn_id={event.turn_id} "
|
||||
f"dominant_emotion={event.snapshot.get('dominant_emotion')!r} "
|
||||
f"pad=({pad.get('pleasure')},{pad.get('arousal')},{pad.get('dominance')})"
|
||||
)
|
||||
else:
|
||||
detail = f"status={event.status} turn_id={event.turn_id}"
|
||||
else: # Text / Thinking handled by counter path; fallback for safety
|
||||
detail = ""
|
||||
return f"[{_ts()}] {kind} sse_id={sid_str} {detail}".rstrip()
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class TuiPresenterState:
|
||||
"""Per-turn presenter state for TUI mode (issue #12).
|
||||
@@ -186,9 +245,6 @@ class TuiPresenterState:
|
||||
"""
|
||||
|
||||
thinking_open: bool = False
|
||||
# v0.6.0: per-turn streaming text buffer. Text deltas accumulate here
|
||||
# and update `current_text` Static in place — no per-token RichLog spam.
|
||||
text_buffer: list[str] = field(default_factory=list)
|
||||
# Thinking-run counter for turn-scoped start/end markers.
|
||||
thinking_run_index: int = 0
|
||||
# v0.7.1: thinking-content accumulator. Worldtree emits Thinking deltas
|
||||
@@ -197,13 +253,29 @@ class TuiPresenterState:
|
||||
# only on `\n` boundaries (one written line per natural paragraph) or
|
||||
# when the run closes (any leftover tail).
|
||||
thinking_chunk_buffer: str = ""
|
||||
# v0.9.0: Text accumulator for live Markdown rendering. Worldtree emits
|
||||
# Text deltas at token granularity; each delta appends to this buffer
|
||||
# and the current_response_widget re-renders Markdown(text_chunk_buffer)
|
||||
# in place. On terminal event the widget is finalized + reference clears.
|
||||
text_chunk_buffer: str = ""
|
||||
# v0.9.0: reference to the Static widget holding the current turn's
|
||||
# response Markdown Renderable. None between turns.
|
||||
current_response_widget: object = None
|
||||
# v0.10.0: per-turn counters for the debug-pane turn-summary line. Text
|
||||
# and Thinking events arrive at token rate; emitting per-delta debug
|
||||
# lines would drown the pane. Instead we count them and surface
|
||||
# aggregated totals when the turn closes.
|
||||
text_delta_count: int = 0
|
||||
text_byte_count: int = 0
|
||||
thinking_delta_count: int = 0
|
||||
thinking_byte_count: int = 0
|
||||
turn_start_ts: float = 0.0
|
||||
|
||||
def render(
|
||||
self,
|
||||
event: Event,
|
||||
*,
|
||||
log: RichLog,
|
||||
current_text: Static,
|
||||
transcript: "VerticalScroll",
|
||||
tools_log: RichLog,
|
||||
debug_log: RichLog,
|
||||
thinking_log: RichLog,
|
||||
@@ -211,18 +283,14 @@ class TuiPresenterState:
|
||||
) -> None:
|
||||
"""Render one Worldtree SSE event with the TUI hierarchy + coalescing.
|
||||
|
||||
v0.6.5 routing:
|
||||
- `log` (transcript) = content only: user-prompt echo (written
|
||||
outside the presenter), terminal labels, post-Done Markdown body.
|
||||
- `current_text` (Static below transcript) = live-streaming Text
|
||||
deltas accumulated into one growing line; cleared on terminal.
|
||||
- `tools_log` = ToolStart + ToolResult.
|
||||
- `debug_log` = WorkerPhase + TextBoundary.
|
||||
- `thinking_log` = streaming Thinking deltas inline (each chunk =
|
||||
one line in the scrollable log). Rule(start)/Rule(end) markers
|
||||
wrap each run. The whole pane scrolls naturally — no separate
|
||||
tail-scrolling Static at the bottom (v0.6.5 removed
|
||||
`thinking-current`).
|
||||
v0.9.0 routing:
|
||||
- `transcript` (VerticalScroll) = chat content: each turn mounts
|
||||
child widgets (turn-header / prompt-echo / response Markdown /
|
||||
done-label). Live Markdown rendering during Text streaming.
|
||||
- `tools_log` (RichLog) = ToolStart + ToolResult.
|
||||
- `debug_log` (RichLog) = WorkerPhase + TextBoundary.
|
||||
- `thinking_log` (RichLog) = streaming Thinking deltas inline
|
||||
(coalesced on `\n`); Rule(start)/Rule(end) wrap each run.
|
||||
|
||||
Exceptions caught at the presenter boundary (INV-009 fallback).
|
||||
"""
|
||||
@@ -230,7 +298,7 @@ class TuiPresenterState:
|
||||
event,
|
||||
(
|
||||
WorkerPhase, Thinking, Text, TextBoundary,
|
||||
ToolStart, ToolResult, Done, Error, Cancelled,
|
||||
ToolStart, ToolResult, Done, Error, Cancelled, AffectUpdate,
|
||||
),
|
||||
)
|
||||
from rich.text import Text as RichText
|
||||
@@ -240,6 +308,37 @@ class TuiPresenterState:
|
||||
return RichText(s, style=_AU_DEMOTED)
|
||||
|
||||
try:
|
||||
# v0.10.0: per-event audit log line to debug pane. Text and
|
||||
# Thinking arrive at token rate, so we count them rather than
|
||||
# emit a line per delta — totals are reported in the turn-
|
||||
# summary on Done/Error/Cancelled. Everything else gets one
|
||||
# debug-pane line per arrival with timestamp + sse_id + a short
|
||||
# event-specific summary, giving the operator a wire-level
|
||||
# timeline of what the server sent.
|
||||
if isinstance(event, Text):
|
||||
if self.text_delta_count == 0:
|
||||
if self.turn_start_ts == 0.0:
|
||||
self.turn_start_ts = _time.monotonic()
|
||||
self.text_delta_count += 1
|
||||
self.text_byte_count += len(event.content)
|
||||
elif isinstance(event, Thinking):
|
||||
if self.thinking_delta_count == 0:
|
||||
if self.turn_start_ts == 0.0:
|
||||
self.turn_start_ts = _time.monotonic()
|
||||
self.thinking_delta_count += 1
|
||||
self.thinking_byte_count += len(event.content)
|
||||
else:
|
||||
if self.turn_start_ts == 0.0:
|
||||
self.turn_start_ts = _time.monotonic()
|
||||
debug_log.write(_dim(_audit_line(event)))
|
||||
# v0.11.0: AffectUpdate is debug-pane-only for now (the audit
|
||||
# line emitted above is the complete handling). Return early
|
||||
# so the event doesn't pass through the thinking-close path
|
||||
# or fall into the unknown-event ValueError branch. A full
|
||||
# persona surface (Persona TabPane, sticky header line, or
|
||||
# similar) is deferred to a later bump pending UX direction.
|
||||
if isinstance(event, AffectUpdate):
|
||||
return
|
||||
# v0.7.1: Thinking deltas coalesce by newline before flushing.
|
||||
# Worldtree emits Thinking events at token granularity; per-delta
|
||||
# RichLog writes produce one visual line per token (per-token-per-
|
||||
@@ -284,51 +383,99 @@ class TuiPresenterState:
|
||||
self.thinking_open = False
|
||||
# Now render the non-thinking event itself.
|
||||
if isinstance(event, Text):
|
||||
# v0.6.0: streaming text accumulates into current_text Static
|
||||
# — one growing live line, NOT per-delta RichLog entries.
|
||||
self.text_buffer.append(event.content)
|
||||
current_text.update("".join(self.text_buffer))
|
||||
# v0.8.1: stream Text deltas into transcript directly,
|
||||
# coalesced on `\n`. Same pattern as Thinking (v0.7.1).
|
||||
# The pre-v0.8.1 #current-text Static is gone — its dock-
|
||||
# bottom growth was overlapping the transcript visually.
|
||||
#
|
||||
# v0.9.0: Text deltas accumulate in text_chunk_buffer and
|
||||
# the current_response_widget renders Markdown(buffer) in
|
||||
# place. First Text delta of the turn mounts a fresh Static
|
||||
# holding the Markdown Renderable; subsequent deltas update
|
||||
# the same widget. Live markdown rendering — no post-Done
|
||||
# re-render needed.
|
||||
from rich.markdown import Markdown
|
||||
|
||||
self.text_chunk_buffer += event.content
|
||||
# --raw bypasses Markdown rendering — useful for debugging
|
||||
# the raw text stream surface, and matches the pre-v0.9.0
|
||||
# --raw semantics (which dropped the post-Done Markdown re-
|
||||
# render). In raw mode the response widget holds plain str.
|
||||
rendered = (
|
||||
self.text_chunk_buffer if raw else Markdown(self.text_chunk_buffer)
|
||||
)
|
||||
if self.current_response_widget is None:
|
||||
self.current_response_widget = Static(
|
||||
rendered, classes="response-md"
|
||||
)
|
||||
transcript.mount(self.current_response_widget)
|
||||
else:
|
||||
self.current_response_widget.update(rendered)
|
||||
transcript.scroll_end(animate=False)
|
||||
return
|
||||
if isinstance(event, (Done, Error, Cancelled)):
|
||||
# Terminal event: clear the streaming Static first so the
|
||||
# live-preview band collapses. Then write the colored label
|
||||
# + (non-raw) Markdown body / (raw) accumulated plain text
|
||||
# to the transcript.
|
||||
accumulated = "".join(self.text_buffer)
|
||||
self.text_buffer.clear()
|
||||
current_text.update("")
|
||||
# Terminal labels tinted per outcome (Aurora green / Dawn red
|
||||
# / Dawn yellow) for at-a-glance scanning.
|
||||
# v0.10.0: emit turn-summary to debug pane before clearing
|
||||
# counters. Aggregates the per-event totals (Text + Thinking
|
||||
# deltas don't get per-event audit lines because they arrive
|
||||
# at token rate; the summary surfaces what was elided).
|
||||
elapsed_ms = (
|
||||
int((_time.monotonic() - self.turn_start_ts) * 1000)
|
||||
if self.turn_start_ts
|
||||
else 0
|
||||
)
|
||||
turn_id = (
|
||||
event.sse_id.turn_id
|
||||
if hasattr(event, "sse_id")
|
||||
else getattr(event, "turn_id", "?")
|
||||
)
|
||||
debug_log.write(_dim(
|
||||
f"[{_ts()}] turn_summary turn_id={turn_id} "
|
||||
f"text_deltas={self.text_delta_count} "
|
||||
f"text_bytes={self.text_byte_count} "
|
||||
f"thinking_deltas={self.thinking_delta_count} "
|
||||
f"thinking_bytes={self.thinking_byte_count} "
|
||||
f"elapsed_ms={elapsed_ms}"
|
||||
))
|
||||
# Terminal event: finalize the response widget (clear ref so
|
||||
# the next turn mounts a fresh one). The accumulated text is
|
||||
# already rendered as Markdown in the widget — no post-Done
|
||||
# re-render, no double-print.
|
||||
self.text_chunk_buffer = ""
|
||||
self.current_response_widget = None
|
||||
# Terminal labels mount as styled Statics. Tinted per outcome
|
||||
# (Aurora green / Dawn red / Dawn yellow) for at-a-glance
|
||||
# scanning.
|
||||
if isinstance(event, Done):
|
||||
log.write(RichText(
|
||||
f"[done] turn_id={event.sse_id.turn_id} model={event.model} "
|
||||
f"duration={_format_duration_ms(event.duration_ms)} "
|
||||
f"usage {_format_usage(event.usage, arrow='→')}",
|
||||
style=_AU_SUCCESS,
|
||||
transcript.mount(Static(
|
||||
RichText(
|
||||
f"[done] turn_id={event.sse_id.turn_id} "
|
||||
f"model={event.model} "
|
||||
f"duration={_format_duration_ms(event.duration_ms)} "
|
||||
f"usage {_format_usage(event.usage, arrow='→')}",
|
||||
style=_AU_SUCCESS,
|
||||
),
|
||||
classes="done-label",
|
||||
))
|
||||
if raw:
|
||||
# Raw mode: emit the accumulated streamed text verbatim
|
||||
# so the operator has a record after the Static clears.
|
||||
if accumulated:
|
||||
log.write(accumulated)
|
||||
else:
|
||||
from rich.markdown import Markdown
|
||||
from rich.rule import Rule
|
||||
|
||||
log.write(Rule(style=_AU_DEMOTED))
|
||||
log.write(Markdown(event.response))
|
||||
elif isinstance(event, Error):
|
||||
log.write(RichText(
|
||||
f"[error] turn_id={event.sse_id.turn_id} code={event.error_code} "
|
||||
f"message={event.message!r}",
|
||||
style=_AU_ERROR,
|
||||
transcript.mount(Static(
|
||||
RichText(
|
||||
f"[error] turn_id={event.sse_id.turn_id} "
|
||||
f"code={event.error_code} message={event.message!r}",
|
||||
style=_AU_ERROR,
|
||||
),
|
||||
classes="error-label",
|
||||
))
|
||||
else: # Cancelled
|
||||
log.write(RichText(
|
||||
f"[cancelled] turn_id={event.turn_id} reason={event.reason!r} "
|
||||
f"partial_message_id={event.partial_message_id}",
|
||||
style=_AU_WARNING,
|
||||
transcript.mount(Static(
|
||||
RichText(
|
||||
f"[cancelled] turn_id={event.turn_id} "
|
||||
f"reason={event.reason!r} "
|
||||
f"partial_message_id={event.partial_message_id}",
|
||||
style=_AU_WARNING,
|
||||
),
|
||||
classes="cancelled-label",
|
||||
))
|
||||
transcript.scroll_end(animate=False)
|
||||
return
|
||||
if isinstance(event, WorkerPhase):
|
||||
# v0.5.0: telemetry → Debug pane, not transcript.
|
||||
@@ -360,22 +507,24 @@ class TuiPresenterState:
|
||||
# the original event AND a render_error line with the class name only
|
||||
# (NO exception message — security clause). Volva F1 fix.
|
||||
#
|
||||
# v0.6.0 routing-under-failure preservation — fallback writes go
|
||||
# to the same destination the successful render would have used:
|
||||
# - ToolStart/ToolResult → tools_log
|
||||
# - Thinking → thinking_log
|
||||
# - WorkerPhase/TextBoundary → debug_log
|
||||
# - everything else → log
|
||||
# v0.9.0 routing-under-failure: panes (RichLog) still write Strip
|
||||
# lines; transcript (VerticalScroll) mounts a Static instead.
|
||||
if isinstance(event, (ToolStart, ToolResult)):
|
||||
target = tools_log
|
||||
tools_log.write(_plain_label(event))
|
||||
tools_log.write(f"[render_error] {type(exc).__name__}")
|
||||
elif isinstance(event, Thinking):
|
||||
target = thinking_log
|
||||
thinking_log.write(_plain_label(event))
|
||||
thinking_log.write(f"[render_error] {type(exc).__name__}")
|
||||
elif isinstance(event, (WorkerPhase, TextBoundary)):
|
||||
target = debug_log
|
||||
debug_log.write(_plain_label(event))
|
||||
debug_log.write(f"[render_error] {type(exc).__name__}")
|
||||
else:
|
||||
target = log
|
||||
target.write(_plain_label(event))
|
||||
target.write(f"[render_error] {type(exc).__name__}")
|
||||
# Transcript-bound event (Text / Done / Error / Cancelled).
|
||||
transcript.mount(Static(_plain_label(event), classes="error-label"))
|
||||
transcript.mount(
|
||||
Static(f"[render_error] {type(exc).__name__}", classes="error-label")
|
||||
)
|
||||
transcript.scroll_end(animate=False)
|
||||
|
||||
|
||||
class AgentPickerApp(App[str | None]):
|
||||
@@ -564,21 +713,43 @@ class RatatoskrApp(App[int]):
|
||||
}
|
||||
/* v0.6.5: thinking-current Static removed; thinking now streams
|
||||
directly into thinking-log so the whole pane scrolls naturally. */
|
||||
#transcript {
|
||||
/* v0.9.0: transcript is a VerticalScroll container holding dynamically
|
||||
mounted Statics + Markdown widgets per turn. Live Markdown rendering
|
||||
replaces the v0.8.x RichLog approach which couldn't render Markdown
|
||||
in-flight (only on Done as a re-render → double-print bug). */
|
||||
#transcript-scroll {
|
||||
height: 1fr;
|
||||
background: $background;
|
||||
padding: 0 1;
|
||||
}
|
||||
/* v0.6.0: streaming-text Static carries in-flight assistant tokens.
|
||||
Replaces per-token RichLog spam — one growing line that updates in
|
||||
place. Cleared on terminal event; final Markdown body lands in the
|
||||
transcript. */
|
||||
#current-text {
|
||||
dock: bottom;
|
||||
/* Per-turn mounted widgets carry id-prefix conventions:
|
||||
- .turn-header "── turn N ──" (dim)
|
||||
- .prompt-echo "❯ user input" (aurora bright cyan)
|
||||
- .response-md Markdown(accumulated_text) — updated live
|
||||
- .done-label "[done] turn_id=…" (aurora green)
|
||||
- .error-label "[error] …" (dawn red)
|
||||
- .cancelled-label "[cancelled] …" (dawn yellow)
|
||||
*/
|
||||
.turn-header {
|
||||
height: auto;
|
||||
padding: 0 1;
|
||||
color: $au-dark-60;
|
||||
}
|
||||
.prompt-echo {
|
||||
height: auto;
|
||||
background: $background;
|
||||
padding: 0 1;
|
||||
}
|
||||
.response-md {
|
||||
height: auto;
|
||||
padding: 0 1;
|
||||
}
|
||||
.done-label, .error-label, .cancelled-label {
|
||||
height: auto;
|
||||
padding: 0 1;
|
||||
}
|
||||
/* v0.8.1: #current-text Static removed. Streaming text now coalesces
|
||||
on `\n` and writes directly to #transcript (same pattern as v0.7.1
|
||||
thinking fix). Eliminates the dock-bottom-growth-overlap bug. */
|
||||
#tools-log, #debug-log, #thinking-log {
|
||||
background: $background;
|
||||
padding: 0 1;
|
||||
@@ -677,8 +848,12 @@ class RatatoskrApp(App[int]):
|
||||
# work without widget-level markup=True.
|
||||
with Horizontal(id="main-row"):
|
||||
with Vertical(id="left-column"):
|
||||
yield RichLog(id="transcript", wrap=True, markup=False, highlight=False)
|
||||
yield Static("", id="current-text")
|
||||
# v0.9.0: transcript is a VerticalScroll holding per-turn
|
||||
# mounted widgets (turn header, prompt echo, response Markdown,
|
||||
# done label). Live Markdown rendering happens via Static
|
||||
# widgets holding `Markdown` Renderables, updated as Text
|
||||
# deltas arrive.
|
||||
yield VerticalScroll(id="transcript-scroll")
|
||||
yield Input(id="prompt", placeholder="Type a message and press Enter")
|
||||
with Vertical(id="right-column"):
|
||||
with TabbedContent(id="side-panes"):
|
||||
@@ -739,22 +914,40 @@ class RatatoskrApp(App[int]):
|
||||
)
|
||||
self.state = "idle"
|
||||
self._set_hint(self.HINT_IDLE)
|
||||
# v0.10.0: startup audit so the debug pane carries a complete
|
||||
# session bootstrap line (server URL, agent, end_user_id, raw flag,
|
||||
# session tail) before the first turn fires.
|
||||
self._audit(
|
||||
f"app_mounted server={self.args.server_url} agent_id={self.agent_id!r} "
|
||||
f"session={self.session_id[-8:]} raw={self.args.raw} "
|
||||
f"end_user_id={getattr(self.args, 'end_user_id', None)!r}"
|
||||
)
|
||||
|
||||
def _write_turn_headers(self, turn_id: int) -> None:
|
||||
"""v0.6.0: Write `── turn N ──` Rule headers across every pane so
|
||||
operators can visually correlate sections during cross-pane
|
||||
debugging. Called from `_stream_turn_worker` on first event of
|
||||
each new turn (idempotent per turn via active_turn_id guard).
|
||||
"""v0.6.0: turn-ID headers across every pane for cross-pane
|
||||
correlation. v0.9.0: transcript is a VerticalScroll; mounts a
|
||||
Static with rule-style text instead of writing a Rule Renderable
|
||||
to RichLog. Other panes still use RichLog.write(Rule).
|
||||
"""
|
||||
from rich.rule import Rule
|
||||
from rich.text import Text as RichText
|
||||
|
||||
title = f"turn {turn_id}"
|
||||
rule = Rule(title=title, style=_AU_DEMOTED)
|
||||
try:
|
||||
self.query_one("#transcript", RichLog).write(rule)
|
||||
# Transcript (VerticalScroll): mount a styled Static.
|
||||
transcript = self.query_one("#transcript-scroll", VerticalScroll)
|
||||
transcript.mount(
|
||||
Static(
|
||||
RichText(f"── turn {turn_id} ──", style=_AU_DEMOTED),
|
||||
classes="turn-header",
|
||||
)
|
||||
)
|
||||
# Other panes (RichLog): write the Rule Renderable.
|
||||
self.query_one("#tools-log", RichLog).write(rule)
|
||||
self.query_one("#debug-log", RichLog).write(rule)
|
||||
self.query_one("#thinking-log", RichLog).write(rule)
|
||||
transcript.scroll_end(animate=False)
|
||||
except Exception:
|
||||
# Defensive: widget tree may be tearing down — never let a
|
||||
# turn-header write block the SSE consumer.
|
||||
@@ -769,13 +962,51 @@ class RatatoskrApp(App[int]):
|
||||
# Widget may be gone during shutdown; ignore.
|
||||
pass
|
||||
|
||||
def _audit(self, line: str) -> None:
|
||||
"""Write a timestamped audit line to the debug pane.
|
||||
|
||||
v0.10.0: shared sink for app-level events that don't pass through
|
||||
the presenter — state transitions, worker spawn/cancel, cancel POST
|
||||
lifecycle, startup probes. The presenter's per-event audit lives at
|
||||
`_audit_line()`; this is its app-side counterpart.
|
||||
"""
|
||||
try:
|
||||
from rich.text import Text as RichText
|
||||
self.query_one("#debug-log", RichLog).write(
|
||||
RichText(f"[{_ts()}] {line}", style=_AU_DEMOTED)
|
||||
)
|
||||
except Exception:
|
||||
# Widget may not exist yet (pre-mount) or be tearing down.
|
||||
pass
|
||||
|
||||
def _transition(
|
||||
self, new_state: Literal["idle", "streaming", "cancelling"], reason: str
|
||||
) -> None:
|
||||
"""Set self.state with debug-pane audit log.
|
||||
|
||||
Every state machine transition flows through here so the debug pane
|
||||
carries a complete idle→streaming→cancelling→idle timeline with the
|
||||
triggering reason. Cheap; safe to call from any context.
|
||||
"""
|
||||
old = self.state
|
||||
self.state = new_state
|
||||
if old != new_state:
|
||||
self._audit(f"state {old} → {new_state} reason={reason}")
|
||||
|
||||
async def on_input_submitted(self, event: Input.Submitted) -> None:
|
||||
"""Echo user prompt, spawn stream worker; busy notice if not idle."""
|
||||
"""Echo user prompt, spawn stream worker; busy notice if not idle.
|
||||
|
||||
v0.9.0: prompt echo mounts as a Static in the transcript VerticalScroll
|
||||
(was log.write to RichLog).
|
||||
"""
|
||||
if event.input.id != "prompt":
|
||||
return
|
||||
log = self.query_one("#transcript", RichLog)
|
||||
transcript = self.query_one("#transcript-scroll", VerticalScroll)
|
||||
if self.state != "idle":
|
||||
log.write("[busy] turn in flight; input ignored")
|
||||
transcript.mount(
|
||||
Static("[busy] turn in flight; input ignored", classes="error-label")
|
||||
)
|
||||
transcript.scroll_end(animate=False)
|
||||
event.input.value = ""
|
||||
return
|
||||
content = event.input.value.strip()
|
||||
@@ -784,37 +1015,53 @@ class RatatoskrApp(App[int]):
|
||||
# v0.4.1 retheme: operator's voice gets Australis bright cyan so it
|
||||
# stands out against the default-foreground assistant text below it.
|
||||
from rich.text import Text as RichText
|
||||
log.write(RichText(f"❯ {content}", style=_AU_USER_ECHO)) # noqa: RUF001
|
||||
transcript.mount(
|
||||
Static(
|
||||
RichText(f"❯ {content}", style=_AU_USER_ECHO), # noqa: RUF001
|
||||
classes="prompt-echo",
|
||||
)
|
||||
)
|
||||
transcript.scroll_end(animate=False)
|
||||
event.input.value = ""
|
||||
self.state = "streaming"
|
||||
self._transition("streaming", "input_submitted")
|
||||
self._audit(f"worker_spawn content_len={len(content)}")
|
||||
self._set_hint(self.HINT_STREAMING)
|
||||
self.stream_worker = self.run_worker(
|
||||
self._stream_turn_worker(content), exclusive=True
|
||||
)
|
||||
|
||||
async def _stream_turn_worker(self, content: str) -> None:
|
||||
"""Drive stream_turn, render events via TuiPresenterState (issue #12)."""
|
||||
"""Drive stream_turn, render events via TuiPresenterState.
|
||||
|
||||
v0.9.0: transcript is a VerticalScroll; the presenter's `transcript`
|
||||
argument is the container, and the presenter mounts Static / Markdown-
|
||||
backed widgets directly. Wire-error labels mount as `error-label`
|
||||
Statics into the transcript-scroll.
|
||||
"""
|
||||
assert self.state == "streaming"
|
||||
assert self.client is not None
|
||||
assert content
|
||||
log = self.query_one("#transcript", RichLog)
|
||||
current_text = self.query_one("#current-text", Static)
|
||||
transcript = self.query_one("#transcript-scroll", VerticalScroll)
|
||||
tools_log = self.query_one("#tools-log", RichLog)
|
||||
debug_log = self.query_one("#debug-log", RichLog)
|
||||
thinking_log = self.query_one("#thinking-log", RichLog)
|
||||
presenter = TuiPresenterState()
|
||||
|
||||
def _mount_wire_error(label: str) -> None:
|
||||
try:
|
||||
transcript.mount(Static(label, classes="error-label"))
|
||||
transcript.scroll_end(animate=False)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
try:
|
||||
async for event in stream_turn(self.client, self.session_id, content):
|
||||
if self.active_turn_id is None:
|
||||
self.active_turn_id = event.sse_id.turn_id
|
||||
# v0.6.0: turn-ID headers across all panes so the
|
||||
# operator can visually correlate sections during
|
||||
# cross-pane debugging.
|
||||
self._write_turn_headers(self.active_turn_id)
|
||||
presenter.render(
|
||||
event,
|
||||
log=log,
|
||||
current_text=current_text,
|
||||
transcript=transcript,
|
||||
tools_log=tools_log,
|
||||
debug_log=debug_log,
|
||||
thinking_log=thinking_log,
|
||||
@@ -823,17 +1070,22 @@ class RatatoskrApp(App[int]):
|
||||
if isinstance(event, (Done, Error, Cancelled)):
|
||||
break
|
||||
except SseConnectFailed as exc:
|
||||
log.write(f"[sse_connect_failed] status={exc.status} body={exc.body!r}")
|
||||
self._audit(f"sse_connect_failed status={exc.status} body={exc.body!r:.120}")
|
||||
_mount_wire_error(f"[sse_connect_failed] status={exc.status} body={exc.body!r}")
|
||||
except SseConnectionDropped as exc:
|
||||
log.write(f"[connection_dropped] last_seen={exc.last_seen_sse_id}")
|
||||
self._audit(f"connection_dropped last_seen={exc.last_seen_sse_id}")
|
||||
_mount_wire_error(f"[connection_dropped] last_seen={exc.last_seen_sse_id}")
|
||||
except MalformedSseId as exc:
|
||||
log.write(f"[malformed_sse_id] raw={exc.raw!r}")
|
||||
self._audit(f"malformed_sse_id raw={exc.raw!r}")
|
||||
_mount_wire_error(f"[malformed_sse_id] raw={exc.raw!r}")
|
||||
except MalformedSseData as exc:
|
||||
log.write(f"[malformed_sse_data] raw={exc.raw!r}")
|
||||
self._audit(f"malformed_sse_data raw={exc.raw!r:.120}")
|
||||
_mount_wire_error(f"[malformed_sse_data] raw={exc.raw!r}")
|
||||
except TurnIdFlip as exc:
|
||||
log.write(f"[turn_id_flip] expected={exc.established} got={exc.got}")
|
||||
self._audit(f"turn_id_flip expected={exc.established} got={exc.got}")
|
||||
_mount_wire_error(f"[turn_id_flip] expected={exc.established} got={exc.got}")
|
||||
finally:
|
||||
self.state = "idle"
|
||||
self._transition("idle", "worker_finally")
|
||||
self.active_turn_id = None
|
||||
self._set_hint(self.HINT_IDLE)
|
||||
|
||||
@@ -845,26 +1097,35 @@ class RatatoskrApp(App[int]):
|
||||
"""Two-stage Ctrl-C state machine per INV-003."""
|
||||
assert self.state in ("idle", "streaming", "cancelling")
|
||||
if self.state == "idle":
|
||||
self._audit("ctrl_c state=idle action=exit code=0")
|
||||
self.exit(0)
|
||||
elif self.state == "streaming":
|
||||
if self.active_turn_id is None:
|
||||
self._audit("ctrl_c state=streaming active_turn_id=None action=force_exit code=3")
|
||||
if self.stream_worker is not None:
|
||||
self.stream_worker.cancel()
|
||||
self.exit(3)
|
||||
return
|
||||
self.state = "cancelling"
|
||||
self._audit(f"ctrl_c state=streaming turn_id={self.active_turn_id} action=cancel_post")
|
||||
self._transition("cancelling", "ctrl_c_cancel_post_issued")
|
||||
self._set_hint(self.HINT_CANCELLING)
|
||||
log = self.query_one("#transcript", RichLog)
|
||||
transcript = self.query_one("#transcript-scroll", VerticalScroll)
|
||||
self.run_worker(
|
||||
_cancel_via_sse(self.client, self.session_id, self.active_turn_id, log=log)
|
||||
_cancel_via_sse(
|
||||
self.client, self.session_id, self.active_turn_id,
|
||||
transcript=transcript,
|
||||
audit=self._audit,
|
||||
)
|
||||
)
|
||||
elif self.state == "cancelling":
|
||||
self._audit("ctrl_c state=cancelling action=force_exit code=3")
|
||||
if self.stream_worker is not None:
|
||||
self.stream_worker.cancel()
|
||||
self.exit(3)
|
||||
|
||||
def action_quit(self) -> None:
|
||||
"""Ctrl-D — immediate exit regardless of state."""
|
||||
self._audit(f"ctrl_d state={self.state} action=exit code=0")
|
||||
if self.stream_worker is not None and not self.stream_worker.is_finished:
|
||||
self.stream_worker.cancel()
|
||||
self.exit(0)
|
||||
@@ -928,6 +1189,14 @@ async def _resolve_then_run(args: ParsedArgs) -> int:
|
||||
# Issue #8: startup agent picker — fetch GET /agents and prompt when
|
||||
# --new is passed without --agent. list_agents errors land on real
|
||||
# stderr before any alt-screen opens (preserves issue #6 INV-001).
|
||||
#
|
||||
# v0.8.0: merge in local tier-3 agent index. Worldtree's GET /agents
|
||||
# doesn't return consumer-defined agents (issue #15 smoke finding);
|
||||
# ratatoskr keeps its own JSON-backed index of agents the operator
|
||||
# defined via `python -m ratatoskr.tier3 define`. Merged here so the
|
||||
# picker shows foundational + local-tier-3 in one list. Dedup by
|
||||
# agent_id (remote wins on conflict, since a server-listed agent
|
||||
# is the authoritative source).
|
||||
chosen_agent_id: str | None = args.agent_id
|
||||
if args.new and args.agent_id is None:
|
||||
try:
|
||||
@@ -940,6 +1209,23 @@ async def _resolve_then_run(args: ParsedArgs) -> int:
|
||||
except (httpx.ConnectError, httpx.ReadTimeout, httpx.TransportError) as exc:
|
||||
sys.stderr.write(f"[network_error] {type(exc).__name__}: {exc}\n")
|
||||
return 21
|
||||
# v0.8.0: append local tier-3 entries not already in the remote list.
|
||||
from ratatoskr.local_agents import load_local_agents
|
||||
|
||||
remote_ids = {a.agent_id for a in agents}
|
||||
for entry in load_local_agents():
|
||||
if entry.agent_id in remote_ids:
|
||||
continue
|
||||
agents.append(AgentInfo(
|
||||
agent_id=entry.agent_id,
|
||||
name=entry.agent_name,
|
||||
description=entry.description,
|
||||
version=None,
|
||||
capabilities=[],
|
||||
supported_models=[],
|
||||
persona_traits={},
|
||||
ui_hints={},
|
||||
))
|
||||
if not agents:
|
||||
sys.stderr.write("[no_agents] server returned empty agent list\n")
|
||||
return 13
|
||||
@@ -980,12 +1266,33 @@ async def _cancel_via_sse(
|
||||
session_id: str,
|
||||
turn_id: int,
|
||||
*,
|
||||
log: RichLog,
|
||||
transcript: VerticalScroll,
|
||||
audit: "Callable[[str], None] | None" = None,
|
||||
) -> None:
|
||||
"""Fire-and-forget cancel; never raises (mirrors cli._cancel_and_log; #3 INV-009)."""
|
||||
"""Fire-and-forget cancel; never raises (mirrors cli._cancel_and_log; #3 INV-009).
|
||||
|
||||
v0.9.0: mounts a `[cancel_failed]` Static into the transcript-scroll
|
||||
container on failure (was log.write to RichLog).
|
||||
v0.10.0: optional `audit` callback (RatatoskrApp._audit) receives one
|
||||
line on POST issue + one on POST result, so the debug pane carries the
|
||||
full cancel lifecycle. Defaults to no-op for legacy callers.
|
||||
"""
|
||||
assert client is not None
|
||||
assert isinstance(turn_id, int) and turn_id > 0
|
||||
if audit is not None:
|
||||
audit(f"cancel_post issued session_id={session_id} turn_id={turn_id}")
|
||||
try:
|
||||
await cancel_turn(client, session_id, turn_id)
|
||||
if audit is not None:
|
||||
audit(f"cancel_post ok turn_id={turn_id}")
|
||||
except (CancelFailed, CancelTurnNotFound, CancelAlreadyCompleted, httpx.RequestError) as exc:
|
||||
log.write(f"[cancel_failed] {type(exc).__name__}: {exc}")
|
||||
if audit is not None:
|
||||
audit(f"cancel_post failed turn_id={turn_id} {type(exc).__name__}: {exc!s:.120}")
|
||||
try:
|
||||
transcript.mount(Static(
|
||||
f"[cancel_failed] {type(exc).__name__}: {exc}",
|
||||
classes="error-label",
|
||||
))
|
||||
transcript.scroll_end(animate=False)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
@@ -0,0 +1,187 @@
|
||||
"""Tests for ratatoskr.local_agents.
|
||||
|
||||
Use ``$RATATOSKR_LOCAL_AGENTS`` env-var override + pytest tmp_path to
|
||||
isolate from the operator's real ``~/.config/ratatoskr/local_agents.json``.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from ratatoskr.local_agents import (
|
||||
LocalAgentEntry,
|
||||
_local_agents_path,
|
||||
add_local_agent,
|
||||
load_local_agents,
|
||||
make_description,
|
||||
remove_local_agent,
|
||||
update_local_agent,
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def local_path(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> Path:
|
||||
"""Point $RATATOSKR_LOCAL_AGENTS at a fresh tmp file for the test."""
|
||||
path = tmp_path / "local_agents.json"
|
||||
monkeypatch.setenv("RATATOSKR_LOCAL_AGENTS", str(path))
|
||||
return path
|
||||
|
||||
|
||||
def _entry(
|
||||
agent_id: str = "ratatoskr:wizard",
|
||||
agent_name: str = "wizard",
|
||||
model: str = "qwen3.6-35-a3b",
|
||||
description: str = "(tier 3) test agent",
|
||||
defined_at: str = "2026-05-25T00:00:00+00:00",
|
||||
) -> LocalAgentEntry:
|
||||
return LocalAgentEntry(
|
||||
agent_id=agent_id,
|
||||
agent_name=agent_name,
|
||||
model=model,
|
||||
description=description,
|
||||
defined_at=defined_at,
|
||||
)
|
||||
|
||||
|
||||
class TestPathResolution:
|
||||
def test_env_override(self, monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.setenv("RATATOSKR_LOCAL_AGENTS", "/tmp/custom-agents.json")
|
||||
assert _local_agents_path() == Path("/tmp/custom-agents.json")
|
||||
|
||||
def test_xdg_config_home(self, monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.delenv("RATATOSKR_LOCAL_AGENTS", raising=False)
|
||||
monkeypatch.setenv("XDG_CONFIG_HOME", "/tmp/xdg-config")
|
||||
assert (
|
||||
_local_agents_path()
|
||||
== Path("/tmp/xdg-config/ratatoskr/local_agents.json")
|
||||
)
|
||||
|
||||
def test_default_home(self, monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.delenv("RATATOSKR_LOCAL_AGENTS", raising=False)
|
||||
monkeypatch.delenv("XDG_CONFIG_HOME", raising=False)
|
||||
path = _local_agents_path()
|
||||
assert path == Path.home() / ".config" / "ratatoskr" / "local_agents.json"
|
||||
|
||||
|
||||
class TestLoadEmpty:
|
||||
def test_missing_file_returns_empty(self, local_path: Path) -> None:
|
||||
assert not local_path.exists()
|
||||
assert load_local_agents() == []
|
||||
|
||||
def test_corrupt_json_returns_empty(self, local_path: Path) -> None:
|
||||
local_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
local_path.write_text("not json at all")
|
||||
assert load_local_agents() == []
|
||||
|
||||
def test_wrong_schema_version_returns_empty(self, local_path: Path) -> None:
|
||||
local_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
local_path.write_text(json.dumps({"version": 999, "agents": []}))
|
||||
assert load_local_agents() == []
|
||||
|
||||
def test_missing_version_key_returns_empty(self, local_path: Path) -> None:
|
||||
local_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
local_path.write_text(json.dumps({"agents": []}))
|
||||
assert load_local_agents() == []
|
||||
|
||||
def test_malformed_row_skipped(self, local_path: Path) -> None:
|
||||
local_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
local_path.write_text(
|
||||
json.dumps(
|
||||
{
|
||||
"version": 1,
|
||||
"agents": [
|
||||
{"agent_id": "incomplete"}, # missing required fields
|
||||
{
|
||||
"agent_id": "ratatoskr:good",
|
||||
"agent_name": "good",
|
||||
"model": "m",
|
||||
"description": "d",
|
||||
"defined_at": "t",
|
||||
},
|
||||
],
|
||||
}
|
||||
)
|
||||
)
|
||||
entries = load_local_agents()
|
||||
assert len(entries) == 1
|
||||
assert entries[0].agent_id == "ratatoskr:good"
|
||||
|
||||
|
||||
class TestAdd:
|
||||
def test_add_one(self, local_path: Path) -> None:
|
||||
add_local_agent(_entry())
|
||||
entries = load_local_agents()
|
||||
assert len(entries) == 1
|
||||
assert entries[0].agent_id == "ratatoskr:wizard"
|
||||
|
||||
def test_add_two_different(self, local_path: Path) -> None:
|
||||
add_local_agent(_entry(agent_id="ratatoskr:a", agent_name="a"))
|
||||
add_local_agent(_entry(agent_id="ratatoskr:b", agent_name="b"))
|
||||
ids = {e.agent_id for e in load_local_agents()}
|
||||
assert ids == {"ratatoskr:a", "ratatoskr:b"}
|
||||
|
||||
def test_add_replaces_same_id(self, local_path: Path) -> None:
|
||||
add_local_agent(_entry(model="old-model"))
|
||||
add_local_agent(_entry(model="new-model"))
|
||||
entries = load_local_agents()
|
||||
assert len(entries) == 1
|
||||
assert entries[0].model == "new-model"
|
||||
|
||||
def test_creates_parent_dirs(
|
||||
self, tmp_path: Path, monkeypatch: pytest.MonkeyPatch
|
||||
) -> None:
|
||||
nested = tmp_path / "deep" / "nested" / "path" / "agents.json"
|
||||
monkeypatch.setenv("RATATOSKR_LOCAL_AGENTS", str(nested))
|
||||
add_local_agent(_entry())
|
||||
assert nested.exists()
|
||||
|
||||
|
||||
class TestUpdate:
|
||||
def test_update_changes_existing(self, local_path: Path) -> None:
|
||||
add_local_agent(_entry(model="v1"))
|
||||
update_local_agent(_entry(model="v2"))
|
||||
entries = load_local_agents()
|
||||
assert len(entries) == 1
|
||||
assert entries[0].model == "v2"
|
||||
|
||||
|
||||
class TestRemove:
|
||||
def test_remove_existing(self, local_path: Path) -> None:
|
||||
add_local_agent(_entry())
|
||||
remove_local_agent("ratatoskr:wizard")
|
||||
assert load_local_agents() == []
|
||||
|
||||
def test_remove_missing_is_noop(self, local_path: Path) -> None:
|
||||
add_local_agent(_entry())
|
||||
remove_local_agent("ratatoskr:doesnotexist")
|
||||
assert len(load_local_agents()) == 1
|
||||
|
||||
|
||||
class TestMakeDescription:
|
||||
def test_first_nonempty_line(self) -> None:
|
||||
prompt = "\n\n# IDENTITY\nYou are a test agent..."
|
||||
desc = make_description(prompt)
|
||||
assert desc.startswith("(tier 3) IDENTITY")
|
||||
|
||||
def test_strips_heading_markers(self) -> None:
|
||||
prompt = "# A nice heading\nMore prompt..."
|
||||
desc = make_description(prompt)
|
||||
assert "(tier 3) A nice heading" == desc
|
||||
|
||||
def test_truncates_long(self) -> None:
|
||||
prompt = "x" * 200
|
||||
desc = make_description(prompt)
|
||||
# 80 char cap including the prefix
|
||||
assert len(desc) == 81 # 80 + ellipsis char
|
||||
assert desc.endswith("…")
|
||||
|
||||
def test_empty_prompt_fallback(self) -> None:
|
||||
desc = make_description("")
|
||||
assert desc == "(tier 3) custom system prompt"
|
||||
|
||||
def test_whitespace_only_fallback(self) -> None:
|
||||
desc = make_description(" \n\n ")
|
||||
assert desc == "(tier 3) custom system prompt"
|
||||
@@ -5,6 +5,7 @@ import pytest
|
||||
import respx
|
||||
|
||||
from ratatoskr.sse_client import (
|
||||
AffectUpdate,
|
||||
CancelAlreadyCompleted,
|
||||
Cancelled,
|
||||
CancelResult,
|
||||
@@ -711,6 +712,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 +878,80 @@ 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)
|
||||
|
||||
+60
-6
@@ -1,5 +1,7 @@
|
||||
"""Tests for ratatoskr.tier3 per docs/contracts/issues/15.contract.md."""
|
||||
|
||||
from pathlib import Path
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
import respx
|
||||
@@ -315,12 +317,29 @@ class TestDeleteAgent:
|
||||
assert exc.value.status == 500
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def _isolated_local_agents(
|
||||
tmp_path: "Path", monkeypatch: pytest.MonkeyPatch
|
||||
) -> "Path":
|
||||
"""Isolate the v0.8.0 local-tier-3 index from the operator's real file."""
|
||||
path = tmp_path / "local_agents.json"
|
||||
monkeypatch.setenv("RATATOSKR_LOCAL_AGENTS", str(path))
|
||||
return path
|
||||
|
||||
|
||||
class TestCli:
|
||||
@respx.mock
|
||||
def test_cli_define_happy(
|
||||
self, capsys: pytest.CaptureFixture[str], monkeypatch: pytest.MonkeyPatch
|
||||
self,
|
||||
capsys: pytest.CaptureFixture[str],
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
_isolated_local_agents: "Path",
|
||||
) -> None:
|
||||
"""cli_define_happy [happy]: argv → 201 mock → stdout confirmation."""
|
||||
"""cli_define_happy [happy]: argv → 201 mock → stdout confirmation;
|
||||
local index updated with the new entry (v0.8.0 hook).
|
||||
"""
|
||||
from ratatoskr.local_agents import load_local_agents
|
||||
|
||||
monkeypatch.setenv("WORLDTREE_API_URL", "https://w.example")
|
||||
monkeypatch.setenv("WORLDTREE_API_KEY", "k")
|
||||
respx.post("https://w.example/agents/define").mock(
|
||||
@@ -335,12 +354,24 @@ class TestCli:
|
||||
out = capsys.readouterr()
|
||||
assert rc == 0
|
||||
assert out.out.strip() == "defined ratatoskr:wizard (qwen3.6-35-a3b)"
|
||||
# v0.8.0: local index now has the new entry.
|
||||
entries = load_local_agents()
|
||||
assert len(entries) == 1
|
||||
assert entries[0].agent_id == "ratatoskr:wizard"
|
||||
assert entries[0].model == "qwen3.6-35-a3b"
|
||||
|
||||
@respx.mock
|
||||
def test_cli_patch_happy(
|
||||
self, capsys: pytest.CaptureFixture[str], monkeypatch: pytest.MonkeyPatch
|
||||
self,
|
||||
capsys: pytest.CaptureFixture[str],
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
_isolated_local_agents: "Path",
|
||||
) -> None:
|
||||
"""cli_patch_happy [happy]: argv → 200 mock → stdout confirmation."""
|
||||
"""cli_patch_happy [happy]: argv → 200 mock → stdout confirmation;
|
||||
local index refreshed with the post-patch state.
|
||||
"""
|
||||
from ratatoskr.local_agents import load_local_agents
|
||||
|
||||
monkeypatch.setenv("WORLDTREE_API_URL", "https://w.example")
|
||||
monkeypatch.setenv("WORLDTREE_API_KEY", "k")
|
||||
respx.patch("https://w.example/agents/ratatoskr:wizard").mock(
|
||||
@@ -350,12 +381,34 @@ class TestCli:
|
||||
out = capsys.readouterr()
|
||||
assert rc == 0
|
||||
assert out.out.strip() == "patched ratatoskr:wizard"
|
||||
entries = load_local_agents()
|
||||
assert len(entries) == 1
|
||||
assert entries[0].agent_id == "ratatoskr:wizard"
|
||||
|
||||
@respx.mock
|
||||
def test_cli_delete_happy(
|
||||
self, capsys: pytest.CaptureFixture[str], monkeypatch: pytest.MonkeyPatch
|
||||
self,
|
||||
capsys: pytest.CaptureFixture[str],
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
_isolated_local_agents: "Path",
|
||||
) -> None:
|
||||
"""cli_delete_happy [happy]: argv → 204 mock → stdout confirmation."""
|
||||
"""cli_delete_happy [happy]: argv → 204 mock → stdout confirmation;
|
||||
local index entry removed (v0.8.0 hook).
|
||||
"""
|
||||
from ratatoskr.local_agents import (
|
||||
LocalAgentEntry,
|
||||
add_local_agent,
|
||||
load_local_agents,
|
||||
)
|
||||
|
||||
# Pre-populate so we can verify removal.
|
||||
add_local_agent(LocalAgentEntry(
|
||||
agent_id="ratatoskr:wizard",
|
||||
agent_name="wizard",
|
||||
model="m",
|
||||
description="d",
|
||||
defined_at="t",
|
||||
))
|
||||
monkeypatch.setenv("WORLDTREE_API_URL", "https://w.example")
|
||||
monkeypatch.setenv("WORLDTREE_API_KEY", "k")
|
||||
respx.delete("https://w.example/agents/ratatoskr:wizard").mock(
|
||||
@@ -365,6 +418,7 @@ class TestCli:
|
||||
out = capsys.readouterr()
|
||||
assert rc == 0
|
||||
assert out.out.strip() == "deleted ratatoskr:wizard"
|
||||
assert load_local_agents() == []
|
||||
|
||||
def test_cli_missing_auth(
|
||||
self, capsys: pytest.CaptureFixture[str], monkeypatch: pytest.MonkeyPatch
|
||||
|
||||
+548
-183
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user