Files
iris_x_hermes/gateway-plugin/tests/ws_probe.py
T
ARIA 678c0344c8 Clean up lint/LSP across gateway, Android, and desktop (alpha -> stable)
Gateway (gateway-plugin/):
- Fix interactive_setup broken imports: print helpers were imported from the
  wrong hermes module (hermes_cli.config instead of hermes_cli.cli_output) plus
  a non-existent print_code; the try/except swallowed the ImportError so
  `hermes gateway setup` for android always bailed out early.
- Fix release_scoped_lock type error (str | None passed where str required).
- Rewrite empty `except: pass` blocks as contextlib.suppress with rationale.
- Restructure two ambiguous ws_server try blocks (hello-auth, frame loop).
- Ruff cleanup: type annotations, import sorting, line wrapping, magic values
  -> named constants, `raise ... from e`, complexity. Add gateway-plugin/ruff.toml.
- Add pyrightconfig.json so the Python LSP resolves hermes-runtime imports.
- Suppress verified false positives inline (parameterized SQL, column-name
  "secrets", hermes-generated media path).

Android (app/androidApp + app/shared):
- Consolidate launcher icons into a single mipmap-anydpi (minSdk 29 >= 26) with
  the monochrome layer; clears ObsoleteSdkInt + MonochromeLauncherIcon.
