From ad0f4522e78a7eb6033b8a17adc475dd63c34e1a Mon Sep 17 00:00:00 2001 From: jaylfc Date: Mon, 5 Oct 2026 02:29:37 +0000 Subject: [PATCH 1/7] tsk-aa5qck: extract avatar slug + dir into agent_avatars module RED-FIRST proof: ```text ==================================== ERRORS ==================================== _________________ ERROR collecting tests/test_agent_avatars.py _________________ ImportError while importing test module '/tmp/exec-tsk-aa5qck/tests/test_agent_avatars.py'. Hint: make sure your test modules/packages have valid Python names. Traceback (most recent call last): File "/home/jay/.local/share/uv/python/cpython-3.14.7-linux-x86_64-gnu/lib/python3.14/importlib/__init__.py", line 88, in _gcd_import module = _bootstrap._gcd_import(name, package, level) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ tests/test_agent_avatars.py:9: in from tinyagentos.agent_avatars import avatar_hash, avatar_source_path E ModuleNotFoundError: No module named 'tinyagentos.agent_avatars' !!!!!!!!!!!!!!!!!!! Interrupted: 1 error during collection !!!!!!!!!!!!!!!!!!!! 1 error in 0.32s ``` Green after fix: ```text ==================================== 4 passed in 0.24s ==================================== ``` Demo mode and lock-widgets tests also green: ```text ==================================== 17 passed in 22.48s ==================================== ... ........................................................................ [ 60%] ............................................... [100%] 119 passed in 43.90s ``` Docs-Reviewed: module extraction only, no route or user-facing behavior change, docs/agent-coordination.md not affected --- .../tsk-aa5qck-agent-avatars-module.md | 3 ++ tests/test_agent_avatars.py | 45 +++++++++++++++++++ tinyagentos/agent_avatars.py | 38 ++++++++++++++++ tinyagentos/routes/auth.py | 25 +---------- 4 files changed, 87 insertions(+), 24 deletions(-) create mode 100644 changelog.d/tsk-aa5qck-agent-avatars-module.md create mode 100644 tests/test_agent_avatars.py create mode 100644 tinyagentos/agent_avatars.py diff --git a/changelog.d/tsk-aa5qck-agent-avatars-module.md b/changelog.d/tsk-aa5qck-agent-avatars-module.md new file mode 100644 index 000000000..997c20826 --- /dev/null +++ b/changelog.d/tsk-aa5qck-agent-avatars-module.md @@ -0,0 +1,3 @@ +### Changed + +- Extracted avatar slug and avatar directory logic from `tinyagentos/routes/auth.py` into a new `tinyagentos/agent_avatars.py` module, and added `avatar_source_path` and `avatar_hash` helpers. diff --git a/tests/test_agent_avatars.py b/tests/test_agent_avatars.py new file mode 100644 index 000000000..00fa37936 --- /dev/null +++ b/tests/test_agent_avatars.py @@ -0,0 +1,45 @@ +"""Tests for tinyagentos.agent_avatars module (avatar_source_path + avatar_hash).""" +from __future__ import annotations + +import hashlib +from pathlib import Path + +import pytest + +from tinyagentos import agent_avatars as aa +from tinyagentos.agent_avatars import avatar_hash, avatar_source_path +from tinyagentos.routes.auth import _avatar_slug + + +class TestAvatarSourcePathAndHash: + + def test_no_jpg_returns_none(self, tmp_path, monkeypatch): + monkeypatch.setattr(aa, "LOCK_AVATAR_DIR", str(tmp_path)) + assert avatar_source_path("Some Agent") is None + assert avatar_hash("Some Agent") is None + + def test_existing_jpg_returns_path_and_hash(self, tmp_path, monkeypatch): + monkeypatch.setattr(aa, "LOCK_AVATAR_DIR", str(tmp_path)) + name = "Some Agent" + slug = _avatar_slug(name) + (tmp_path / f"{slug}.jpg").write_bytes(b"one") + assert avatar_source_path(name) == tmp_path / f"{slug}.jpg" + h = avatar_hash(name) + assert isinstance(h, str) + assert len(h) == 16 + assert h == hashlib.sha256(b"one").hexdigest()[:16] + + def test_hash_changes_when_file_changes(self, tmp_path, monkeypatch): + monkeypatch.setattr(aa, "LOCK_AVATAR_DIR", str(tmp_path)) + name = "Some Agent" + slug = _avatar_slug(name) + p = tmp_path / f"{slug}.jpg" + p.write_bytes(b"one") + first = avatar_hash(name) + p.write_bytes(b"two") + second = avatar_hash(name) + assert first != second + + def test_imported_slug_matches_module_slug(self): + name = "Some Agent" + assert _avatar_slug(name) == aa._avatar_slug(name) diff --git a/tinyagentos/agent_avatars.py b/tinyagentos/agent_avatars.py new file mode 100644 index 000000000..f8ad5182b --- /dev/null +++ b/tinyagentos/agent_avatars.py @@ -0,0 +1,38 @@ +from __future__ import annotations + +import hashlib +import os +from pathlib import Path + +LOCK_AVATAR_DIR = os.environ.get("TAOS_LOCK_AVATAR_DIR", "/var/lib/taos/lock-avatars") + + +def _avatar_slug(name: str) -> str: + """Slug for an agent name, restricted to characters that cannot traverse. + + Anything outside [a-z0-9-] is dropped rather than escaped: this value is + used to build a filesystem path, so a conservative whitelist is the control + that keeps "../" and absolute paths out, not a sanitiser that tries to spot + bad input. + """ + out = [] + for ch in name.strip().lower(): + if ch.isalnum() and ch.isascii(): + out.append(ch) + elif out and out[-1] != "-": + out.append("-") + return "".join(out).strip("-") + + +def avatar_source_path(name: str) -> Path | None: + """Path to the avatar image file for *name*, or None if it is not installed.""" + path = Path(LOCK_AVATAR_DIR) / f"{_avatar_slug(name)}.jpg" + return path if path.is_file() else None + + +def avatar_hash(name: str) -> str | None: + """SHA-256 hex of the avatar image bytes, first 16 chars, or None.""" + path = avatar_source_path(name) + if path is None: + return None + return hashlib.sha256(path.read_bytes()).hexdigest()[:16] diff --git a/tinyagentos/routes/auth.py b/tinyagentos/routes/auth.py index 4d3720908..7182d3e37 100644 --- a/tinyagentos/routes/auth.py +++ b/tinyagentos/routes/auth.py @@ -27,6 +27,7 @@ Response, StreamingResponse, ) +from tinyagentos.agent_avatars import LOCK_AVATAR_DIR, _avatar_slug from tinyagentos.auth import ( PIN_MAX_LEN, PIN_MIN_LEN, @@ -9603,30 +9604,6 @@ def _attach(agent: dict) -> None: -#: Where lock-screen agent avatars are read from. One flat directory of -#: ".jpg" files, slug being the agent name lowercased with non-alphanumerics -#: collapsed to "-". Overridable so a packaged install can point it at its own -#: data dir rather than this default. -LOCK_AVATAR_DIR = os.environ.get("TAOS_LOCK_AVATAR_DIR", "/var/lib/taos/lock-avatars") - - -def _avatar_slug(name: str) -> str: - """Slug for an agent name, restricted to characters that cannot traverse. - - Anything outside [a-z0-9-] is dropped rather than escaped: this value is - used to build a filesystem path, so a conservative whitelist is the control - that keeps "../" and absolute paths out, not a sanitiser that tries to spot - bad input. - """ - out = [] - for ch in name.strip().lower(): - if ch.isalnum() and ch.isascii(): - out.append(ch) - elif out and out[-1] != "-": - out.append("-") - return "".join(out).strip("-") - - @router.get("/lock-avatar/{slug}") async def lock_avatar(slug: str, request: Request): """Serve one lock-screen avatar. Console-only, same reasoning as the widgets. From 4477ec5f05361c52f65fa4ff5d3ee646f346b792 Mon Sep 17 00:00:00 2001 From: jaylfc Date: Mon, 5 Oct 2026 03:57:02 +0000 Subject: [PATCH 2/7] feat(device): add GET /api/device/v1/state (agents:read) Extract async def assemble_lock_agents(request, owner_id) from lock_widgets in auth.py. The new assembler builds the agent list with container status, demo placeholders, device-live agents, and pending decisions, owner-filtered when requested. Create tinyagentos/routes/device_state.py: device-bearer only, required scope agents:read. Response shape: agents list with name, status, framework, avatar {hue, hash}, attention, last_recap, decision {id, question, options} | null; plus server.version, time, demo bool. Server-side string caps (name 48, status 120, last_recap 180, question 280, option 40) with trailing ellipsis. avatar.hash is the first 16 hex chars of the avatar image SHA-256, or null when no image is installed. Demo flag comes from _demo_enabled(request); demo decision id is "" (unanswerable). Wire device_state_router into routes/__init__.py and add /api/device/v1/state to the device-bearer allowlist in auth_middleware.py (AGENTS_READ scope). Tests: red-first tests/test_device_v1_state.py with 6 cases (owner filtering, caps, scope gating, demo flag, settings switch, avatar hash). Captured FAIL block in RED-PROOF.md before implementation. Docs: extend docs/routes.d/16-device-v1.md, docs/agent-coordination.md. Changelog fragment: changelog.d/tsk-ypsu2z-device-v1-state.md. Closes: tsk-ypsu2z # RED-PROOF - GET /api/device/v1/state Captured before implementation so the initial test run is guaranteed to fail. ```text ============================= test session starts ============================= collected 6 items tests/test_device_v1_state.py FFFFFF [100%] =================================== FAILURES =================================== _____________________ test_state_returns_only_owner_agents _____________________ vapp = @pytest.mark.asyncio async def test_state_returns_only_owner_agents(vapp): app = vapp app.state.config.agents = [ {"name": "alice-agent", "framework": "openclaw", "user_id": "owner-1"}, {"name": "bob-agent", "framework": "hermes", "user_id": "owner-2"}, ] tok1 = await _device(app, user_id="owner-1", scopes=("agents:read",)) tok2 = await _device(app, user_id="owner-2", scopes=("agents:read",)) async with _client(app) as c: r1 = await c.get("/api/device/v1/state", headers=_bearer(tok1)) r2 = await c.get("/api/device/v1/state", headers=_bearer(tok2)) > assert r1.status_code == 200, r1.text E AssertionError: {"error":"Authentication required"} E assert 401 == 200 E + where 401 = .status_code tests/test_device_v1_state.py:58: AssertionError ------------------------------ Captured log setup ------------------------------ WARNING tinyagentos.containers.backend:backend.py:320 No container backend detected (Incus / Docker / Podman / Apple / Native). Cluster features and worker containers will be disabled. Install one (e.g. 'sudo apt install incus' on Debian/Debian, 'sudo dnf install incus' on Fedora) and restart taOS. _________________________ test_state_caps_long_strings ________________________ vapp = @pytest.mark.asyncio async def test_state_caps_long_strings(vapp): app = vapp long_status = "x" * 500 long_recap = "r" * 500 app.state.config.agents = [ {"name": "cap-agent", "framework": "openclaw", "user_id": "u1", "status": long_status}, ] await app.state.agent_messages.send( from_agent="cap-agent", to_agent="cap-agent", message=long_recap ) tok = await _device(app, user_id="u1", scopes=("agents:read",)) async with _client(app) as c: r = await c.get("/api/device/v1/state", headers=_bearer(tok)) > assert r.status_code == 200, r.text E AssertionError: {"error":"Authentication required"} E assert 401 == 200 E + where 401 = .status_code tests/test_device_v1_state.py:87: AssertionError ------------------------------ Captured log setup ------------------------------ WARNING tinyagentos.containers.backend:backend.py:320 No container backend detected (Incus / Docker / Podman / Apple / Native). Cluster features and worker containers will be disabled. Install one (e.g. 'sudo apt install incus' on Debian/Debian, 'sudo dnf install incus' on Fedora) and restart taOS. ___________________ test_state_requires_agents_read ___________________ vapp = @pytest.mark.asyncio async def test_state_requires_agents_read(vapp): app = vapp tok = await _device(app, user_id="u1", scopes=("chat:send",)) async with _client(app) as c: r = await c.get("/api/device/v1/state", headers=_bearer(tok)) > assert r.status_code == 403, r.text E AssertionError: {"error":"Authentication required"} E assert 401 == 403 E + where 401 = .status_code tests/test_device_v1_state.py:105: AssertionError ------------------------------ Captured log setup ------------------------------ WARNING tinyagentos.containers.backend:backend.py:320 No container backend detected (Incus / Docker / Podman / Apple / Native). Cluster features and worker containers will be disabled. Install one (e.g. 'sudo apt install incus' on Debian/Debian, 'sudo dnf install incus' on Fedora) and restart taOS. ________________ test_state_demo_flag_and_unanswerable_decision ________________ vapp = monkeypatch = <_pytest.monkeypatch.MonkeyPatch object at 0x701d07916750> @pytest.mark.asyncio async def test_state_demo_flag_and_unanswerable_decision(vapp, monkeypatch): from tinyagentos.demo_mode import write_demo_mode app = vapp monkeypatch.setenv("TAOS_LOCK_DEMO_AGENTS", "DemoA:openclaw:Drafting") monkeypatch.setenv("TAOS_LOCK_DEMO_DECISION", "Ship it?") monkeypatch.setenv("TAOS_LOCK_DEMO_DECISION_AGENT", "DemoA") write_demo_mode(app.state.data_dir, True) app.state.config.agents = [ {"name": "real-agent", "framework": "openclaw", "user_id": "u1"}, ] tok = await _device(app, user_id="u1", scopes=("agents:read",)) async with _client(app) as c: r_state = await c.get("/api/device/v1/state", headers=_bearer(tok)) r_widgets = await c.get("/auth/lock-widgets") > assert r_state.status_code == 200, r_state.text E AssertionError: {"error":"Authentication required"} E assert 401 == 200 E + where 401 = .status_code tests/test_device_v1_state.py:130: AssertionError ------------------------------ Captured log setup ------------------------------ WARNING tinyagentos.containers.backend:backend.py:320 No container backend detected (Incus / Docker / Podman / Apple / Native). Cluster features and worker containers will be disabled. Install one (e.g. 'sudo apt install incus' on Debian/Debian, 'sudo dnf install incus' on Fedora) and restart taOS. __________________ test_settings_demo_switch_takes_state_down __________________ vapp = @pytest.mark.asyncio async def test_settings_demo_switch_takes_state_down(vapp): from tinyagentos.demo_mode import write_demo_mode app = vapp monkeypatch = pytest.MonkeyPatch() monkeypatch.setenv("TAOS_LOCK_DEMO_AGENTS", "DemoB:openclaw:Drafting") monkeypatch.setenv("TAOS_LOCK_DEMO_DECISION", "Ship it?") monkeypatch.setenv("TAOS_LOCK_DEMO_DECISION_AGENT", "DemoB") write_demo_mode(app.state.data_dir, True) tok = await _device(app, user_id="u1", scopes=("agents:read",)) async with _client(app) as c: r_on = await c.get("/api/device/v1/state", headers=_bearer(tok)) w_on = await c.get("/auth/lock-widgets") > assert r_on.json()["demo"] is True ^^^^^^^^^^^^^^^^^^^ E KeyError: 'demo' tests/test_device_v1_state.py:173: KeyError ------------------------------ Captured log setup ------------------------------ WARNING tinyagentos.containers.backend:backend.py:320 No container backend detected (Incus / Docker / Podman / Apple / Native). Cluster features and worker containers will be disabled. Install one (e.g. 'sudo apt install incus' on Debian/Debian, 'sudo dnf install incus' on Fedora) and restart taOS. ___________________________ test_state_avatar_hash _____________________________ vapp = monkeypatch = <_pytest.monkeypatch.MonkeyPatch object at 0x701d153ab350> @pytest.mark.asyncio async def test_state_avatar_hash(vapp, monkeypatch): from tinyagentos import agent_avatars as avatars app = vapp app.state.config.agents = [ {"name": "avatar-agent", "framework": "openclaw", "user_id": "u1"}, ] with tempfile.TemporaryDirectory() as tmpdir: monkeypatch.setattr(avatars, "LOCK_AVATAR_DIR", tmpdir) tok = await _device(app, user_id="u1", scopes=("agents:read",)) async with _client(app) as c: r = await c.get("/api/device/v1/state", headers=_bearer(tok)) > assert r.status_code == 200, r.text E AssertionError: {"error":"Authentication required"} E assert 401 == 200 E + where 401 = .status_code tests/test_device_v1_state.py:205: AssertionError ------------------------------ Captured log setup ------------------------------ WARNING tinyagentos.containers.backend:backend.py:320 No container backend detected (Incus / Docker / Podman / Apple / Native). Cluster features and worker containers will be disabled. Install one (e.g. 'sudo apt install incus' on Debian/Debian, 'sudo dnf install incus' on Fedora) and restart taOS. =========================== short test summary info ============================ FAILED tests/test_device_v1_state.py::test_state_returns_only_owner_agents FAILED tests/test_device_v1_state.py::test_state_caps_long_strings FAILED tests/test_device_v1_state.py::test_state_requires_agents_read FAILED tests/test_device_v1_state.py::test_state_demo_flag_and_unanswerable_decision FAILED tests/test_device_v1_state.py::test_settings_demo_switch_takes_state_down FAILED tests/test_device_v1_state.py::test_state_avatar_hash 6 failed in 14.84s ``` --- changelog.d/tsk-ypsu2z-device-v1-state.md | 7 + docs/agent-coordination.md | 12 ++ docs/routes.d/16-device-v1.md | 42 ++++ tests/test_device_v1_state.py | 232 ++++++++++++++++++++++ tinyagentos/auth_middleware.py | 3 + tinyagentos/routes/__init__.py | 5 + tinyagentos/routes/auth.py | 186 +++++++---------- tinyagentos/routes/device_state.py | 113 +++++++++++ 8 files changed, 489 insertions(+), 111 deletions(-) create mode 100644 changelog.d/tsk-ypsu2z-device-v1-state.md create mode 100644 tests/test_device_v1_state.py create mode 100644 tinyagentos/routes/device_state.py diff --git a/changelog.d/tsk-ypsu2z-device-v1-state.md b/changelog.d/tsk-ypsu2z-device-v1-state.md new file mode 100644 index 000000000..00f4e9365 --- /dev/null +++ b/changelog.d/tsk-ypsu2z-device-v1-state.md @@ -0,0 +1,7 @@ +### Added + +- GET `/api/device/v1/state` (device bearer, scope `agents:read`): returns the owner-filtered agent list for the paired device, plus `server.version`, `time`, and `demo` flag. Server-side string caps (name 48, status 120, `last_recap` 180, question 280, option 40) with trailing ellipsis. `avatar.hash` is the first 16 hex chars of the avatar image SHA-256, or null when no image is installed. + +### Security + +- `/api/device/v1/state` is registered in the device-bearer allowlist (`_DEVICE_BEARER_PATHS`) and is CSRF-exempt like the other device-bearer routes. Only a `taosdev_...` bearer with the `agents:read` scope may reach it; missing scope returns 403 with `device_scope_missing`. diff --git a/docs/agent-coordination.md b/docs/agent-coordination.md index 00b817a2a..cf9ae8459 100644 --- a/docs/agent-coordination.md +++ b/docs/agent-coordination.md @@ -993,6 +993,18 @@ approval channel for privileged grants. Device scoped tokens do not expire and cannot be self-rotated; the only revocation is `DELETE /api/devices/{id}` from a session. +The device-bearer self-service allowlist in `tinyagentos/auth_middleware.py` +covers: + +- `PATCH /api/devices/{id}/push-token` — scope `push:register` +- `GET /api/decisions`, `GET /api/decisions/{id}`, `GET /api/decisions/{id}/history` — scope `agents:read` +- `POST /api/decisions/{id}/answer` — scope `decisions:answer` +- `POST /api/library/ingest` — scope `library:ingest` +- `POST /api/projects/{slug}/files/upload` — scope `files:upload` +- `POST /api/chat/messages` — scope `chat:send` +- `POST /api/device/v1/voice`, `POST /api/device/v1/voice/tts` — scope `voice:stt` / `voice:tts` +- `GET /api/device/v1/state` — scope `agents:read` (owner-filtered agent list) + ## Share destinations (device bearer) `GET /api/share/destinations` lets a paired device DISCOVER share destinations. It is diff --git a/docs/routes.d/16-device-v1.md b/docs/routes.d/16-device-v1.md index 50b9eb881..f4269e856 100644 --- a/docs/routes.d/16-device-v1.md +++ b/docs/routes.d/16-device-v1.md @@ -5,3 +5,45 @@ are bearer-authenticated. The following routes are public (no bearer token): - `POST /api/devices/pair-requests` - `GET /api/devices/pair-requests/{pair_request_id}` + +### GET /api/device/v1/state + +**Scope:** `agents:read` (device bearer). + +Returns the owner-filtered agent list for the paired device, plus the server +version, current time, and demo flag. + +**Response 200:** + +```json +{ + "agents": [ + { + "name": "alice-agent", + "status": "running", + "framework": "openclaw", + "avatar": { + "hue": 123, + "hash": "abc123def4567890" + }, + "attention": true, + "last_recap": "recapped message...", + "decision": { + "id": "", + "question": "Ship it?", + "options": ["Approve", "Deny"] + } + } + ], + "server": { + "version": "1.0.0-beta.55" + }, + "time": 1727640000.123, + "demo": false +} +``` + +**Error codes:** + +- `401` -- missing or invalid device bearer. +- `403` -- `{"error": "device_scope_missing", "scope": "agents:read"}` when the device token lacks the required scope. diff --git a/tests/test_device_v1_state.py b/tests/test_device_v1_state.py new file mode 100644 index 000000000..2daaacf80 --- /dev/null +++ b/tests/test_device_v1_state.py @@ -0,0 +1,232 @@ +"""S2b: GET /api/device/v1/state (device bearer, scope agents:read). + +RED-FIRST: this module is written BEFORE the route exists so the first run +must FAIL with 404. The FAIL block is captured in RED-PROOF.md. +""" +from __future__ import annotations + +import hashlib +import json +import os +import tempfile +from pathlib import Path + +import pytest +import pytest_asyncio +from httpx import ASGITransport, AsyncClient + +PLAIN = "http://testserver:6969" + + +def _bearer(tok): + return {"Authorization": f"Bearer {tok}"} + + +def _client(app, base=PLAIN): + return AsyncClient(transport=ASGITransport(app=app), base_url=base) + + +@pytest_asyncio.fixture +async def vapp(app, client): + return app + + +async def _device(app, user_id="u1", platform="ios", scopes=("agents:read",)): + st = app.state.device_store + d = await st.register(user_id=user_id, platform=platform) + if scopes is not None: + await st.set_scopes(d["device_id"], list(scopes)) + return d["scoped_token"] + + +# (a) two owners, each with agents; a device paired to owner 1 sees only owner 1's agents. +@pytest.mark.asyncio +async def test_state_returns_only_owner_agents(vapp): + app = vapp + app.state.config.agents = [ + {"name": "alice-agent", "framework": "openclaw", "user_id": "owner-1"}, + {"name": "bob-agent", "framework": "hermes", "user_id": "owner-2"}, + ] + + tok1 = await _device(app, user_id="owner-1", scopes=("agents:read",)) + tok2 = await _device(app, user_id="owner-2", scopes=("agents:read",)) + + async with _client(app) as c: + r1 = await c.get("/api/device/v1/state", headers=_bearer(tok1)) + r2 = await c.get("/api/device/v1/state", headers=_bearer(tok2)) + + assert r1.status_code == 200, r1.text + assert r2.status_code == 200, r2.text + + names1 = {a["name"] for a in r1.json()["agents"]} + names2 = {a["name"] for a in r2.json()["agents"]} + + assert names1 == {"alice-agent"} + assert names2 == {"bob-agent"} + + +# (b) a 500-char recap arrives with length <= 180 and ends with the ellipsis. +@pytest.mark.asyncio +async def test_state_caps_long_strings(vapp): + app = vapp + long_status = "x" * 500 + long_recap = "r" * 500 + app.state.config.agents = [ + {"name": "cap-agent", "framework": "openclaw", "user_id": "u1", "status": long_status}, + ] + + await app.state.agent_messages.send( + from_agent="cap-agent", to_agent="cap-agent", message=long_recap + ) + + tok = await _device(app, user_id="u1", scopes=("agents:read",)) + + async with _client(app) as c: + r = await c.get("/api/device/v1/state", headers=_bearer(tok)) + + assert r.status_code == 200, r.text + data = r.json() + agent = next(a for a in data["agents"] if a["name"] == "cap-agent") + assert len(agent["status"]) <= 120 + assert len(agent["last_recap"]) <= 180 + assert agent["status"].endswith("\u2026") + assert agent["last_recap"].endswith("\u2026") + + +# (c) device token without agents:read -> 403, detail error device_scope_missing. +@pytest.mark.asyncio +async def test_state_requires_agents_read(vapp): + app = vapp + tok = await _device(app, user_id="u1", scopes=("chat:send",)) + + async with _client(app) as c: + r = await c.get("/api/device/v1/state", headers=_bearer(tok)) + + assert r.status_code == 403, r.text + assert r.json()["detail"] == {"error": "device_scope_missing", "scope": "agents:read"} + + +# (g) demo flag + unanswerable decision + lock-widgets name parity. +@pytest.mark.asyncio +async def test_state_demo_flag_and_unanswerable_decision(vapp, monkeypatch): + from tinyagentos.demo_mode import write_demo_mode + + app = vapp + monkeypatch.setenv("TAOS_LOCK_DEMO_AGENTS", "DemoA:openclaw:Drafting") + monkeypatch.setenv("TAOS_LOCK_DEMO_DECISION", "Ship it?") + monkeypatch.setenv("TAOS_LOCK_DEMO_DECISION_AGENT", "DemoA") + write_demo_mode(app.state.data_dir, True) + + app.state.config.agents = [ + {"name": "real-agent", "framework": "openclaw", "user_id": "u1"}, + ] + + tok = await _device(app, user_id="u1", scopes=("agents:read",)) + + async with _client(app) as c: + r_state = await c.get("/api/device/v1/state", headers=_bearer(tok)) + r_widgets = await c.get("/auth/lock-widgets") + + assert r_state.status_code == 200, r_state.text + assert r_widgets.status_code == 200, r_widgets.text + + state_data = r_state.json() + widgets_data = r_widgets.json() + + assert state_data["demo"] is True + state_names = {a["name"] for a in state_data["agents"]} + widgets_names = {a["name"] for a in widgets_data["agents"] if not a.get("system")} + assert state_names == widgets_names + + for a in state_data["agents"]: + dec = a.get("decision") + if dec and dec.get("id") == "": + assert dec["id"] == "" + break + + write_demo_mode(app.state.data_dir, False) + + async with _client(app) as c: + r_off = await c.get("/api/device/v1/state", headers=_bearer(tok)) + + assert r_off.json()["demo"] is False + + +# (h) Settings demo switch takes state down. +@pytest.mark.asyncio +async def test_settings_demo_switch_takes_state_down(vapp): + from tinyagentos.demo_mode import write_demo_mode + + app = vapp + monkeypatch = pytest.MonkeyPatch() + monkeypatch.setenv("TAOS_LOCK_DEMO_AGENTS", "DemoB:openclaw:Drafting") + monkeypatch.setenv("TAOS_LOCK_DEMO_DECISION", "Ship it?") + monkeypatch.setenv("TAOS_LOCK_DEMO_DECISION_AGENT", "DemoB") + write_demo_mode(app.state.data_dir, True) + + tok = await _device(app, user_id="u1", scopes=("agents:read",)) + + async with _client(app) as c: + r_on = await c.get("/api/device/v1/state", headers=_bearer(tok)) + w_on = await c.get("/auth/lock-widgets") + + assert r_on.json()["demo"] is True + state_names_on = {a["name"] for a in r_on.json()["agents"]} + widgets_names_on = {a["name"] for a in w_on.json()["agents"] if not a.get("system")} + assert state_names_on == widgets_names_on + + write_demo_mode(app.state.data_dir, False) + + async with _client(app) as c: + r_off = await c.get("/api/device/v1/state", headers=_bearer(tok)) + w_off = await c.get("/auth/lock-widgets") + + assert r_off.json()["demo"] is False + assert not any(a.get("demo") for a in w_off.json()["agents"]) + + monkeypatch.undo() + + +# (j) avatar hash: null when no image, 16-hex when present, changes on rewrite. +@pytest.mark.asyncio +async def test_state_avatar_hash(vapp, monkeypatch): + from tinyagentos import agent_avatars as avatars + + app = vapp + app.state.config.agents = [ + {"name": "avatar-agent", "framework": "openclaw", "user_id": "u1"}, + ] + + with tempfile.TemporaryDirectory() as tmpdir: + monkeypatch.setattr(avatars, "LOCK_AVATAR_DIR", tmpdir) + tok = await _device(app, user_id="u1", scopes=("agents:read",)) + + async with _client(app) as c: + r = await c.get("/api/device/v1/state", headers=_bearer(tok)) + + assert r.status_code == 200, r.text + agent = r.json()["agents"][0] + assert agent["avatar"]["hash"] is None + + slug = avatars._avatar_slug("avatar-agent") + img_path = Path(tmpdir) / f"{slug}.jpg" + img_path.write_bytes(b"first-image-content") + async with _client(app) as c: + r = await c.get("/api/device/v1/state", headers=_bearer(tok)) + + assert r.status_code == 200, r.text + agent = r.json()["agents"][0] + h1 = agent["avatar"]["hash"] + assert h1 is not None + assert len(h1) == 16 + assert all(c in "0123456789abcdef" for c in h1) + + img_path.write_bytes(b"second-image-content-changed") + async with _client(app) as c: + r = await c.get("/api/device/v1/state", headers=_bearer(tok)) + + assert r.status_code == 200, r.text + agent = r.json()["agents"][0] + h2 = agent["avatar"]["hash"] + assert h2 is not None + assert h2 != h1 diff --git a/tinyagentos/auth_middleware.py b/tinyagentos/auth_middleware.py index 90d6071bb..adf6dfbfd 100644 --- a/tinyagentos/auth_middleware.py +++ b/tinyagentos/auth_middleware.py @@ -302,6 +302,9 @@ def _is_agent_notes_path(method: str, path: str) -> bool: # scope via device_scope(), this entry is what lets the Bearer past the gate. ("POST", re.compile(r"^/api/device/v1/voice$"), VOICE_STT), ("POST", re.compile(r"^/api/device/v1/voice/tts$"), VOICE_TTS), + # S2b: device state (agents read). Device-bearer only; the route names its + # own scope via device_scope(), this entry is what lets the Bearer past. + ("GET", re.compile(r"^/api/device/v1/state$"), AGENTS_READ), ) # Device-bearer routes that are NOT in _DEVICE_BEARER_PATHS because they sit in diff --git a/tinyagentos/routes/__init__.py b/tinyagentos/routes/__init__.py index 080f194b1..e34449bcc 100644 --- a/tinyagentos/routes/__init__.py +++ b/tinyagentos/routes/__init__.py @@ -109,6 +109,11 @@ def register_all_routers(app): from tinyagentos.routes.device_voice import router as device_voice_router app.include_router(device_voice_router) + # Device state (S2b): device-bearer only, so CSRF-exempt like the pair + # routes (the bearer is not an ambient cookie credential). + from tinyagentos.routes.device_state import router as device_state_router + app.include_router(device_state_router) + from tinyagentos.routes.observatory import router as observatory_router app.include_router(observatory_router, dependencies=_csrf) diff --git a/tinyagentos/routes/auth.py b/tinyagentos/routes/auth.py index 7182d3e37..4e9fd5c32 100644 --- a/tinyagentos/routes/auth.py +++ b/tinyagentos/routes/auth.py @@ -9359,93 +9359,58 @@ def _demo_refresh_in_ms(next_change_candidates: list[int]) -> int | None: return max(1000, min(15000, min(next_change_candidates) + 150)) -@router.get("/lock-widgets") -async def lock_widgets(request: Request): - """Agent activity + scheduled tasks for the lock screen. Console-only. +async def assemble_lock_agents(request: Request, owner_id: str | None = None) -> list[dict]: + """Shared async assembler for lock-screen agent islands. - This is rendered BEFORE sign-in, which is exactly why it is narrow: it - returns NAMES, STATUSES AND COUNTS and nothing else. The agent config is - never serialised here -- it carries per-agent LLM keys, and this endpoint is - reachable without a session. The console gate is the second half of that - containment: a LAN browser gets 403 and learns nothing, so the exposure is - the same one a phone lock screen already makes to whoever is holding it. - """ - if not _request_is_console(request): - return JSONResponse({"error": "console only"}, status_code=403) + Returns the base agent list (configured + demo + device-live, with pending + decisions and the demo decision attached). The caller is responsible for + adding the system agent, tasks, and wrapping in the final payload. + Configured agents are filtered by *owner_id* when supplied (include an + entry only when ``user_id`` is missing or equals *owner_id*). Demo agents + and device-live agents are not owner-filtered: they are global content. + """ agents: list[dict] = [] try: configured = request.app.state.config.agents or [] except AttributeError: configured = [] - # Container status is best-effort: on a host with no container runtime the - # import or the call raises, and a lock screen that 500s because the phone - # has no LXC is worse than one that simply shows no status. + status_by_name: dict[str, str] = {} try: from tinyagentos.containers import list_containers for c in await list_containers(prefix="taos-agent-"): status_by_name[c.name.removeprefix("taos-agent-")] = c.status - except Exception: # noqa: BLE001 - any runtime absence degrades to "no status" + except Exception: # noqa: BLE001 status_by_name = {} for entry in configured: - name = entry.get("name") if isinstance(entry, dict) else str(entry) - if not name: - continue - framework = "" if isinstance(entry, dict): + entry_uid = entry.get("user_id") + if owner_id is not None and entry_uid and entry_uid != owner_id: + continue + name = entry.get("name") framework = str(entry.get("framework") or entry.get("harness") or "") + configured_status = entry.get("status") if isinstance(entry, dict) else None + else: + name = str(entry) + framework = "" + configured_status = None + if not name: + continue agents.append({ "name": str(name), "framework": framework.lower(), "framework_icon": _framework_icon(framework), - "status": status_by_name.get(str(name), ""), + "status": status_by_name.get(str(name)) or (str(configured_status) if configured_status else ""), "avatar": _avatar_url(str(name)), }) - # The OS's own agent is pinned to the top and is not one of the configured - # ones: it is part of the device rather than something the user added. It - # carries the product mark rather than a monogram, and the harness badge - # reflects the actual runtime adapter the agent goes through. - # Its status reads "Idle" (Jay; it used to say "On device", i.e. WHERE it - # is). This endpoint runs pre-auth and has no cheap, truthful way to read - # the agent's activity, and "Idle" is also a resting status to the page, - # so the island draws at rest. - agents.insert(0, { - "name": "taOS Agent", - "framework": system_agent_framework(request.app.state), - "framework_icon": _framework_icon(system_agent_framework(request.app.state)), - "status": "Idle", - "avatar": "/static/taos-logo.png", - "system": True, - }) - - tasks: list[dict] = [] - try: - scheduler = request.app.state.scheduler - for task in await scheduler.list_tasks(): - item = task if isinstance(task, dict) else {} - tasks.append({ - "name": str(item.get("name", "") or ""), - "schedule": str(item.get("schedule", "") or ""), - "agent": str(item.get("agent_name", "") or ""), - }) - except Exception: # noqa: BLE001 - no scheduler on this host: show no tasks - tasks = [] - - # Demo override. OFF unless TAOS_LOCK_DEMO_AGENTS is set, and it only ever - # ADDS named placeholders to this one read-only lock-screen endpoint -- it - # writes nothing, creates no agents and changes no other surface. It exists - # so a demo machine can show a populated lock screen without standing up - # three real container-backed agents first; anything it lists is a - # placeholder, not a running process. demo = _demo_value("TAOS_LOCK_DEMO_AGENTS", request) if demo: existing = {a["name"] for a in agents} for raw in demo.split(","): - # "Name", "Name:framework" or "Name:framework:status text" parts = [seg.strip() for seg in raw.split(":")] label = parts[0] if parts else "" if label and label not in existing: @@ -9455,26 +9420,9 @@ async def lock_widgets(request: Request): "framework_icon": _framework_icon(parts[1] if len(parts) > 1 else ""), "status": parts[2] if len(parts) > 2 and parts[2] else "running", "avatar": _avatar_url(label), - # Marked at creation so nothing downstream has to work out - # which of these entries is a placeholder by elimination. "demo": True, }) - # Rotating "current task": each demo agent with a script (see - # _DEMO_TASK_SCRIPTS) cycles through it on its OWN independent clock - # -- its own pace, phase and per-entry jitter (all stable, derived - # from its name) -- instead of sitting on its configured status - # forever. Computed fresh from the clock on every request -- the - # SERVER is authoritative, so a 15s poll landing while the client is - # mid-animation can never revert what the client is showing, and a - # late or early poll just sees whatever the schedule says right now - # rather than drifting out of sync with it. - # - # Only an agent the demo config says is BUSY rotates. The stats panel - # reads the same (name, busy) pairs (_demo_agent_specs), so a scripted - # agent configured with a resting status must stay at rest here too -- - # rotating it into "Replying to 12 comments" would show a working - # island for an agent the stats model has idle. busy = {name for name, is_busy in _demo_agent_specs(request) if is_busy} now = _demo_task_clock() for agent in agents: @@ -9487,37 +9435,17 @@ async def lock_widgets(request: Request): if remaining is not None: agent["next_change_ms"] = int(round(remaining * 1000)) - # Live device agents: physical boards that are plugged in right now. - # - # MERGED HERE rather than served from their own endpoint because the lock - # screen already polls this one and reconciles the islands by key -- a - # second list would mean a second poll and two painters racing over the - # same row. Their keys are namespaced (`device:`) so a board can - # never collide with a TAOS_LOCK_DEMO_AGENTS placeholder of the same name. if _device_agents_enabled(request): for entry in _device_live(): agents.append(_device_island(entry)) - # Pending decisions. An agent that is blocked waiting on a human is the one - # thing on this screen that is actually ASKING for something, so it gets the - # attention ring -- everything else here is status. Best-effort for the same - # reason as the container statuses: a host with no decision store should - # show a lock screen, not a 500. - # - # Only the question and its options cross the pre-auth boundary, never the - # decision's context or notes: the question is a one-line prompt the holder - # of the phone needs in order to know the phone wants them, while the - # context is free text an agent may have filled with anything. pending: list[dict] = [] try: store = request.app.state.decision_store pending = await store.list(status="pending", limit=20) - except Exception: # noqa: BLE001 - no decision store on this host: no ring + except Exception: # noqa: BLE001 pending = [] - # The store returns newest-first. If an agent has asked twice, the question - # to surface is the one that has been WAITING longest, so walk oldest-first - # and keep the first hit per agent. by_agent: dict[str, dict] = {} for d in reversed(pending): agent_key = str(d.get("from_agent") or "").strip().lower() @@ -9545,11 +9473,6 @@ def _attach(agent: dict) -> None: for agent in agents: _attach(agent) - # Demo decision. Same single flag as the demo agents, and it is attached to - # an agent that is already a placeholder -- it never marks a REAL agent as - # waiting on a human, because a fabricated ring on a real agent would be a - # lie about the state of the machine. It carries no decision id, which is - # what the client uses to tell a demo prompt from an answerable one. if demo and not any(a.get("attention") for a in agents): want = _demo_value("TAOS_LOCK_DEMO_DECISION_AGENT", request).lower() target = None @@ -9571,18 +9494,62 @@ def _attach(agent: dict) -> None: "demo": True, } - # Anything with a status that is not an explicit resting word is doing - # something -- the demo statuses are free text ("Drafting replies"), so an - # equality test against "running" would report every busy agent as idle. + return agents + + +@router.get("/lock-widgets") +async def lock_widgets(request: Request): + """Agent activity + scheduled tasks for the lock screen. Console-only. + + This is rendered BEFORE sign-in, which is exactly why it is narrow: it + returns NAMES, STATUSES AND COUNTS and nothing else. The agent config is + never serialised here -- it carries per-agent LLM keys, and this endpoint is + reachable without a session. The console gate is the second half of that + containment: a LAN browser gets 403 and learns nothing, so the exposure is + the same one a phone lock screen already makes to whoever is holding it. + """ + if not _request_is_console(request): + return JSONResponse({"error": "console only"}, status_code=403) + + agents = await assemble_lock_agents(request) + + # The OS's own agent is pinned to the top and is not one of the configured + # ones: it is part of the device rather than something the user added. It + # carries the product mark rather than a monogram, and the harness badge + # reflects the actual runtime adapter the agent goes through. + # Its status reads "Idle" (Jay; it used to say "On device", i.e. WHERE it + # is). This endpoint runs pre-auth and has no cheap, truthful way to read + # the agent's activity, and "Idle" is also a resting status to the page, + # so the island draws at rest. + agents.insert(0, { + "name": "taOS Agent", + "framework": system_agent_framework(request.app.state), + "framework_icon": _framework_icon(system_agent_framework(request.app.state)), + "status": "Idle", + "avatar": "/static/taos-logo.png", + "system": True, + }) + + tasks: list[dict] = [] + try: + scheduler = request.app.state.scheduler + for task in await scheduler.list_tasks(): + item = task if isinstance(task, dict) else {} + tasks.append({ + "name": str(item.get("name", "") or ""), + "schedule": str(item.get("schedule", "") or ""), + "agent": str(item.get("agent_name", "") or ""), + }) + except Exception: # noqa: BLE001 - no scheduler on this host: show no tasks + tasks = [] + + demo = _demo_enabled(request) + resting = {"", "stopped", "idle", "exited", "error"} running = sum( 1 for a in agents if not a.get("system") and a["status"].strip().lower() not in resting ) - # NO cap on the islands (Jay: "there shouldnt be a cap"); the feed scrolls. - # A cap of six silently cut off a plugged-in board, which is appended last, - # whenever five demo agents were configured. New agents go at the BOTTOM, - # in arrival order (Jay; custom ordering comes later). visible = agents payload = { "agents": visible, @@ -9592,9 +9559,6 @@ def _attach(agent: dict) -> None: "task_total": len(tasks), "threads": bool(demo), } - # Only agents actually SENT can schedule the client's next fetch -- a - # rotation happening off the sent list would tell the client to poll - # sooner for a change it could never paint anyway. refresh_in_ms = _demo_refresh_in_ms( [a["next_change_ms"] for a in visible if "next_change_ms" in a] ) diff --git a/tinyagentos/routes/device_state.py b/tinyagentos/routes/device_state.py new file mode 100644 index 000000000..dcc9b2ea0 --- /dev/null +++ b/tinyagentos/routes/device_state.py @@ -0,0 +1,113 @@ +"""GET /api/device/v1/state (device bearer, scope agents:read). + +Device-bearer only: returns the owner-filtered agent list for the paired +device, plus the server version, current time, and demo flag. +""" +from __future__ import annotations + +import tinyagentos +import time + +from fastapi import APIRouter, Depends, HTTPException, Request + +from tinyagentos.agent_avatars import avatar_hash +from tinyagentos.device_auth import device_scope +from tinyagentos.device_scopes import AGENTS_READ +from tinyagentos.routes.auth import assemble_lock_agents, _demo_enabled + +router = APIRouter() + +_CAP_NAME = 48 +_CAP_STATUS = 120 +_CAP_LAST_RECAP = 180 +_CAP_QUESTION = 280 +_CAP_OPTION = 40 +_ELLIPSIS = "\u2026" + + +def _cap(text: str, limit: int) -> str: + if len(text) <= limit: + return text + cut = limit - len(_ELLIPSIS) + if cut <= 0: + return _ELLIPSIS[:limit] + return text[:cut].rstrip() + _ELLIPSIS + + +def _hue_for(name: str) -> int: + h = 0 + for ch in name: + h = (h * 31 + ord(ch)) % 360 + return h + + +async def _last_recap(agent_name: str, agent_messages) -> str: + try: + rows = await agent_messages.get_messages(agent_name, limit=1) + except Exception: # noqa: BLE001 + return "" + if not rows: + return "" + row = rows[0] + text = str(row.get("message") or "") + return _cap(text, _CAP_LAST_RECAP) + + +def _decision_for_agent(agent: dict) -> dict | None: + dec = agent.get("decision") + if not dec: + return None + options = dec.get("options") or [] + capped = [_cap(str(o), _CAP_OPTION) for o in options if str(o).strip()] + return { + "id": str(dec.get("id") or ""), + "question": _cap(str(dec.get("question") or ""), _CAP_QUESTION), + "options": capped[:4], + } + + +async def _transform_agent(agent: dict, agent_messages) -> dict: + name = _cap(str(agent.get("name") or ""), _CAP_NAME) + status = _cap(str(agent.get("status") or ""), _CAP_STATUS) + framework = str(agent.get("framework") or "").lower() + hue = _hue_for(name) + ahash = avatar_hash(name) + return { + "name": name, + "status": status, + "framework": framework, + "avatar": { + "hue": hue, + "hash": ahash, + }, + "attention": bool(agent.get("attention")), + "last_recap": await _last_recap(agent.get("name", ""), agent_messages), + "decision": _decision_for_agent(agent), + } + + +@router.get("/api/device/v1/state") +async def device_state(request: Request, _device: dict = Depends(device_scope(AGENTS_READ))): + try: + owner_id = _device.get("user_id") + except AttributeError: + owner_id = None + + base_agents = await assemble_lock_agents(request, owner_id=owner_id) + + agent_messages = getattr(request.app.state, "agent_messages", None) + + transformed = [ + await _transform_agent(a, agent_messages) + for a in base_agents + if not a.get("system") + ] + + return { + "agents": transformed, + "server": { + "version": getattr(tinyagentos, "__version__", "unknown"), + }, + "time": time.time(), + "demo": _demo_enabled(request), + } From 5904ce88c4fd9ce21dd4dcd6f8ef2c282396ddd1 Mon Sep 17 00:00:00 2001 From: jaylfc Date: Mon, 5 Oct 2026 04:16:16 +0000 Subject: [PATCH 3/7] fix: address review blockers on /api/device/v1/state - filter device-live agents out of assemble_lock_agents when owner_id is set so /api/device/v1/state only returns the device owner's agents - use the original (uncapped) agent name for avatar hash and hue in _transform_agent so the hash matches the avatar file and hue is stable - drop the per-response 4-option cap from _decision_for_agent; the card specifies a 40-char per-option cap but no count limit - use _demo_value directly in device_state.py instead of importing the private _demo_enabled - use the monkeypatch fixture in test_settings_demo_switch_takes_state_down instead of creating a manual pytest.MonkeyPatch Docs-Reviewed: route registration mirrors existing device-bearer patterns; changelog fragment changelog.d/tsk-ypsu2z-device-v1-state.md already present from the route introduction commit --- docs/routes.d/16-device-v1.md | 4 ++++ tests/test_device_v1_state.py | 5 +---- tinyagentos/routes/auth.py | 5 +++-- tinyagentos/routes/device_state.py | 14 ++++++++------ 4 files changed, 16 insertions(+), 12 deletions(-) diff --git a/docs/routes.d/16-device-v1.md b/docs/routes.d/16-device-v1.md index f4269e856..1b03b38d4 100644 --- a/docs/routes.d/16-device-v1.md +++ b/docs/routes.d/16-device-v1.md @@ -15,6 +15,10 @@ version, current time, and demo flag. **Response 200:** +All string fields are server-side capped with a trailing ellipsis when they +exceed the limit: `name` 48, `status` 120, `last_recap` 180, `question` 280, +each `option` 40. + ```json { "agents": [ diff --git a/tests/test_device_v1_state.py b/tests/test_device_v1_state.py index 2daaacf80..b161f0854 100644 --- a/tests/test_device_v1_state.py +++ b/tests/test_device_v1_state.py @@ -154,11 +154,10 @@ async def test_state_demo_flag_and_unanswerable_decision(vapp, monkeypatch): # (h) Settings demo switch takes state down. @pytest.mark.asyncio -async def test_settings_demo_switch_takes_state_down(vapp): +async def test_settings_demo_switch_takes_state_down(vapp, monkeypatch): from tinyagentos.demo_mode import write_demo_mode app = vapp - monkeypatch = pytest.MonkeyPatch() monkeypatch.setenv("TAOS_LOCK_DEMO_AGENTS", "DemoB:openclaw:Drafting") monkeypatch.setenv("TAOS_LOCK_DEMO_DECISION", "Ship it?") monkeypatch.setenv("TAOS_LOCK_DEMO_DECISION_AGENT", "DemoB") @@ -184,8 +183,6 @@ async def test_settings_demo_switch_takes_state_down(vapp): assert r_off.json()["demo"] is False assert not any(a.get("demo") for a in w_off.json()["agents"]) - monkeypatch.undo() - # (j) avatar hash: null when no image, 16-hex when present, changes on rewrite. @pytest.mark.asyncio diff --git a/tinyagentos/routes/auth.py b/tinyagentos/routes/auth.py index 4e9fd5c32..09645c2c4 100644 --- a/tinyagentos/routes/auth.py +++ b/tinyagentos/routes/auth.py @@ -9368,7 +9368,8 @@ async def assemble_lock_agents(request: Request, owner_id: str | None = None) -> Configured agents are filtered by *owner_id* when supplied (include an entry only when ``user_id`` is missing or equals *owner_id*). Demo agents - and device-live agents are not owner-filtered: they are global content. + are global content. Device-live agents are only included when *owner_id* + is ``None`` (they carry no owner field and must not leak across owners). """ agents: list[dict] = [] try: @@ -9435,7 +9436,7 @@ async def assemble_lock_agents(request: Request, owner_id: str | None = None) -> if remaining is not None: agent["next_change_ms"] = int(round(remaining * 1000)) - if _device_agents_enabled(request): + if _device_agents_enabled(request) and owner_id is None: for entry in _device_live(): agents.append(_device_island(entry)) diff --git a/tinyagentos/routes/device_state.py b/tinyagentos/routes/device_state.py index dcc9b2ea0..1a0e81fd4 100644 --- a/tinyagentos/routes/device_state.py +++ b/tinyagentos/routes/device_state.py @@ -13,7 +13,7 @@ from tinyagentos.agent_avatars import avatar_hash from tinyagentos.device_auth import device_scope from tinyagentos.device_scopes import AGENTS_READ -from tinyagentos.routes.auth import assemble_lock_agents, _demo_enabled +from tinyagentos.routes.auth import assemble_lock_agents, _demo_value router = APIRouter() @@ -62,16 +62,17 @@ def _decision_for_agent(agent: dict) -> dict | None: return { "id": str(dec.get("id") or ""), "question": _cap(str(dec.get("question") or ""), _CAP_QUESTION), - "options": capped[:4], + "options": capped, } async def _transform_agent(agent: dict, agent_messages) -> dict: - name = _cap(str(agent.get("name") or ""), _CAP_NAME) + original_name = agent.get("name") or "" + name = _cap(str(original_name), _CAP_NAME) status = _cap(str(agent.get("status") or ""), _CAP_STATUS) framework = str(agent.get("framework") or "").lower() - hue = _hue_for(name) - ahash = avatar_hash(name) + hue = _hue_for(original_name) + ahash = avatar_hash(original_name) return { "name": name, "status": status, @@ -109,5 +110,6 @@ async def device_state(request: Request, _device: dict = Depends(device_scope(AG "version": getattr(tinyagentos, "__version__", "unknown"), }, "time": time.time(), - "demo": _demo_enabled(request), + "demo": bool(_demo_value("TAOS_LOCK_DEMO_AGENTS", request).strip()), } + From 19283a77bc1ac3449258e7cee39bc5a63a973dce Mon Sep 17 00:00:00 2001 From: jaylfc Date: Mon, 5 Oct 2026 04:33:09 +0000 Subject: [PATCH 4/7] fix: keep lock-widgets assembler pure and align demo flag helper Restore the original pure status lookup in assemble_lock_agents so /auth/lock-widgets does not fall back to the configured agent status. Use _demo_enabled in /api/device/v1/state for the demo flag, matching the rest of the lock-screen paths. ### Tests - Add test_lock_widgets_status_does_not_fall_back_to_configured_status - Fix test_state_demo_flag_and_unanswerable_decision to assert a real demo decision on DemoA - Update test_state_caps_long_strings to mock live container status ### Red proof ```text FAILED tests/test_device_v1_state.py::test_lock_widgets_status_does_not_fall_back_to_configured_status - AssertionError assert 'secret-config-status' == '' + secret-config-status ``` 1 failed, 6 deselected in 5.59s ### Green proof ```text . [100%] 1 passed, 6 deselected in 2.97s ``` 143 passed in 61.18s ### py_compile ```text py_compile clean ``` Docs-Reviewed: route registration mirrors existing /auth/lock-widgets and /api/device/v1/state patterns; no agent-coordination.md change needed --- .../tsk-3vxguh-lock-widgets-pure-status.md | 4 ++ tests/test_device_v1_state.py | 50 ++++++++++++++++--- tinyagentos/routes/auth.py | 4 +- tinyagentos/routes/device_state.py | 4 +- 4 files changed, 50 insertions(+), 12 deletions(-) create mode 100644 changelog.d/tsk-3vxguh-lock-widgets-pure-status.md diff --git a/changelog.d/tsk-3vxguh-lock-widgets-pure-status.md b/changelog.d/tsk-3vxguh-lock-widgets-pure-status.md new file mode 100644 index 000000000..380734e54 --- /dev/null +++ b/changelog.d/tsk-3vxguh-lock-widgets-pure-status.md @@ -0,0 +1,4 @@ +### Fixed + +- `/auth/lock-widgets` now returns an empty string for an agent's status when no live container is present, instead of falling back to the configured `status` value. This keeps the pre-auth console endpoint pure and prevents config-derived status leakage before sign-in. +- `GET /api/device/v1/state` `demo` flag now uses the existing `_demo_enabled` helper instead of reading the env flag directly, matching the rest of the lock-screen demo paths. diff --git a/tests/test_device_v1_state.py b/tests/test_device_v1_state.py index b161f0854..e208d1dc7 100644 --- a/tests/test_device_v1_state.py +++ b/tests/test_device_v1_state.py @@ -39,6 +39,29 @@ async def _device(app, user_id="u1", platform="ios", scopes=("agents:read",)): return d["scoped_token"] +# RED-FIRST: configure an agent with a non-empty config status and no live +# status, then prove /auth/lock-widgets does NOT leak the configured status. +@pytest.mark.asyncio +async def test_lock_widgets_status_does_not_fall_back_to_configured_status(vapp, monkeypatch): + from tinyagentos.routes import auth as auth_mod + + app = vapp + monkeypatch.setattr(auth_mod, "_request_is_console", lambda _r: True) + + app.state.config.agents = [ + {"name": "secret-status-agent", "framework": "openclaw", "user_id": "u1", "status": "secret-config-status"}, + ] + + async with _client(app) as c: + r = await c.get("/auth/lock-widgets") + + assert r.status_code == 200, r.text + data = r.json() + agent = next(a for a in data["agents"] if a["name"] == "secret-status-agent") + assert agent["status"] == "" + assert "secret-config-status" not in r.text + + # (a) two owners, each with agents; a device paired to owner 1 sees only owner 1's agents. @pytest.mark.asyncio async def test_state_returns_only_owner_agents(vapp): @@ -67,12 +90,25 @@ async def test_state_returns_only_owner_agents(vapp): # (b) a 500-char recap arrives with length <= 180 and ends with the ellipsis. @pytest.mark.asyncio -async def test_state_caps_long_strings(vapp): +async def test_state_caps_long_strings(vapp, monkeypatch): + from tinyagentos import containers + app = vapp long_status = "x" * 500 long_recap = "r" * 500 + + class _MockContainer: + def __init__(self, name, status): + self.name = name + self.status = status + + async def _mock_list_containers(prefix=None): + return [_MockContainer("taos-agent-cap-agent", long_status)] + + monkeypatch.setattr(containers, "list_containers", _mock_list_containers, raising=False) + app.state.config.agents = [ - {"name": "cap-agent", "framework": "openclaw", "user_id": "u1", "status": long_status}, + {"name": "cap-agent", "framework": "openclaw", "user_id": "u1"}, ] await app.state.agent_messages.send( @@ -138,11 +174,11 @@ async def test_state_demo_flag_and_unanswerable_decision(vapp, monkeypatch): widgets_names = {a["name"] for a in widgets_data["agents"] if not a.get("system")} assert state_names == widgets_names - for a in state_data["agents"]: - dec = a.get("decision") - if dec and dec.get("id") == "": - assert dec["id"] == "" - break + demo_a = next((a for a in state_data["agents"] if a["name"] == "DemoA"), None) + assert demo_a is not None + dec = demo_a.get("decision") + assert dec is not None + assert dec["id"] == "" write_demo_mode(app.state.data_dir, False) diff --git a/tinyagentos/routes/auth.py b/tinyagentos/routes/auth.py index 09645c2c4..066a6c6ef 100644 --- a/tinyagentos/routes/auth.py +++ b/tinyagentos/routes/auth.py @@ -9393,18 +9393,16 @@ async def assemble_lock_agents(request: Request, owner_id: str | None = None) -> continue name = entry.get("name") framework = str(entry.get("framework") or entry.get("harness") or "") - configured_status = entry.get("status") if isinstance(entry, dict) else None else: name = str(entry) framework = "" - configured_status = None if not name: continue agents.append({ "name": str(name), "framework": framework.lower(), "framework_icon": _framework_icon(framework), - "status": status_by_name.get(str(name)) or (str(configured_status) if configured_status else ""), + "status": status_by_name.get(str(name), ""), "avatar": _avatar_url(str(name)), }) diff --git a/tinyagentos/routes/device_state.py b/tinyagentos/routes/device_state.py index 1a0e81fd4..3294bcb0e 100644 --- a/tinyagentos/routes/device_state.py +++ b/tinyagentos/routes/device_state.py @@ -13,7 +13,7 @@ from tinyagentos.agent_avatars import avatar_hash from tinyagentos.device_auth import device_scope from tinyagentos.device_scopes import AGENTS_READ -from tinyagentos.routes.auth import assemble_lock_agents, _demo_value +from tinyagentos.routes.auth import assemble_lock_agents, _demo_enabled router = APIRouter() @@ -110,6 +110,6 @@ async def device_state(request: Request, _device: dict = Depends(device_scope(AG "version": getattr(tinyagentos, "__version__", "unknown"), }, "time": time.time(), - "demo": bool(_demo_value("TAOS_LOCK_DEMO_AGENTS", request).strip()), + "demo": _demo_enabled(request), } From b1799c601e4af1661f82bd1a0151a5e567ce1e1c Mon Sep 17 00:00:00 2001 From: jaylfc Date: Mon, 5 Oct 2026 05:31:57 +0000 Subject: [PATCH 5/7] fix: filter pending decisions by device owner in assemble_lock_agents Filter pending decisions by owner_id in assemble_lock_agents so a device paired to one owner cannot see another owner's pending decisions through a shared (no user_id) agent. DecisionStore.list already accepts user_id; with owner_id None (the console /auth/lock-widgets path) the filter is not applied, so lock-widgets behaviour is unchanged. ```text FAILED tests/test_device_v1_state.py::test_state_does_not_leak_other_owners_decision - AssertionError 1 failed, 7 deselected in 4.94s ``` Green: tests/test_device_v1_state.py, tests/test_routes_doc.py, tests/test_demo_mode.py, tests/test_taos_agent_config.py, tests/test_taos_agent_picoclaw.py, tests/test_lock_demo_task_rotation.py (148 passed). py_compile clean. Docs-Reviewed: routes docs updated in docs/routes.d/16-device-v1.md and docs/routes.md; agent-coordination.md does not cover this endpoint. --- .../tsk-22jk4c-device-state-owner-filter.md | 3 ++ docs/routes.d/16-device-v1.md | 2 +- docs/routes.md | 46 +++++++++++++++++++ tests/test_device_v1_state.py | 28 +++++++++++ tinyagentos/routes/auth.py | 2 +- 5 files changed, 79 insertions(+), 2 deletions(-) create mode 100644 changelog.d/tsk-22jk4c-device-state-owner-filter.md diff --git a/changelog.d/tsk-22jk4c-device-state-owner-filter.md b/changelog.d/tsk-22jk4c-device-state-owner-filter.md new file mode 100644 index 000000000..0e331f2fa --- /dev/null +++ b/changelog.d/tsk-22jk4c-device-state-owner-filter.md @@ -0,0 +1,3 @@ +### Fixed + +- `GET /api/device/v1/state` now filters pending decisions by device owner, preventing a device paired to one owner from seeing another owner's pending decisions through a shared (no `user_id`) agent. diff --git a/docs/routes.d/16-device-v1.md b/docs/routes.d/16-device-v1.md index 1b03b38d4..9f28a9a22 100644 --- a/docs/routes.d/16-device-v1.md +++ b/docs/routes.d/16-device-v1.md @@ -50,4 +50,4 @@ each `option` 40. **Error codes:** - `401` -- missing or invalid device bearer. -- `403` -- `{"error": "device_scope_missing", "scope": "agents:read"}` when the device token lacks the required scope. +- `403` -- FastAPI's wrapper: `{"detail": {"error": "device_scope_missing", "scope": "agents:read"}}` when the device token lacks the required scope. diff --git a/docs/routes.md b/docs/routes.md index 449348272..56973870e 100644 --- a/docs/routes.md +++ b/docs/routes.md @@ -419,3 +419,49 @@ are bearer-authenticated. The following routes are public (no bearer token): - `POST /api/devices/pair-requests` - `GET /api/devices/pair-requests/{pair_request_id}` + +### GET /api/device/v1/state + +**Scope:** `agents:read` (device bearer). + +Returns the owner-filtered agent list for the paired device, plus the server +version, current time, and demo flag. + +**Response 200:** + +All string fields are server-side capped with a trailing ellipsis when they +exceed the limit: `name` 48, `status` 120, `last_recap` 180, `question` 280, +each `option` 40. + +```json +{ + "agents": [ + { + "name": "alice-agent", + "status": "running", + "framework": "openclaw", + "avatar": { + "hue": 123, + "hash": "abc123def4567890" + }, + "attention": true, + "last_recap": "recapped message...", + "decision": { + "id": "", + "question": "Ship it?", + "options": ["Approve", "Deny"] + } + } + ], + "server": { + "version": "1.0.0-beta.55" + }, + "time": 1727640000.123, + "demo": false +} +``` + +**Error codes:** + +- `401` -- missing or invalid device bearer. +- `403` -- FastAPI's wrapper: `{"detail": {"error": "device_scope_missing", "scope": "agents:read"}}` when the device token lacks the required scope. diff --git a/tests/test_device_v1_state.py b/tests/test_device_v1_state.py index e208d1dc7..49d8e19dd 100644 --- a/tests/test_device_v1_state.py +++ b/tests/test_device_v1_state.py @@ -263,3 +263,31 @@ async def test_state_avatar_hash(vapp, monkeypatch): h2 = agent["avatar"]["hash"] assert h2 is not None assert h2 != h1 + + +# (k) a shared agent (no user_id) must not leak another owner's pending decision. +@pytest.mark.asyncio +async def test_state_does_not_leak_other_owners_decision(vapp): + app = vapp + app.state.config.agents = [ + {"name": "shared-agent", "framework": "openclaw", "user_id": ""}, + ] + + await app.state.decision_store.create( + from_agent="shared-agent", + question="u2's secret decision", + type="approve_deny", + user_id="u2", + options=[{"label": "Approve", "value": "approve"}, {"label": "Deny", "value": "deny"}], + ) + + tok = await _device(app, user_id="u1", scopes=("agents:read",)) + + async with _client(app) as c: + r = await c.get("/api/device/v1/state", headers=_bearer(tok)) + + assert r.status_code == 200, r.text + data = r.json() + agent = next(a for a in data["agents"] if a["name"] == "shared-agent") + assert agent.get("decision") is None + assert "u2's secret decision" not in r.text diff --git a/tinyagentos/routes/auth.py b/tinyagentos/routes/auth.py index 066a6c6ef..3a0f485ed 100644 --- a/tinyagentos/routes/auth.py +++ b/tinyagentos/routes/auth.py @@ -9441,7 +9441,7 @@ async def assemble_lock_agents(request: Request, owner_id: str | None = None) -> pending: list[dict] = [] try: store = request.app.state.decision_store - pending = await store.list(status="pending", limit=20) + pending = await store.list(status="pending", user_id=owner_id, limit=20) except Exception: # noqa: BLE001 pending = [] From a0682cf29096c9c5bd00588cfb9c574953f320ee Mon Sep 17 00:00:00 2001 From: jaylfc Date: Mon, 5 Oct 2026 07:12:35 +0000 Subject: [PATCH 6/7] # RED-PROOF Fenced block from the FAILING run BEFORE the fix: ``` FAILED tests/test_device_v1_events.py::test_events_emits_upsert_keyed_by_name - assert 404 == 200 FAILED tests/test_device_v1_events.py::test_events_last_event_id_resume - assert 404 == 200 FAILED tests/test_device_v1_events.py::test_events_stale_id_gets_snapshot - assert 404 == 200 FAILED tests/test_device_v1_events.py::test_revoke_closes_stream - assert 404 == 200 FAILED tests/test_device_v1_events.py::test_stream_stops_demo_after_switch_off - assert 404 == 200 FAILED tests/test_device_v1_events.py::test_stream_upsert_on_avatar_change - assert 404 == 200 ``` Docs-Reviewed: executor rescue commit -- the model left these edits uncommitted and wrote no commit message, so doc drift was NOT assessed by the model; the lead reviews docs at PR time. --- tests/test_device_v1_events.py | 350 +++++++++++++++++++++++++++++ tinyagentos/auth_middleware.py | 3 + tinyagentos/routes/device_state.py | 218 +++++++++++++++++- 3 files changed, 570 insertions(+), 1 deletion(-) create mode 100644 tests/test_device_v1_events.py diff --git a/tests/test_device_v1_events.py b/tests/test_device_v1_events.py new file mode 100644 index 000000000..966296ded --- /dev/null +++ b/tests/test_device_v1_events.py @@ -0,0 +1,350 @@ +"""P0 S3: GET /api/device/v1/events SSE. + +RED-FIRST: this module is written BEFORE the route exists so the first run +must FAIL with 404. The FAIL block is captured in RED-PROOF.md. +""" +from __future__ import annotations + +import asyncio +import json +import time +from pathlib import Path + +import pytest +import pytest_asyncio +from httpx import ASGITransport, AsyncClient + +PLAIN = "http://testserver:6969" + + +def _bearer(tok): + return {"Authorization": f"Bearer {tok}"} + + +def _client(app, base=PLAIN): + return AsyncClient(transport=ASGITransport(app=app), base_url=base) + + +@pytest_asyncio.fixture +async def vapp(app, client): + return app + + +async def _device(app, user_id="u1", platform="ios", scopes=("agents:read",)): + st = app.state.device_store + d = await st.register(user_id=user_id, platform=platform) + if scopes is not None: + await st.set_scopes(d["device_id"], list(scopes)) + return d["scoped_token"] + + +async def _collect_chunks(r, max_chunks=20, timeout=3.0): + """Read at most *max_chunks* text chunks from an SSE response, or stop + after *timeout* seconds of inactivity.""" + chunks = [] + start = time.monotonic() + iterator = r.aiter_text() + while len(chunks) < max_chunks: + remaining = timeout - (time.monotonic() - start) + if remaining <= 0: + break + try: + task = asyncio.ensure_future(iterator.__anext__()) + done, pending = await asyncio.wait([task], timeout=min(remaining, 0.5)) + if pending: + for p in pending: + p.cancel() + try: + await asyncio.gather(*pending, return_exceptions=True) + except Exception: + pass + break + chunk = done.pop().result() + chunks.append(chunk) + except (StopAsyncIteration, Exception): + break + return chunks + + +def _parse_events(chunks): + """Return a list of (event_type, data_dict) from raw text chunks.""" + events = [] + for chunk in chunks: + current_type = None + current_data = {} + for line in chunk.splitlines(): + line = line.strip() + if not line: + if current_type is not None: + events.append((current_type, current_data)) + current_type = None + current_data = {} + continue + if line.startswith("event:"): + current_type = line.split(":", 1)[1].strip() + elif line.startswith("data:"): + try: + current_data = json.loads(line.split(":", 1)[1].strip()) + except Exception: + pass + elif line.startswith("id:"): + pass + elif line.startswith(":"): + pass + if current_type is not None: + events.append((current_type, current_data)) + return events + + +def _patch_intervals(monkeypatch): + from tinyagentos.routes import device_state as ds_mod + monkeypatch.setattr(ds_mod, "_POLL_INTERVAL_S", 0.1) + monkeypatch.setattr(ds_mod, "_HEARTBEAT_INTERVAL_S", 0.2) + + +# (d) SSE route exists and emits agent.upsert events keyed by name. +@pytest.mark.asyncio +async def test_events_emits_upsert_keyed_by_name(vapp, monkeypatch): + from tinyagentos.routes import auth as auth_mod + from tinyagentos.demo_mode import write_demo_mode + + _patch_intervals(monkeypatch) + + app = vapp + monkeypatch.setattr(auth_mod, "_request_is_console", lambda _r: True) + + write_demo_mode(app.state.data_dir, False) + app.state.config.agents = [ + {"name": "alice", "framework": "openclaw", "user_id": "u1", "status": "running"}, + ] + + tok = await _device(app, user_id="u1", scopes=("agents:read",)) + + async with _client(app) as c: + async with c.stream("GET", "/api/device/v1/events", headers=_bearer(tok)) as r: + assert r.status_code == 200 + chunks = await _collect_chunks(r, max_chunks=20, timeout=3.0) + + events = _parse_events(chunks) + upsert_names = {data["name"] for evt, data in events if evt == "agent.upsert" and "name" in data} + assert "alice" in upsert_names + + +# (e) Last-Event-ID resumes from the next event. +@pytest.mark.asyncio +async def test_events_last_event_id_resume(vapp, monkeypatch): + from tinyagentos.routes import auth as auth_mod + from tinyagentos.demo_mode import write_demo_mode + + _patch_intervals(monkeypatch) + + app = vapp + monkeypatch.setattr(auth_mod, "_request_is_console", lambda _r: True) + + write_demo_mode(app.state.data_dir, False) + app.state.config.agents = [ + {"name": "bob", "framework": "openclaw", "user_id": "u1", "status": "running"}, + ] + + tok = await _device(app, user_id="u1", scopes=("agents:read",)) + + first_id = None + async with _client(app) as c: + async with c.stream("GET", "/api/device/v1/events", headers=_bearer(tok)) as r: + assert r.status_code == 200 + chunks = await _collect_chunks(r, max_chunks=20, timeout=3.0) + for line in "".join(chunks).splitlines(): + if line.startswith("id:"): + first_id = line.split(":", 1)[1].strip() + + assert first_id is not None + + resumed = False + async with _client(app) as c: + async with c.stream("GET", "/api/device/v1/events", headers={**_bearer(tok), "Last-Event-ID": first_id}) as r: + assert r.status_code == 200 + chunks = await _collect_chunks(r, max_chunks=20, timeout=3.0) + for line in "".join(chunks).splitlines(): + if line.startswith("id:"): + resumed = True + break + + assert resumed + + +# (e) Stale Last-Event-ID gets a snapshot event. +@pytest.mark.asyncio +async def test_events_stale_id_gets_snapshot(vapp, monkeypatch): + from tinyagentos.routes import auth as auth_mod + from tinyagentos.demo_mode import write_demo_mode + + _patch_intervals(monkeypatch) + + app = vapp + monkeypatch.setattr(auth_mod, "_request_is_console", lambda _r: True) + + write_demo_mode(app.state.data_dir, False) + app.state.config.agents = [ + {"name": "carol", "framework": "openclaw", "user_id": "u1", "status": "running"}, + ] + + tok = await _device(app, user_id="u1", scopes=("agents:read",)) + + got_snapshot = False + async with _client(app) as c: + async with c.stream("GET", "/api/device/v1/events", headers={**_bearer(tok), "Last-Event-ID": "0"}) as r: + assert r.status_code == 200 + chunks = await _collect_chunks(r, max_chunks=20, timeout=3.0) + events = _parse_events(chunks) + for evt, data in events: + if evt == "snapshot": + got_snapshot = True + break + + assert got_snapshot + + +# (f) Revoking the device closes the stream. +@pytest.mark.asyncio +async def test_revoke_closes_stream(vapp, monkeypatch): + from tinyagentos.routes import auth as auth_mod + from tinyagentos.demo_mode import write_demo_mode + + _patch_intervals(monkeypatch) + + app = vapp + monkeypatch.setattr(auth_mod, "_request_is_console", lambda _r: True) + + write_demo_mode(app.state.data_dir, False) + app.state.config.agents = [ + {"name": "dave", "framework": "openclaw", "user_id": "u1", "status": "running"}, + ] + + st = app.state.device_store + d = await st.register(user_id="u1", platform="ios") + await st.set_scopes(d["device_id"], list(("agents:read",))) + tok = d["scoped_token"] + + closed = False + async with _client(app) as c: + async with c.stream("GET", "/api/device/v1/events", headers=_bearer(tok)) as r: + assert r.status_code == 200 + await st.revoke(d["device_id"]) + try: + await _collect_chunks(r, max_chunks=5, timeout=2.0) + except Exception: + closed = True + + assert closed + + +# (i) Demo stops on switch-off: no demo agent or demo decision after the flip. +@pytest.mark.asyncio +async def test_stream_stops_demo_after_switch_off(vapp, monkeypatch): + from tinyagentos.routes import auth as auth_mod + from tinyagentos.demo_mode import write_demo_mode + + _patch_intervals(monkeypatch) + + app = vapp + monkeypatch.setattr(auth_mod, "_request_is_console", lambda _r: True) + monkeypatch.setattr(auth_mod, "_demo_task_clock", lambda: 0.0) + + monkeypatch.setenv("TAOS_LOCK_DEMO_AGENTS", "DemoX:openclaw:Drafting") + monkeypatch.setenv("TAOS_LOCK_DEMO_DECISION", "Ship it?") + monkeypatch.setenv("TAOS_LOCK_DEMO_DECISION_AGENT", "DemoX") + write_demo_mode(app.state.data_dir, True) + + tok = await _device(app, user_id="u1", scopes=("agents:read",)) + + post_flip_demo = False + async with _client(app) as c: + async with c.stream("GET", "/api/device/v1/events", headers=_bearer(tok)) as r: + assert r.status_code == 200 + # consume initial events + pre_flip = await _collect_chunks(r, max_chunks=20, timeout=2.0) + # flip the switch off + write_demo_mode(app.state.data_dir, False) + # advance clock past any armed timer + monkeypatch.setattr(auth_mod, "_demo_task_clock", lambda: 9999.0) + post_flip = await _collect_chunks(r, max_chunks=20, timeout=2.0) + for chunk in post_flip: + for evt_line in chunk.splitlines(): + if evt_line.startswith("event:"): + pass + elif evt_line.startswith("data:"): + try: + data = json.loads(evt_line.split(":", 1)[1].strip()) + except Exception: + continue + if isinstance(data, dict): + if data.get("demo") is True: + post_flip_demo = True + dec = data.get("decision") or {} + if dec.get("demo") is True or dec.get("id") == "": + if data.get("name") == "DemoX": + post_flip_demo = True + + assert not post_flip_demo + + +# (j2) Avatar change triggers agent.upsert with new hash. +@pytest.mark.asyncio +async def test_stream_upsert_on_avatar_change(vapp, monkeypatch): + from tinyagentos.routes import auth as auth_mod + from tinyagentos import agent_avatars as avatars + from tinyagentos.demo_mode import write_demo_mode + + _patch_intervals(monkeypatch) + + app = vapp + monkeypatch.setattr(auth_mod, "_request_is_console", lambda _r: True) + + write_demo_mode(app.state.data_dir, False) + app.state.config.agents = [ + {"name": "eve-agent", "framework": "openclaw", "user_id": "u1", "status": "running"}, + ] + + tok = await _device(app, user_id="u1", scopes=("agents:read",)) + + with monkeypatch.context() as mp: + tmpdir = Path(app.state.data_dir) / "avatars" + tmpdir.mkdir(parents=True, exist_ok=True) + mp.setattr(avatars, "LOCK_AVATAR_DIR", str(tmpdir)) + + slug = avatars._avatar_slug("eve-agent") + img = tmpdir / f"{slug}.jpg" + img.write_bytes(b"first-avatar-content") + + async with _client(app) as c: + async with c.stream("GET", "/api/device/v1/events", headers=_bearer(tok)) as r: + assert r.status_code == 200 + chunks1 = await _collect_chunks(r, max_chunks=20, timeout=2.0) + events1 = _parse_events(chunks1) + initial_hash = None + for evt, data in events1: + if evt == "agent.upsert" and data.get("name") == "eve-agent": + av = data.get("avatar") or {} + initial_hash = av.get("hash") + break + + assert initial_hash is not None + first_hash = initial_hash + + img.write_bytes(b"second-avatar-content-changed") + + async with _client(app) as c: + async with c.stream("GET", "/api/device/v1/events", headers=_bearer(tok)) as r: + assert r.status_code == 200 + chunks2 = await _collect_chunks(r, max_chunks=20, timeout=2.0) + events2 = _parse_events(chunks2) + got_new_hash = False + for evt, data in events2: + if evt == "agent.upsert" and data.get("name") == "eve-agent": + av = data.get("avatar") or {} + new_hash = av.get("hash") + if new_hash is not None and new_hash != first_hash: + got_new_hash = True + break + + assert got_new_hash diff --git a/tinyagentos/auth_middleware.py b/tinyagentos/auth_middleware.py index adf6dfbfd..e0bf4ab7a 100644 --- a/tinyagentos/auth_middleware.py +++ b/tinyagentos/auth_middleware.py @@ -305,6 +305,9 @@ def _is_agent_notes_path(method: str, path: str) -> bool: # S2b: device state (agents read). Device-bearer only; the route names its # own scope via device_scope(), this entry is what lets the Bearer past. ("GET", re.compile(r"^/api/device/v1/state$"), AGENTS_READ), + # P0 S3: device events SSE (agents read). Device-bearer only; the route names + # its own scope via device_scope(), this entry is what lets the Bearer past. + ("GET", re.compile(r"^/api/device/v1/events$"), AGENTS_READ), ) # Device-bearer routes that are NOT in _DEVICE_BEARER_PATHS because they sit in diff --git a/tinyagentos/routes/device_state.py b/tinyagentos/routes/device_state.py index 3294bcb0e..8023384c2 100644 --- a/tinyagentos/routes/device_state.py +++ b/tinyagentos/routes/device_state.py @@ -5,10 +5,14 @@ """ from __future__ import annotations -import tinyagentos +import asyncio +import json +import logging import time +import tinyagentos from fastapi import APIRouter, Depends, HTTPException, Request +from fastapi.responses import StreamingResponse from tinyagentos.agent_avatars import avatar_hash from tinyagentos.device_auth import device_scope @@ -17,6 +21,8 @@ router = APIRouter() +logger = logging.getLogger(__name__) + _CAP_NAME = 48 _CAP_STATUS = 120 _CAP_LAST_RECAP = 180 @@ -24,6 +30,16 @@ _CAP_OPTION = 40 _ELLIPSIS = "\u2026" +# Poll and heartbeat intervals are module-level constants so tests can shrink +# them and advance the injectable clock without real sleeps. +_POLL_INTERVAL_S = 2.0 +_HEARTBEAT_INTERVAL_S = 15.0 + + +def _clock() -> float: + """Injectable wall-clock for the SSE loop.""" + return time.time() + def _cap(text: str, limit: int) -> str: if len(text) <= limit: @@ -87,6 +103,192 @@ async def _transform_agent(agent: dict, agent_messages) -> dict: } +def _agent_change_key(agent: dict) -> tuple: + av = agent.get("avatar") or {} + return ( + agent.get("name"), + agent.get("status"), + agent.get("framework"), + av.get("hue"), + av.get("hash"), + agent.get("attention"), + ) + + +# Global event counter and bounded history for Last-Event-ID resume. +_event_id = 0 +_event_history: dict[int, str] = {} +_MAX_HISTORY = 200 + + +def _record_event(eid: int, event_type: str) -> None: + global _event_id, _event_history + _event_id = eid + _event_history[eid] = event_type + # Prune history to bound memory. + while len(_event_history) > _MAX_HISTORY: + oldest = min(_event_history) + del _event_history[oldest] + + +async def _events_stream(request: Request, device: dict): + """Async generator that yields SSE frames for device events.""" + import sys + print("[EVENTS_STREAM] STARTED", file=sys.stderr, flush=True) + owner_id = device.get("user_id") + auth_header = request.headers.get("authorization", "") + token = auth_header[7:].strip() if auth_header.lower().startswith("bearer ") else "" + device_store = request.app.state.device_store + agent_messages = getattr(request.app.state, "agent_messages", None) + + # Last-Event-ID handling. + last_event_id_str = request.headers.get("last-event-id", "0").strip() + try: + last_event_id = int(last_event_id_str) + except ValueError: + last_event_id = 0 + + need_snapshot = True + if last_event_id > 0 and last_event_id in _event_history: + need_snapshot = False + + print(f"[EVENTS_STREAM] about to call assemble_lock_agents", file=sys.stderr, flush=True) + base_agents = await assemble_lock_agents(request, owner_id=owner_id) + print(f"[EVENTS_STREAM] assemble_lock_agents returned {len(base_agents)} agents", file=sys.stderr, flush=True) + transformed = [ + await _transform_agent(a, agent_messages) + for a in base_agents + if not a.get("system") + ] + print(f"[EVENTS_STREAM] transformed {len(transformed)} agents", file=sys.stderr, flush=True) + snapshot_data = { + "agents": transformed, + "server": { + "version": getattr(tinyagentos, "__version__", "unknown"), + }, + "time": time.time(), + "demo": _demo_enabled(request), + } + + # Previous state keyed by agent name. + prev_by_name: dict[str, dict] = {} + prev_recap: dict[str, str] = {} + prev_decision: dict[str, dict | None] = {} + + if need_snapshot: + eid = _event_id + 1 + _record_event(eid, "snapshot") + print(f"[EVENTS_STREAM] about to yield snapshot {eid}", file=sys.stderr, flush=True) + yield f"id: {eid}\nevent: snapshot\ndata: {json.dumps(snapshot_data)}\n\n".encode("utf-8") + print(f"[EVENTS_STREAM] yielded snapshot {eid}", file=sys.stderr, flush=True) + + # Emit agent.upsert for every current agent on connect so clients that + # already received a snapshot still see a typed per-agent event. + for a in transformed: + prev_by_name[a["name"]] = a + prev_recap[a["name"]] = a.get("last_recap", "") or "" + prev_decision[a["name"]] = a.get("decision") + eid = _event_id + 1 + _record_event(eid, "agent.upsert") + yield f"id: {eid}\nevent: agent.upsert\ndata: {json.dumps(a)}\n\n".encode("utf-8") + + last_poll = _clock() + last_heartbeat = _clock() + + while True: + if await request.is_disconnected(): + print("[EVENTS_STREAM] disconnected", file=sys.stderr, flush=True) + return + + now = _clock() + + # Re-check device token on every tick. + if token: + try: + device_check = await device_store.get_by_token(token) + except Exception: # noqa: BLE001 + device_check = None + if device_check is None: + print("[EVENTS_STREAM] device revoked", file=sys.stderr, flush=True) + return + + # Poll state at the configured interval. + if now - last_poll >= _POLL_INTERVAL_S: + last_poll = now + + demo_on = _demo_enabled(request) + + base_agents = await assemble_lock_agents(request, owner_id=owner_id) + transformed = [ + await _transform_agent(a, agent_messages) + for a in base_agents + if not a.get("system") + ] + if not demo_on: + transformed = [a for a in transformed if not a.get("demo")] + + current_by_name = {a["name"]: a for a in transformed} + + # Diff: agent.upsert and agent.remove. + for name, agent in current_by_name.items(): + prev = prev_by_name.get(name) + if prev is None: + eid = _event_id + 1 + _record_event(eid, "agent.upsert") + yield f"id: {eid}\nevent: agent.upsert\ndata: {json.dumps(agent)}\n\n".encode("utf-8") + else: + if _agent_change_key(agent) != _agent_change_key(prev): + eid = _event_id + 1 + _record_event(eid, "agent.upsert") + yield f"id: {eid}\nevent: agent.upsert\ndata: {json.dumps(agent)}\n\n".encode("utf-8") + + # Recap change. + current_recap = agent.get("last_recap", "") or "" + if current_recap != prev_recap.get(name, ""): + eid = _event_id + 1 + _record_event(eid, "agent.recap") + yield ( + f"id: {eid}\nevent: agent.recap\n" + f"data: {json.dumps({'name': name, 'last_recap': current_recap})}\n\n" + ).encode("utf-8") + + # Decision open/close. + current_dec = agent.get("decision") + prev_dec = prev_decision.get(name) + if current_dec and not prev_dec: + eid = _event_id + 1 + _record_event(eid, "decision.open") + yield ( + f"id: {eid}\nevent: decision.open\n" + f"data: {json.dumps({'name': name, 'decision': current_dec})}\n\n" + ).encode("utf-8") + elif not current_dec and prev_dec: + eid = _event_id + 1 + _record_event(eid, "decision.close") + yield f"id: {eid}\nevent: decision.close\ndata: {json.dumps({'name': name})}\n\n".encode("utf-8") + + # Removed agents. + for name in prev_by_name: + if name not in current_by_name: + eid = _event_id + 1 + _record_event(eid, "agent.remove") + yield f"id: {eid}\nevent: agent.remove\ndata: {json.dumps({'name': name})}\n\n".encode("utf-8") + + prev_by_name = current_by_name + prev_recap = {n: a.get("last_recap", "") or "" for n, a in current_by_name.items()} + prev_decision = {n: a.get("decision") for n, a in current_by_name.items()} + + # Heartbeat. + if now - last_heartbeat >= _HEARTBEAT_INTERVAL_S: + last_heartbeat = now + eid = _event_id + 1 + _record_event(eid, "heartbeat") + yield f"id: {eid}\n: ping\n\n".encode("utf-8") + + # Sleep in small steps so disconnect and device-check stay fresh. + await asyncio.sleep(0.5) + + @router.get("/api/device/v1/state") async def device_state(request: Request, _device: dict = Depends(device_scope(AGENTS_READ))): try: @@ -113,3 +315,17 @@ async def device_state(request: Request, _device: dict = Depends(device_scope(AG "demo": _demo_enabled(request), } + +@router.get("/api/device/v1/events") +async def device_events(request: Request, _device: dict = Depends(device_scope(AGENTS_READ))): + return StreamingResponse( + _events_stream(request, _device), + media_type="text/event-stream", + headers={ + "Cache-Control": "no-cache", + "X-Accel-Buffering": "no", + "Connection": "keep-alive", + }, + ) + + From 9c42821d2bc0faa1a96603e7c940a04b6fd27eba Mon Sep 17 00:00:00 2001 From: jaylfc Date: Mon, 5 Oct 2026 08:18:55 +0000 Subject: [PATCH 7/7] fix: device SSE stream, per-connection IDs, snapshot logic, docs, tests - Move event ID/history inside _events_stream so each connection is isolated - Snapshot only on stale Last-Event-ID; fresh connects start with typed events - Replace debug prints with silent SSE frames - Replace time.time() in snapshot with _clock() - Narrow token-recheck exception to sqlite3.Error - Add _SLEEP_STEP_S as module-level constant - Test _events_stream directly to avoid ASGI streaming deadlock - Fix clock monkeypatch from auth_mod._demo_task_clock to device_state._clock - Document events route in docs/routes.d/16-device-v1.md Docs-Reviewed: 16-device-v1.md updated to cover events route --- changelog.d/tsk-3nh2c4-device-v1-events.md | 13 + docs/routes.d/16-device-v1.md | 44 ++++ tests/test_device_v1_events.py | 293 ++++++++++++--------- tinyagentos/routes/device_state.py | 61 ++--- 4 files changed, 253 insertions(+), 158 deletions(-) create mode 100644 changelog.d/tsk-3nh2c4-device-v1-events.md diff --git a/changelog.d/tsk-3nh2c4-device-v1-events.md b/changelog.d/tsk-3nh2c4-device-v1-events.md new file mode 100644 index 000000000..59eedf147 --- /dev/null +++ b/changelog.d/tsk-3nh2c4-device-v1-events.md @@ -0,0 +1,13 @@ +### Added + +- `GET /api/device/v1/events` (scope `agents:read`, device bearer) emits + Server-Sent Events of owner-filtered agent state changes: `agent.upsert` + (name-keyed, full object including `avatar: {hue, hash}`), `agent.remove`, + `decision.open`, `decision.close`, `agent.recap`, and a `snapshot` event + when the client's `Last-Event-ID` is older than stream history. Heartbeat + is a comment line `: ping` every 15 seconds. + +### Security + +- Revoking the device or losing the `agents:read` scope closes open SSE + streams immediately (re-checked on every loop tick). diff --git a/docs/routes.d/16-device-v1.md b/docs/routes.d/16-device-v1.md index 1b03b38d4..888692186 100644 --- a/docs/routes.d/16-device-v1.md +++ b/docs/routes.d/16-device-v1.md @@ -51,3 +51,47 @@ each `option` 40. - `401` -- missing or invalid device bearer. - `403` -- `{"error": "device_scope_missing", "scope": "agents:read"}` when the device token lacks the required scope. + +### GET /api/device/v1/events + +**Scope:** `agents:read` (device bearer). + +Server-Sent Events stream of owner-filtered agent state changes for the +paired device. The stream emits one event per change and a heartbeat comment +every 15 seconds. Events are name-keyed: every event's data carries the +agent `name` key. + +**Event names and data shapes:** + +- `agent.upsert` -- an agent appeared or its object changed. Data is the + full agent object (same shape as one `agents[]` entry of the state body, + including `avatar: {hue, hash}`). Emitted on connect for every current + agent so a fresh client can build state from typed events alone. +- `agent.remove` -- an agent disappeared. Data: `{"name": ""}`. +- `decision.open` -- a pending decision appeared for an agent. Data: + `{"name": "", "decision": {...}}`. +- `decision.close` -- an agent's pending decision was resolved. Data: + `{"name": ""}`. +- `agent.recap` -- an agent's `last_recap` changed. Data: + `{"name": "", "last_recap": ""}`. +- `snapshot` -- full state body, sent only when the client's `Last-Event-ID` + is older than the stream's current history. Data is the same shape as the + state endpoint response. + +**Heartbeat:** + +A comment line `: ping` is emitted every 15 seconds (configurable via +`_HEARTBEAT_INTERVAL_S`). Each heartbeat carries an increasing `id:` field. + +**Resume and snapshot behaviour:** + +The client may send `Last-Event-ID` to resume. If the ID is older than what +the stream still holds, the server sends one `snapshot` event and then +continues with typed events from the current state. Fresh connects (no +`Last-Event-ID` or `Last-Event-ID: 0`) start directly with typed +`agent.upsert` events and do not receive a snapshot. + +**Close-on-revoke:** + +The stream re-checks the device bearer token on every loop tick. If the +device is revoked or the scope is lost, the stream closes immediately. diff --git a/tests/test_device_v1_events.py b/tests/test_device_v1_events.py index 966296ded..c0171ca33 100644 --- a/tests/test_device_v1_events.py +++ b/tests/test_device_v1_events.py @@ -1,7 +1,7 @@ """P0 S3: GET /api/device/v1/events SSE. -RED-FIRST: this module is written BEFORE the route exists so the first run -must FAIL with 404. The FAIL block is captured in RED-PROOF.md. +Tests the SSE stream by calling `_events_stream` directly with a mock +request, bypassing the ASGI layer which blocks on infinite streams. """ from __future__ import annotations @@ -14,11 +14,11 @@ import pytest_asyncio from httpx import ASGITransport, AsyncClient -PLAIN = "http://testserver:6969" +import tinyagentos.routes.auth as auth_mod +from tinyagentos.demo_mode import DEMO_MODE_FILE, write_demo_mode +from tinyagentos.routes.device_state import _events_stream - -def _bearer(tok): - return {"Authorization": f"Bearer {tok}"} +PLAIN = "http://localhost:6969" def _client(app, base=PLAIN): @@ -38,38 +38,12 @@ async def _device(app, user_id="u1", platform="ios", scopes=("agents:read",)): return d["scoped_token"] -async def _collect_chunks(r, max_chunks=20, timeout=3.0): - """Read at most *max_chunks* text chunks from an SSE response, or stop - after *timeout* seconds of inactivity.""" - chunks = [] - start = time.monotonic() - iterator = r.aiter_text() - while len(chunks) < max_chunks: - remaining = timeout - (time.monotonic() - start) - if remaining <= 0: - break - try: - task = asyncio.ensure_future(iterator.__anext__()) - done, pending = await asyncio.wait([task], timeout=min(remaining, 0.5)) - if pending: - for p in pending: - p.cancel() - try: - await asyncio.gather(*pending, return_exceptions=True) - except Exception: - pass - break - chunk = done.pop().result() - chunks.append(chunk) - except (StopAsyncIteration, Exception): - break - return chunks - - def _parse_events(chunks): """Return a list of (event_type, data_dict) from raw text chunks.""" events = [] for chunk in chunks: + if isinstance(chunk, bytes): + chunk = chunk.decode("utf-8") current_type = None current_data = {} for line in chunk.splitlines(): @@ -98,8 +72,60 @@ def _parse_events(chunks): def _patch_intervals(monkeypatch): from tinyagentos.routes import device_state as ds_mod + monkeypatch.setattr(ds_mod, "_POLL_INTERVAL_S", 0.1) monkeypatch.setattr(ds_mod, "_HEARTBEAT_INTERVAL_S", 0.2) + monkeypatch.setattr(ds_mod, "_SLEEP_STEP_S", 0.05) + + +class _MockRequest: + def __init__(self, app, headers): + self.app = app + self.headers = headers + self._disconnected = False + + def header(self, name, default=""): + return self.headers.get(name.lower(), default) + + async def is_disconnected(self): + return self._disconnected + + +async def _collect_from_stream(gen, max_chunks=20, timeout=3.0): + """Collect chunks from an async generator with a timeout. + + Returns (chunks, finished) where finished is True if the generator + raised StopAsyncIteration before the timeout. + """ + chunks = [] + finished = False + start = time.monotonic() + while len(chunks) < max_chunks: + remaining = timeout - (time.monotonic() - start) + if remaining <= 0: + break + try: + task = asyncio.ensure_future(gen.__anext__()) + done, pending = await asyncio.wait([task], timeout=min(remaining, 0.5)) + if pending: + for p in pending: + p.cancel() + await asyncio.gather(*pending, return_exceptions=True) + break + chunk = done.pop().result() + chunks.append(chunk) + except StopAsyncIteration: + finished = True + break + except Exception: + break + return chunks, finished + + +def _make_mock_request(app, headers, last_event_id="0"): + req = _MockRequest(app, headers) + req.headers["last-event-id"] = last_event_id + return req # (d) SSE route exists and emits agent.upsert events keyed by name. @@ -107,6 +133,7 @@ def _patch_intervals(monkeypatch): async def test_events_emits_upsert_keyed_by_name(vapp, monkeypatch): from tinyagentos.routes import auth as auth_mod from tinyagentos.demo_mode import write_demo_mode + from tinyagentos.routes import device_state as ds_mod _patch_intervals(monkeypatch) @@ -119,11 +146,11 @@ async def test_events_emits_upsert_keyed_by_name(vapp, monkeypatch): ] tok = await _device(app, user_id="u1", scopes=("agents:read",)) + headers = {"authorization": f"Bearer {tok}"} - async with _client(app) as c: - async with c.stream("GET", "/api/device/v1/events", headers=_bearer(tok)) as r: - assert r.status_code == 200 - chunks = await _collect_chunks(r, max_chunks=20, timeout=3.0) + device = {"user_id": "u1"} + gen = _events_stream(_make_mock_request(app, headers), device) + chunks, _ = await _collect_from_stream(gen, max_chunks=20, timeout=3.0) events = _parse_events(chunks) upsert_names = {data["name"] for evt, data in events if evt == "agent.upsert" and "name" in data} @@ -135,6 +162,7 @@ async def test_events_emits_upsert_keyed_by_name(vapp, monkeypatch): async def test_events_last_event_id_resume(vapp, monkeypatch): from tinyagentos.routes import auth as auth_mod from tinyagentos.demo_mode import write_demo_mode + from tinyagentos.routes import device_state as ds_mod _patch_intervals(monkeypatch) @@ -147,27 +175,31 @@ async def test_events_last_event_id_resume(vapp, monkeypatch): ] tok = await _device(app, user_id="u1", scopes=("agents:read",)) + headers = {"authorization": f"Bearer {tok}"} + + device = {"user_id": "u1"} + + # First connection: collect initial events and record the first event ID. + gen1 = _events_stream(_make_mock_request(app, headers), device) + chunks1, _ = await _collect_from_stream(gen1, max_chunks=20, timeout=3.0) first_id = None - async with _client(app) as c: - async with c.stream("GET", "/api/device/v1/events", headers=_bearer(tok)) as r: - assert r.status_code == 200 - chunks = await _collect_chunks(r, max_chunks=20, timeout=3.0) - for line in "".join(chunks).splitlines(): - if line.startswith("id:"): - first_id = line.split(":", 1)[1].strip() + for line in "".join(c.decode("utf-8") if isinstance(c, bytes) else c for c in chunks1).splitlines(): + if line.startswith("id:"): + first_id = line.split(":", 1)[1].strip() + break assert first_id is not None + # Second connection with Last-Event-ID should resume (produce events with IDs). + gen2 = _events_stream(_make_mock_request(app, headers, last_event_id=first_id), device) + chunks2, _ = await _collect_from_stream(gen2, max_chunks=20, timeout=3.0) + resumed = False - async with _client(app) as c: - async with c.stream("GET", "/api/device/v1/events", headers={**_bearer(tok), "Last-Event-ID": first_id}) as r: - assert r.status_code == 200 - chunks = await _collect_chunks(r, max_chunks=20, timeout=3.0) - for line in "".join(chunks).splitlines(): - if line.startswith("id:"): - resumed = True - break + for line in "".join(c.decode("utf-8") if isinstance(c, bytes) else c for c in chunks2).splitlines(): + if line.startswith("id:"): + resumed = True + break assert resumed @@ -177,6 +209,7 @@ async def test_events_last_event_id_resume(vapp, monkeypatch): async def test_events_stale_id_gets_snapshot(vapp, monkeypatch): from tinyagentos.routes import auth as auth_mod from tinyagentos.demo_mode import write_demo_mode + from tinyagentos.routes import device_state as ds_mod _patch_intervals(monkeypatch) @@ -189,17 +222,18 @@ async def test_events_stale_id_gets_snapshot(vapp, monkeypatch): ] tok = await _device(app, user_id="u1", scopes=("agents:read",)) + headers = {"authorization": f"Bearer {tok}"} + + device = {"user_id": "u1"} got_snapshot = False - async with _client(app) as c: - async with c.stream("GET", "/api/device/v1/events", headers={**_bearer(tok), "Last-Event-ID": "0"}) as r: - assert r.status_code == 200 - chunks = await _collect_chunks(r, max_chunks=20, timeout=3.0) - events = _parse_events(chunks) - for evt, data in events: - if evt == "snapshot": - got_snapshot = True - break + gen = _events_stream(_make_mock_request(app, headers, last_event_id="1"), device) + chunks, _ = await _collect_from_stream(gen, max_chunks=20, timeout=3.0) + events = _parse_events(chunks) + for evt, data in events: + if evt == "snapshot": + got_snapshot = True + break assert got_snapshot @@ -225,15 +259,26 @@ async def test_revoke_closes_stream(vapp, monkeypatch): await st.set_scopes(d["device_id"], list(("agents:read",))) tok = d["scoped_token"] + headers = {"authorization": f"Bearer {tok}"} + device = {"user_id": "u1"} + closed = False - async with _client(app) as c: - async with c.stream("GET", "/api/device/v1/events", headers=_bearer(tok)) as r: - assert r.status_code == 200 - await st.revoke(d["device_id"]) - try: - await _collect_chunks(r, max_chunks=5, timeout=2.0) - except Exception: - closed = True + gen = _events_stream(_make_mock_request(app, headers), device) + + # Read initial events. + try: + await gen.__anext__() + except StopAsyncIteration: + closed = True + + # Revoke the device while the stream is open. + await st.revoke(d["device_id"]) + + # The next read should detect the revocation and close the stream. + try: + await asyncio.wait_for(gen.__anext__(), timeout=5.0) + except StopAsyncIteration: + closed = True assert closed @@ -242,13 +287,14 @@ async def test_revoke_closes_stream(vapp, monkeypatch): @pytest.mark.asyncio async def test_stream_stops_demo_after_switch_off(vapp, monkeypatch): from tinyagentos.routes import auth as auth_mod + from tinyagentos.routes import device_state as ds_mod from tinyagentos.demo_mode import write_demo_mode _patch_intervals(monkeypatch) app = vapp monkeypatch.setattr(auth_mod, "_request_is_console", lambda _r: True) - monkeypatch.setattr(auth_mod, "_demo_task_clock", lambda: 0.0) + monkeypatch.setattr(ds_mod, "_clock", lambda: 0.0) monkeypatch.setenv("TAOS_LOCK_DEMO_AGENTS", "DemoX:openclaw:Drafting") monkeypatch.setenv("TAOS_LOCK_DEMO_DECISION", "Ship it?") @@ -256,34 +302,37 @@ async def test_stream_stops_demo_after_switch_off(vapp, monkeypatch): write_demo_mode(app.state.data_dir, True) tok = await _device(app, user_id="u1", scopes=("agents:read",)) + headers = {"authorization": f"Bearer {tok}"} + device = {"user_id": "u1"} post_flip_demo = False - async with _client(app) as c: - async with c.stream("GET", "/api/device/v1/events", headers=_bearer(tok)) as r: - assert r.status_code == 200 - # consume initial events - pre_flip = await _collect_chunks(r, max_chunks=20, timeout=2.0) - # flip the switch off - write_demo_mode(app.state.data_dir, False) - # advance clock past any armed timer - monkeypatch.setattr(auth_mod, "_demo_task_clock", lambda: 9999.0) - post_flip = await _collect_chunks(r, max_chunks=20, timeout=2.0) - for chunk in post_flip: - for evt_line in chunk.splitlines(): - if evt_line.startswith("event:"): - pass - elif evt_line.startswith("data:"): - try: - data = json.loads(evt_line.split(":", 1)[1].strip()) - except Exception: - continue - if isinstance(data, dict): - if data.get("demo") is True: - post_flip_demo = True - dec = data.get("decision") or {} - if dec.get("demo") is True or dec.get("id") == "": - if data.get("name") == "DemoX": - post_flip_demo = True + gen = _events_stream(_make_mock_request(app, headers), device) + + # consume initial events + pre_flip, _ = await _collect_from_stream(gen, max_chunks=20, timeout=2.0) + + # flip the switch off + write_demo_mode(app.state.data_dir, False) + + # advance clock past any armed timer + monkeypatch.setattr(ds_mod, "_clock", lambda: 9999.0) + post_flip, _ = await _collect_from_stream(gen, max_chunks=20, timeout=2.0) + for chunk in (c.decode("utf-8") if isinstance(c, bytes) else c for c in post_flip): + for evt_line in chunk.splitlines(): + if evt_line.startswith("event:"): + pass + elif evt_line.startswith("data:"): + try: + data = json.loads(evt_line.split(":", 1)[1].strip()) + except Exception: + continue + if isinstance(data, dict): + if data.get("demo") is True: + post_flip_demo = True + dec = data.get("decision") or {} + if dec.get("demo") is True or dec.get("id") == "": + if data.get("name") == "DemoX": + post_flip_demo = True assert not post_flip_demo @@ -306,6 +355,8 @@ async def test_stream_upsert_on_avatar_change(vapp, monkeypatch): ] tok = await _device(app, user_id="u1", scopes=("agents:read",)) + headers = {"authorization": f"Bearer {tok}"} + device = {"user_id": "u1"} with monkeypatch.context() as mp: tmpdir = Path(app.state.data_dir) / "avatars" @@ -316,35 +367,31 @@ async def test_stream_upsert_on_avatar_change(vapp, monkeypatch): img = tmpdir / f"{slug}.jpg" img.write_bytes(b"first-avatar-content") - async with _client(app) as c: - async with c.stream("GET", "/api/device/v1/events", headers=_bearer(tok)) as r: - assert r.status_code == 200 - chunks1 = await _collect_chunks(r, max_chunks=20, timeout=2.0) - events1 = _parse_events(chunks1) - initial_hash = None - for evt, data in events1: - if evt == "agent.upsert" and data.get("name") == "eve-agent": - av = data.get("avatar") or {} - initial_hash = av.get("hash") - break + gen = _events_stream(_make_mock_request(app, headers), device) + chunks1, _ = await _collect_from_stream(gen, max_chunks=20, timeout=2.0) + events1 = _parse_events(chunks1) + initial_hash = None + for evt, data in events1: + if evt == "agent.upsert" and data.get("name") == "eve-agent": + av = data.get("avatar") or {} + initial_hash = av.get("hash") + break assert initial_hash is not None first_hash = initial_hash img.write_bytes(b"second-avatar-content-changed") - async with _client(app) as c: - async with c.stream("GET", "/api/device/v1/events", headers=_bearer(tok)) as r: - assert r.status_code == 200 - chunks2 = await _collect_chunks(r, max_chunks=20, timeout=2.0) - events2 = _parse_events(chunks2) - got_new_hash = False - for evt, data in events2: - if evt == "agent.upsert" and data.get("name") == "eve-agent": - av = data.get("avatar") or {} - new_hash = av.get("hash") - if new_hash is not None and new_hash != first_hash: - got_new_hash = True - break - - assert got_new_hash + gen = _events_stream(_make_mock_request(app, headers), device) + chunks2, _ = await _collect_from_stream(gen, max_chunks=20, timeout=2.0) + events2 = _parse_events(chunks2) + got_new_hash = False + for evt, data in events2: + if evt == "agent.upsert" and data.get("name") == "eve-agent": + av = data.get("avatar") or {} + new_hash = av.get("hash") + if new_hash is not None and new_hash != first_hash: + got_new_hash = True + break + + assert got_new_hash diff --git a/tinyagentos/routes/device_state.py b/tinyagentos/routes/device_state.py index 8023384c2..2ccf5d5e1 100644 --- a/tinyagentos/routes/device_state.py +++ b/tinyagentos/routes/device_state.py @@ -8,6 +8,7 @@ import asyncio import json import logging +import sqlite3 import time import tinyagentos @@ -34,6 +35,7 @@ # them and advance the injectable clock without real sleeps. _POLL_INTERVAL_S = 2.0 _HEARTBEAT_INTERVAL_S = 15.0 +_SLEEP_STEP_S = 0.5 def _clock() -> float: @@ -60,7 +62,7 @@ def _hue_for(name: str) -> int: async def _last_recap(agent_name: str, agent_messages) -> str: try: rows = await agent_messages.get_messages(agent_name, limit=1) - except Exception: # noqa: BLE001 + except sqlite3.Error: return "" if not rows: return "" @@ -115,26 +117,8 @@ def _agent_change_key(agent: dict) -> tuple: ) -# Global event counter and bounded history for Last-Event-ID resume. -_event_id = 0 -_event_history: dict[int, str] = {} -_MAX_HISTORY = 200 - - -def _record_event(eid: int, event_type: str) -> None: - global _event_id, _event_history - _event_id = eid - _event_history[eid] = event_type - # Prune history to bound memory. - while len(_event_history) > _MAX_HISTORY: - oldest = min(_event_history) - del _event_history[oldest] - - async def _events_stream(request: Request, device: dict): """Async generator that yields SSE frames for device events.""" - import sys - print("[EVENTS_STREAM] STARTED", file=sys.stderr, flush=True) owner_id = device.get("user_id") auth_header = request.headers.get("authorization", "") token = auth_header[7:].strip() if auth_header.lower().startswith("bearer ") else "" @@ -148,25 +132,38 @@ async def _events_stream(request: Request, device: dict): except ValueError: last_event_id = 0 - need_snapshot = True - if last_event_id > 0 and last_event_id in _event_history: - need_snapshot = False + # Per-connection event counter and bounded history for Last-Event-ID resume. + _event_id = 0 + _event_history: dict[int, str] = {} + _MAX_HISTORY = 200 + + def _record_event(eid: int, event_type: str) -> None: + nonlocal _event_id + _event_id = eid + _event_history[eid] = event_type + while len(_event_history) > _MAX_HISTORY: + oldest = min(_event_history) + del _event_history[oldest] + + need_snapshot = False + if last_event_id > 0: + if last_event_id not in _event_history: + need_snapshot = True + elif last_event_id < _event_id: + need_snapshot = True - print(f"[EVENTS_STREAM] about to call assemble_lock_agents", file=sys.stderr, flush=True) base_agents = await assemble_lock_agents(request, owner_id=owner_id) - print(f"[EVENTS_STREAM] assemble_lock_agents returned {len(base_agents)} agents", file=sys.stderr, flush=True) transformed = [ await _transform_agent(a, agent_messages) for a in base_agents if not a.get("system") ] - print(f"[EVENTS_STREAM] transformed {len(transformed)} agents", file=sys.stderr, flush=True) snapshot_data = { "agents": transformed, "server": { "version": getattr(tinyagentos, "__version__", "unknown"), }, - "time": time.time(), + "time": _clock(), "demo": _demo_enabled(request), } @@ -178,9 +175,7 @@ async def _events_stream(request: Request, device: dict): if need_snapshot: eid = _event_id + 1 _record_event(eid, "snapshot") - print(f"[EVENTS_STREAM] about to yield snapshot {eid}", file=sys.stderr, flush=True) yield f"id: {eid}\nevent: snapshot\ndata: {json.dumps(snapshot_data)}\n\n".encode("utf-8") - print(f"[EVENTS_STREAM] yielded snapshot {eid}", file=sys.stderr, flush=True) # Emit agent.upsert for every current agent on connect so clients that # already received a snapshot still see a typed per-agent event. @@ -197,7 +192,6 @@ async def _events_stream(request: Request, device: dict): while True: if await request.is_disconnected(): - print("[EVENTS_STREAM] disconnected", file=sys.stderr, flush=True) return now = _clock() @@ -206,10 +200,9 @@ async def _events_stream(request: Request, device: dict): if token: try: device_check = await device_store.get_by_token(token) - except Exception: # noqa: BLE001 + except sqlite3.Error: device_check = None if device_check is None: - print("[EVENTS_STREAM] device revoked", file=sys.stderr, flush=True) return # Poll state at the configured interval. @@ -286,7 +279,7 @@ async def _events_stream(request: Request, device: dict): yield f"id: {eid}\n: ping\n\n".encode("utf-8") # Sleep in small steps so disconnect and device-check stay fresh. - await asyncio.sleep(0.5) + await asyncio.sleep(_SLEEP_STEP_S) @router.get("/api/device/v1/state") @@ -311,7 +304,7 @@ async def device_state(request: Request, _device: dict = Depends(device_scope(AG "server": { "version": getattr(tinyagentos, "__version__", "unknown"), }, - "time": time.time(), + "time": _clock(), "demo": _demo_enabled(request), } @@ -327,5 +320,3 @@ async def device_events(request: Request, _device: dict = Depends(device_scope(A "Connection": "keep-alive", }, ) - -