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/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/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/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/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..f2ac033d0 100644 --- a/docs/routes.d/16-device-v1.md +++ b/docs/routes.d/16-device-v1.md @@ -5,3 +5,93 @@ 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. + +### 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/docs/routes.md b/docs/routes.md index 449348272..dac010ac2 100644 --- a/docs/routes.md +++ b/docs/routes.md @@ -419,3 +419,93 @@ 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. + +### 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_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/tests/test_device_v1_events.py b/tests/test_device_v1_events.py new file mode 100644 index 000000000..c0171ca33 --- /dev/null +++ b/tests/test_device_v1_events.py @@ -0,0 +1,397 @@ +"""P0 S3: GET /api/device/v1/events SSE. + +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 + +import asyncio +import json +import time +from pathlib import Path + +import pytest +import pytest_asyncio +from httpx import ASGITransport, AsyncClient + +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 + +PLAIN = "http://localhost:6969" + + +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"] + + +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(): + 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) + 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. +@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 + from tinyagentos.routes import device_state as ds_mod + + _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",)) + headers = {"authorization": f"Bearer {tok}"} + + 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} + 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 + from tinyagentos.routes import device_state as ds_mod + + _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",)) + 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 + 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 + 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 + + +# (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 + from tinyagentos.routes import device_state as ds_mod + + _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",)) + headers = {"authorization": f"Bearer {tok}"} + + device = {"user_id": "u1"} + + got_snapshot = False + 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 + + +# (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"] + + headers = {"authorization": f"Bearer {tok}"} + device = {"user_id": "u1"} + + closed = False + 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 + + +# (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.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(ds_mod, "_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",)) + headers = {"authorization": f"Bearer {tok}"} + device = {"user_id": "u1"} + + post_flip_demo = False + 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 + + +# (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",)) + headers = {"authorization": f"Bearer {tok}"} + device = {"user_id": "u1"} + + 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") + + 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") + + 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/tests/test_device_v1_state.py b/tests/test_device_v1_state.py new file mode 100644 index 000000000..49d8e19dd --- /dev/null +++ b/tests/test_device_v1_state.py @@ -0,0 +1,293 @@ +"""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"] + + +# 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): + 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, 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"}, + ] + + 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 + + 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) + + 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, monkeypatch): + from tinyagentos.demo_mode import write_demo_mode + + app = vapp + 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"]) + + +# (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 + + +# (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/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/auth_middleware.py b/tinyagentos/auth_middleware.py index 90d6071bb..e0bf4ab7a 100644 --- a/tinyagentos/auth_middleware.py +++ b/tinyagentos/auth_middleware.py @@ -302,6 +302,12 @@ 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), + # 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/__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 4d3720908..3a0f485ed 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, @@ -9358,44 +9359,45 @@ 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 + 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: 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 "") + else: + name = str(entry) + framework = "" + if not name: + continue agents.append({ "name": str(name), "framework": framework.lower(), @@ -9404,47 +9406,10 @@ async def lock_widgets(request: Request): "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: @@ -9454,26 +9419,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: @@ -9486,37 +9434,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): + if _device_agents_enabled(request) and owner_id is None: 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 + pending = await store.list(status="pending", user_id=owner_id, limit=20) + 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() @@ -9544,11 +9472,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 @@ -9570,18 +9493,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, @@ -9591,9 +9558,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] ) @@ -9603,30 +9567,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. diff --git a/tinyagentos/routes/device_state.py b/tinyagentos/routes/device_state.py new file mode 100644 index 000000000..2ccf5d5e1 --- /dev/null +++ b/tinyagentos/routes/device_state.py @@ -0,0 +1,322 @@ +"""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 asyncio +import json +import logging +import sqlite3 +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 +from tinyagentos.device_scopes import AGENTS_READ +from tinyagentos.routes.auth import assemble_lock_agents, _demo_enabled + +router = APIRouter() + +logger = logging.getLogger(__name__) + +_CAP_NAME = 48 +_CAP_STATUS = 120 +_CAP_LAST_RECAP = 180 +_CAP_QUESTION = 280 +_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 +_SLEEP_STEP_S = 0.5 + + +def _clock() -> float: + """Injectable wall-clock for the SSE loop.""" + return time.time() + + +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 sqlite3.Error: + 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, + } + + +async def _transform_agent(agent: dict, agent_messages) -> dict: + 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(original_name) + ahash = avatar_hash(original_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), + } + + +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"), + ) + + +async def _events_stream(request: Request, device: dict): + """Async generator that yields SSE frames for device events.""" + 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 + + # 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 + + 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") + ] + snapshot_data = { + "agents": transformed, + "server": { + "version": getattr(tinyagentos, "__version__", "unknown"), + }, + "time": _clock(), + "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") + yield f"id: {eid}\nevent: snapshot\ndata: {json.dumps(snapshot_data)}\n\n".encode("utf-8") + + # 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(): + return + + now = _clock() + + # Re-check device token on every tick. + if token: + try: + device_check = await device_store.get_by_token(token) + except sqlite3.Error: + device_check = None + if device_check is None: + 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(_SLEEP_STEP_S) + + +@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": _clock(), + "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", + }, + )