- Bump core-splashscreen 1.0.1 -> 1.2.0; pin targetSdk 34 (deliberate).
- Suppress verified findings inline (LAN ws:// default, correct GCM IV usage).

Desktop (app/desktopApp):
- Move the desktop to a Java 21 runtime (org.gradle.java.home) and set the
  desktop jvmTarget to 21 (Android stays JVM 17 / minSdk 29). Fixes the startup
  UnsupportedClassVersionError and restores Markdown renderer 0.44.0.

Tooling/config:
- .pi-lens.json: disable verified-noisy heuristics (documented in docs).
- .gitleaks.toml: allowlist git-ignored false-positive paths.
- docs/18-code-review.md: full findings + verification.

Verified: ruff clean, pyright 0 errors, 64/64 gateway tests, all Kotlin tests,
Android lint 0 issues, Android installed+launched on device, desktop launches
on JDK 21.
2026-08-21 18:47:03 +02:00

753 lines
31 KiB
Python

#!/usr/bin/env python3
"""WS test-client harness (docs/13-testing.md §13.2).
Connects to the REAL running gateway and drives a turn, printing every
frame. This is how we empirically confirm the exact frame shapes before /
while building the Kotlin client.
Usage::
hermes gateway & # with the android plugin
python gateway-plugin/tests/ws_probe.py --token <ANDROID_TOKEN> \
--send "hello"
Options:
--url ws://host:port/ws (default ws://127.0.0.1:8790/ws)
--token ANDROID_TOKEN (default: $ANDROID_TOKEN)
--device device_id (default: probe-<rand>)
--send TEXT send this message after pairing (default: "hello")
--upload F M4: upload F (chunked media.upload) and attach it to the
message.send via media_refs
--pull-offer M4: when a media.offer arrives during the turn, pull the
media (chunked) and verify the byte count. Offers are
emitted right AFTER the final message (MEDIA: tag
extraction runs post-turn), so after the final the probe
keeps listening for --offer-grace seconds for one.
--sync C M5: after pairing, send sync {cursor: C} and print the
replay + sync.done (no turn is driven)
--fcm-token M5: attach this FCM token to the hello payload
--fcm-reg M5: after pairing, send fcm.register with --fcm-token
--timeout S seconds to wait for the final reply (default 120)
--authfail expect an auth rejection (wrong token) and exit 0 on it
Assertion modes (checked after the turn; see exit codes below):
--assert-turn M2: the turn produced message.start -> >=1
message.update -> message.stop (scenario 2)
--assert-reasoning M2: the final message.stop carries a non-empty
reasoning field (scenario 3)
--assert-tools M2: >=1 tool.start with a matching tool.end
(matched by index; scenario 4)
--assert-commentary M2: >=1 commentary frame (scenario 5)
--assert-read-receipt M7: a read.receipt frame arrives after the sent
message. SKIPs (exit 0) when the frame never
arrives (old gateway without the M7 frame).
--assert-status M7: a status frame is received. SKIPs (exit 0)
when the frame never arrives.
Request modes (no turn driven unless --send/--upload also given):
--search Q [--scope all|chat] [--chat-id C]
M3: send search {query, scope, limit} and assert >=1 hit
in search.results (scenario 8). With --send, the turn is
driven first, then the search runs.
--channel-create NAME M3: send channel.create, print the new chat_id
("== channel created: <chat_id>"), exit
--channel-delete CHAT M3: send channel.delete, assert channel.deleted
--channel-list M3: send channel.list, print the directory
--watch CHAT_ID wait up to --timeout for a message to land in
CHAT_ID (cron delivery E2E, scenario 7)
Exit codes:
0 ok (incl. SKIP for absent M7 frames)
2 connect failed
3 no hello.ack
4 expected hello.ack, got something else
5 --authfail but the token was accepted
6 timeout waiting for the final message
7 no final assistant message
8 upload/sync failed
9 media pull failed
10 --assert-turn failed (no ordered start/update/stop segment)
11 --assert-reasoning failed (final has no non-empty reasoning)
12 --assert-tools failed (no tool.start with a matching tool.end)
13 --assert-commentary failed (no commentary frame)
14 --search failed (error or zero hits)
15 --channel-create / --channel-list failed
16 --channel-delete failed
17 --watch timed out (no message landed in the channel)
18 --assert-read-receipt failed (frame arrived before the sent message)
19 --assert-status failed (status frame arrived with an empty payload)
"""
import argparse
import asyncio
import hashlib
import json
import mimetypes
import os
import sys
import time
import uuid
try:
import websockets
except ImportError: # pragma: no cover
sys.stderr.write("websockets is required (hermes core dep); run inside the hermes venv\n")
raise
def _print_frame(raw):
try:
data = json.loads(raw)
except (json.JSONDecodeError, TypeError):
print(f" <- {raw!r}")
return None
ftype = data.get("type", "?")
chat = data.get("chat_id")
fid = data.get("id")
payload = data.get("payload", {})
# Compact one-line summary + full payload for the interesting frames.
extra = ""
if ftype == "message":
text = (payload.get("text") or "")
extra = f" role={payload.get('role')} id={payload.get('message_id')} text={text[:120]!r}"
if payload.get("reasoning"):
extra += f" reasoning={payload['reasoning'][:80]!r}"
elif ftype == "message.start":
extra = f" id={payload.get('message_id')} role={payload.get('role')}"
elif ftype == "message.update":
text = (payload.get("text") or "")
extra = f" id={payload.get('message_id')} text={text[:100]!r}"
elif ftype == "message.stop":
text = (payload.get("final_text") or "")
extra = f" id={payload.get('message_id')} text={text[:120]!r}"
if payload.get("reasoning"):
extra += f" reasoning={payload['reasoning'][:80]!r}"
elif ftype == "tool.start":
extra = (f" idx={payload.get('index')} name={payload.get('name')!r} "
f"preview={str(payload.get('preview'))[:80]!r}")
elif ftype == "tool.progress":
extra = (
f" idx={payload.get('index')} name={payload.get('name')!r} "
f"note={payload.get('note')!r}"
)
elif ftype == "tool.end":
extra = (f" idx={payload.get('index')} name={payload.get('name')!r} "
f"ok={payload.get('ok')} dur={payload.get('duration')}")
elif ftype == "commentary":
extra = f" id={payload.get('message_id')} text={(payload.get('text') or '')[:120]!r}"
elif ftype == "hello.ack":
extra = f" caps={payload.get('server_caps')} cursor={payload.get('sync_cursor')}"
elif ftype == "error":
extra = f" code={payload.get('code')} msg={payload.get('message')!r}"
elif ftype == "typing":
extra = f" on={payload.get('on')}"
elif ftype == "pong":
extra = ""
elif ftype == "media.offer":
extra = (f" media_id={payload.get('media_id')} kind={payload.get('kind')} "
f"mime={payload.get('mime')} size={payload.get('size')} "
f"file={payload.get('filename')!r} msg={payload.get('message_id')}")
elif ftype == "media.upload.ack":
extra = f" ok={payload.get('ok')} ref={payload.get('media_ref')}"
elif ftype == "media.pull.end":
extra = f" ok={payload.get('ok')}"
elif ftype == "notification":
extra = (f" kind={payload.get('kind')} title={payload.get('title')!r} "
f"body={(payload.get('body') or '')[:100]!r}")
elif ftype in {"sync", "sync.done"}:
extra = f" cursor={payload.get('cursor')}"
elif ftype == "search.results":
hits = payload.get("hits") or []
extra = f" query={payload.get('query')!r} scope={payload.get('scope')} hits={len(hits)}"
elif ftype == "channel.created":
extra = f" chat_id={payload.get('chat_id')} name={payload.get('name')!r}"
elif ftype == "channel.deleted":
extra = f" chat_id={payload.get('chat_id')}"
elif ftype == "channel.list":
extra = f" channels={len(payload.get('channels') or [])}"
elif ftype in {"read.receipt", "status"}:
extra = f" payload={ {k: payload[k] for k in list(payload)[:4]} }"
scope = f" chat={chat}" if chat else ""
idpart = f" id={fid}" if fid is not None else ""
print(f" <- {ftype}{idpart}{scope}{extra}")
return data
def _kind_for_path(path: str) -> str:
mime, _ = mimetypes.guess_type(path)
mime = mime or "application/octet-stream"
if mime.startswith("image/"):
return "image"
if mime.startswith("video/"):
return "video"
if mime.startswith("audio/"):
return "audio"
return "document"
async def upload_file(ws, path: str, media_ref: str, next_id: int) -> int:
"""Drive media.upload.start -> binary chunks -> media.upload.end.
Returns the next free request id; raises on a non-ack terminal frame.
"""
with open(path, "rb") as f:
data = f.read()
mime, _ = mimetypes.guess_type(path)
await ws.send(json.dumps({
"v": 1, "id": next_id, "type": "media.upload.start",
"payload": {
"media_ref": media_ref,
"kind": _kind_for_path(path),
"mime": mime or "application/octet-stream",
"size": len(data),
"filename": os.path.basename(path),
},
}))
print(f" -> media.upload.start id={next_id} ref={media_ref} size={len(data)}")
chunk = 256 * 1024
for off in range(0, len(data), chunk):
await ws.send(data[off:off + chunk])
await ws.send(json.dumps({
"v": 1, "id": next_id + 1, "type": "media.upload.end",
"payload": {"media_ref": media_ref, "sha256": hashlib.sha256(data).hexdigest()},
}))
print(f" -> media.upload.end id={next_id + 1} ref={media_ref}")
while True:
raw = await asyncio.wait_for(ws.recv(), timeout=60)
data_frame = _print_frame(raw)
if data_frame is None:
continue
if data_frame.get("type") == "media.upload.ack":
if not data_frame["payload"].get("ok"):
raise RuntimeError(f"upload rejected: {data_frame['payload']}")
return next_id + 2
if data_frame.get("type") == "error":
raise RuntimeError(f"upload failed: {data_frame['payload']}")
async def pull_media(ws, media_id: str, request_id: int, expected_size: int | None) -> None:
"""media.pull -> binary frames -> media.pull.end; verifies the size."""
await ws.send(json.dumps({
"v": 1, "id": request_id, "type": "media.pull",
"payload": {"media_id": media_id},
}))
print(f" -> media.pull id={request_id} media_id={media_id}")
total = 0
while True:
raw = await asyncio.wait_for(ws.recv(), timeout=120)
if isinstance(raw, (bytes, bytearray)):
total += len(raw)
continue
data = _print_frame(raw)
if data is None:
continue
if data.get("type") == "media.pull.end":
if not data["payload"].get("ok"):
raise RuntimeError(f"pull failed: {data['payload']}")
if expected_size is not None and total != expected_size:
raise RuntimeError(f"pull size mismatch: got {total}, want {expected_size}")
print(f"== pulled {total} bytes (sha256 of stream verified by size match)")
return
if data.get("type") == "error":
raise RuntimeError(f"pull failed: {data['payload']}")
class _TurnState:
"""Assertion-relevant facts collected while driving a turn."""
def __init__(self):
self.seq: list[tuple[str, str | None]] = [] # (type, message_id)
self.tool_starts: set[int] = set()
self.tool_ends: set[int] = set()
self.commentary = 0
self.final_stop_reasoning: str | None = None
self.final_message_reasoning: str | None = None
self.user_echo_seen = False
self.read_receipt: bool | None = None # None = never arrived
self.status_seen = False
self.status_empty = False
self.pulled = False
def track(self, ftype: str, payload: dict) -> None:
if ftype in ("message.start", "message.update", "message.stop"):
self.seq.append((ftype, payload.get("message_id")))
if ftype == "message.stop":
r = payload.get("reasoning")
if isinstance(r, str) and r.strip():
self.final_stop_reasoning = r
if ftype == "message" and payload.get("role") == "assistant":
r = payload.get("reasoning")
if isinstance(r, str) and r.strip():
self.final_message_reasoning = r
if ftype == "message" and payload.get("role") == "user":
self.user_echo_seen = True
if ftype == "tool.start" and isinstance(payload.get("index"), int):
self.tool_starts.add(payload["index"])
if ftype == "tool.end" and isinstance(payload.get("index"), int):
self.tool_ends.add(payload["index"])
if ftype == "commentary":
self.commentary += 1
if ftype == "read.receipt":
self.read_receipt = self.user_echo_seen
if ftype == "status":
self.status_seen = True
if not payload:
self.status_empty = True
def _evaluate_assertions(args, st: _TurnState) -> list[tuple[int, bool, str]]:
"""Evaluate the enabled assertion modes. Returns (exit_code, ok, message)
per failed-or-passed assertion; SKIPs are printed here and not returned."""
results: list[tuple[int, bool, str]] = []
if args.assert_turn:
ok = False
for mid in {m for _, m in st.seq if m is not None}:
events = [t for t, m in st.seq if m == mid]
if "message.start" in events and "message.stop" in events:
i_start = events.index("message.start")
i_stop = events.index("message.stop")
if any(i_start < i < i_stop
for i, e in enumerate(events) if e == "message.update"):
ok = True
break
results.append((10, ok,
"assert-turn: no message.start -> >=1 message.update -> message.stop"))
if args.assert_reasoning:
reasoning = st.final_stop_reasoning or st.final_message_reasoning
results.append((11, bool(reasoning),
"assert-reasoning: final message has no non-empty reasoning"))
if args.assert_tools:
ok = bool(st.tool_starts) and bool(st.tool_starts & st.tool_ends)
results.append((12, ok,
"assert-tools: no tool.start with a matching tool.end"))
if args.assert_commentary:
results.append((13, st.commentary >= 1,
"assert-commentary: no commentary frame"))
if args.assert_read_receipt:
if st.read_receipt is None:
print("== SKIP: no read.receipt frame (M7 frame not live on this gateway)")
elif not st.read_receipt:
results.append((18, False,
"assert-read-receipt: read.receipt arrived before the sent message"))
if args.assert_status:
if not st.status_seen:
print("== SKIP: no status frame (M7 frame not live on this gateway)")
elif st.status_empty:
results.append((19, False,
"assert-status: status frame arrived with an empty payload"))
return results
async def _recv_frames(ws, timeout: float):
"""Yield parsed frames (dicts) until *timeout* seconds elapse."""
deadline = time.time() + timeout
while time.time() < deadline:
try:
raw = await asyncio.wait_for(ws.recv(), timeout=deadline - time.time())
except asyncio.TimeoutError:
return
if isinstance(raw, (bytes, bytearray)):
continue
data = _print_frame(raw)
if data is not None:
yield data
async def _search_mode(ws, args, next_id: int) -> int:
"""M3: send a search frame, wait for search.results, assert >=1 hit."""
req_id = next_id
payload = {"query": args.search, "scope": args.scope, "limit": 20}
if args.scope == "chat":
payload["chat_id"] = args.chat_id
await ws.send(json.dumps({"v": 1, "id": req_id, "type": "search", "payload": payload}))
print(f" -> search id={req_id} query={args.search!r} scope={args.scope}")
async for data in _recv_frames(ws, timeout=30):
if data.get("type") == "search.results" and data.get("id") == req_id:
hits = (data.get("payload") or {}).get("hits") or []
print(f"== search: {len(hits)} hit(s)")
for h in hits[:10]:
print(f" hit chat={h.get('chat_id')} role={h.get('role')} "
f"snippet={str(h.get('snippet'))[:100]!r}")
if hits:
return 0
print("!! search: no hits")
return 14
if data.get("type") == "error":
print(f"!! search failed: {data.get('payload')}")
return 14
print("!! search: no search.results within 30s")
return 14
async def _channel_create_mode(ws, args) -> int:
"""M3: channel.create -> channel.created; print the new chat_id."""
req_id = 1
await ws.send(json.dumps({
"v": 1, "id": req_id, "type": "channel.create",
"payload": {"name": args.channel_create},
}))
print(f" -> channel.create id={req_id} name={args.channel_create!r}")
async for data in _recv_frames(ws, timeout=30):
if data.get("type") == "channel.created" and data.get("id") == req_id:
chat_id = (data.get("payload") or {}).get("chat_id")
print(f"== channel created: {chat_id}")
await ws.close()
return 0
if data.get("type") == "error":
print(f"!! channel.create failed: {data.get('payload')}")
await ws.close()
return 15
print("!! channel.create: no channel.created within 30s")
await ws.close()
return 15
async def _channel_delete_mode(ws, args) -> int:
"""M3: channel.delete -> channel.deleted."""
req_id = 1
await ws.send(json.dumps({
"v": 1, "id": req_id, "type": "channel.delete",
"payload": {"chat_id": args.channel_delete},
}))
print(f" -> channel.delete id={req_id} chat_id={args.channel_delete!r}")
async for data in _recv_frames(ws, timeout=30):
if data.get("type") == "channel.deleted" and data.get("id") == req_id:
print(f"== channel deleted: {args.channel_delete}")
await ws.close()
return 0
if data.get("type") == "error":
print(f"!! channel.delete failed: {data.get('payload')}")
await ws.close()
return 16
print("!! channel.delete: no channel.deleted within 30s")
await ws.close()
return 16
async def _channel_list_mode(ws, args) -> int:
"""M3: channel.list -> print the directory."""
req_id = 1
await ws.send(json.dumps({"v": 1, "id": req_id, "type": "channel.list", "payload": {}}))
print(" -> channel.list")
async for data in _recv_frames(ws, timeout=30):
if data.get("type") == "channel.list" and data.get("id") == req_id:
for c in (data.get("payload") or {}).get("channels") or []:
print(f"== channel: {c.get('chat_id')} name={c.get('name')!r} "
f"default={bool(c.get('is_default'))}")
await ws.close()
return 0
if data.get("type") == "error":
print(f"!! channel.list failed: {data.get('payload')}")
await ws.close()
return 15
print("!! channel.list: no response within 30s")
await ws.close()
return 15
async def _watch_mode(ws, args) -> int:
"""Wait up to --timeout for a message to land in args.watch (cron E2E)."""
print(f"== watching {args.watch} for a message (timeout {args.timeout:.0f}s)")
deadline = time.time() + args.timeout
while time.time() < deadline:
try:
raw = await asyncio.wait_for(ws.recv(), timeout=deadline - time.time())
except asyncio.TimeoutError:
print(f"!! timeout after {args.timeout:.0f}s watching {args.watch}")
await ws.close()
return 17
if isinstance(raw, (bytes, bytearray)):
continue
data = _print_frame(raw)
if data is None:
continue
if data.get("chat_id") != args.watch:
continue
ftype = data.get("type")
payload = data.get("payload") or {}
if ftype == "message" and payload.get("role") in ("assistant", "cron"):
print(f"== message landed in {args.watch}: {str(payload.get('text'))[:120]!r}")
await ws.close()
return 0
print(f"!! no message landed in {args.watch}")
await ws.close()
return 17
async def run(args) -> int:
url = args.url
token = args.token
device_id = args.device
print(f"== ws_probe: connecting {url} device={device_id}")
try:
ws = await websockets.connect(url, open_timeout=10)
except Exception as e:
print(f"!! connect failed: {e}")
return 2
hello = {
"v": 1,
"type": "hello",
"payload": {
"token": token,
"device_id": device_id,
"device_name": "ws-probe",
"caps": {"min_protocol": 1},
},
}
if args.fcm_token:
hello["payload"]["fcm_token"] = args.fcm_token
await ws.send(json.dumps(hello))
print(" -> hello" + (f" fcm_token={args.fcm_token[:12]}…" if args.fcm_token else ""))
# First response must be hello.ack (or an auth error).
try:
first = await asyncio.wait_for(ws.recv(), timeout=10)
except asyncio.TimeoutError:
print("!! no hello.ack within 10s")
await ws.close()
return 3
data = _print_frame(first)
if data is None or data.get("type") != "hello.ack":
if args.authfail:
print("== auth rejected as expected")
await ws.close()
return 0
print("!! expected hello.ack")
await ws.close()
return 4
if args.authfail:
print("!! expected auth rejection but got hello.ack")
await ws.close()
return 5
# M5: optional fcm.register after pairing.
if args.fcm_reg:
reg_token = args.fcm_token or f"probe-{uuid.uuid4().hex[:12]}"
await ws.send(json.dumps({
"v": 1, "type": "fcm.register",
"payload": {"fcm_token": reg_token},
}))
print(f" -> fcm.register fcm_token={reg_token[:12]}…")
# M5: sync catch-up mode (no turn driven).
if args.sync is not None:
await ws.send(json.dumps({
"v": 1, "id": 1, "type": "sync", "payload": {"cursor": args.sync},
}))
print(f" -> sync cursor={args.sync}")
while True:
raw = await asyncio.wait_for(ws.recv(), timeout=30)
data = _print_frame(raw)
if data is None:
continue
if data.get("type") == "sync.done":
print(f"== sync done at cursor {data['payload'].get('cursor')}")
await ws.close()
return 0
if data.get("type") == "error":
print(f"!! sync failed: {data['payload']}")
await ws.close()
return 8
# M3: request modes (no turn driven).
if args.channel_create:
return await _channel_create_mode(ws, args)
if args.channel_delete:
return await _channel_delete_mode(ws, args)
if args.channel_list:
return await _channel_list_mode(ws, args)
if args.watch:
return await _watch_mode(ws, args)
if not args.send and not args.upload and not args.search:
print("== paired OK (no --send/--upload/--search; exiting)")
await ws.close()
return 0
# M4: optional inbound upload before the turn.
media_refs: list[str] = []
next_id = 1
if args.upload:
media_ref = f"mu_probe_{uuid.uuid4().hex[:8]}"
try:
next_id = await upload_file(ws, args.upload, media_ref, next_id)
except Exception as e:
print(f"!! upload failed: {e}")
await ws.close()
return 8
media_refs.append(media_ref)
# Drive a turn (if --send or --upload).
st = _TurnState()
got_final = False
if args.send or args.upload:
msg_id = next_id
send_payload: dict = {"text": args.send or ""}
if media_refs:
send_payload["media_refs"] = media_refs
send_frame = {
"v": 1,
"id": msg_id,
"type": "message.send",
"chat_id": "android:default",
"payload": send_payload,
}
await ws.send(json.dumps(send_frame))
print(f" -> message.send id={msg_id} text={args.send!r} media_refs={media_refs}")
deadline = time.time() + args.timeout
seen_final_frame = False
while time.time() < deadline:
try:
raw = await asyncio.wait_for(ws.recv(), timeout=deadline - time.time())
except asyncio.TimeoutError:
print(f"!! timeout after {args.timeout}s waiting for final message")
await ws.close()
return 6
data = _print_frame(raw)
if data is None:
continue
ftype = data.get("type")
payload = data.get("payload") or {}
st.track(ftype, payload)
# M4: fetch offered media live (outbound direction).
if ftype == "media.offer" and args.pull_offer and payload.get("media_id"):
try:
await pull_media(
ws, payload["media_id"], next_id, payload.get("size")
)
next_id += 1
st.pulled = True
except Exception as e:
print(f"!! pull failed: {e}")
await ws.close()
return 9
# A standalone assistant `message` (non-streaming) is immediately final.
if ftype == "message" and payload.get("role") == "assistant":
got_final = True
break
# A `message.stop` finalizes a streaming segment; the turn is done
# once typing stops afterwards (multi-segment turns have several
# stops).
if ftype == "message.stop":
seen_final_frame = True
if ftype == "typing" and payload.get("on") is False and seen_final_frame:
got_final = True
break
# M4: media offers are emitted right AFTER the final message (the
# MEDIA: tag is extracted post-turn); give them a grace window.
if got_final and args.pull_offer and not st.pulled:
grace_deadline = time.time() + args.offer_grace
while time.time() < grace_deadline:
try:
raw = await asyncio.wait_for(
ws.recv(), timeout=grace_deadline - time.time()
)
except asyncio.TimeoutError:
break
if isinstance(raw, (bytes, bytearray)):
continue
data = _print_frame(raw)
if data is None:
continue
is_offer = data.get("type") == "media.offer"
if is_offer and (data.get("payload") or {}).get("media_id"):
try:
await pull_media(
ws, data["payload"]["media_id"], next_id,
data["payload"].get("size"),
)
next_id += 1
st.pulled = True
except Exception as e:
print(f"!! pull failed: {e}")
await ws.close()
return 9
break
if not st.pulled:
print(f"== no media.offer within {args.offer_grace:.0f}s grace")
if not got_final:
print("!! no final assistant message")
await ws.close()
return 7
print("== final assistant message received")
# M3: optional search (standalone, or after the turn).
if args.search:
rc = await _search_mode(ws, args, next_id)
await ws.close()
return rc
await ws.close()
for code, ok, msg in _evaluate_assertions(args, st):
if not ok:
print(f"!! {msg}")
return code
return 0
def main() -> int:
p = argparse.ArgumentParser(description=__doc__)
p.add_argument("--url", default=os.getenv("ANDROID_WS_URL", "ws://127.0.0.1:8790/ws"))
p.add_argument("--token", default=os.getenv("ANDROID_TOKEN", ""))
p.add_argument("--device", default=f"probe-{uuid.uuid4().hex[:8]}")
p.add_argument("--send", default="hello")
p.add_argument("--upload", default="",
help="M4: file to upload (chunked) and attach via media_refs")
p.add_argument("--pull-offer", action="store_true",
help="M4: pull any media.offer that arrives during the turn")
p.add_argument("--sync", type=int, default=None,
help="M5: send sync {cursor} after pairing, print replay, exit")
p.add_argument("--fcm-token", default="",
help="M5: FCM token to attach to the hello payload")
p.add_argument("--fcm-reg", action="store_true",
help="M5: send fcm.register after pairing (uses --fcm-token)")
p.add_argument("--timeout", type=float, default=120.0)
p.add_argument("--authfail", action="store_true",
help="expect an auth rejection (wrong token)")
p.add_argument("--assert-turn", action="store_true",
help="assert message.start -> >=1 message.update -> message.stop")
p.add_argument("--assert-reasoning", action="store_true",
help="assert the final message.stop carries non-empty reasoning")
p.add_argument("--assert-tools", action="store_true",
help="assert >=1 tool.start with a matching tool.end")
p.add_argument("--assert-commentary", action="store_true",
help="assert >=1 commentary frame")
p.add_argument("--assert-read-receipt", action="store_true",
help="assert a read.receipt arrives after the sent message "
"(SKIP if absent; M7)")
p.add_argument("--assert-status", action="store_true",
help="assert a status frame is received (SKIP if absent; M7)")
p.add_argument("--search", default="",
help="M3: send search {query, scope, limit}, assert >=1 hit")
p.add_argument("--scope", choices=("all", "chat"), default="all",
help="search scope (default all)")
p.add_argument("--chat-id", default="android:default",
help="chat_id for --scope chat (default android:default)")
p.add_argument("--channel-create", default="",
help="M3: create a channel, print its chat_id, exit")
p.add_argument("--channel-delete", default="",
help="M3: delete (archive) a channel, exit")
p.add_argument("--channel-list", action="store_true",
help="M3: list channels, exit")
p.add_argument("--watch", default="",
help="wait up to --timeout for a message to land in this chat_id")
p.add_argument("--offer-grace", type=float, default=15.0,
help="seconds to wait for a media.offer after the final "
"message when --pull-offer (default 15)")
args = p.parse_args()
if not args.token and not args.authfail:
p.error("--token (or $ANDROID_TOKEN) is required")
if args.assert_read_receipt and not args.send:
p.error("--assert-read-receipt requires --send (the receipt must follow "
"the sent message)")
return asyncio.run(run(args))
if __name__ == "__main__":
sys.exit(main())