feat(relay): outbound buffering — queue agent→phone messages for an offline phone

Closes the gap the device test exposed: an agent answer pushed while the phone
had dropped its (connection-scoped) subscription used to 503 and be lost.

- ProactiveChannel.push() now queues on no-subscriber (bounded deque, drop-oldest,
  24h staleness TTL) and returns {delivered:false, queued:true, buffered:N} instead
  of raising 503. _flush_outbound delivers the backlog FIFO on the next
  proactive.subscribe; stale entries are pruned, and a socket that dies mid-flush
  re-buffers the remainder.
- Observe/cancel: peek_outbound()/cancel_outbound() + loopback routes
  GET /phone/outbound (count + summaries) and DELETE /phone/outbound[?message_id=]
  (cancel all / one) — the enabling layer for a host-side "queued + cancel" UI.
- Adapter send() unchanged (a 200 queued reads as success). Handler's no-subscriber
  503 path removed (now queues); ProactiveError is left only for a live write fail.
- Tests: 54 proactive+phone pass (flush-on-subscribe FIFO, bounded drop-oldest,
  stale-drop, cancel one/all, close clears).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Bailey Dixon
2026-06-29 16:30:01 -04:00
co-authored by Claude Opus 4.8
parent b5d30287a9
commit ebb4f041dd
5 changed files with 278 additions and 38 deletions
+2
View File
@@ -12,6 +12,8 @@
**Verification.** `python -m unittest plugin.tests.test_phone_platform plugin.tests.test_proactive_channel` — 49 pass (incl. the new connect-contract guard). Two-way reply round-trip confirmed end-to-end on-device: agent → phone notification → inline reply → drained through `/phone/replies` → `handle_message` → agent answer back in the *same* thread. One observed gap: when the phone has dropped its (connection-scoped) proactive subscription, the agent's answer `503`s and is lost — tracked as **outbound buffering** in TODO.
**Follow-on — outbound buffering (`plugin/relay/channels/proactive.py`, `server.py`).** Closes that gap. When no phone is subscribed, `ProactiveChannel.push()` now *queues* the message in a bounded deque (drop-oldest, 24 h staleness TTL) and returns `{delivered: false, queued: true, …}` instead of raising 503; `_flush_outbound` delivers the backlog FIFO on the next `proactive.subscribe` (stale entries pruned, socket-died-mid-flush re-buffers the remainder). Queued messages are inspectable + cancelable before they flush via `peek_outbound`/`cancel_outbound` and new loopback routes `GET /phone/outbound` (count + summaries) and `DELETE /phone/outbound[?message_id=…]` (cancel all / one). The adapter's `send()` is unchanged — a 200 queued reads as success. 54 proactive + phone tests pass (new: flush-on-subscribe FIFO, bounded drop-oldest, stale-drop, cancel one/all, close clears). UI surfacing of the queued state (host-side `relay` view + per-message status in the threaded surface) is specced in TODO.
## 2026-06-28 — Phone platform (Phase 2c: two-way reply — the inbound leg)
**Why.** Proactive messaging was push-only: the agent could message the phone, but the user couldn't answer. The phone was registered as a Hermes *platform* but only the outbound half (`send()`) was wired; its inbound path was a no-op, so a reply never reached the agent. Phase 2c wires the inbound leg so a reply becomes an inbound platform message the agent processes on the `phone` channel and answers over the existing `send()` — closing the loop into a conversation.
+1 -1
View File
@@ -26,7 +26,7 @@ Phase 1 (end-to-end spine) shipped on `Codename-11/phone-platform` — `send_mes
**Decision (2026-06-29): "separate lanes, unified surface."** The phone/agent conversation stays its own **gateway-platform lane** — distinct from the Standard Chat tab, which must keep working on vanilla upstream Hermes with no plugin — but is surfaced as a **first-class chat-style thread** that reuses the chat UI and sits alongside Chat. NOT a Chat "transport": a transport is an interchangeable pipe for the *same* user-chat conversation; the phone platform is a *different* conversation (agent-initiated, own session store/attribution, relay auth), so treating it as a transport miscategorizes it and couples a standard surface to a relay-only capability.
- **Outbound buffering (highest value — do first).** Mirror the inbound reply buffer for the *outbound* direction. Today when no phone is subscribed, `ProactiveChannel.push()` raises → HTTP 503 and the agent's answer is **lost** (observed live: 6× 503 after the app backgrounded). Buffer outbound messages in a bounded deque (drop-oldest, optional staleness TTL) and flush on `proactive.subscribe`; return a "queued" result instead of 503 and have the adapter `send()` treat queued as success.
- **Outbound buffering — ✅ relay-side DONE (2026-06-29).** `ProactiveChannel.push()` now queues agent→phone messages in a bounded deque (drop-oldest, 24 h TTL) when no phone is subscribed and returns `{queued: true}` (not 503); `_flush_outbound` delivers FIFO on the next subscribe (stale pruned). Inspect/cancel via `peek_outbound`/`cancel_outbound` + loopback `GET`/`DELETE /phone/outbound`. **Remaining — UI surfacing of the queued state** (the queued state exists while the phone is OFFLINE, so it's naturally a *host-side* surface, not the phone): (a) wire a desktop CLI `relay queue` / `relay queue --clear` over the new endpoints (and/or a dashboard Relay-tab view) to show "N queued for offline phone" + cancel; (b) in the threaded agent surface, mark messages that arrived-while-away, and show the user's OWN pending replies (the Phase 3 reply queue) with a sending/Cancel affordance — that's where phone-side "queued + cancel" belongs.
- **Threaded "agent" surface.** Render the `phone` platform session as a real chat thread in the app (history + reply from a chat composer — the `proactive.reply` path already exists) instead of notification-by-notification + discrete inbox cards. Needs a "read phone session history" path (relay-exposed, or read the gateway session store) + an app surface that reuses chat rendering but is gated behind relay+pairing and visually distinct from the Standard Chat tab. Verify first how hermes-desktop renders a *plugin-registered* (non-bundled) platform's sessions for attribution parity.
- **Per-thread `chat_id`.** Everything is hardcoded `chat_id="phone"` (one thread) today; the adapter already plumbs `chat_id`, so varying it yields multiple threads (per topic, or the agent opening distinct conversations). Ties into the threaded surface.
- **Message status + delivery state.** Surface sent / delivered / queued / failed per message in the thread (depends on outbound buffering's queued state) so the user knows whether the agent actually reached them.
+145 -24
View File
@@ -15,7 +15,12 @@ Flow (outbound — agent → phone):
phone platform adapter POSTs loopback to the relay's ``/phone/message``
route → the HTTP handler calls :meth:`push`.
4. :meth:`push` sends a ``phone.message`` envelope over the latched
phone WebSocket. The app raises a notification / inbox entry.
phone WebSocket. The app raises a notification / inbox entry. If no phone
is subscribed, :meth:`push` instead *buffers* the message (bounded,
drop-oldest, stale-pruned) and :meth:`_flush_outbound` delivers it on the
next ``proactive.subscribe`` — so an answer to a backgrounded phone is
queued, not lost. A queued message can be cancelled (``cancel_outbound``)
or inspected (``peek_outbound``) before it flushes.
Flow (inbound — phone → agent, Phase 2c):
5. The user replies (notification inline-reply or inbox reply box). The app
@@ -77,6 +82,13 @@ logger = logging.getLogger(__name__)
# is only useful fresh, so overflow drops the OLDEST (deque maxlen semantics).
MAX_BUFFERED_REPLIES = 100
# Bound the OUTBOUND buffer — agent→phone messages queued while no phone is
# subscribed (backgrounded / offline). Overflow drops the OLDEST (deque maxlen).
# A queued message older than the TTL is dropped on flush rather than delivered
# stale (a day-old "build is green" is noise, not signal).
MAX_BUFFERED_OUTBOUND = 50
OUTBOUND_TTL_SECONDS = 24 * 3600
class ProactiveError(Exception):
"""Raised when a proactive message cannot be delivered to a phone."""
@@ -119,6 +131,14 @@ class ProactiveChannel:
self._reply_event: asyncio.Event = asyncio.Event()
self.reply_count: int = 0
# Outbound buffer (agent → phone): messages pushed while no phone is
# subscribed are parked here and flushed FIFO on the next
# ``proactive.subscribe``, so the agent's answer survives a
# backgrounded/offline phone instead of being dropped with a 503.
# Bounded + drop-oldest; stale entries pruned on flush.
self._outbound: deque[dict[str, Any]] = deque(maxlen=MAX_BUFFERED_OUTBOUND)
self.queued_count: int = 0
# ── Envelope dispatch (inbound from phone) ───────────────────────────
async def handle(
@@ -155,6 +175,8 @@ class ProactiveChannel:
)
except Exception as exc: # pragma: no cover - best-effort ack
logger.debug("proactive: failed to ack subscribe: %s", exc)
# Deliver anything that queued while no phone was subscribed.
await self._flush_outbound(ws)
async def _handle_unsubscribe(self, ws: web.WebSocketResponse) -> None:
"""Release ``ws`` if it is the current subscriber."""
@@ -165,26 +187,9 @@ class ProactiveChannel:
# ── Outbound push (called from the HTTP handler) ─────────────────────
async def push(self, payload: dict[str, Any]) -> dict[str, Any]:
"""Send a ``phone.message`` envelope to the subscribed phone.
``payload`` is the body POSTed by the phone platform adapter
(``chat_id``, ``text``, ``title``, ``surfacing``, ``reply_to``,
``metadata``). Returns ``{delivered, message_id}``.
Raises :class:`ProactiveError` if no phone is subscribed (→ HTTP
503) or the socket write fails (→ HTTP 502).
"""
ws = self.phone_ws
if ws is None or ws.closed:
raise ProactiveError(
"No phone subscribed. Open the Hermes app and enable "
"'Let Hermes message me'."
)
message_id = str(payload.get("message_id") or uuid.uuid4().hex[:12])
sent_at = time.time()
out_payload = {
def _build_out_payload(self, payload: dict[str, Any], message_id: str) -> dict[str, Any]:
"""Build the ``phone.message`` payload that is sent / buffered."""
return {
"message_id": message_id,
"chat_id": payload.get("chat_id"),
"text": payload.get("text", ""),
@@ -192,16 +197,54 @@ class ProactiveChannel:
"surfacing": payload.get("surfacing"),
"reply_to": payload.get("reply_to"),
"metadata": payload.get("metadata") or {},
"sent_at": int(sent_at * 1000),
"sent_at": int(time.time() * 1000),
}
async def push(self, payload: dict[str, Any]) -> dict[str, Any]:
"""Send a ``phone.message`` to the subscribed phone, or queue it.
``payload`` is the body POSTed by the phone platform adapter
(``chat_id``, ``text``, ``title``, ``surfacing``, ``reply_to``,
``metadata``).
Returns ``{delivered: True, message_id}`` on a live send. When **no
phone is subscribed**, the message is parked in a bounded buffer and
``{delivered: False, queued: True, message_id, buffered}`` is returned
— it flushes on the next ``proactive.subscribe`` so the agent's answer
survives a backgrounded/offline phone (it used to be dropped with 503).
A queued message can be cancelled before it flushes via
:meth:`cancel_outbound`.
Raises :class:`ProactiveError` only when a *live* socket write fails
(→ HTTP 502); a missing subscriber is no longer an error.
"""
message_id = str(payload.get("message_id") or uuid.uuid4().hex[:12])
out_payload = self._build_out_payload(payload, message_id)
ws = self.phone_ws
if ws is None or ws.closed:
self._outbound.append(out_payload)
self.queued_count += 1
logger.info(
"proactive >>> queued (no subscriber) message_id=%s chat=%s (buffered=%d)",
message_id,
out_payload["chat_id"],
len(self._outbound),
)
return {
"delivered": False,
"queued": True,
"message_id": message_id,
"buffered": len(self._outbound),
}
try:
await ws.send_str(_envelope("phone.message", out_payload, message_id))
except Exception as exc:
logger.error("proactive: failed to push message: %s", exc)
raise ProactiveError(f"Failed to send message to phone: {exc}") from exc
self.last_push_at = sent_at
self.last_push_at = out_payload["sent_at"] / 1000.0
self.push_count += 1
logger.info(
"proactive >>> message_id=%s chat=%s len=%d",
@@ -211,6 +254,79 @@ class ProactiveChannel:
)
return {"delivered": True, "message_id": message_id}
async def _flush_outbound(self, ws: web.WebSocketResponse) -> None:
"""Deliver queued outbound messages to a freshly-subscribed phone.
Drains the outbound buffer FIFO. Entries older than
:data:`OUTBOUND_TTL_SECONDS` are dropped rather than delivered stale.
If the socket dies mid-flush, the remaining (incl. the one that failed)
are re-buffered for the next subscribe instead of lost.
"""
if not self._outbound:
return
pending = list(self._outbound)
self._outbound.clear()
now_ms = int(time.time() * 1000)
ttl_ms = OUTBOUND_TTL_SECONDS * 1000
flushed = 0
for i, out_payload in enumerate(pending):
sent_at = out_payload.get("sent_at") or 0
if ttl_ms and sent_at and (now_ms - sent_at) > ttl_ms:
continue # stale — drop rather than deliver late
try:
await ws.send_str(
_envelope("phone.message", out_payload, out_payload.get("message_id"))
)
flushed += 1
except Exception as exc:
logger.warning(
"proactive: outbound flush interrupted (%s) — re-buffering %d",
exc,
len(pending) - i,
)
for leftover in pending[i:]:
self._outbound.append(leftover)
break
if flushed:
self.last_push_at = time.time()
self.push_count += flushed
logger.info(
"proactive: flushed %d queued message(s) to phone on subscribe", flushed
)
def peek_outbound(self) -> list[dict[str, Any]]:
"""Return a UI-friendly summary of queued outbound messages (FIFO).
Text is truncated so a status view / `relay` CLI can list what's
waiting without dumping full bodies. Read-only — does not drain.
"""
return [
{
"message_id": m.get("message_id"),
"chat_id": m.get("chat_id"),
"text": (m.get("text") or "")[:200],
"sent_at": m.get("sent_at"),
}
for m in self._outbound
]
def cancel_outbound(self, message_id: str | None = None) -> int:
"""Cancel queued outbound messages before they flush.
``message_id=None`` clears the whole queue; otherwise removes just that
message. Returns the number cancelled. No effect on already-delivered
messages (the buffer only holds the not-yet-delivered).
"""
if message_id is None:
n = len(self._outbound)
self._outbound.clear()
return n
before = len(self._outbound)
kept = [m for m in self._outbound if m.get("message_id") != message_id]
self._outbound.clear()
self._outbound.extend(kept)
return before - len(self._outbound)
# ── Inbound reply buffer (phone → agent) ─────────────────────────────
def _handle_reply(self, payload: dict[str, Any]) -> None:
@@ -276,6 +392,10 @@ class ProactiveChannel:
"""Number of replies currently waiting for the gateway poller."""
return len(self._replies)
def buffered_outbound_count(self) -> int:
"""Number of agent→phone messages queued for the next subscribe."""
return len(self._outbound)
# ── Lifecycle ────────────────────────────────────────────────────────
def is_phone_subscribed(self) -> bool:
@@ -292,8 +412,9 @@ class ProactiveChannel:
)
async def close(self) -> None:
"""Server shutdown — drop the subscriber reference + reply buffer."""
"""Server shutdown — drop the subscriber reference + buffers."""
self.phone_ws = None
self.subscribed_at = None
self._replies.clear()
self._reply_event.clear()
self._outbound.clear()
+34 -5
View File
@@ -3193,11 +3193,12 @@ async def handle_phone_message(request: web.Request) -> web.Response:
exposed to the LAN.
POST /phone/message {chat_id, text, title?, surfacing?, reply_to?, metadata?}
→ 200 {"delivered": true, "message_id": "..."}
→ 200 {"delivered": true, "message_id": "..."} (sent to a live phone)
→ 200 {"delivered": false, "queued": true, "message_id": "...", "buffered": N}
(no phone subscribed — queued, flushed on the next subscribe)
→ 400 invalid body / empty text
→ 403 non-loopback caller
→ 502 socket write failed
→ 503 no phone subscribed
"""
remote = request.remote or ""
if remote not in ("127.0.0.1", "::1"):
@@ -3218,13 +3219,38 @@ async def handle_phone_message(request: web.Request) -> web.Response:
try:
result = await server.proactive.push(payload)
except ProactiveError as exc:
msg = str(exc)
status = 503 if "subscrib" in msg.lower() or "no phone" in msg.lower() else 502
return web.json_response({"error": msg}, status=status)
# push() now buffers when no phone is subscribed, so the only
# ProactiveError path left is a live socket write that failed → 502.
return web.json_response({"error": str(exc)}, status=502)
return web.json_response(result, status=200)
async def handle_phone_outbound(request: web.Request) -> web.Response:
"""Inspect or cancel the queued agent→phone outbound messages.
Loopback-only (same rationale as ``/phone/message``). Lets a status view /
`relay` CLI / dashboard show what's waiting for an offline phone and cancel
it before it flushes on the next subscribe.
GET /phone/outbound → {queued, messages:[{message_id, chat_id, text, sent_at}]}
DELETE /phone/outbound → cancel ALL → {cancelled}
DELETE /phone/outbound?message_id=<id> → cancel one → {cancelled}
"""
remote = request.remote or ""
if remote not in ("127.0.0.1", "::1"):
return web.json_response({"error": "loopback only"}, status=403)
server: RelayServer = request.app["server"]
if request.method == "DELETE":
mid = request.query.get("message_id") or None
cancelled = server.proactive.cancel_outbound(mid)
return web.json_response({"cancelled": cancelled}, status=200)
messages = server.proactive.peek_outbound()
return web.json_response({"queued": len(messages), "messages": messages}, status=200)
# Cap the server-side long-poll hold so a wedged gateway poller can't pin a
# request open forever; the adapter re-polls. Generous vs. the adapter's own
# read timeout so a normal empty poll returns from here, not from a client
@@ -4049,6 +4075,9 @@ def create_app(config: RelayConfig) -> web.Application:
# Inbound reply leg (Phase 2c) — the gateway adapter long-polls here to
# drain ``proactive.reply`` envelopes buffered by the ProactiveChannel.
app.router.add_get("/phone/replies", handle_phone_replies)
# Inspect / cancel the queued agent→phone outbound buffer (offline phone).
app.router.add_get("/phone/outbound", handle_phone_outbound)
app.router.add_delete("/phone/outbound", handle_phone_outbound)
# Relay-owned agent context audit.
app.router.add_get("/context/injected", handle_context_injected)
+96 -8
View File
@@ -5,7 +5,9 @@ Validated surface:
``proactive.subscribed``.
* :meth:`ProactiveChannel.push` sends a ``phone.message`` envelope over the
latched WS and returns ``{delivered, message_id}``.
* push with no subscriber raises :class:`ProactiveError`.
* push with no subscriber buffers the message; it flushes FIFO on the next
``proactive.subscribe`` (bounded, drop-oldest, stale-pruned) and can be
cancelled (``cancel_outbound``) before it flushes.
* push over a closed/failing socket raises :class:`ProactiveError`.
* ``proactive.unsubscribe`` / ``detach_ws`` release the subscriber.
* a phone must not originate ``phone.message`` (ignored).
@@ -139,23 +141,27 @@ class ProactiveChannelTests(unittest.TestCase):
_run(run())
def test_push_without_subscriber_raises(self) -> None:
def test_push_without_subscriber_queues(self) -> None:
async def run() -> None:
ch = ProactiveChannel()
with self.assertRaises(ProactiveError) as ctx:
await ch.push({"text": "hi"})
self.assertIn("subscrib", str(ctx.exception).lower())
result = await ch.push({"text": "hi", "chat_id": "phone"})
self.assertFalse(result["delivered"])
self.assertTrue(result["queued"])
self.assertEqual(result["buffered"], 1)
self.assertTrue(result["message_id"])
self.assertEqual(ch.buffered_outbound_count(), 1)
_run(run())
def test_push_over_closed_ws_raises(self) -> None:
def test_push_over_closed_ws_queues(self) -> None:
async def run() -> None:
ch = ProactiveChannel()
ws = _FakeWs()
await ch.handle(ws, {"type": "proactive.subscribe"})
ws.closed = True
with self.assertRaises(ProactiveError):
await ch.push({"text": "hi"})
result = await ch.push({"text": "hi"})
self.assertTrue(result["queued"])
self.assertEqual(ch.buffered_outbound_count(), 1)
_run(run())
@@ -171,6 +177,88 @@ class ProactiveChannelTests(unittest.TestCase):
_run(run())
# ── Outbound buffering (queue when no phone, flush on subscribe) ──────
def test_queued_outbound_flushes_on_subscribe(self) -> None:
async def run() -> None:
ch = ProactiveChannel()
# Pushed with no phone subscribed → queued, FIFO.
await ch.push({"text": "first", "chat_id": "phone"})
await ch.push({"text": "second", "chat_id": "phone"})
self.assertEqual(ch.buffered_outbound_count(), 2)
# Phone subscribes → ack, then both queued messages delivered.
ws = _FakeWs()
await ch.handle(ws, {"type": "proactive.subscribe"})
types = [m["type"] for m in ws.sent]
self.assertEqual(
types, ["proactive.subscribed", "phone.message", "phone.message"]
)
texts = [m["payload"]["text"] for m in ws.sent if m["type"] == "phone.message"]
self.assertEqual(texts, ["first", "second"])
self.assertEqual(ch.buffered_outbound_count(), 0)
_run(run())
def test_outbound_buffer_bounded_drops_oldest(self) -> None:
async def run() -> None:
from plugin.relay.channels import proactive as mod
ch = ProactiveChannel()
overflow = mod.MAX_BUFFERED_OUTBOUND + 5
for i in range(overflow):
await ch.push({"text": f"m{i}", "chat_id": "phone"})
self.assertEqual(ch.buffered_outbound_count(), mod.MAX_BUFFERED_OUTBOUND)
ws = _FakeWs()
await ch.handle(ws, {"type": "proactive.subscribe"})
texts = [m["payload"]["text"] for m in ws.sent if m["type"] == "phone.message"]
self.assertEqual(texts[0], "m5") # oldest 5 dropped
self.assertEqual(texts[-1], f"m{overflow - 1}")
_run(run())
def test_stale_outbound_dropped_on_flush(self) -> None:
async def run() -> None:
from plugin.relay.channels import proactive as mod
ch = ProactiveChannel()
await ch.push({"text": "fresh", "chat_id": "phone"})
stale = ch._build_out_payload({"text": "stale", "chat_id": "phone"}, "old")
stale["sent_at"] -= (mod.OUTBOUND_TTL_SECONDS + 60) * 1000
ch._outbound.appendleft(stale) # oldest
ws = _FakeWs()
await ch.handle(ws, {"type": "proactive.subscribe"})
texts = [m["payload"]["text"] for m in ws.sent if m["type"] == "phone.message"]
self.assertEqual(texts, ["fresh"]) # stale dropped, fresh delivered
_run(run())
def test_cancel_outbound_one_and_all(self) -> None:
async def run() -> None:
ch = ProactiveChannel()
r1 = await ch.push({"text": "a", "chat_id": "phone"})
await ch.push({"text": "b", "chat_id": "phone"})
self.assertEqual(ch.buffered_outbound_count(), 2)
# peek summary is read-only.
self.assertEqual([p["text"] for p in ch.peek_outbound()], ["a", "b"])
self.assertEqual(ch.buffered_outbound_count(), 2)
# cancel one by id, then cancel the rest.
self.assertEqual(ch.cancel_outbound(r1["message_id"]), 1)
self.assertEqual([p["text"] for p in ch.peek_outbound()], ["b"])
self.assertEqual(ch.cancel_outbound(), 1)
self.assertEqual(ch.buffered_outbound_count(), 0)
_run(run())
def test_close_clears_outbound_buffer(self) -> None:
async def run() -> None:
ch = ProactiveChannel()
await ch.push({"text": "hi", "chat_id": "phone"})
self.assertEqual(ch.buffered_outbound_count(), 1)
await ch.close()
self.assertEqual(ch.buffered_outbound_count(), 0)
_run(run())
# ── Hardening ────────────────────────────────────────────────────────
def test_phone_message_from_phone_ignored(self) -> None: