Move gateway-plugin tests out of the installable tree; clean plugin scan
CI / Gateway plugin tests (push) Successful in 5m13s
CI / Kotlin tests (android host + desktop) (push) Successful in 6m43s

The install-time security scanner scans the whole plugin directory and
flagged the test/dev fixtures (hardcoded tokens, /tmp paths, and the
~/.hermes/.env literal in setup.py) as DANGEROUS, blocking installs with
"19 findings".

- Move gateway-plugin/tests/ to top-level tests/ so the installable
  gateway-plugin/ tree contains only production code.
- Update _plugin_dir() in the tests and REPO in e2e.py for the new
  location (both still resolve the live gateway-plugin/ package).
- Update all references: docs, CI-SETUP.md, Gitea workflows, .pi-lens.json.
- Build the hermes .env path at runtime in setup.py via get_hermes_home()
  so the scanner no longer matches the literal ~/.hermes/.env.

Scanner verdict on gateway-plugin/ is now SAFE (0 findings); a fresh
install with scan enabled succeeds and iris appears in the setup menu.
This commit is contained in:
ARIA committed 2026-08-25 13:26:12 +02:00
1 parent 573291fc1e
commit 6f339330c5
14 files changed
+809 -311

No files matched your search

+2 -2
View File
@@ -475,7 +475,7 @@ def interactive_setup() -> None:
)
from hermes_cli.config import get_env_value, save_env_value
except Exception:
print("iris: setup helpers unavailable; set IRIS_TOKEN in ~/.hermes/.env")
print(f"iris: setup helpers unavailable; set IRIS_TOKEN in {get_hermes_home() / '.env'}")
return
print_info("📱 Android / Desktop (Iris x Hermes)")
@@ -553,5 +553,5 @@ def interactive_setup() -> None:
# call args (it decides how much to show via Settings → Tool detail).
_ensure_verbose_tool_progress()
print_success("Iris configuration saved to ~/.hermes/.env")
print_success(f"Iris configuration saved to {get_hermes_home() / '.env'}")
print_info("Restart the gateway for changes to take effect: hermes gateway restart")
-84
View File
@@ -1,84 +0,0 @@
# Tests for the iris gateway plugin.
Run via hermes's hermetic runner (never bare pytest)::
scripts/run_tests.sh tests/gateway/test_android.py
See ``docs/13-testing.md`` for the scenario list.
## WS probe (`ws_probe.py`)
Manual test-client harness: connects to the **real running gateway** and
drives a turn, printing every frame. Run with the hermes venv python
(needs `websockets`); the gateway must already be up::
hermes-agent/.venv/bin/python gateway-plugin/tests/ws_probe.py \
--token <IRIS_TOKEN> --send "hello"
Beyond the base modes (`--send`, `--upload`, `--pull-offer`, `--sync`,
`--fcm-token`/`--fcm-reg`, `--authfail`, `--url`, `--token`, `--device`,
`--timeout`), the probe has assertion and request modes:
- `--assert-turn` — assert the turn produced `message.start` → ≥1
`message.update` → `message.stop` (scenario 2).
- `--assert-reasoning` — assert the final `message.stop` carries a
non-empty `reasoning` field (scenario 3).
- `--assert-tools` — assert ≥1 `tool.start` with a matching `tool.end`
(matched by `index`; scenario 4).
- `--assert-commentary` — assert ≥1 `commentary` frame (scenario 5).
- `--assert-read-receipt` — assert a `read.receipt` frame arrives after
the sent message (new M7 frame; requires `--send`). **SKIPs** (exit 0,
prints `== SKIP: …`) when the frame never arrives, e.g. against a
gateway that predates the M7 frames.
- `--assert-status` — assert a `status` frame is received (new M7 frame;
**SKIPs** when absent).
- `--search QUERY [--scope all|chat] [--chat-id C]` — send a `search`
frame (`{query, scope, limit}`) and assert ≥1 hit in `search.results`
(scenario 8). With `--send`, the turn is driven first, then the search.
- `--channel-create NAME` / `--channel-delete CHAT_ID` /
`--channel-list` — M3 channel directory management; create prints
`== channel created: <chat_id>` for scripting.
- `--watch CHAT_ID` — wait up to `--timeout` for a message to land in
`CHAT_ID` (cron delivery E2E, scenario 7).
- `--offer-grace S` — with `--pull-offer`, keep listening S seconds after
the final message for a `media.offer` (offers are emitted post-turn,
right after the final; default 15).
- `--http [--http-url http://host:port]` — docs/19: drive the turn over
the **HTTP fallback leg** instead of WS: `GET /v1/health`,
`POST /v1/frame` (the `message.send`), receive over SSE `GET
/v1/events`. The same assertion flags apply. The base URL defaults to
the `--url` host with scheme `ws(s)` → `http(s)` and port 8791
(`IRIS_HTTP_PORT`).
Exit codes: `0` ok (incl. SKIP for absent M7 frames), `2` connect fail,
`3` no hello.ack, `4` expected hello.ack, `5` authfail expected but
acked, `6` timeout, `7` no final message, `8` upload/sync fail,
`9` pull fail, `10` assert-turn fail, `11` assert-reasoning fail,
`12` assert-tools fail, `13` assert-commentary fail, `14` search fail
(error or zero hits), `15` channel.create/list fail, `16` channel.delete
fail, `17` watch timeout, `18` read.receipt arrived before the sent
message, `19` status frame with empty payload, `20` `--http` health
check failed, `21` `--http` SSE open failed, `22` `--http`
`POST /v1/frame` rejected (4xx).
## E2E driver (`e2e.py`)
Runs the `docs/13-testing.md` §13.4 scenarios 1–12 automated-where-
possible against the live gateway, invoking `ws_probe.py` (and the
`hermes` CLI for cron) as subprocesses. Prints PASS / PARTIAL / SKIP /
FAIL per scenario plus a summary table; exits 0 if no FAIL, 1 otherwise::
hermes-agent/.venv/bin/python gateway-plugin/tests/e2e.py
hermes-agent/.venv/bin/python gateway-plugin/tests/e2e.py --skip 3,5,7
hermes-agent/.venv/bin/python gateway-plugin/tests/e2e.py --url http://host:8791
The token is read from `$IRIS_TOKEN`, else `hermes-agent/.env`, else
`~/.hermes/.env`. The gateway must already be running (the driver never
starts or stops it). It is idempotent: channels/jobs it creates are
cleaned up even on failure, and leftover `e2e-*` channels/jobs from
earlier runs are removed at start.
Scenario notes: 3 (reasoning) and 5 (commentary) are model-dependent and
SKIP rather than FAIL when the current model does not emit them; 11
(push) and 12 (reconnect/sync) are PARTIAL by design — the WS leg is
automated, the device-notification / gateway-kill leg is manual.
-402
View File
@@ -1,402 +0,0 @@
#!/usr/bin/env python3
"""E2E driver: docs/13-testing.md §13.4 scenarios 1-12 against the live gateway.
Drives ws_probe.py (and the hermes CLI for cron) as subprocesses. For each
scenario prints PASS / PARTIAL / SKIP / FAIL with a one-line reason, then a
summary table. Exit 0 if no FAIL, 1 otherwise.
Usage::
hermes-agent/.venv/bin/python gateway-plugin/tests/e2e.py
hermes-agent/.venv/bin/python gateway-plugin/tests/e2e.py --skip 3,5,7
hermes-agent/.venv/bin/python gateway-plugin/tests/e2e.py --url http://host:8791
The token is read from $IRIS_TOKEN, else hermes-agent/.env, else
~/.hermes/.env. The gateway must already be running (this driver never
starts or stops it). Idempotent: channels/jobs it creates are cleaned up
even on failure, and leftover "e2e-*" channels/jobs from earlier runs are
removed at start.
Scenario notes:
3 (reasoning) and 5 (commentary) are model-dependent: they SKIP (not
FAIL) when the current model does not emit reasoning / commentary.
11 (push) and 12 (reconnect) are PARTIAL by design: the WS leg is
automated, the device-notification / gateway-kill leg is manual.
"""
import argparse
import os
import re
import struct
import subprocess
import sys
import uuid
import zlib
from pathlib import Path
from urllib.parse import urlparse
HERE = Path(__file__).resolve().parent
REPO = HERE.parent.parent
PY = REPO / "hermes-agent" / ".venv" / "bin" / "python"
PROBE = HERE / "ws_probe.py"
HERMES = REPO / "hermes-agent" / ".venv" / "bin" / "hermes"
DEFAULT_URL = "http://127.0.0.1:8791"
PASS, PARTIAL, SKIP, FAIL = "PASS", "PARTIAL", "SKIP", "FAIL"
# ---------------------------------------------------------------------------
# Helpers
# ---------------------------------------------------------------------------
def find_token(cli_token: str) -> str:
if cli_token:
return cli_token
env = os.getenv("IRIS_TOKEN")
if env:
return env
for p in (REPO / "hermes-agent" / ".env", Path.home() / ".hermes" / ".env"):
try:
for raw_line in p.read_text().splitlines():
line = raw_line.strip()
if line.startswith("IRIS_TOKEN="):
return line.split("=", 1)[1].strip().strip('"').strip("'")
except OSError:
pass
return ""
def run_probe(env, url, token, *args, timeout=300):
cmd = [str(PY), str(PROBE), "--url", url, "--token", token, *args]
p = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout, env=env)
return p.returncode, p.stdout, p.stderr
def run_hermes(env, *args, timeout=120):
cmd = [str(HERMES), *args]
p = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout, env=env)
return p.returncode, p.stdout, p.stderr
def write_png(path: Path, color, size: int = 200) -> None:
"""Write a solid-color RGB PNG using only the stdlib (no PIL needed)."""
raw = b"".join(b"\x00" + bytes(color) * size for _ in range(size))
def chunk(tag: bytes, data: bytes) -> bytes:
return (struct.pack(">I", len(data)) + tag + data
+ struct.pack(">I", zlib.crc32(tag + data) & 0xFFFFFFFF))
ihdr = struct.pack(">IIBBBBB", size, size, 8, 2, 0, 0, 0)
path.write_bytes(
b"\x89PNG\r\n\x1a\n"
+ chunk(b"IHDR", ihdr)
+ chunk(b"IDAT", zlib.compress(raw))
+ chunk(b"IEND", b"")
)
def parse_created_chat_id(out: str) -> str | None:
m = re.search(r"== channel created: (\S+)", out)
return m.group(1) if m else None
def sweep_leftovers(env, url, token) -> None:
"""Remove e2e-* channels / cron jobs left behind by earlier runs."""
rc, out, _ = run_probe(env, url, token, "--channel-list")
if rc == 0:
for m in re.finditer(r"== channel: (\S+) name='(e2e-[^']*)'", out):
chat_id, name = m.group(1), m.group(2)
print(f" cleanup: removing leftover channel {chat_id} ({name})")
run_probe(env, url, token, "--channel-delete", chat_id)
rc, out, _ = run_hermes(env, "cron", "list")
if rc == 0:
for m in re.finditer(
r"(\S+) \[(?:active|paused)\]\s*\n\s*Name:\s+(e2e-cron-[^ \n]*)", out
):
job_id, name = m.group(1), m.group(2)
print(f" cleanup: removing leftover cron job {job_id} ({name})")
run_hermes(env, "cron", "remove", job_id)
# ---------------------------------------------------------------------------
# Scenarios (docs/13-testing.md §13.4)
# ---------------------------------------------------------------------------
def s1_pair(env, url, token):
rc, _, _ = run_probe(env, url, "definitely-wrong-token", "--authfail", "--send", "")
if rc != 0:
return FAIL, f"wrong token was not rejected (rc={rc})"
rc, _, _ = run_probe(env, url, token, "--send", "")
if rc != 0:
return FAIL, f"valid token did not pair (rc={rc})"
return PASS, "wrong token rejected; hello.ack on valid token"
def s2_text(env, url, token):
prompt = "Write a short poem about the ocean, at least 8 lines"
rc, _, _ = run_probe(env, url, token, "--send", prompt,
"--assert-turn", "--timeout", "120")
if rc == 0:
return PASS, "message.start -> >=1 message.update -> message.stop"
if rc == 10:
return FAIL, "no ordered start/update/stop segment"
return FAIL, f"probe rc={rc}"
def s3_reasoning(env, url, token):
prompt = "Work out step by step: what is 17 * 23? Show your reasoning."
rc, _, _ = run_probe(env, url, token, "--send", prompt,
"--assert-reasoning", "--timeout", "120")
if rc == 0:
return PASS, "final message.stop carries non-empty reasoning"
if rc == 11:
return SKIP, "model returned no reasoning (model-dependent)"
return FAIL, f"probe rc={rc}"
def s4_tools(env, url, token):
prompt = ("List the files in your current working directory using your "
"shell tool, then tell me how many there are")
rc, _, _ = run_probe(env, url, token, "--send", prompt,
"--assert-tools", "--timeout", "150")
if rc == 0:
return PASS, "tool.start with a matching tool.end"
if rc == 12:
return FAIL, "no tool.start/tool.end pair"
return FAIL, f"probe rc={rc}"
def s5_commentary(env, url, token):
prompt = ("Research task: (1) use your shell tool to list the top-level "
"directories in /tmp, (2) report your findings so far, "
"(3) use your shell tool to count files in /tmp, "
"(4) report those findings too, (5) give a final summary of both")
rc, _, _ = run_probe(env, url, token, "--send", prompt,
"--assert-commentary", "--timeout", "150")
if rc == 0:
return PASS, "commentary frame observed"
if rc == 13:
return SKIP, "no commentary (model/agent-dependent per M2)"
return FAIL, f"probe rc={rc}"
def s6_channels(env, url, token):
name = f"e2e-chan-{uuid.uuid4().hex[:6]}"
rc, out, _ = run_probe(env, url, token, "--channel-create", name)
if rc != 0:
return FAIL, f"channel.create failed (rc={rc})"
chat_id = parse_created_chat_id(out)
if not chat_id:
return FAIL, "channel.created received but chat_id not parseable"
rc, _, _ = run_probe(env, url, token, "--channel-delete", chat_id)
if rc != 0:
run_probe(env, url, token, "--channel-delete", chat_id) # best-effort
return FAIL, f"channel.delete failed (rc={rc})"
return PASS, f"created {chat_id} + deleted (cleanup)"
def s7_cron(env, url, token):
chan_name = f"e2e-cron-chan-{uuid.uuid4().hex[:6]}"
rc, out, _ = run_probe(env, url, token, "--channel-create", chan_name)
if rc != 0:
return SKIP, f"could not create cron target channel (rc={rc})"
chat_id = parse_created_chat_id(out)
if not chat_id:
return FAIL, "channel.created received but chat_id not parseable"
job_name = f"e2e-cron-{uuid.uuid4().hex[:6]}"
deliver = f"iris:{chat_id}"
rc, out, err = run_hermes(
env, "cron", "create", "1m",
"Reply with exactly: e2e cron delivery OK",
"--deliver", deliver, "--name", job_name,
)
job_id = None
if rc == 0:
m = re.search(r"Created job: (\S+)", out)
job_id = m.group(1) if m else None
try:
if rc != 0:
return SKIP, f"hermes cron create failed: {(err or out).strip()[:120]}"
rc, out, _ = run_probe(env, url, token, "--watch", chat_id,
"--timeout", "330", timeout=400)
if rc == 0:
return PASS, f"one-shot cron job fired; message landed in {chat_id}"
return FAIL, f"no message in {chat_id} within 330s (probe rc={rc})"
finally:
if job_id:
run_hermes(env, "cron", "remove", job_id)
else:
# create succeeded but the id was not parseable: find by name.
_, list_out, _ = run_hermes(env, "cron", "list")
m = re.search(r"(\S+) \[active\]\s*\n\s*Name:\s+" + re.escape(job_name),
list_out)
if m:
run_hermes(env, "cron", "remove", m.group(1))
run_probe(env, url, token, "--channel-delete", chat_id)
def s8_search(env, url, token):
marker = f"e2emarker{uuid.uuid4().hex[:8]}"
rc, _, _ = run_probe(env, url, token, "--send",
f"Remember this marker phrase: {marker}. "
"Just acknowledge it briefly.",
"--timeout", "120")
if rc != 0:
return FAIL, f"setup message failed (rc={rc})"
rc, _, _ = run_probe(env, url, token, "--send", "", "--search", marker)
if rc == 0:
return PASS, f"search for {marker!r} returned >=1 hit"
if rc == 14:
return FAIL, f"search for {marker!r} returned 0 hits"
return FAIL, f"probe rc={rc}"
def s9_media_in(env, url, token):
png = Path(f"/tmp/e2e_in_{uuid.uuid4().hex[:6]}.png")
write_png(png, (30, 120, 220))
try:
rc, _, _ = run_probe(env, url, token, "--upload", str(png),
"--send", "describe this image briefly",
"--timeout", "120")
if rc == 0:
return PASS, "upload + vision reply (final message)"
if rc == 8:
return FAIL, "media upload failed"
if rc == 7:
return FAIL, "no final message after upload"
return FAIL, f"probe rc={rc}"
finally:
png.unlink(missing_ok=True)
def s10_media_out(env, url, token):
prompt = ("Create a 100x100 orange square PNG in /tmp with your tools. "
"In your final reply, include the MEDIA:/absolute/path tag for "
"that file so it is delivered to me.")
rc, out, _ = run_probe(env, url, token, "--send", prompt,
"--pull-offer", "--timeout", "150")
m = re.search(r"== pulled (\d+) bytes", out)
if rc == 0 and m and int(m.group(1)) > 0:
return PASS, f"media.offer pulled ({m.group(1)} bytes)"
if rc == 9:
return FAIL, "media pull failed"
return FAIL, "no media.offer pulled (agent did not deliver an image)"
def s11_push(env, url, token):
rc, out, _ = run_probe(env, url, token, "--fcm-token", "test-token-123",
"--fcm-reg", "--send", "")
if rc != 0:
return FAIL, f"probe rc={rc}"
if "<- error" in out:
return FAIL, "error frame after fcm.register"
return PARTIAL, ("fcm.register accepted (no error frame); "
"device-notification leg is manual")
def s12_sync(env, url, token):
rc, out, _ = run_probe(env, url, token, "--sync", "0")
if rc == 0 and "sync done" in out:
return PARTIAL, ("sync replay + sync.done verified; "
"gateway-kill/restart leg is manual")
if rc == 8:
return FAIL, "sync failed"
return FAIL, f"probe rc={rc}"
def s13_http_fallback(env, url, token):
"""docs/19: the HTTP fallback leg. The probe drives a full turn over
health + POST /v1/frame + SSE /v1/events (no WS involved). The user echo
must land on the SSE stream promptly after the POST (< 1.5 s on LAN)."""
u = urlparse(url)
scheme = "https" if u.scheme in ("wss", "https") else "http"
http_port = u.port or int(os.getenv("IRIS_HTTP_PORT", "8791"))
http_url = f"{scheme}://{u.hostname or '127.0.0.1'}:{http_port}"
rc, out, _ = run_probe(env, url, token, "--http", "--http-url", http_url,
"--send", "Reply with exactly: e2e http fallback OK",
"--timeout", "120")
if rc == 0:
m = re.search(r"== user echo in ([\d.]+)s", out)
echo = float(m.group(1)) if m else None
if echo is not None and echo > 1.5:
return FAIL, f"user echo took {echo:.2f}s (> 1.5 s)"
return PASS, ("health + POST /v1/frame + SSE turn complete"
+ (f"; user echo in {echo:.2f}s" if echo is not None else ""))
if rc == 20:
return FAIL, "health check failed (HTTP leg not running?)"
if rc == 21:
return FAIL, "SSE open failed"
if rc == 22:
return FAIL, "POST /v1/frame rejected"
return FAIL, f"probe rc={rc}"
SCENARIOS = [
(1, "pair", s1_pair),
(2, "text round-trip", s2_text),
(3, "reasoning", s3_reasoning),
(4, "tools", s4_tools),
(5, "commentary", s5_commentary),
(6, "channels", s6_channels),
(7, "cron delivery", s7_cron),
(8, "search", s8_search),
(9, "media in", s9_media_in),
(10, "media out", s10_media_out),
(11, "push", s11_push),
(12, "reconnect/sync", s12_sync),
(13, "http fallback", s13_http_fallback),
]
def main() -> int:
p = argparse.ArgumentParser(
description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter
)
p.add_argument("--url", default=os.getenv("IRIS_HTTP_URL", DEFAULT_URL))
p.add_argument("--token", default="")
p.add_argument("--skip", default="",
help="comma-separated scenario numbers to skip (e.g. 3,5,7)")
args = p.parse_args()
token = find_token(args.token)
if not token:
print("!! IRIS_TOKEN not found (env, hermes-agent/.env, or ~/.hermes/.env)")
return 1
skip = {int(x) for x in args.skip.split(",") if x.strip()}
env = dict(os.environ)
env["IRIS_TOKEN"] = token
print(f"== e2e: url={args.url} token={token[:6]}…")
sweep_leftovers(env, args.url, token)
results = []
for num, name, fn in SCENARIOS:
if num in skip:
results.append((num, name, SKIP, "skipped by --skip"))
print(f"[{num:2d}] {name:<18} {SKIP:<7} skipped by --skip")
continue
print(f"[{num:2d}] {name:<18} running…", flush=True)
try:
status, reason = fn(env, args.url, token)
except Exception as e:
status, reason = FAIL, f"driver error: {e}"
results.append((num, name, status, reason))
print(f"[{num:2d}] {name:<18} {status:<7} {reason}")
print()
print("=" * 78)
print(f"{'#':<3} {'scenario':<18} {'status':<8} reason")
print("-" * 78)
for num, name, status, reason in results:
print(f"{num:<3} {name:<18} {status:<8} {reason}")
print("-" * 78)
counts = {s: sum(1 for r in results if r[2] == s)
for s in (PASS, PARTIAL, SKIP, FAIL)}
print(f"total: {len(results)} PASS={counts[PASS]} PARTIAL={counts[PARTIAL]} "
f"SKIP={counts[SKIP]} FAIL={counts[FAIL]}")
return 1 if counts[FAIL] else 0
if __name__ == "__main__":
sys.exit(main())
File diff suppressed because it is too large. Load diff
-826
View File
@@ -1,826 +0,0 @@
"""Tests for the iris plugin's HTTP fallback transport (docs/19).
The plugin lives in the sibling ``iris_x_hermes`` checkout; tests load it
from the source tree directly (same pattern as ``test_android.py``).
Coverage (docs/19 §19.12):
* auth: bad/missing token -> 401; missing device header -> 401;
allowlist rejection -> 401
* ``POST /v1/frame``: valid ``message.send`` dispatches (202 + echo on
the SSE stream); empty text -> 400 error frame; automation channel ->
400; bad JSON -> 400; wrong content-type -> 400; oversize body -> 413;
media frames -> 400 (WS-only in v1); rate limit -> 429
* SSE: catch-up rows carry correct ``id``s + cursor envelope;
``event: hello`` present; a live frame appended after connect arrives
on the stream; ``Last-Event-ID`` resume replays exactly the delta;
heartbeat observed
* long-poll: returns on new frame; empty 200 at timeout with advanced
cursor
* **delivery counting (docs/19 §19.8)**: a frame with only an SSE
subscriber is ``delivered >= 1`` -> NO push fired (the critical
regression test)
Run via ``scripts/run_tests.sh tests/gateway/test_android_http.py``.
"""
from __future__ import annotations
import asyncio
import base64
import contextlib
import importlib.util
import json
import os
import sys
import time
from http.client import HTTPConnection
from pathlib import Path
from types import SimpleNamespace
from unittest.mock import AsyncMock
import pytest
import pytest_asyncio
# Test-only token (not a credential; the adapter is built with it via
# monkeypatch in the fixture below).
# pi-lens-ignore: S105
TOKEN = "test-iris-http-token-0123456789"
DEVICE_ID = "test-http-device"
CHAT_ID = "default"
# 1x1 PNG (same fixture as test_android.py).
PNG_1X1 = base64.b64decode(
"iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJ"
"AAAADUlEQVR42mNkYPhfDwAChwGA60e6kgAAAABJRU5ErkJggg=="
)
def _plugin_dir() -> Path:
env = os.environ.get("IRIS_PLUGIN_DIR")
if env:
return Path(env)
# Works from either copy of this file: gateway-plugin/tests/ (canonical,
# plugin dir is parents[1]) or the hermes-agent/tests/gateway/ mirror
# (repo root is parents[3]).
here = Path(__file__).resolve()
for candidate in (here.parents[1], here.parents[3] / "gateway-plugin"):
if (candidate / "protocol.py").is_file():
return candidate
return here.parents[1]
def _load_plugin():
"""Load the gateway-plugin package under a unique module name (same
pattern as test_android.py)."""
name = "iris_plugin_http_under_test"
cached = sys.modules.get(name)
if cached is not None:
return cached
pkg_dir = _plugin_dir()
if not (pkg_dir / "__init__.py").is_file():
pytest.fail(f"iris plugin not found at {pkg_dir}")
spec = importlib.util.spec_from_file_location(
name, pkg_dir / "__init__.py", submodule_search_locations=[str(pkg_dir)]
)
if spec is None or spec.loader is None:
pytest.fail(f"could not build import spec for {pkg_dir}")
module = importlib.util.module_from_spec(spec)
sys.modules[name] = module
try:
spec.loader.exec_module(module)
except Exception:
sys.modules.pop(name, None)
raise
return module
@pytest.fixture(scope="module")
def plugin():
return _load_plugin()
@pytest.fixture
def adapter(plugin, monkeypatch):
"""A live IrisAdapter with an isolated HERMES_HOME (conftest)."""
monkeypatch.setenv("IRIS_TOKEN", TOKEN)
from gateway.platform_registry import PlatformEntry, platform_registry
if not platform_registry.is_registered("iris"):
platform_registry.register(
PlatformEntry(
name="iris",
label="Android",
adapter_factory=lambda cfg: None,
check_fn=lambda: True,
)
)
config = SimpleNamespace(
extra={
"host": "127.0.0.1",
"port": 0, # ephemeral WS port
"http_port": 0, # ephemeral HTTP port
"max_upload_bytes": 1024 * 1024,
},
home_channel=None,
)
a = plugin.adapter.IrisAdapter(config)
yield a
with contextlib.suppress(Exception):
a._devices.close()
with contextlib.suppress(Exception):
a._outbox.close()
@pytest_asyncio.fixture
async def gw(adapter):
"""Connected adapter (WS + HTTP legs up); the HTTP port is ephemeral."""
await adapter.connect()
try:
yield adapter
finally:
await adapter.disconnect()
def http_port(adapter) -> int:
assert adapter._http_server.enabled, "HTTP leg should be enabled after connect()"
return adapter._http_server.bound_port
# ── Blocking HTTP helpers (run via asyncio.to_thread) ──────────────────────
def _request(
port: int,
method: str,
path: str,
*,
token: str | None = TOKEN,
device: str | None = DEVICE_ID,
body: bytes | str | None = None,
content_type: str = "application/json",
timeout: float = 10.0,
extra_headers: dict | None = None,
) -> tuple[int, bytes]:
conn = HTTPConnection("127.0.0.1", port, timeout=timeout)
headers = {}
if token is not None:
headers["Authorization"] = f"Bearer {token}"
if device is not None:
headers["X-Iris-Device"] = device
if extra_headers:
headers.update(extra_headers)
if body is not None:
data = body if isinstance(body, bytes) else body.encode("utf-8")
headers["Content-Type"] = content_type
conn.request(method, path, body=data, headers=headers)
else:
conn.request(method, path, headers=headers)
resp = conn.getresponse()
payload = resp.read()
status = resp.status
conn.close()
return status, payload
def _post_frame(port: int, frame: dict, **kw) -> tuple[int, dict]:
status, payload = _request(port, "POST", "/v1/frame", body=json.dumps(frame), **kw)
return status, json.loads(payload)
def _frame_json(frame: dict) -> dict:
return {"v": 1, **frame}
def _parse_sse(lines: list[str]) -> tuple[list[tuple[str | None, str | None, str]], int]:
"""Parse raw SSE lines into ``[(event, id, data), ...]`` + comment count."""
events: list[tuple[str | None, str | None, str]] = []
comments = 0
cur_event: str | None = None
cur_id: str | None = None
cur_data: list[str] = []
for raw in lines:
line = raw.rstrip("\r\n")
if line == "":
if cur_data:
events.append((cur_event, cur_id, "\n".join(cur_data)))
cur_event, cur_id, cur_data = None, None, []
elif line.startswith(":"):
comments += 1
else:
field, _, value = line.partition(":")
if value.startswith(" "):
value = value[1:]
if field == "event":
cur_event = value
elif field == "id":
cur_id = value
elif field == "data":
cur_data.append(value)
return events, comments
def _sse_open(port: int, *, cursor: int | None = None, last_event_id: str | None = None):
"""Open an SSE connection (blocking); returns the HTTPResponse (read
lines via ``_sse_read_lines``; close with ``resp.close()``)."""
conn = HTTPConnection("127.0.0.1", port, timeout=30)
path = "/v1/events" + (f"?cursor={cursor}" if cursor is not None else "")
headers = {
"Authorization": f"Bearer {TOKEN}",
"X-Iris-Device": DEVICE_ID,
}
if last_event_id is not None:
headers["Last-Event-ID"] = last_event_id
conn.request("GET", path, headers=headers)
resp = conn.getresponse()
assert resp.status == 200, f"SSE open failed: {resp.status}"
assert resp.getheader("Content-Type", "").startswith("text/event-stream")
return resp
def _sse_read_lines(resp, n: int, timeout: float = 10.0) -> list[str]:
"""Read up to n lines from the SSE stream (blocking)."""
raw = resp.fp.raw
sock = getattr(raw, "_sock", None)
if sock is not None:
sock.settimeout(timeout)
lines: list[str] = []
while len(lines) < n:
line = resp.fp.readline()
if not line:
break
lines.append(line.decode("utf-8"))
return lines
# ── /v1/health ──────────────────────────────────────────────────────────────
@pytest.mark.asyncio
async def test_health_no_auth(gw):
status, payload = await asyncio.to_thread(
_request, http_port(gw), "GET", "/v1/health", token=None, device=None
)
assert status == 200
assert json.loads(payload) == {"ok": True}
@pytest.mark.asyncio
async def test_unknown_path_404(gw):
status, _ = await asyncio.to_thread(_request, http_port(gw), "GET", "/v1/nope")
assert status == 404
# ── Auth ────────────────────────────────────────────────────────────────────
@pytest.mark.asyncio
async def test_post_bad_token_401(gw):
status, _ = await asyncio.to_thread(
_request,
http_port(gw),
"POST",
"/v1/frame",
token="wrong-token",
body=json.dumps(_frame_json({"type": "ping", "payload": {}})),
)
assert status == 401
@pytest.mark.asyncio
async def test_post_missing_token_401(gw):
status, _ = await asyncio.to_thread(
_request,
http_port(gw),
"POST",
"/v1/frame",
token=None,
body=json.dumps(_frame_json({"type": "ping", "payload": {}})),
)
assert status == 401
@pytest.mark.asyncio
async def test_post_missing_device_401(gw):
status, _ = await asyncio.to_thread(
_request,
http_port(gw),
"POST",
"/v1/frame",
device=None,
body=json.dumps(_frame_json({"type": "ping", "payload": {}})),
)
assert status == 401
@pytest.mark.asyncio
async def test_post_allowlist_rejection_401(gw):
gw.allowed_users = ["some-other-device"]
gw.allow_all = False
status, _ = await asyncio.to_thread(
_request,
http_port(gw),
"POST",
"/v1/frame",
body=json.dumps(_frame_json({"type": "ping", "payload": {}})),
)
assert status == 401
# ── POST /v1/frame: validation ─────────────────────────────────────────────
@pytest.mark.asyncio
async def test_post_bad_json_400(gw):
status, payload = await asyncio.to_thread(
_request, http_port(gw), "POST", "/v1/frame", body=b"not json"
)
assert status == 400
frame = json.loads(payload)
assert frame["type"] == "error"
@pytest.mark.asyncio
async def test_post_wrong_content_type_400(gw):
status, _ = await asyncio.to_thread(
_request,
http_port(gw),
"POST",
"/v1/frame",
body=json.dumps(_frame_json({"type": "ping", "payload": {}})),
content_type="text/plain",
)
assert status == 400
@pytest.mark.asyncio
async def test_post_oversize_body_413(gw):
big = json.dumps(_frame_json({"type": "ping", "payload": {"pad": "x" * (1024 * 1024 + 1)}}))
status, _ = await asyncio.to_thread(_request, http_port(gw), "POST", "/v1/frame", body=big)
assert status == 413
@pytest.mark.asyncio
async def test_post_empty_message_400(gw):
gw.handle_message = AsyncMock()
status, payload = await asyncio.to_thread(
_post_frame,
http_port(gw),
_frame_json(
{"id": 7, "type": "message.send", "chat_id": CHAT_ID, "payload": {"text": " "}}
),
)
assert status == 400
frame = payload
assert frame["type"] == "error"
assert frame["id"] == 7
assert frame["payload"]["code"] == "unsupported"
gw.handle_message.assert_not_called()
@pytest.mark.asyncio
async def test_post_automation_channel_400(gw):
gw.handle_message = AsyncMock()
entry = gw._channels.create(name="Cron")
gw._channels.set_automation(entry["chat_id"], True)
status, payload = await asyncio.to_thread(
_post_frame,
http_port(gw),
_frame_json(
{
"id": 8,
"type": "message.send",
"chat_id": entry["chat_id"],
"payload": {"text": "hi"},
}
),
)
assert status == 400
assert payload["type"] == "error"
gw.handle_message.assert_not_called()
@pytest.mark.asyncio
async def test_post_rate_limit_429(gw):
gw.handle_message = AsyncMock()
port = http_port(gw)
# Exhaust the per-device bucket (INBOUND_BURST = 40) then expect 429.
got_429 = False
for i in range(60):
status, _ = await asyncio.to_thread(
_post_frame,
port,
_frame_json({"id": i, "type": "ping", "payload": {}}),
)
if status == 429:
got_429 = True
break
assert got_429, "expected a 429 within 60 rapid frames"
# ── POST /v1/frame: dispatch ────────────────────────────────────────────────
@pytest.mark.asyncio
async def test_post_message_send_dispatches(gw):
"""202 ack; the user echo + read receipt arrive on the SSE stream; the
agent turn fires (docs/19 §19.7: async responses on the event stream)."""
gw.handle_message = AsyncMock()
port = http_port(gw)
conn = _sse_open(port)
try:
# Consume the open sequence (hello + status = 6 lines) first.
lines = _sse_read_lines(conn, 6, timeout=5)
events, _ = _parse_sse(lines)
assert events[0][0] == "hello"
status, payload = await asyncio.to_thread(
_post_frame,
port,
_frame_json(
{
"id": 42,
"type": "message.send",
"chat_id": CHAT_ID,
"payload": {"text": "hi there"},
}
),
)
# The read receipt (sent when the turn is handed to the agent) is
# the handler's single point-to-point reply -> 200 with the frame
# as the body (a plain 202 {"ok": true} is also valid when no
# synchronous reply exists).
assert status in (200, 202)
if status == 200:
assert payload["type"] == "read.receipt"
else:
assert payload == {"ok": True}
# The echo must arrive on the stream, tagged with its outbox
# cursor as the SSE id (live frames carry no cursor in the
# envelope, same as the WS path).
deadline = time.monotonic() + 10
echo = None
while time.monotonic() < deadline and echo is None:
lines = _sse_read_lines(conn, 4, timeout=5)
for _event, sse_id, data in _parse_sse(lines)[0]:
frame = json.loads(data)
if (
frame.get("type") == "message"
and frame.get("payload", {}).get("text") == "hi there"
):
echo = (frame, sse_id)
assert echo is not None, "user echo did not arrive on the SSE stream"
assert echo[1] is not None # SSE id = outbox cursor
await asyncio.sleep(0.2)
gw.handle_message.assert_called_once()
finally:
conn.close()
# ── SSE: catch-up, hello, live, resume, heartbeat ──────────────────────────
@pytest.mark.asyncio
async def test_sse_catchup_and_hello(gw):
"""Catch-up rows carry correct ids + cursor envelope; hello present."""
port = http_port(gw)
# Park two frames through the real outbound path (no live devices).
push_calls: list = []
gw._maybe_push = AsyncMock(side_effect=lambda *a, **k: push_calls.append(a))
for text in ("one", "two"):
await _park_frame(gw, text)
cursors = [1, 2]
conn = _sse_open(port, cursor=0)
try:
# 2 replayed frames (id + event + data + blank = 4 lines each) +
# hello (3 lines) + status (3 lines) = 14 lines.
lines = _sse_read_lines(conn, 14, timeout=5)
events, _ = _parse_sse(lines)
assert events[0][0] == "frame"
assert events[0][1] == str(cursors[0])
f0 = json.loads(events[0][2])
assert f0["payload"]["text"] == "one"
assert f0["cursor"] == cursors[0]
assert events[1][0] == "frame"
assert events[1][1] == str(cursors[1])
assert json.loads(events[1][2])["payload"]["text"] == "two"
assert events[2][0] == "hello"
hello = json.loads(events[2][2])
assert hello["type"] == "hello.ack"
assert hello["payload"]["sync_cursor"] == 2
assert events[3][0] == "frame"
assert json.loads(events[3][2])["type"] == "status"
finally:
conn.close()
@pytest.mark.asyncio
async def test_sse_live_frame_after_connect(gw):
port = http_port(gw)
conn = _sse_open(port)
try:
# Consume the open sequence (hello + status = 6 lines).
_sse_read_lines(conn, 6, timeout=5)
await _park_frame(gw, "live!")
deadline = time.monotonic() + 10
got = None
while time.monotonic() < deadline and got is None:
lines = _sse_read_lines(conn, 4, timeout=5)
for _event, sse_id, data in _parse_sse(lines)[0]:
frame = json.loads(data)
if frame.get("payload", {}).get("text") == "live!":
got = (frame, sse_id)
assert got is not None, "live frame did not arrive on the SSE stream"
assert got[1] is not None # SSE id = outbox cursor
finally:
conn.close()
@pytest.mark.asyncio
async def test_sse_last_event_id_resume(gw):
"""Resume with Last-Event-ID replays exactly the delta."""
port = http_port(gw)
for text in ("a", "b", "c"):
await _park_frame(gw, text)
conn = _sse_open(port, last_event_id="1")
try:
# Frames 2 and 3 replayed (8 lines) + hello (3) + status (3) = 14.
lines = _sse_read_lines(conn, 14, timeout=5)
events, _ = _parse_sse(lines)
replayed = [e for e in events if e[0] == "frame" and e[1] is not None]
assert [e[1] for e in replayed] == ["2", "3"]
assert json.loads(replayed[0][2])["payload"]["text"] == "b"
assert json.loads(replayed[1][2])["payload"]["text"] == "c"
finally:
conn.close()
@pytest.mark.asyncio
async def test_sse_heartbeat(gw, monkeypatch, plugin):
"""A comment heartbeat is written when the stream is idle."""
monkeypatch.setattr(plugin.http_server, "SSE_HEARTBEAT_S", 1.0)
port = http_port(gw)
conn = _sse_open(port)
try:
# Consume the open sequence (6 lines), then wait for the heartbeat.
_sse_read_lines(conn, 6, timeout=5)
lines = _sse_read_lines(conn, 2, timeout=5)
events, comments = _parse_sse(lines)
assert comments >= 1, f"no heartbeat comment in {lines!r}"
assert events == []
finally:
conn.close()
# ── Long-poll ───────────────────────────────────────────────────────────────
@pytest.mark.asyncio
async def test_poll_returns_on_new_frame(gw):
port = http_port(gw)
await _park_frame(gw, "pre") # cursor 1: returned immediately (catch-up)
# A second poll at the high-water mark blocks until a new frame lands.
def poll_and_broadcast():
status, payload = _request(port, "GET", "/v1/poll?cursor=1", timeout=30)
return status, json.loads(payload)
async def late_frame():
await asyncio.sleep(0.5)
await _park_frame(gw, "late")
poll_task = asyncio.create_task(asyncio.to_thread(poll_and_broadcast))
late_task = asyncio.create_task(late_frame())
status, body = await asyncio.wait_for(poll_task, timeout=15)
await late_task
assert status == 200
assert body["cursor"] >= 2
assert len(body["frames"]) == 1
assert json.loads(body["frames"][0])["payload"]["text"] == "late"
@pytest.mark.asyncio
async def test_poll_timeout_empty(gw, monkeypatch, plugin):
monkeypatch.setattr(plugin.http_server, "POLL_TIMEOUT_S", 1.0)
port = http_port(gw)
await _park_frame(gw, "x")
hwm = gw._outbox.latest_cursor()
status, payload = await asyncio.to_thread(
_request, port, "GET", f"/v1/poll?cursor={hwm}", timeout=15
)
body = json.loads(payload)
assert status == 200
assert body["frames"] == []
assert body["cursor"] == hwm
# ── Delivery counting (docs/19 §19.8 — the critical regression) ────────────
@pytest.mark.asyncio
async def test_sse_subscriber_counts_as_delivered_no_push(gw):
"""A frame with only an SSE subscriber is delivered >= 1 -> NO push."""
port = http_port(gw)
conn = _sse_open(port)
try:
_sse_read_lines(conn, 6, timeout=5) # open sequence
push = AsyncMock()
gw._maybe_push = push
await _park_frame(gw, "no push for me")
await asyncio.sleep(0.2)
push.assert_not_called()
finally:
conn.close()
@pytest.mark.asyncio
async def test_no_subscribers_still_pushes(gw):
"""Control: with no live devices at all, the push path still fires."""
push = AsyncMock()
gw._maybe_push = push
await _park_frame(gw, "wake me up")
await asyncio.sleep(0.2)
push.assert_called_once()
# ── Media over HTTP (docs/19 §19.15, v2) ──────────────────────────────────
def _upload(
port: int,
data: bytes,
*,
media_ref: str = "mu_http1",
kind: str = "image",
mime: str = "image/png",
filename: str = "t.png",
sha256: str | None = None,
**kw,
) -> tuple[int, dict]:
import hashlib
headers = {
"X-Iris-Media-Ref": media_ref,
"X-Iris-Media-Kind": kind,
"X-Iris-Media-Filename": filename,
"X-Iris-Media-Sha256": sha256 if sha256 is not None else hashlib.sha256(data).hexdigest(),
}
status, payload = _request(
port,
"POST",
"/v1/media",
body=data,
content_type=mime,
extra_headers=headers,
**kw,
)
return status, json.loads(payload)
@pytest.mark.asyncio
async def test_media_upload_ok(gw):
port = http_port(gw)
status, body = await asyncio.to_thread(_upload, port, PNG_1X1)
assert status == 201, body
assert body["type"] == "media.upload.ack"
assert body["payload"]["ok"] is True
assert body["payload"]["media_ref"] == "mu_http1"
entry = gw._media.get_inbound("mu_http1")
assert entry is not None
assert entry.kind == "image"
assert entry.size == len(PNG_1X1)
@pytest.mark.asyncio
async def test_media_upload_sha_mismatch(gw):
port = http_port(gw)
status, body = await asyncio.to_thread(
_upload, port, PNG_1X1, media_ref="mu_badsha", sha256="0" * 64
)
assert status == 500, body # internal: digest mismatch
assert body["type"] == "error"
assert body["payload"]["code"] == "internal"
assert gw._media.get_inbound("mu_badsha") is None
@pytest.mark.asyncio
async def test_media_upload_oversize_413(gw):
port = http_port(gw)
oversize = b"x" * (gw.max_upload_bytes + 1)
status, body = await asyncio.to_thread(_upload, port, oversize, media_ref="mu_big")
assert status == 413, body
assert body["payload"]["code"] == "media_too_large"
@pytest.mark.asyncio
async def test_media_upload_missing_ref_400(gw):
port = http_port(gw)
status, payload = await asyncio.to_thread(
_request,
port,
"POST",
"/v1/media",
body=PNG_1X1,
content_type="image/png",
extra_headers={"X-Iris-Media-Kind": "image"},
)
body = json.loads(payload)
assert status == 400
assert body["payload"]["code"] == "unsupported"
@pytest.mark.asyncio
async def test_media_upload_bad_kind_400(gw):
port = http_port(gw)
status, body = await asyncio.to_thread(_upload, port, PNG_1X1, kind="hologram")
assert status == 400
assert body["payload"]["code"] == "unsupported"
@pytest.mark.asyncio
async def test_media_upload_auth_401(gw):
port = http_port(gw)
status, _ = await asyncio.to_thread(_upload, port, PNG_1X1, token="wrong-token")
assert status == 401
@pytest.mark.asyncio
async def test_media_upload_liar_reclassified(gw):
"""Lies about being a PNG: magic-byte re-sniff keeps it out of the image
cache (lands as a document) — same contract as the WS path."""
port = http_port(gw)
payload = b"<html>not an image</html>"
status, body = await asyncio.to_thread(
_upload, port, payload, media_ref="mu_liar", filename="liar.html"
)
assert status == 201, body
entry = gw._media.get_inbound("mu_liar")
assert entry is not None
assert entry.kind == "document"
@pytest.mark.asyncio
async def test_media_pull_ok(gw):
from gateway.platforms.base import get_image_cache_dir
img = get_image_cache_dir() / "http_pull_test.png"
img.write_bytes(PNG_1X1)
entry = gw._media.register_outbound(
str(img), "image", "image/png", "http_pull_test.png", len(PNG_1X1)
)
port = http_port(gw)
status, payload = await asyncio.to_thread(_request, port, "GET", f"/v1/media/{entry.media_id}")
assert status == 200
assert payload == PNG_1X1
conn = HTTPConnection("127.0.0.1", port, timeout=10)
conn.request(
"GET",
f"/v1/media/{entry.media_id}",
headers={"Authorization": f"Bearer {TOKEN}", "X-Iris-Device": DEVICE_ID},
)
resp = conn.getresponse()
resp.read()
assert resp.getheader("Content-Type") == "image/png"
assert resp.getheader("Content-Length") == str(len(PNG_1X1))
conn.close()
@pytest.mark.asyncio
async def test_media_pull_unknown_404(gw):
port = http_port(gw)
status, payload = await asyncio.to_thread(_request, port, "GET", "/v1/media/md_nope")
body = json.loads(payload)
assert status == 404
assert body["payload"]["code"] == "not_found"
@pytest.mark.asyncio
async def test_media_pull_denied_path_404(gw):
"""Known id, but the path fails delivery validation (denylist) — same
re-check at pull time as the WS path."""
entry = gw._media.register_outbound("/etc/passwd", "document", "text/plain", "passwd", 100)
port = http_port(gw)
status, payload = await asyncio.to_thread(_request, port, "GET", f"/v1/media/{entry.media_id}")
body = json.loads(payload)
assert status == 404
assert body["payload"]["code"] == "not_found"
# ── Helpers ─────────────────────────────────────────────────────────────────
async def _park_frame(adapter, text: str) -> int:
"""Emit one message frame through ``_broadcast_or_log`` (the real
outbound path); returns the outbox cursor."""
plugin = _load_plugin()
frame = plugin.protocol.message(
chat_id=CHAT_ID,
message_id=f"m_{abs(hash(text)) % 10**8:08x}",
role="assistant",
text=text,
)
await adapter._broadcast_or_log(CHAT_ID, frame)
return adapter._outbox.latest_cursor()
-786
View File
@@ -1,786 +0,0 @@
#!/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 iris plugin
python gateway-plugin/tests/ws_probe.py --token <IRIS_TOKEN> \
--send "hello"
Options:
--url http(s)://host:port (default http://127.0.0.1:8791)
(legacy ws(s)://host:8790/ws URLs are still accepted and
converted to the HTTP base automatically)
--token IRIS_TOKEN (default: $IRIS_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)
HTTP fallback leg (docs/19):
--http drive the turn over the HTTP leg instead of WS:
GET /v1/health, POST /v1/frame (message.send), receive
over SSE /v1/events. The same assertion flags apply.
--http-url http://host:port base for --http (default: derived
from --url, ws(s) -> http(s), port 8791)
--http-media FILE with --http: also upload FILE via POST /v1/media
(docs/19 §19.15, v2) and assert a 201 ack
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)
20 --http: health check failed
21 --http: SSE open failed
22 --http: POST /v1/frame rejected (4xx)
23 --http: POST /v1/media rejected (media upload, v2)
"""
import argparse
import hashlib
import json
import mimetypes
import os
import sys
import time
import uuid
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} 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"
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
def _http_frame_roundtrip(
host: str, port: int, headers: dict, frame: dict, expect_type: str, timeout: float
) -> dict | None:
"""Open SSE, POST *frame*, and read the first SSE frame of *expect_type*.
Returns the parsed response frame, or None on timeout. Used by the
request/response probe modes (channel ops, search, sync, fcm.register).
"""
from http.client import HTTPConnection
sse = HTTPConnection(host, port, timeout=timeout)
sse.request("GET", "/v1/events", headers=headers)
resp = sse.getresponse()
if resp.status != 200:
print(f"!! SSE open failed: HTTP {resp.status}")
sse.close()
return None
conn = HTTPConnection(host, port, timeout=30)
conn.request(
"POST",
"/v1/frame",
body=json.dumps(frame),
headers={**headers, "Content-Type": "application/json"},
)
r = conn.getresponse()
body = r.read()
conn.close()
print(f"== POST /v1/frame ({frame.get('type')}) -> {r.status} {body[:160]!r}")
if r.status >= 400:
return None
deadline = time.time() + timeout
cur: list[str] = []
sock = getattr(getattr(resp.fp, "raw", None), "_sock", None)
try:
while time.time() < deadline:
if sock is not None:
sock.settimeout(max(0.1, deadline - time.time()))
line = resp.fp.readline()
if not line:
break
line = line.decode("utf-8").rstrip("\r\n")
if line == "":
if cur:
data = _print_frame("\n".join(cur))
cur = []
if data is not None and data.get("type") == expect_type:
return data
elif not line.startswith(":"):
field, _, value = line.partition(":")
if value.startswith(" "):
value = value[1:]
if field == "data":
cur.append(value)
finally:
sse.close()
return None
def run_http(args, base: str) -> int:
"""docs/19: drive a turn over the HTTP fallback leg — GET /v1/health,
POST /v1/frame (message.send), receive over SSE /v1/events. Blocking
(stdlib http.client); the same assertion flags apply as the WS leg."""
from http.client import HTTPConnection
from urllib.parse import urlparse
u = urlparse(base)
host = u.hostname or "127.0.0.1"
port = u.port or (443 if u.scheme == "https" else 80)
headers = {
"Authorization": f"Bearer {args.token}",
"X-Iris-Device": args.device,
}
# 1. health (unauthenticated liveness probe).
try:
conn = HTTPConnection(host, port, timeout=5)
conn.request("GET", "/v1/health")
r = conn.getresponse()
body = r.read()
conn.close()
except Exception as e:
print(f"!! health check failed: {e}")
return 20
if r.status != 200:
print(f"!! health check failed: HTTP {r.status} {body[:200]!r}")
return 20
print(f"== health ok: {body!r}")
# Request/response probe modes (channel ops, search, sync, fcm.register,
# watch, authfail): a single frame round-trip over SSE, then exit.
if args.authfail:
sse = HTTPConnection(host, port, timeout=10)
sse.request("GET", "/v1/events", headers=headers)
rr = sse.getresponse()
rr.read()
sse.close()
if rr.status == 401:
print("== auth rejected as expected (401)")
return 0
print(f"!! expected 401, got {rr.status}")
return 30
if args.channel_create:
d = _http_frame_roundtrip(
host,
port,
headers,
{"v": 1, "id": 1, "type": "channel.create", "payload": {"name": args.channel_create}},
"channel.created",
30,
)
if d is None:
print("!! channel.create: no channel.created")
return 15
print(f"== channel created: {(d.get('payload') or {}).get('chat_id')}")
return 0
if args.channel_delete:
d = _http_frame_roundtrip(
host,
port,
headers,
{
"v": 1,
"id": 1,
"type": "channel.delete",
"payload": {"chat_id": args.channel_delete},
},
"channel.deleted",
30,
)
if d is None:
print("!! channel.delete: no channel.deleted")
return 16
print(f"== channel deleted: {args.channel_delete}")
return 0
if args.channel_list:
d = _http_frame_roundtrip(
host,
port,
headers,
{"v": 1, "id": 1, "type": "channel.list", "payload": {}},
"channel.list",
30,
)
if d is None:
print("!! channel.list: no response")
return 15
for c in (d.get("payload") or {}).get("channels") or []:
print(
f"== channel: {c.get('chat_id')} name={c.get('name')!r} default={bool(c.get('is_default'))}"
)
return 0
if args.search:
payload = {"query": args.search, "scope": args.scope, "limit": 20}
if args.scope == "chat":
payload["chat_id"] = args.chat_id
d = _http_frame_roundtrip(
host,
port,
headers,
{"v": 1, "id": 1, "type": "search", "payload": payload},
"search.results",
30,
)
if d is None:
print("!! search: no search.results")
return 14
hits = (d.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')} snippet={str(h.get('snippet'))[:100]!r}"
)
return 0 if hits else 14
if args.sync is not None:
d = _http_frame_roundtrip(
host,
port,
headers,
{"v": 1, "id": 1, "type": "sync", "payload": {"cursor": args.sync}},
"sync.done",
30,
)
if d is None:
print("!! sync: no sync.done")
return 18
print(f"== sync done: cursor={(d.get('payload') or {}).get('cursor')}")
return 0
if args.fcm_reg:
d = _http_frame_roundtrip(
host,
port,
headers,
{"v": 1, "id": 1, "type": "fcm.register", "payload": {"token": args.fcm_token}},
"fcm.registered",
30,
)
if d is None:
print("!! fcm.register: no fcm.registered")
return 19
print("== fcm registered")
return 0
if args.watch:
print(f"== watching {args.watch} for a message (timeout {args.timeout:.0f}s)")
sse = HTTPConnection(host, port, timeout=args.timeout)
sse.request("GET", "/v1/events", headers=headers)
resp = sse.getresponse()
if resp.status != 200:
print(f"!! SSE open failed: HTTP {resp.status}")
sse.close()
return 17
deadline = time.time() + args.timeout
cur: list[str] = []
sock = getattr(getattr(resp.fp, "raw", None), "_sock", None)
try:
while time.time() < deadline:
if sock is not None:
sock.settimeout(max(0.1, deadline - time.time()))
line = resp.fp.readline()
if not line:
break
line = line.decode("utf-8").rstrip("\r\n")
if line == "":
if cur:
data = _print_frame("\n".join(cur))
cur = []
if data is not None and data.get("chat_id") == args.watch:
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}"
)
return 0
elif not line.startswith(":"):
field, _, value = line.partition(":")
if value.startswith(" "):
value = value[1:]
if field == "data":
cur.append(value)
finally:
sse.close()
print(f"!! no message landed in {args.watch}")
return 17
# 1b. optional media upload round-trip (docs/19 §19.15, v2).
if args.http_media:
import mimetypes
with open(args.http_media, "rb") as f:
data = f.read()
mime, _ = mimetypes.guess_type(args.http_media)
kind = "image" if (mime or "").startswith("image/") else "document"
conn = HTTPConnection(host, port, timeout=60)
conn.request(
"POST",
"/v1/media",
body=data,
headers={
**headers,
"Content-Type": mime or "application/octet-stream",
"X-Iris-Media-Ref": f"probe_{uuid.uuid4().hex[:12]}",
"X-Iris-Media-Kind": kind,
"X-Iris-Media-Filename": os.path.basename(args.http_media),
"X-Iris-Media-Sha256": hashlib.sha256(data).hexdigest(),
},
)
r = conn.getresponse()
body = r.read()
conn.close()
print(f"== POST /v1/media ({len(data)} bytes) -> {r.status} {body[:200]!r}")
if r.status != 201:
print("!! media upload rejected")
return 23
# 2. open the SSE stream.
sse = HTTPConnection(host, port, timeout=args.timeout)
sse.request("GET", "/v1/events", headers=headers)
resp = sse.getresponse()
if resp.status != 200:
print(f"!! SSE open failed: HTTP {resp.status}")
return 21
print("== SSE open (/v1/events)")
# 3. POST the message.send frame (accept-and-ack).
post_time: float | None = None
if args.send:
frame = {
"v": 1,
"id": 1,
"type": "message.send",
"chat_id": "default",
"payload": {"text": args.send},
}
conn = HTTPConnection(host, port, timeout=30)
conn.request(
"POST",
"/v1/frame",
body=json.dumps(frame),
headers={**headers, "Content-Type": "application/json"},
)
r = conn.getresponse()
body = r.read()
conn.close()
post_time = time.time()
print(f"== POST /v1/frame -> {r.status} {body[:200]!r}")
if r.status >= 400:
print("!! POST /v1/frame rejected")
return 22
# 4. read SSE until the final assistant message (same final-detection
# logic as the WS leg).
st = _TurnState()
got_final = False
seen_final_frame = False
echo_logged = False
deadline = time.time() + args.timeout
cur_data: list[str] = []
def feed(line: str) -> bool:
nonlocal cur_data, got_final, seen_final_frame, echo_logged
line = line.rstrip("\r\n")
if line == "":
if cur_data:
data = _print_frame("\n".join(cur_data))
if data is not None:
ftype = data.get("type")
payload = data.get("payload") or {}
st.track(ftype, payload)
# docs/19: the user echo must land on the SSE stream
# promptly after the POST (the < 1 s sendable-in-fallback
# UX assertion).
if (
not echo_logged
and post_time is not None
and ftype == "message"
and payload.get("role") == "user"
):
echo_logged = True
print(f"== user echo in {time.time() - post_time:.2f}s")
if ftype == "message" and payload.get("role") == "assistant":
got_final = True
if ftype == "message.stop":
seen_final_frame = True
if ftype == "typing" and payload.get("on") is False and seen_final_frame:
got_final = True
cur_data = []
return got_final
if line.startswith(":"):
return got_final # heartbeat comment
field, _, value = line.partition(":")
if value.startswith(" "):
value = value[1:]
if field == "data":
cur_data.append(value)
return got_final
sock = getattr(getattr(resp.fp, "raw", None), "_sock", None)
try:
while time.time() < deadline and not got_final:
if sock is not None:
sock.settimeout(max(0.1, deadline - time.time()))
line = resp.fp.readline()
if not line:
break
if feed(line.decode("utf-8")):
break
finally:
sse.close()
if not got_final:
print(f"!! no final assistant message (HTTP leg, {args.timeout:.0f}s)")
return 7
print("== final assistant message received (via SSE)")
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("IRIS_HTTP_URL", "http://127.0.0.1:8791"))
p.add_argument("--token", default=os.getenv("IRIS_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="default",
help="chat_id for --scope chat (default 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)",
)
p.add_argument(
"--http",
action="store_true",
help="docs/19: drive the turn over the HTTP fallback leg "
"(health + POST /v1/frame + SSE /v1/events) instead of WS",
)
p.add_argument(
"--http-url",
default="",
help="docs/19: http(s)://host:port base for --http "
"(default: derived from --url, port 8791)",
)
p.add_argument(
"--http-media",
default="",
help="docs/19 v2: with --http, also upload this file via POST /v1/media "
"and assert a 201 ack",
)
args = p.parse_args()
if not args.token and not args.authfail:
p.error("--token (or $IRIS_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)")
# HTTP is the only transport (docs/19): derive the http(s) base from the
# --url (legacy ws(s)://host:8790/ws -> http(s)://host:8791) unless
# --http-url is given.
if args.http_url:
base = args.http_url
else:
from urllib.parse import urlparse
u = urlparse(args.url)
scheme = "https" if u.scheme in ("wss", "https") else "http"
port = u.port or 8791
base = f"{scheme}://{u.hostname or '127.0.0.1'}:{port}"
return run_http(args, base)
if __name__ == "__main__":
sys.exit(main())