diff --git a/docs/PROTOCOL.md b/docs/PROTOCOL.md index 6680675..17ab02e 100644 --- a/docs/PROTOCOL.md +++ b/docs/PROTOCOL.md @@ -59,7 +59,7 @@ Every error uses one shape with a machine-readable code: Codes: `unauthorized`, `auth_timeout`, `version_mismatch`, `bad_request`, `invalid_json`, `unknown_type`, `not_found`, `conversation_not_found`, `approval_not_found`, `clarify_not_found`, `too_large`, `rate_limited`, -`internal_error`, `not_implemented`. +`run_active`, `internal_error`, `not_implemented`. Limits: chat text ≤ 64,000 chars; inbound WS frame ≤ 2 MiB (violations get `too_large`); conversation history fetch ≤ 500 messages. @@ -98,13 +98,20 @@ Limits: chat text ≤ 64,000 chars; inbound WS frame ≤ 2 MiB (violations get ← { "type": "conversation.history", "conversation_id": "a1b2…", "request_id": "r3", "messages": [ { "message_id": "m12", "role": "user", "text": "hey", "ts": "…|null" }, - { "message_id": "m13", "role": "assistant", "text": "hi!", "ts": "…|null" } ] } + { "message_id": "m13", "role": "assistant", "text": "hi!", "ts": "…|null" } ], + "attachments": [], + "run": null, + "tools": [], + "approvals": [], + "clarifications": [] } ``` History is the authoritative Hermes transcript (`role` is always `user` or -`assistant`). On reconnect, re-open the last-open conversations and resume — -no client-side message cache is needed for correctness. Unknown conversation -→ `conversation_not_found` error. +`assistant`). The other fields form a recoverable snapshot: unexpired +attachments, the active run (if any), latest structured tool states, and +pending approval/clarification requests. On reconnect, replace local state +with this snapshot, then consume new live events. Unknown conversation → +`conversation_not_found` error. ### `conversation.rename` @@ -168,14 +175,15 @@ Run end: ```json ← { "type": "run.finished", "conversation_id": "a1b2…", "run_id": "8c1f…", - "status": "completed" | "cancelled" | "failed" | "idle", + "status": "completed" | "cancelled" | "failed", "error": "only on failure" } ``` State machine per assistant turn: `run.accepted → message.start → (message.delta)* → message.complete → run.finished`. -A turn with no streaming skips `message.start`/`message.delta`. Never infer -state from text — use these events. +A non-streaming turn may skip `message.delta`; `message.start` is always sent +when the run is accepted. Only one run may be active per conversation; +another `chat.send` receives `run_active`. Never infer state from text. ### `run.cancel` — stop an active run @@ -186,8 +194,8 @@ state from text — use these events. Cancellation uses Hermes's supported interrupt mechanism (agent interrupt + run-generation invalidation) — the conversation stays consistent and -resumable. Cancelling with no active run returns `run.finished` -`status:"idle"`. +resumable. A stale `run_id`, missing run, or run that Hermes can no longer +interrupt returns `not_found`; success emits exactly one cancelled event. --- @@ -200,21 +208,20 @@ Tool activity arrives as `tool.event` messages, completely separate from ```json { "type": "tool.event", "conversation_id": "a1b2…", + "run_id": "8c1f…", "tool_call_id": "t-1a2b3c4d5e6f", "tool_name": "web_search", - "status": "running" | "completed" | "failed", - "description": "cats — short preview from the agent (may be null)", + "status": "running" | "completed" | "failed" | "cancelled", "args_redacted": { "query": "cats" }, // only on "running"; secret-looking keys redacted "duration_ms": 1234, // only on completion/failure (may be null) "error": "only on failed, truncated", // only on "failed" "ts": "..." } ``` -Correlate `running` → `completed`/`failed` by `tool_call_id`. Note: the -running event's ID comes from the adapter and the completion event from -Hermes's `post_tool_call` hook; when they differ, correlate by -`(tool_name, conversation)` as a fallback and prefer the completion event's -ID going forward. +Correlate `running` → `completed`/`failed` by `tool_call_id`. Both events use +Hermes's authoritative call ID from the pre/post tool hooks and include the +active `run_id` when one exists. Secret-looking argument keys are redacted +recursively before leaving the server. No fake "Searching the web…" text is ever injected into `message.*` events. @@ -228,6 +235,8 @@ When Hermes pauses for a human decision on a dangerous action: { "type": "approval.request", "approval_id": "3d4e5f6070a1", "session_key": "agent:main:pheby:dm:a1b2…", + "conversation_id": "a1b2…", + "run_id": "8c1f…", "command": "rm -rf /tmp/build-output", "description": "Destructive shell command (rm -rf)", "choices": ["once", "session", "always", "deny"], @@ -240,11 +249,13 @@ Respond: → { "type": "approval.respond", "approval_id": "3d4e5f6070a1", "choice": "once" | "session" | "always" | "deny", "reason": "optional free text with deny", "request_id": "r8" } -← { "type": "approval.resolved", "approval_id": "3d4e5f6070a1", "choice": "once", "request_id": "r8" } -// broadcast confirmation (also informs other tabs): -← { "type": "approval.resolved", "approval_id": "3d4e5f6070a1", "choice": "once", "accepted": true } +← { "type": "approval.resolved", "approval_id": "3d4e5f6070a1", + "choice": "once", "accepted": true, "request_id": "r8" } ``` +The single resolution event is broadcast to every connected client; the +responding client's `request_id` is included on that broadcast. + Choices map to Hermes semantics: `once` (approve this action), `session` (approve pattern for this conversation), `always` (also persist), `deny` (decline; the agent is told NOT to retry). Unknown/stale ID → @@ -259,6 +270,8 @@ own timeout, so a silently-closed socket can't leave a zombie gate. { "type": "clarify.request", "clarify_id": "c1a2b3d4e5", "session_key": "agent:main:pheby:dm:a1b2…", + "conversation_id": "a1b2…", + "run_id": "8c1f…", "question": "Deploy to staging or production?", "choices": ["staging", "production"], // null ⇒ free text only "allow_free_text": true, @@ -269,13 +282,13 @@ Respond (either a choice value or free text): ```json → { "type": "clarify.respond", "clarify_id": "c1a2b3d4e5", "response": "production", "request_id": "r9" } -← { "type": "clarify.resolved", "clarify_id": "c1a2b3d4e5", "request_id": "r9" } -← { "type": "clarify.resolved", "clarify_id": "c1a2b3d4e5", "accepted": true } // broadcast +← { "type": "clarify.resolved", "clarify_id": "c1a2b3d4e5", + "accepted": true, "request_id": "r9" } ``` -`accepted:false` on the broadcast means Hermes had already resolved/timed out -the prompt. Always render an "Other" affordance — Hermes clarifications -accept free text. +The single resolution event is broadcast. A stale/timed-out ID instead gets +`clarify_not_found`. Always render an "Other" affordance — Hermes +clarifications accept free text. --- @@ -330,7 +343,7 @@ Authorization: Bearer (or ApiKey , or X-Pheby-Secre ### List providers + models ```json -→ { "type": "models.list", "request_id": "r10" } +→ { "type": "models.list", "conversation_id": "a1b2…", "request_id": "r10" } ← { "type": "models.snapshot", "request_id": "r10", "providers": [ { "slug": "openrouter", "name": "OpenRouter", "is_current": true, @@ -343,13 +356,17 @@ Authorization: Bearer (or ApiKey , or X-Pheby-Secre ``` Lists come from Hermes's own credential-aware picker data — nothing is -hardcoded. Models are exactly what the configured providers expose. +hardcoded. Models are exactly what the configured providers expose. The +optional `conversation_id` makes `current_model`, `current_provider`, and +`scope` reflect that conversation's override. ### Read / change current model ```json -→ { "type": "models.current", "request_id": "r11" } -← { "type": "model.current", "request_id": "r11", "model": "z-ai/glm-5.3-flash", "provider": "openrouter", "ts": "..." } +→ { "type": "models.current", "conversation_id": "a1b2…", "request_id": "r11" } +← { "type": "model.current", "request_id": "r11", "model": "z-ai/glm-5.3-flash", + "provider": "openrouter", "scope": "conversation", + "conversation_id": "a1b2…", "ts": "..." } → { "type": "model.set", "model": "anthropic/claude-sonnet-4", "provider": "anthropic", // optional @@ -369,10 +386,11 @@ conversation the default moved. ## Reasoning effort ```json -→ { "type": "reasoning.current", "request_id": "r13" } +→ { "type": "reasoning.current", "conversation_id": "a1b2…", "request_id": "r13" } ← { "type": "reasoning.snapshot", "request_id": "r13", "effort": "medium", // current effective effort (may be null = provider default) "enabled": true, // false ⇒ thinking disabled + "scope": "conversation", "conversation_id": "a1b2…", "supported_efforts": ["none","minimal","low","medium","high","xhigh","max","ultra"], "ts": "..." } @@ -399,8 +417,9 @@ notifications) is pushed as normal `message.*` / `run.*` events even when it is not a reply to your last request. Reconnection procedure for clients: 1. Reconnect WS, redo `hello`. -2. Re-`conversation.open` the conversations you show; replace local state - with `conversation.history` (authoritative). +2. Re-`conversation.open` the conversations you show; replace local history, + attachments, run/tool state, and pending decisions with its authoritative + snapshot. 3. Re-`models.current` / `reasoning.current` if those views are visible. 4. Live events continue from there. diff --git a/docs/README.md b/docs/README.md index 6fd8737..5955be9 100644 --- a/docs/README.md +++ b/docs/README.md @@ -29,7 +29,7 @@ Pheby plugin (in Hermes gateway process) ├─ aiohttp server: GET /ws, GET /attachments/{id}, GET /health ├─ protocol layer: JSON message types, auth, limits, errors ├─ PhebyAdapter (BasePlatformAdapter subclass) - │ outbound: send/edit → chat events, format_tool_event → tool.event, + │ outbound: send/edit → chat events, tool hooks → tool.event, │ send_clarify → clarify.request, send_document/… → attachments │ inbound: WS messages → MessageEvent → gateway pipeline └─ hermes_bridge: sessions, approvals (tools.approval), @@ -44,9 +44,12 @@ Key properties: - **Hermes is authoritative.** Conversations are Hermes sessions keyed `agent:main:pheby:dm:`; history/titles live in Hermes `state.db`. Pheby keeps only a thin, rebuildable name index. -- **The full gateway pipeline works unchanged** — auth/pairing, tool +- **The full gateway pipeline works unchanged** — tool approval, clarify, deliverables, streaming, cron delivery — because inbound - messages are ordinary `MessageEvent`s on a registered platform. + messages are ordinary `MessageEvent`s on a registered platform. The adapter + marks its authenticated WebSocket as the upstream authorization boundary, + so a valid Pheby secret is not rejected by Hermes's separate messaging- + platform allowlist layer. - **Single-user, shared-secret.** One `PHEBY_SECRET` gates every WS connection and attachment download (constant-time compare). No accounts, no registration — but IDs and message shapes are multi-client friendly. @@ -99,8 +102,7 @@ platforms: Env vars override YAML: `PHEBY_SECRET`, `PHEBY_BIND_HOST`, `PHEBY_PORT`, `PHEBY_DEBUG`, `PHEBY_LOG_CHAT_CONTENT`. Optional: `PHEBY_HOME_CHANNEL` -(conversation ID receiving cron deliveries), `PHEBY_ALLOWED_USERS`, -`PHEBY_ALLOW_ALL_USERS`. +(conversation ID receiving cron deliveries). ### 4. Start Hermes with the gateway @@ -134,7 +136,7 @@ pheby.example.com { # WebSocket + API reverse_proxy /ws 127.0.0.1:8620 - reverse_proxy /attachments 127.0.0.1:8620 + reverse_proxy /attachments/* 127.0.0.1:8620 reverse_proxy /health 127.0.0.1:8620 # Optionally restrict by source when on a public VPS: @@ -153,10 +155,10 @@ pheby.example.com { Notes: - Caddy provides HTTPS + automatic certificates; the plugin never sees TLS. -- `X-Forwarded-For` is not used for auth decisions (the lockout key is the - direct peer address — behind Caddy that is Caddy itself, so lockout is - effectively global; that is acceptable for a single-user deployment and - still stops brute force). +- Authentication always depends on the shared secret. For lockout accounting, + `X-Forwarded-For` is accepted only when the direct peer is loopback (the + documented local-Caddy setup). A directly exposed non-loopback client cannot + spoof that header to evade rate limiting. - The Android client connects to `wss://pheby.example.com/ws` and downloads attachments from `https://pheby.example.com/attachments/{id}` with `Authorization: Bearer `. @@ -243,23 +245,21 @@ around. None require core modifications; all are handled cleanly. valid effort levels. Pheby exposes Hermes's canonical level set (`none, minimal, low, medium, high, xhigh, max, ultra`) and documents that an unsupported level surfaces as a provider error on the next turn. -3. **Conversation delete is a documented approximation.** Hermes's - SessionStore has no public per-routing-key delete; Pheby deletes the - authoritative transcript row (`SessionDB.delete_session`) and resets the - routing entry, which yields the same user-visible behavior. -4. **Tool `running` → `completed` ID correlation.** Start events get - adapter-generated IDs; completion events carry Hermes's authoritative - `tool_call_id` from the `post_tool_call` hook. The protocol documents - correlating by `(tool_name, conversation)` when IDs differ. (Gateway - tool-start events don't carry Hermes's call ID yet.) -5. **Standalone cron delivery.** Cron jobs targeting `pheby` are delivered +3. **Conversation routing deletion uses a private Hermes detail.** Hermes's + SessionStore has no public per-routing-key delete. Pheby deletes the + authoritative transcript through `SessionDB.delete_session`, then removes + the exact routing entry under SessionStore's own lock/save discipline. It + does not call `reset_session` (which would recreate the deleted chat). This + is isolated in `hermes_bridge.py` but may need adjustment after a Hermes + SessionStore refactor. +4. **Standalone cron delivery.** Cron jobs targeting `pheby` are delivered in-process with the gateway. A `standalone_sender_fn` hook exists but cannot push to a WS server it isn't hosting; out-of-process cron delivery to Pheby is not supported (documented, fail-loud). -6. **Reconnect recovery is state-based, not event-replay.** Missed events are - recovered by re-opening conversations (authoritative history), not by - replaying a server-side event log. This is the spec's preferred approach - and keeps the protocol simple. +5. **Reconnect recovery is state-based, not event-replay.** Re-opening a + conversation returns authoritative history plus unexpired attachments and + the current run, tool, approval, and clarification snapshot. The server + does not retain a replay log of every transient delta. --- @@ -270,9 +270,9 @@ python3 -m venv .venv && .venv/bin/pip install pytest pytest-asyncio aiohttp pyy .venv/bin/python -m pytest tests/ -o addopts= -q --asyncio-mode=auto ``` -Tests are gateway-free (fakes; real Hermes primitives exercised in-process -where safe — no LLM calls). Two live smoke tests bind an ephemeral localhost -port. +The 46 tests are gateway-free (fakes; real Hermes primitives exercised +in-process where safe — no LLM calls). Two live smoke tests bind an ephemeral +localhost port. ``` repo layout diff --git a/plugin/pheby/__init__.py b/plugin/pheby/__init__.py index cac8bdd..da95d2d 100644 --- a/plugin/pheby/__init__.py +++ b/plugin/pheby/__init__.py @@ -9,12 +9,11 @@ set ``PHEBY_SECRET`` in ``~/.hermes/.env``, and restart the gateway. from __future__ import annotations import logging -import os from typing import Any logger = logging.getLogger(__name__) -__version__ = "1.0.0" +__version__ = "1.0.1" def register(ctx: Any) -> None: @@ -47,8 +46,6 @@ def register(ctx: Any) -> None: env_enablement_fn=env_enablement, # Home channel for cron / notification delivery when configured. cron_deliver_env_var="PHEBY_HOME_CHANNEL", - allowed_users_env="PHEBY_ALLOWED_USERS", - allow_all_env="PHEBY_ALLOW_ALL_USERS", emoji="🐱", pii_safe=True, # single-user private platform; no PII in routing IDs allow_update_command=True, @@ -61,15 +58,21 @@ def register(ctx: Any) -> None: ), ) - # post_tool_call observer → structured tool-result events. Registered - # against the plugin context so it loads with the plugin, before any - # adapter is constructed (the hook is a no-op until the adapter serves). + # Tool observers provide authoritative session_id + tool_call_id values. + # The adapter's presentation-only format_tool_event hook deliberately eats + # Hermes's generic tool chrome to avoid duplicate/ambiguous events. + def _pre_tool_call(**kwargs: Any) -> None: + adapter = adapter_holder.get("adapter") + if adapter is not None: + adapter.on_pre_tool_call(**kwargs) + def _post_tool_call(**kwargs: Any) -> None: adapter = adapter_holder.get("adapter") if adapter is not None: adapter.on_post_tool_call(**kwargs) try: + ctx.register_hook("pre_tool_call", _pre_tool_call) ctx.register_hook("post_tool_call", _post_tool_call) except Exception: logger.debug("[pheby] post_tool_call hook registration failed", diff --git a/plugin/pheby/adapter.py b/plugin/pheby/adapter.py index ced1c20..ffb9aed 100644 --- a/plugin/pheby/adapter.py +++ b/plugin/pheby/adapter.py @@ -6,7 +6,7 @@ clarify, deliverables, streaming — works unchanged on the Pheby platform. Outbound mapping: * ``send`` / ``edit_message`` → chat draft events (S_MESSAGE_DELTA etc.) -* ``format_tool_event`` → structured S_TOOL_EVENT JSON (never fake text) +* pre/post tool hooks → structured S_TOOL_EVENT JSON (never fake text) * ``send_clarify`` → structured S_CLARIFY_REQUEST * ``send_document`` etc. → attachment registration + S_ATTACHMENT_ADDED @@ -18,25 +18,15 @@ to any other platform's messages. from __future__ import annotations import asyncio +import importlib.util import logging -import mimetypes -import os -import time -from pathlib import Path -from typing import Any, Dict, List, Optional, Tuple +from typing import Any, Dict, Optional -try: - import aiohttp - from aiohttp import web - AIOHTTP_AVAILABLE = True -except ImportError: # pragma: no cover - AIOHTTP_AVAILABLE = False +AIOHTTP_AVAILABLE = importlib.util.find_spec("aiohttp") is not None from gateway.config import Platform, PlatformConfig from gateway.platforms.base import ( BasePlatformAdapter, - MessageEvent, - MessageType, SendResult, ) @@ -46,9 +36,6 @@ from .config import PhebyConfig, load_config logger = logging.getLogger(__name__) -_MEDIA_TAG_RE = None # populated lazily from base module helpers - - class PhebyAdapter(BasePlatformAdapter): """Serve the Pheby WebSocket/HTTP protocol and map it onto the gateway.""" @@ -74,8 +61,13 @@ class PhebyAdapter(BasePlatformAdapter): super().__init__(config=config, platform=platform) self._pcfg: PhebyConfig = load_config(config.extra or {}) self._server: Any = None + self._loop: Optional[asyncio.AbstractEventLoop] = None self._drafts: Dict[str, Dict[str, Any]] = {} # conv → draft state - self._typing: Dict[str, float] = {} + + @property + def authorization_is_upstream(self) -> bool: + """The authenticated WebSocket transport already authorized input.""" + return True # ── connection lifecycle ───────────────────────────────────────────── async def connect(self, *, is_reconnect: bool = False) -> bool: @@ -90,11 +82,15 @@ class PhebyAdapter(BasePlatformAdapter): retryable=False) return False from .server import PhebyServer + self._loop = asyncio.get_running_loop() hermes_bridge.set_adapter(self) self._server = PhebyServer(self._pcfg, adapter=self) hermes_bridge.set_server(self._server) ok = await self._server.start() if not ok: + self._server = None + hermes_bridge.clear_services(self) + self._loop = None return False self._mark_connected() logger.info("[pheby] adapter connected (protocol v%d)", @@ -106,6 +102,8 @@ class PhebyAdapter(BasePlatformAdapter): if self._server is not None: await self._server.stop() self._server = None + hermes_bridge.clear_services(self) + self._loop = None self._mark_disconnected() logger.info("[pheby] adapter disconnected") @@ -123,29 +121,61 @@ class PhebyAdapter(BasePlatformAdapter): ) -> SendResult: """Deliver assistant text (final response, commentary, or notices). - The stream consumer calls ``send`` for the first streamed chunk and - the gateway calls it for the final response; both land as - ``S_MESSAGE_COMPLETE``. Streamed deltas ride ``edit_message``. + A streaming preview (``expect_edits``) remains a draft; only the final + send or a finalizing edit emits ``S_MESSAGE_COMPLETE`` and closes the + run. Streamed cumulative updates ride ``edit_message``. """ if self._server is None: return SendResult(success=False, error="server not running") conversation_id = self._conv_from_chat_id(chat_id) - message_id = f"m-{proto.new_id()[:12]}" + metadata = metadata or {} - # A draft exists while the turn streams; the final text supersedes - # the draft and closes it out. - draft = self._drafts.pop(conversation_id, None) + # Commentary and gateway notices are complete standalone timeline + # items. They must not close the primary assistant draft/run. + if metadata.get("_interim_send") or metadata.get("non_conversational"): + message_id = f"m-{proto.new_id()[:12]}" + event = { + "type": proto.S_MESSAGE_COMPLETE, + "conversation_id": conversation_id, + "message_id": message_id, + "text": content, + "kind": ("notice" if metadata.get("non_conversational") + else "commentary"), + "ts": proto.now_iso(), + } + await self._server.broadcast(event) + return SendResult(success=True, message_id=message_id) + + draft = self._drafts.setdefault(conversation_id, { + "message_id": f"draft-{proto.new_id()[:12]}", "text": ""}) + draft["text"] = content + + # Hermes uses send(expect_edits=True) for the first visible streaming + # preview. It is a cumulative draft update, not a completed message. + if metadata.get("expect_edits"): + await self._server.broadcast({ + "type": proto.S_MESSAGE_DELTA, + "conversation_id": conversation_id, + "run_id": hermes_bridge.active_run_id(conversation_id), + "message_id": draft["message_id"], + "text": content, + "ts": proto.now_iso(), + }) + return SendResult(success=True, message_id=draft["message_id"]) + + # Fresh/fallback final sends carry notify=True. A non-streaming final + # may arrive without metadata, so any ordinary send that reaches this + # path is also treated as a completed assistant message. + draft = self._drafts.pop(conversation_id, draft) + message_id = draft["message_id"] event = { "type": proto.S_MESSAGE_COMPLETE, "conversation_id": conversation_id, - "message_id": (draft or {}).get("message_id", message_id), + "run_id": hermes_bridge.active_run_id(conversation_id), + "message_id": message_id, "text": content, "ts": proto.now_iso(), } - if metadata and metadata.get("non_conversational"): - # Gateway lifecycle/status notices — deliver as a system note so - # the client can render them differently (or ignore). - event["kind"] = "notice" await self._server.broadcast(event) hermes_bridge.note_run_finished(conversation_id, "completed") return SendResult(success=True, message_id=message_id) @@ -162,8 +192,8 @@ class PhebyAdapter(BasePlatformAdapter): """Streaming path: GatewayStreamConsumer edits the in-place draft. The stream-consumer contract requires concrete adapters to accept - ``finalize=`` even when ignored (it's False during progressive - edits; the final content always arrives via ``send()``). + ``finalize=``. It is false during progressive edits and true when the + stream consumer itself owns final delivery. """ if self._server is None: return SendResult(success=False, error="server not running") @@ -173,13 +203,19 @@ class PhebyAdapter(BasePlatformAdapter): "text": "", }) draft["text"] = content # consumer sends cumulative text - await self._server.broadcast({ - "type": proto.S_MESSAGE_DELTA, + event = { + "type": (proto.S_MESSAGE_COMPLETE if finalize + else proto.S_MESSAGE_DELTA), "conversation_id": conversation_id, + "run_id": hermes_bridge.active_run_id(conversation_id), "message_id": draft["message_id"], "text": content, "ts": proto.now_iso(), - }) + } + await self._server.broadcast(event) + if finalize: + self._drafts.pop(conversation_id, None) + hermes_bridge.note_run_finished(conversation_id, "completed") return SendResult(success=True, message_id=draft["message_id"]) # ── structured stream events ───────────────────────────────────────── @@ -191,80 +227,75 @@ class PhebyAdapter(BasePlatformAdapter): S_TOOL_EVENT broadcast and return None so the gateway's text queue stays clean. (The dispatcher treats None as "adapter ate the event".) """ - try: - conversation_id = self._active_conversation_id() - if not conversation_id or self._server is None: - return None - if isinstance(event, ToolCallShim): - return None # never used at runtime; type-safety shim only - from gateway.stream_events import ToolCallChunk, ToolCallFinished - tool_event: Dict[str, Any] - if isinstance(event, ToolCallChunk): - tool_id = f"t-{proto.new_id()[:12]}" - args = event.args if isinstance(event.args, dict) else None - self._remember_tool(tool_id, conversation_id, event.tool_name) - tool_event = { - "type": proto.S_TOOL_EVENT, - "conversation_id": conversation_id, - "tool_call_id": tool_id, - "tool_name": event.tool_name, - "status": "running", - "description": proto.safe_str(event.preview, 300) - if event.preview else None, - "args_redacted": _redact_args(args), - "ts": proto.now_iso(), - } - elif isinstance(event, ToolCallFinished): - tool_id = self._lookup_tool(event.tool_name, conversation_id) - tool_event = { - "type": proto.S_TOOL_EVENT, - "conversation_id": conversation_id, - "tool_call_id": tool_id, - "tool_name": event.tool_name, - "status": "completed" if event.ok else "failed", - "duration_ms": int(event.duration * 1000) - if event.duration else None, - "ts": proto.now_iso(), - } - else: - return None - asyncio.ensure_future(self._server.broadcast(tool_event)) - except Exception: - logger.debug("[pheby] tool event translation failed", - exc_info=True) - return None # never render tool chrome as chat text - - def _tool_state(self) -> Dict[str, Any]: - if not hasattr(self, "_tool_calls"): - self._tool_calls: Dict[Tuple[str, str], str] = {} - self._tool_order: List[Tuple[str, str]] = [] - return {"calls": self._tool_calls, "order": self._tool_order} - - def _remember_tool(self, tool_id: str, conversation_id: str, - tool_name: str) -> None: - state = self._tool_state() - state["calls"][(tool_name, conversation_id)] = tool_id - state["order"].append((tool_name, conversation_id)) - if len(state["order"]) > 200: - old = state["order"].pop(0) - state["calls"].pop(old, None) - - def _lookup_tool(self, tool_name: str, conversation_id: str) -> str: - state = self._tool_state() - return state["calls"].get((tool_name, conversation_id), - f"t-{proto.new_id()[:12]}") + # Authoritative structured events come from the pre/post_tool_call + # hooks, which include session_id + tool_call_id. Eating this display + # event avoids duplicate/ambiguously-routed tool chrome. + return None # -- Hermes plugin hooks (registered in __init__.py register()) -------- + def schedule_broadcast(self, payload: Dict[str, Any]) -> None: + """Thread-safe hook → aiohttp-loop delivery.""" + loop, server = self._loop, self._server + if loop is None or server is None or loop.is_closed(): + return + + def _spawn() -> None: + asyncio.create_task(server.broadcast(dict(payload))) + try: + if asyncio.get_running_loop() is loop: + _spawn() + return + except RuntimeError: + pass + loop.call_soon_threadsafe(_spawn) + + def _conversation_for_session_id(self, session_id: str) -> Optional[str]: + runner = getattr(self, "gateway_runner", None) + store = getattr(runner, "session_store", None) if runner else None + if store is None or not session_id: + return None + try: + for entry in store.list_sessions(): + if str(getattr(entry, "session_id", "")) != str(session_id): + continue + origin = getattr(entry, "origin", None) + if getattr(getattr(origin, "platform", None), "value", "") == "pheby": + return str(getattr(origin, "chat_id", "") or "") or None + except Exception: + logger.debug("[pheby] tool session lookup failed", exc_info=True) + return None + + def on_pre_tool_call(self, **kwargs: Any) -> None: + conversation_id = self._conversation_for_session_id( + str(kwargs.get("session_id") or "")) + if not conversation_id: + return + tool_id = str(kwargs.get("tool_call_id") or + f"t-{proto.new_id()[:12]}") + event = { + "type": proto.S_TOOL_EVENT, + "conversation_id": conversation_id, + "run_id": hermes_bridge.active_run_id(conversation_id), + "tool_call_id": tool_id, + "tool_name": str(kwargs.get("tool_name") or "tool"), + "status": "running", + "args_redacted": _redact_args(kwargs.get("args")), + "ts": proto.now_iso(), + } + hermes_bridge.record_tool_event(conversation_id, event) + self.schedule_broadcast(event) + def on_post_tool_call(self, **kwargs: Any) -> None: """Observer for the ``post_tool_call`` plugin hook. Hermes fires this after every tool execution with the authoritative tool_call_id, status, duration, and result. We relay it as a structured ``S_TOOL_EVENT`` so the client can settle the matching - "running" event emitted by ``format_tool_event``. + running event emitted by ``on_pre_tool_call``. """ try: - conversation_id = self._active_conversation_id() + conversation_id = self._conversation_for_session_id( + str(kwargs.get("session_id") or "")) if not conversation_id or self._server is None: return tool_name = str(kwargs.get("tool_name") or "tool") @@ -273,12 +304,14 @@ class PhebyAdapter(BasePlatformAdapter): event = { "type": proto.S_TOOL_EVENT, "conversation_id": conversation_id, + "run_id": hermes_bridge.active_run_id(conversation_id), "tool_call_id": str(kwargs.get("tool_call_id") - or self._lookup_tool(tool_name, - conversation_id)), + or f"t-{proto.new_id()[:12]}"), "tool_name": tool_name, "status": "completed" if status in ("ok", "success", "") - else "failed" if status == "error" else status or "completed", + else "cancelled" if status == "cancelled" + else "failed" if status in ("error", "blocked") + else status or "completed", "duration_ms": int(duration_ms) if duration_ms else None, # Result summaries are intentionally NOT included by default: # tool results can embed file paths/host details. The client @@ -288,47 +321,43 @@ class PhebyAdapter(BasePlatformAdapter): error_message = kwargs.get("error_message") if error_message and event["status"] == "failed": event["error"] = proto.safe_str(error_message, 200) - asyncio.ensure_future(self._server.broadcast(event)) + hermes_bridge.record_tool_event(conversation_id, event) + self.schedule_broadcast(event) except Exception: logger.debug("[pheby] post_tool_call relay failed", exc_info=True) - - def _active_conversation_id(self) -> Optional[str]: - """Best-effort current conversation for adapter-level callbacks.""" - if not self._active_sessions: - return None - # Most recent active session wins (single-user platform). - key = sorted(self._active_sessions.keys())[-1] - # Session keys end with :dm: - return key.rsplit(":", 1)[-1] if ":" in key else None - # ── typing indicator → run activity ────────────────────────────────── async def send_typing(self, chat_id: str, metadata=None) -> None: # Pheby clients show their own activity UI from run/tool events. return # ── approvals ──────────────────────────────────────────────────────── - async def send_approval_prompt(self, session_key: str, - approval_data: Dict[str, Any]) -> None: - """Called from the approval notify callback (agent thread → here).""" + async def send_exec_approval( + self, + chat_id: str, + command: str, + session_key: str, + description: str = "dangerous command", + metadata: Optional[Dict[str, Any]] = None, + allow_permanent: bool = True, + allow_session: bool = True, + smart_denied: bool = False, + ) -> SendResult: + """Hermes's native structured-approval extension point.""" + if self._server is None: + return SendResult(success=False, error="server not running") try: - await hermes_bridge.push_approval(approval_data, session_key) - except Exception: + await hermes_bridge.push_approval({ + "command": command, + "description": description, + "allow_permanent": allow_permanent and not smart_denied, + "allow_session": allow_session and not smart_denied, + }, session_key) + return SendResult(success=True, + message_id=f"approval-{proto.new_id()[:12]}") + except Exception as exc: logger.error("[pheby] approval push failed", exc_info=True) - - def register_approval_notify(self, session_key: str) -> None: - """Wire tools.approval's per-session notify callback to Pheby.""" - from tools.approval import register_gateway_notify, \ - unregister_gateway_notify - loop = asyncio.get_event_loop() - # The callback runs on the agent's worker thread; bridge to the loop. - def _notify(approval_data: Dict[str, Any]) -> None: - asyncio.run_coroutine_threadsafe( - self.send_approval_prompt(session_key, approval_data), loop) - register_gateway_notify(session_key, _notify) - self._approval_notify_sessions = getattr( - self, "_approval_notify_sessions", set()) - self._approval_notify_sessions.add(session_key) + return SendResult(success=False, error=str(exc)) # ── clarification ──────────────────────────────────────────────────── async def send_clarify( @@ -461,25 +490,32 @@ class PhebyAdapter(BasePlatformAdapter): return _send -class ToolCallShim: - """Marker type for internal typing only — never instantiated.""" - - def _redact_args(args: Optional[Dict[str, Any]]) -> Optional[Dict[str, Any]]: - """Strip likely-secret values from tool args before sending to client.""" + """Recursively strip likely-secret values before sending args to client.""" if not isinstance(args, dict): return None sensitive = ("key", "token", "secret", "password", "credential", "auth") - out: Dict[str, Any] = {} - for k, v in args.items(): - k_l = str(k).lower() - if any(s in k_l for s in sensitive): - out[str(k)] = "[redacted]" - elif isinstance(v, str) and len(v) > 500: - out[str(k)] = v[:500] + "…" - else: - out[str(k)] = v - return out + + def _clean(value: Any, depth: int = 0) -> Any: + if depth > 6: + return "[truncated]" + if isinstance(value, dict): + cleaned: Dict[str, Any] = {} + for key, nested in list(value.items())[:100]: + key_s = str(key) + cleaned[key_s] = ("[redacted]" if any( + marker in key_s.lower() for marker in sensitive) + else _clean(nested, depth + 1)) + return cleaned + if isinstance(value, (list, tuple)): + return [_clean(item, depth + 1) for item in list(value)[:100]] + if isinstance(value, str): + return value[:500] + ("…" if len(value) > 500 else "") + if value is None or isinstance(value, (bool, int, float)): + return value + return proto.safe_str(value, 500) + + return _clean(args) __all__ = ["PhebyAdapter", "AIOHTTP_AVAILABLE"] diff --git a/plugin/pheby/attachments.py b/plugin/pheby/attachments.py index 715e004..3358afc 100644 --- a/plugin/pheby/attachments.py +++ b/plugin/pheby/attachments.py @@ -201,7 +201,7 @@ class AttachmentStore: def describe(self, attachment_id: str) -> Optional[Dict[str, Any]]: """Public descriptor for an attachment (no server paths).""" meta = self._meta.get(attachment_id) - if not meta: + if not meta or self._is_expired(meta): return None return { "attachment_id": attachment_id, diff --git a/plugin/pheby/hermes_bridge.py b/plugin/pheby/hermes_bridge.py index a0f8bda..691a217 100644 --- a/plugin/pheby/hermes_bridge.py +++ b/plugin/pheby/hermes_bridge.py @@ -9,9 +9,9 @@ so the module can be imported by unit tests without a Hermes install. from __future__ import annotations import asyncio -import contextvars import logging import threading +import time import uuid from typing import Any, Dict, List, Optional, Tuple @@ -27,6 +27,15 @@ _PENDING_CLARIFIES: Dict[str, Dict[str, Any]] = {} _RUN_LOCK = threading.Lock() _ACTIVE_RUNS: Dict[str, Dict[str, Any]] = {} # conversation_id → run info +_TOOL_EVENTS: Dict[str, Dict[str, Dict[str, Any]]] = {} + +# The adapter and server are process services, not request-local values. +# ContextVars lose their values when Hermes calls plugin hooks from agent +# worker threads, which made approvals/tool events disappear. Access is +# guarded because hook callbacks can arrive from multiple workers. +_SERVICE_LOCK = threading.RLock() +_ADAPTER: Any = None +_SERVER: Any = None def _runner() -> Any: @@ -35,16 +44,24 @@ def _runner() -> Any: return getattr(adapter, "gateway_runner", None) if adapter else None -_ADAPTER_CTX: contextvars.ContextVar = contextvars.ContextVar( - "pheby_adapter", default=None) - - def _current_adapter() -> Any: - return _ADAPTER_CTX.get() + with _SERVICE_LOCK: + return _ADAPTER def set_adapter(adapter: Any) -> None: - _ADAPTER_CTX.set(adapter) + global _ADAPTER + with _SERVICE_LOCK: + _ADAPTER = adapter + + +def clear_services(adapter: Any = None) -> None: + """Release process-wide references when the owning adapter disconnects.""" + global _ADAPTER, _SERVER + with _SERVICE_LOCK: + if adapter is None or _ADAPTER is adapter: + _ADAPTER = None + _SERVER = None def _session_store() -> Any: @@ -99,15 +116,45 @@ def _session_key_for(conversation_id: str) -> str: return ConversationRouter.session_key_for(conversation_id) +def _conversation_from_session_key(session_key: str) -> str: + """Extract the opaque chat ID from a Pheby DM session key.""" + marker = ":pheby:dm:" + if marker in str(session_key): + return str(session_key).split(marker, 1)[1] + return str(session_key).rsplit(":", 1)[-1] + + # ═══════════════════════════════════════════════════════════════════════════ # Conversations # ═══════════════════════════════════════════════════════════════════════════ +async def create_conversation(server: Any, name: Optional[str]) -> str: + """Create the Pheby ID and an empty Hermes session routing entry.""" + cid = await server.router.new_conversation(name) + store = _session_store() + if store is not None: + try: + await asyncio.to_thread( + store.get_or_create_session, _source_for(cid), False, False) + if name: + await rename_conversation(cid, name) + except Exception: + # Keep the router entry so the empty conversation remains usable; + # the first message can create its Hermes session normally. + logger.warning("[pheby] empty Hermes session creation failed", + exc_info=True) + return cid + + async def list_conversations(server: Any = None) -> List[Dict[str, Any]]: """Enumerate conversations known to the router + Hermes session store.""" server = server or _current_server() router = server.router if server else None out: List[Dict[str, Any]] = [] seen: set = set() + router_names: Dict[str, str] = {} + if router is not None: + for cid in await router.known_ids(): + router_names[cid] = await router.get_name(cid) or cid # 1. Sessions Hermes already tracks for the pheby platform. store = _session_store() @@ -126,8 +173,10 @@ async def list_conversations(server: Any = None) -> List[Dict[str, Any]]: seen.add(cid) out.append({ "conversation_id": cid, - "name": (getattr(entry, "display_name", None) - or _router_name(router, cid) or cid), + # An explicit Pheby rename wins over Hermes's initial + # source-derived display name. + "name": (router_names.get(cid) + or getattr(entry, "display_name", None) or cid), "session_id": getattr(entry, "session_id", None), "last_active": _iso(getattr(entry, "updated_at", None)), "source": "hermes", @@ -155,12 +204,6 @@ async def list_conversations(server: Any = None) -> List[Dict[str, Any]]: return out -def _router_name(router: Any, cid: str) -> Optional[str]: - if router is None: - return None - return router._names.get(cid) - - def _iso(value: Any) -> Optional[str]: try: return value.isoformat() if value else None @@ -182,7 +225,6 @@ async def conversation_history(conversation_id: str, limit: int session_id: Optional[str] = None if store is not None: try: - source = _source_for(conversation_id) entry = await asyncio.to_thread(store.peek_session_id, _session_key_for(conversation_id)) if entry: @@ -195,7 +237,8 @@ async def conversation_history(conversation_id: str, limit: int if db is not None and session_id: try: rows = await asyncio.to_thread( - db.get_messages_as_conversation, session_id) + db.get_messages_as_conversation, session_id, + include_row_ids=True) for row in rows[-limit:]: role = row.get("role") if role not in ("user", "assistant"): @@ -207,8 +250,8 @@ async def conversation_history(conversation_id: str, limit: int if not text.strip() and role == "assistant": continue messages.append({ - "message_id": f"m{row.get('id', len(messages))}" - if isinstance(row.get("id"), (int, str)) else None, + "message_id": f"m{row.get('_row_id')}" + if isinstance(row.get("_row_id"), (int, str)) else None, "role": role, "text": text, "ts": row.get("timestamp") if isinstance( @@ -229,44 +272,88 @@ async def conversation_history(conversation_id: str, limit: int return messages, found -async def delete_conversation(conversation_id: str) -> bool: - """Delete a conversation from the router + Hermes (best effort on DB). - - Hermes limitation: the SessionStore has no public per-key delete; the - authoritative delete is ``SessionDB.delete_session`` on the current - session id. The routing entry is also reset so the next message starts - a fresh session. Documented approximation — see README limitations. - """ +async def rename_conversation(conversation_id: str, name: str) -> bool: + """Rename the Pheby index and the live Hermes routing entry.""" server = _current_server() router = server.router if server else None if router is None: return False - if not await router.forget(conversation_id): + if not await router.rename(conversation_id, name): + return False + + store = _session_store() + if store is not None: + session_key = _session_key_for(conversation_id) + + def _rename_route() -> None: + with store._lock: + store._ensure_loaded_locked() + entry = store._entries.get(session_key) + if entry is not None: + entry.display_name = name + store._save() + + try: + await asyncio.to_thread(_rename_route) + except Exception: + logger.debug("[pheby] Hermes display-name update failed", + exc_info=True) + return True + + +async def delete_conversation(conversation_id: str) -> bool: + """Delete a transcript and remove its live routing entry. + + Hermes currently has no public per-key removal method. We therefore use + the same lock/save discipline as SessionStore's own pruning code. Calling + ``reset_session`` here would create a replacement entry and make the + deleted conversation immediately reappear. + """ + server = _current_server() + router = server.router if server else None + if router is None or active_run(conversation_id) is not None: return False store = _session_store() db = _session_db() + session_key = _session_key_for(conversation_id) session_id = None if store is not None: try: session_id = await asyncio.to_thread( - store.peek_session_id, _session_key_for(conversation_id)) + store.peek_session_id, session_key) except Exception: session_id = None + router_known = await router.get_name(conversation_id) is not None + if not router_known and not session_id: + return False + if session_id and db is not None: try: - await asyncio.to_thread(db.delete_session, session_id) + deleted = await asyncio.to_thread(db.delete_session, session_id) + if deleted is False: + return False except Exception: - logger.debug("[pheby] session db delete failed", exc_info=True) - if store is not None: - try: - await asyncio.to_thread(store.reset_session, - _session_key_for(conversation_id), - None) - except Exception: - logger.debug("[pheby] store reset failed", exc_info=True) + logger.error("[pheby] session db delete failed", exc_info=True) + return False - _ACTIVE_RUNS.pop(conversation_id, None) + if store is not None: + def _remove_route() -> None: + with store._lock: + store._ensure_loaded_locked() + if store._entries.pop(session_key, None) is not None: + store._save() + + try: + await asyncio.to_thread(_remove_route) + except Exception: + logger.error("[pheby] routing removal failed", exc_info=True) + return False + + await router.forget(conversation_id) + with _RUN_LOCK: + _ACTIVE_RUNS.pop(conversation_id, None) + _TOOL_EVENTS.pop(conversation_id, None) return True @@ -301,10 +388,19 @@ async def send_chat(server: Any, conversation_id: str, text: str, metadata={"pheby_run_id": run_id}, ) - _ACTIVE_RUNS[conversation_id] = { - "run_id": run_id, - "started": asyncio.get_event_loop().time(), - } + with _RUN_LOCK: + already_active = conversation_id in _ACTIVE_RUNS + if not already_active: + _ACTIVE_RUNS[conversation_id] = { + "run_id": run_id, + "started": asyncio.get_running_loop().time(), + } + _TOOL_EVENTS[conversation_id] = {} + if already_active: + await client.send_json(proto.error_event( + proto.ERR_RUN_ACTIVE, + "A run is already active for this conversation", request_id)) + return await client.send_json({ "type": proto.S_RUN_ACCEPTED, "conversation_id": conversation_id, @@ -327,14 +423,73 @@ async def send_chat(server: Any, conversation_id: str, text: str, # The base adapter's handle_message() spawns background tasks and # returns quickly; the eventual reply arrives through adapter.send(). - await adapter.handle_message(event) + try: + await adapter.handle_message(event) + except Exception: + with _RUN_LOCK: + current = _ACTIVE_RUNS.get(conversation_id) + if current and current.get("run_id") == run_id: + _ACTIVE_RUNS.pop(conversation_id, None) + if getattr(adapter, "_drafts", None) is not None: + adapter._drafts.pop(conversation_id, None) + await server.broadcast({ + "type": proto.S_RUN_FINISHED, + "conversation_id": conversation_id, + "run_id": run_id, + "status": "failed", + "error": "Gateway rejected the message", + }) + raise + + +def active_run(conversation_id: str) -> Optional[Dict[str, Any]]: + with _RUN_LOCK: + run = _ACTIVE_RUNS.get(conversation_id) + return dict(run) if run else None + + +def active_run_id(conversation_id: str) -> Optional[str]: + run = active_run(conversation_id) + return str(run["run_id"]) if run else None + + +def record_tool_event(conversation_id: str, event: Dict[str, Any]) -> None: + tool_id = str(event.get("tool_call_id") or "") + if not tool_id: + return + with _RUN_LOCK: + bucket = _TOOL_EVENTS.setdefault(conversation_id, {}) + bucket[tool_id] = dict(event) + if len(bucket) > 200: + bucket.pop(next(iter(bucket))) + + +def runtime_snapshot(conversation_id: str) -> Dict[str, Any]: + """Recoverable transient state included by ``conversation.open``.""" + session_key = _session_key_for(conversation_id) + with _RUN_LOCK: + run = _ACTIVE_RUNS.get(conversation_id) + tools = list(_TOOL_EVENTS.get(conversation_id, {}).values()) + approvals = [dict(p["event"]) for p in _PENDING_APPROVALS.values() + if p.get("session_key") == session_key] + clarifications = [dict(p["event"]) for p in _PENDING_CLARIFIES.values() + if p.get("session_key") == session_key] + return { + "run": dict(run) if run else None, + "tools": tools, + "approvals": approvals, + "clarifications": clarifications, + } def note_run_finished(conversation_id: str, status: str = "completed", error: Optional[str] = None) -> None: """Called by the adapter when a turn completes/fails.""" - run = _ACTIVE_RUNS.pop(conversation_id, None) - run_id = run["run_id"] if run else None + with _RUN_LOCK: + run = _ACTIVE_RUNS.pop(conversation_id, None) + if run is None: + return + run_id = run["run_id"] server = _current_server() if server is None: return @@ -348,9 +503,9 @@ def note_run_finished(conversation_id: str, status: str = "completed", if error: payload["error"] = proto.safe_str(error, 300) try: - loop = asyncio.get_event_loop() - if loop.is_running(): - asyncio.ensure_future(server.broadcast(payload)) + adapter = _current_adapter() + if adapter is not None: + adapter.schedule_broadcast(payload) except RuntimeError: pass @@ -358,8 +513,10 @@ def note_run_finished(conversation_id: str, status: str = "completed", async def cancel_run(conversation_id: str, run_id: Optional[str]) -> bool: """Cancel an active run via Hermes's supported interrupt path.""" adapter = _current_adapter() - run = _ACTIVE_RUNS.get(conversation_id) - if run and run_id and run["run_id"] != run_id: + run = active_run(conversation_id) + if run is None: + return False + if run_id and run["run_id"] != run_id: return False # stale run id — nothing to cancel session_key = _session_key_for(conversation_id) @@ -369,10 +526,10 @@ async def cancel_run(conversation_id: str, run_id: Optional[str]) -> bool: # Preferred: gateway's own /stop dispatch (cancels task + drains). running = getattr(runner, "_running_agents", {}).get(session_key) agent = running if running is not None else None - if agent is not None and agent is not getattr( - type(runner), "_AGENT_PENDING_SENTINEL", object()): + interrupt = getattr(agent, "interrupt", None) + if callable(interrupt): try: - agent.interrupt("Cancelled by Pheby client") + interrupt("Cancelled by Pheby client") invalidate = getattr( runner, "_invalidate_session_run_generation", None) if callable(invalidate): @@ -382,13 +539,16 @@ async def cancel_run(conversation_id: str, run_id: Optional[str]) -> bool: logger.debug("[pheby] agent interrupt failed", exc_info=True) if not interrupted and adapter is not None: try: - await adapter.interrupt_session_activity( - session_key, conversation_id) - interrupted = True + event = getattr(adapter, "_active_sessions", {}).get(session_key) + if event is not None: + await adapter.interrupt_session_activity( + session_key, conversation_id) + interrupted = True except Exception: logger.debug("[pheby] adapter interrupt failed", exc_info=True) - note_run_finished(conversation_id, - "cancelled" if interrupted else "idle") + if interrupted: + with _RUN_LOCK: + _ACTIVE_RUNS.pop(conversation_id, None) return interrupted @@ -399,17 +559,20 @@ async def push_approval(approval_data: Dict[str, Any], session_key: str) -> None: """Adapter callback: a dangerous action needs a human decision.""" approval_id = uuid.uuid4().hex[:12] - from gateway.run import _redact_approval_command - command = _redact_approval_command(approval_data.get("command", "")) + # gateway.run redacts the command before calling send_exec_approval. + command = approval_data.get("command", "") choices: List[str] = ["once", "deny"] if approval_data.get("allow_session", True): choices.insert(1, "session") if approval_data.get("allow_permanent", True): choices.insert(-1, "always") - event = { + conversation_id = _conversation_from_session_key(session_key) + event: Dict[str, Any] = { "type": proto.S_APPROVAL_REQUEST, "approval_id": approval_id, "session_key": session_key, + "conversation_id": conversation_id, + "run_id": active_run_id(conversation_id), "command": proto.safe_str(command, 2000), "description": proto.safe_str( approval_data.get("description", ""), 1000), @@ -418,7 +581,8 @@ async def push_approval(approval_data: Dict[str, Any], } _PENDING_APPROVALS[approval_id] = { "session_key": session_key, - "created": asyncio.get_event_loop().time(), + "created": time.monotonic(), + "event": event, } server = _current_server() if server is not None: @@ -428,11 +592,11 @@ async def push_approval(approval_data: Dict[str, Any], async def resolve_approval(approval_id: str, choice: str, reason: Optional[str]) -> bool: """Forward an approval decision to Hermes (tools.approval primitives).""" - pending = _PENDING_APPROVALS.pop(approval_id, None) - if pending is None: - return False if choice not in ("once", "session", "always", "deny"): return False + pending = _PENDING_APPROVALS.get(approval_id) + if pending is None: + return False try: from tools.approval import resolve_gateway_approval count = await asyncio.to_thread( @@ -442,20 +606,14 @@ async def resolve_approval(approval_id: str, choice: str, except Exception: logger.error("[pheby] approval resolve failed", exc_info=True) ok = False - server = _current_server() - if server is not None: - await server.broadcast({ - "type": proto.S_APPROVAL_RESOLVED, - "approval_id": approval_id, - "choice": choice, - "accepted": ok, - }) - return True + if ok: + _PENDING_APPROVALS.pop(approval_id, None) + return ok def fail_stale_approvals(max_age: float = 3600.0) -> None: """Drop approval IDs whose Hermes-side gate has surely timed out.""" - now = asyncio.get_event_loop().time() + now = time.monotonic() for aid in [a for a, p in _PENDING_APPROVALS.items() if now - p["created"] > max_age]: _PENDING_APPROVALS.pop(aid, None) @@ -467,10 +625,13 @@ def fail_stale_approvals(max_age: float = 3600.0) -> None: async def push_clarify(clarify_id: str, session_key: str, question: str, choices: Optional[List[str]]) -> None: """Adapter callback: the agent needs the user to choose.""" - event = { + conversation_id = _conversation_from_session_key(session_key) + event: Dict[str, Any] = { "type": proto.S_CLARIFY_REQUEST, "clarify_id": clarify_id, "session_key": session_key, + "conversation_id": conversation_id, + "run_id": active_run_id(conversation_id), "question": proto.safe_str(question, 2000), "choices": [proto.safe_str(c, 300) for c in choices] if choices else None, @@ -479,7 +640,8 @@ async def push_clarify(clarify_id: str, session_key: str, question: str, } _PENDING_CLARIFIES[clarify_id] = { "session_key": session_key, - "created": asyncio.get_event_loop().time(), + "created": time.monotonic(), + "event": event, } server = _current_server() if server is not None: @@ -488,7 +650,7 @@ async def push_clarify(clarify_id: str, session_key: str, question: str, async def resolve_clarify(clarify_id: str, response: str) -> bool: """Forward a clarification answer to Hermes's clarify primitive.""" - pending = _PENDING_CLARIFIES.pop(clarify_id, None) + pending = _PENDING_CLARIFIES.get(clarify_id) if pending is None: return False try: @@ -506,32 +668,33 @@ async def resolve_clarify(clarify_id: str, response: str) -> bool: except Exception: logger.error("[pheby] clarify resolve failed", exc_info=True) ok = False - server = _current_server() - if server is not None: - await server.broadcast({ - "type": proto.S_CLARIFY_RESOLVED, - "clarify_id": clarify_id, - "accepted": bool(ok), - }) + if ok: + _PENDING_CLARIFIES.pop(clarify_id, None) return bool(ok) # ═══════════════════════════════════════════════════════════════════════════ # Models & reasoning # ═══════════════════════════════════════════════════════════════════════════ -async def models_snapshot() -> Dict[str, Any]: +async def models_snapshot(conversation_id: Optional[str] = None) -> Dict[str, Any]: """Providers + models Hermes currently exposes (credential-aware).""" def _collect() -> Dict[str, Any]: + from hermes_cli.config import get_compatible_custom_providers from hermes_cli.model_switch import list_picker_providers cfg = _load_cfg() model_cfg = (cfg.get("model") or {}) if isinstance(cfg, dict) else {} current_model = str(model_cfg.get("default", "") or "") current_provider = str(model_cfg.get("provider", "openrouter") or "") + custom_providers = get_compatible_custom_providers(cfg) + excluded = (cfg.get("model_catalog") or {}).get( + "excluded_providers", []) providers = list_picker_providers( current_provider=current_provider, + current_base_url=str(model_cfg.get("base_url", "") or ""), current_model=current_model, user_providers=cfg.get("providers") if isinstance(cfg, dict) else None, - probe_custom_providers=False, # don't block on offline endpoints + custom_providers=custom_providers, + excluded_providers=excluded if isinstance(excluded, list) else [], ) return {"providers": providers, "current_model": current_model, "current_provider": current_provider} @@ -542,11 +705,18 @@ async def models_snapshot() -> Dict[str, Any]: data = {"providers": [], "current_model": "", "current_provider": "", "error": "Model catalog unavailable"} data["supported_reasoning_efforts"] = list(proto.REASONING_EFFORTS) + if conversation_id: + current = await current_model_snapshot(conversation_id) + data["current_model"] = current.get("model", data["current_model"]) + data["current_provider"] = current.get( + "provider", data["current_provider"]) + data["scope"] = current.get("scope", "global") data["ts"] = proto.now_iso() return data -async def current_model_snapshot() -> Dict[str, Any]: +async def current_model_snapshot( + conversation_id: Optional[str] = None) -> Dict[str, Any]: def _collect() -> Dict[str, Any]: cfg = _load_cfg() model_cfg = (cfg.get("model") or {}) if isinstance(cfg, dict) else {} @@ -556,6 +726,19 @@ async def current_model_snapshot() -> Dict[str, Any]: data = await asyncio.to_thread(_collect) except Exception: data = {"model": "", "provider": "", "error": "Config unavailable"} + data["scope"] = "global" + store = _session_store() + if conversation_id and store is not None: + try: + override = await asyncio.to_thread( + store.get_model_override, _session_key_for(conversation_id)) + if override: + data["model"] = override.get("model", data["model"]) + data["provider"] = override.get("provider", data["provider"]) + data["scope"] = "conversation" + except Exception: + logger.debug("[pheby] model override read failed", exc_info=True) + data["conversation_id"] = conversation_id data["ts"] = proto.now_iso() return data @@ -567,6 +750,7 @@ async def set_model(model: str, provider: Optional[str], return {"ok": False, "code": proto.ERR_BAD_REQUEST, "message": "model is required"} try: + from hermes_cli.config import get_compatible_custom_providers from hermes_cli.model_switch import switch_model cfg = _load_cfg() model_cfg = (cfg.get("model") or {}) if isinstance(cfg, dict) else {} @@ -580,7 +764,7 @@ async def set_model(model: str, provider: Optional[str], False, # is_global → session-scoped when conversation given provider or "", cfg.get("providers") if isinstance(cfg, dict) else None, - None, + get_compatible_custom_providers(cfg), ) except Exception as exc: logger.error("[pheby] switch_model failed", exc_info=True) @@ -602,6 +786,11 @@ async def set_model(model: str, provider: Optional[str], store = _session_store() if conversation_id and store is not None: try: + if await asyncio.to_thread( + store.peek_session_id, + _session_key_for(conversation_id)) is None: + return {"ok": False, "code": proto.ERR_CONVERSATION_NOT_FOUND, + "message": "Conversation not found"} await asyncio.to_thread(store.set_model_override, _session_key_for(conversation_id), override) @@ -624,13 +813,14 @@ async def set_model(model: str, provider: Optional[str], def _save_global_model(model: str, provider: str) -> None: - from hermes_cli.config import load_config, save_config_value + from hermes_cli.config import save_config_value save_config_value("model.default", model) if provider: save_config_value("model.provider", provider) -async def reasoning_snapshot() -> Dict[str, Any]: +async def reasoning_snapshot( + conversation_id: Optional[str] = None) -> Dict[str, Any]: def _collect() -> Dict[str, Any]: from hermes_constants import resolve_reasoning_config cfg = _load_cfg() @@ -647,6 +837,22 @@ async def reasoning_snapshot() -> Dict[str, Any]: except Exception: data = {"effort": None, "enabled": None, "error": "Config unavailable"} + data["scope"] = "global" + runner = _runner() + if conversation_id and runner is not None: + try: + cfg = await asyncio.to_thread( + runner._resolve_session_reasoning_config, + session_key=_session_key_for(conversation_id), model="") + if cfg is not None: + data = ({"effort": "none", "enabled": False} + if cfg.get("enabled") is False else + {"effort": cfg.get("effort"), "enabled": True}) + data["scope"] = "conversation" + except Exception: + logger.debug("[pheby] reasoning override read failed", + exc_info=True) + data["conversation_id"] = conversation_id data["supported_efforts"] = ["none"] + list(proto.REASONING_EFFORTS) data["ts"] = proto.now_iso() return data @@ -664,6 +870,12 @@ async def set_reasoning(effort: str, runner = _runner() if runner is not None and conversation_id: try: + store = _session_store() + if store is None or await asyncio.to_thread( + store.peek_session_id, + _session_key_for(conversation_id)) is None: + return {"ok": False, "code": proto.ERR_CONVERSATION_NOT_FOUND, + "message": "Conversation not found"} await asyncio.to_thread( runner._set_session_reasoning_override, _session_key_for(conversation_id), parsed) @@ -691,23 +903,24 @@ def _load_cfg() -> Dict[str, Any]: # ═══════════════════════════════════════════════════════════════════════════ -# Server context +# Process service references # ═══════════════════════════════════════════════════════════════════════════ -_SERVER_CTX: contextvars.ContextVar = contextvars.ContextVar( - "pheby_server", default=None) - - def set_server(server: Any) -> None: - _SERVER_CTX.set(server) + global _SERVER + with _SERVICE_LOCK: + _SERVER = server def _current_server() -> Any: - return _SERVER_CTX.get() + with _SERVICE_LOCK: + return _SERVER __all__ = [ - "set_adapter", "set_server", "list_conversations", "conversation_history", + "set_adapter", "set_server", "clear_services", "create_conversation", + "list_conversations", "conversation_history", "rename_conversation", "delete_conversation", "send_chat", "cancel_run", "note_run_finished", + "active_run", "active_run_id", "record_tool_event", "runtime_snapshot", "push_approval", "resolve_approval", "fail_stale_approvals", "push_clarify", "resolve_clarify", "models_snapshot", "current_model_snapshot", "set_model", "reasoning_snapshot", diff --git a/plugin/pheby/plugin.yaml b/plugin/pheby/plugin.yaml index a72a126..1110760 100644 --- a/plugin/pheby/plugin.yaml +++ b/plugin/pheby/plugin.yaml @@ -1,7 +1,7 @@ name: pheby label: Pheby kind: platform -version: 1.0.0 +version: 1.0.1 description: > Pheby platform adapter for Hermes Agent — serves a WebSocket + HTTPS protocol for the Pheby native Android client (Kotlin/Compose) behind a @@ -37,11 +37,3 @@ optional_env: description: "Conversation ID receiving cron/scheduled deliveries by default" prompt: "Home conversation ID (or empty)" password: false - - name: PHEBY_ALLOWED_USERS - description: "Comma-separated allowlist treated as user IDs by the gateway (optional)" - prompt: "Allowed user IDs (or empty)" - password: false - - name: PHEBY_ALLOW_ALL_USERS - description: "Allow any authenticated client (dev only)" - prompt: "Allow all users? (true/false)" - password: false diff --git a/plugin/pheby/protocol.py b/plugin/pheby/protocol.py index 8fe1359..afdbdef 100644 --- a/plugin/pheby/protocol.py +++ b/plugin/pheby/protocol.py @@ -48,6 +48,7 @@ ERR_NOT_FOUND = "not_found" ERR_CONVERSATION_NOT_FOUND = "conversation_not_found" ERR_APPROVAL_NOT_FOUND = "approval_not_found" ERR_CLARIFY_NOT_FOUND = "clarify_not_found" +ERR_RUN_ACTIVE = "run_active" ERR_TOO_LARGE = "too_large" ERR_RATE_LIMITED = "rate_limited" ERR_INTERNAL = "internal_error" @@ -172,7 +173,7 @@ __all__ = [ "ERR_UNAUTHORIZED", "ERR_AUTH_TIMEOUT", "ERR_VERSION_MISMATCH", "ERR_BAD_REQUEST", "ERR_INVALID_JSON", "ERR_UNKNOWN_TYPE", "ERR_NOT_FOUND", "ERR_CONVERSATION_NOT_FOUND", "ERR_APPROVAL_NOT_FOUND", - "ERR_CLARIFY_NOT_FOUND", "ERR_TOO_LARGE", "ERR_RATE_LIMITED", + "ERR_CLARIFY_NOT_FOUND", "ERR_RUN_ACTIVE", "ERR_TOO_LARGE", "ERR_RATE_LIMITED", "ERR_INTERNAL", "ERR_NOT_IMPLEMENTED", "C_HELLO", "C_PING", "C_CONVERSATION_LIST", "C_CONVERSATION_OPEN", "C_CONVERSATION_CREATE", "C_CONVERSATION_RENAME", "C_CONVERSATION_DELETE", diff --git a/plugin/pheby/server.py b/plugin/pheby/server.py index dd7d57c..941ac82 100644 --- a/plugin/pheby/server.py +++ b/plugin/pheby/server.py @@ -14,8 +14,10 @@ Hermes integration (runs, approvals, clarifications, models) lives in from __future__ import annotations import asyncio +import ipaddress import logging import time +from urllib.parse import quote from pathlib import Path from typing import Any, Dict, List, Optional @@ -153,13 +155,17 @@ class PhebyServer: "message": "Attachment unavailable"}}, status=404) desc = self.store.describe(attachment_id) or {} - safe_name = desc.get("filename", "file.bin") + safe_name = str(desc.get("filename", "file.bin")) + ascii_name = safe_name.encode("ascii", "replace").decode("ascii") \ + .replace('"', "_").replace("\\", "_") + disposition = (f'attachment; filename="{ascii_name}"; ' + f"filename*=UTF-8''{quote(safe_name)}") logger.info("[pheby] attachment download: id=%s bytes=%s", attachment_id, desc.get("size")) return web.FileResponse( blob, headers={ - "Content-Disposition": f'attachment; filename="{safe_name}"', + "Content-Disposition": disposition, "Content-Type": desc.get("mime_type", "application/octet-stream"), }, @@ -173,7 +179,7 @@ class PhebyServer: self._conn_counter += 1 conn_id = f"c{self._conn_counter}" - peer = request.remote or "unknown" + peer = self._auth_peer(request) if self._is_locked_out(peer): logger.warning("[pheby] auth lockout active for %s — refusing", peer) @@ -224,6 +230,21 @@ class PhebyServer: def _record_auth_failure(self, peer: str) -> None: self._auth_failures.setdefault(peer, []).append(time.time()) + @staticmethod + def _auth_peer(request: web.Request) -> str: + """Use Caddy's client IP only when the direct peer is loopback.""" + direct = request.remote or "unknown" + try: + if not ipaddress.ip_address(direct).is_loopback: + return direct + except ValueError: + return direct + forwarded = request.headers.get("X-Forwarded-For", "").split(",", 1)[0].strip() + try: + return str(ipaddress.ip_address(forwarded)) if forwarded else direct + except ValueError: + return direct + async def _authenticate(self, client: ClientConnection, peer: str) -> bool: """Wait for the hello frame and validate the shared secret.""" @@ -253,7 +274,15 @@ class PhebyServer: proto.ERR_UNAUTHORIZED, "Invalid secret")) return False requested = message.get("protocol_version") - if requested is not None and int(requested) != proto.PROTOCOL_VERSION: + try: + requested_version = (proto.PROTOCOL_VERSION if requested is None + else int(requested)) + except (TypeError, ValueError): + await client.send_json(proto.error_event( + proto.ERR_VERSION_MISMATCH, + "protocol_version must be an integer")) + return False + if requested_version != proto.PROTOCOL_VERSION: await client.send_json(proto.error_event( proto.ERR_VERSION_MISMATCH, f"Protocol version mismatch: server={proto.PROTOCOL_VERSION}, " @@ -343,17 +372,20 @@ class PhebyServer: proto.ERR_CONVERSATION_NOT_FOUND, "Conversation not found", request_id)) return + runtime = self.bridge.runtime_snapshot(conversation_id) await client.send_json({ "type": proto.S_CONVERSATION_HISTORY, "conversation_id": conversation_id, "messages": history, + "attachments": self.store.list_for_conversation(conversation_id), + **runtime, **({"request_id": request_id} if request_id else {}), }) async def _handle_conversation_create(self, client, message, request_id): name = message.get("name") - cid = await self.router.new_conversation( - str(name) if name else None) + cid = await self.bridge.create_conversation( + self, str(name) if name else None) await client.send_json({ "type": proto.S_CONVERSATION_CREATED, "conversation_id": cid, @@ -375,7 +407,7 @@ class PhebyServer: proto.ERR_BAD_REQUEST, "conversation_id and name (≤200 chars) required", request_id)) return - ok = await self.router.rename(conversation_id, name) + ok = await self.bridge.rename_conversation(conversation_id, name) if not ok: await client.send_json(proto.error_event( proto.ERR_CONVERSATION_NOT_FOUND, "Conversation not found", @@ -428,17 +460,30 @@ class PhebyServer: proto.ERR_TOO_LARGE, f"text exceeds {proto.MAX_TEXT_CHARS} chars", request_id)) return - await self.bridge.send_chat(conversation_id, text, client, request_id) + await self.bridge.send_chat( + self, conversation_id, text, client, request_id) async def _handle_run_cancel(self, client, message, request_id): conversation_id = str(message.get("conversation_id", "")) - run_id = message.get("run_id") - ok = await self.bridge.cancel_run(conversation_id, run_id) - await client.send_json({ - "type": proto.S_RUN_FINISHED if ok else proto.S_ERROR, - **({"run_id": run_id, "status": "cancelled"} - if ok else {"error": {"code": proto.ERR_NOT_FOUND, - "message": "No active run"}}), + if not ConversationRouter.is_valid_conversation_id(conversation_id): + await client.send_json(proto.error_event( + proto.ERR_BAD_REQUEST, "Invalid conversation_id", request_id)) + return + requested_run_id = message.get("run_id") + active = self.bridge.active_run(conversation_id) + run_id = active.get("run_id") if active else requested_run_id + ok = await self.bridge.cancel_run(conversation_id, requested_run_id) + if not ok: + await client.send_json(proto.error_event( + proto.ERR_NOT_FOUND, "No matching active run", request_id)) + return + if self.adapter is not None: + self.adapter._drafts.pop(conversation_id, None) + await self.broadcast({ + "type": proto.S_RUN_FINISHED, + "conversation_id": conversation_id, + "run_id": run_id, + "status": "cancelled", **({"request_id": request_id} if request_id else {}), }) @@ -453,10 +498,11 @@ class PhebyServer: proto.ERR_APPROVAL_NOT_FOUND, "Unknown or already-resolved approval", request_id)) return - await client.send_json({ + await self.broadcast({ "type": proto.S_APPROVAL_RESOLVED, "approval_id": approval_id, "choice": choice, + "accepted": True, **({"request_id": request_id} if request_id else {}), }) @@ -470,21 +516,26 @@ class PhebyServer: proto.ERR_CLARIFY_NOT_FOUND, "Unknown or already-resolved clarification", request_id)) return - await client.send_json({ + await self.broadcast({ "type": proto.S_CLARIFY_RESOLVED, "clarify_id": clarify_id, + "accepted": True, **({"request_id": request_id} if request_id else {}), }) async def _handle_models_list(self, client, message, request_id): - snapshot = await self.bridge.models_snapshot() + conversation_id = message.get("conversation_id") + snapshot = await self.bridge.models_snapshot( + str(conversation_id) if conversation_id else None) snapshot["type"] = proto.S_MODELS_SNAPSHOT if request_id: snapshot["request_id"] = request_id await client.send_json(snapshot) async def _handle_model_current(self, client, message, request_id): - snapshot = await self.bridge.current_model_snapshot() + conversation_id = message.get("conversation_id") + snapshot = await self.bridge.current_model_snapshot( + str(conversation_id) if conversation_id else None) snapshot["type"] = proto.S_MODEL_CURRENT_SNAPSHOT if request_id: snapshot["request_id"] = request_id @@ -514,7 +565,9 @@ class PhebyServer: if k != "request_id"}) async def _handle_reasoning_current(self, client, message, request_id): - snapshot = await self.bridge.reasoning_snapshot() + conversation_id = message.get("conversation_id") + snapshot = await self.bridge.reasoning_snapshot( + str(conversation_id) if conversation_id else None) snapshot["type"] = proto.S_REASONING_SNAPSHOT if request_id: snapshot["request_id"] = request_id diff --git a/tests/test_pheby.py b/tests/test_pheby.py index cf3e4e1..cd9c27f 100644 --- a/tests/test_pheby.py +++ b/tests/test_pheby.py @@ -13,13 +13,12 @@ from __future__ import annotations import asyncio import json -import os +import threading import time from pathlib import Path from typing import Any, Dict, List, Optional import pytest -from aiohttp import web # Make the plugin package importable regardless of install layout. import sys @@ -156,7 +155,6 @@ def make_server(tmp_path: Path, **overrides) -> PhebyServer: **overrides, }) cfg.secret = overrides.get("secret", cfg.secret or "test-secret-abc123") - from hermes_constants import get_hermes_home # conftest redirects home root = Path(tmp_path) / "attachments" server = PhebyServer(cfg, adapter=FakeAdapter()) server.store = AttachmentStore(root=root, retention_days=7) @@ -248,6 +246,19 @@ class TestAuth: ok = await server._authenticate(client, "peer4") assert not ok and not client.authenticated + @pytest.mark.asyncio + async def test_malformed_version_is_structured_error(self, tmp_path): + server = make_server(tmp_path) + client = FakeClientConnection() + client.ws.inbox.put_nowait(proto.encode_message({ + "type": proto.C_HELLO, "secret": "test-secret-abc123", + "protocol_version": {"not": "an integer"}})) + ok = await server._authenticate(client, "peer5") + assert not ok and not client.authenticated + error = client.ws.events()[-1] + assert error["type"] == proto.S_ERROR + assert error["error"]["code"] == proto.ERR_VERSION_MISMATCH + def test_constant_time_equals(self): assert constant_time_equals("abc", "abc") assert not constant_time_equals("abc", "abd") @@ -414,7 +425,8 @@ class TestAttachments: # ═══════════════════════════════════════════════════════════════════════════ class TestBridge: @pytest.mark.asyncio - async def test_tool_start_event_is_structured_not_text(self, tmp_path): + async def test_tool_start_event_is_structured_not_text(self, tmp_path, + monkeypatch): server = make_server(tmp_path) client = FakeClientConnection() client.authenticated = True @@ -428,8 +440,10 @@ class TestBridge: enabled=True, extra={"secret": "test-secret-abc123"})) real._pcfg = server.config real._server = server + real._loop = asyncio.get_running_loop() real._active_sessions = {"agent:main:pheby:dm:conv1": asyncio.Event()} - # Tool-call dedup state is created lazily via _tool_state() + monkeypatch.setattr(real, "_conversation_for_session_id", + lambda _sid: "conv1") from gateway.stream_events import ToolCallChunk marker = real.format_tool_event( @@ -437,21 +451,26 @@ class TestBridge: args={"query": "cats"}, index=0), mode="all") assert marker is None # never rendered as chat text - await asyncio.sleep(0) # let ensure_future broadcast run + # The display event is intentionally eaten. The authoritative hook + # carries the real Hermes session and tool-call IDs. + real.on_pre_tool_call( + session_id="session-1", tool_name="web_search", + tool_call_id="call-1", args={"query": "cats"}) + await asyncio.sleep(0) events = client.ws.events() tool_events = [e for e in events if e["type"] == proto.S_TOOL_EVENT] assert len(tool_events) == 1 ev = tool_events[0] assert ev["tool_name"] == "web_search" assert ev["status"] == "running" - assert ev["tool_call_id"] + assert ev["tool_call_id"] == "call-1" assert ev["conversation_id"] == "conv1" # No fake prose leaked into a message event assert not any(e.get("type") == proto.S_MESSAGE_COMPLETE for e in events) @pytest.mark.asyncio - async def test_post_tool_call_completion(self, tmp_path): + async def test_post_tool_call_completion(self, tmp_path, monkeypatch): server = make_server(tmp_path) client = FakeClientConnection() client.authenticated = True @@ -462,10 +481,13 @@ class TestBridge: real = PhebyAdapter(PlatformConfig( enabled=True, extra={"secret": "test-secret-abc123"})) real._server = server + real._loop = asyncio.get_running_loop() real._active_sessions = {"agent:main:pheby:dm:conv1": asyncio.Event()} + monkeypatch.setattr(real, "_conversation_for_session_id", + lambda _sid: "conv1") real.on_post_tool_call( - tool_name="terminal", tool_call_id="call_9", + session_id="session-1", tool_name="terminal", tool_call_id="call_9", status="ok", duration_ms=1234) await asyncio.sleep(0) # let ensure_future run ev = [e for e in client.ws.events() @@ -475,7 +497,8 @@ class TestBridge: assert ev["duration_ms"] == 1234 @pytest.mark.asyncio - async def test_approval_push_and_resolve_roundtrip(self, tmp_path): + async def test_approval_push_and_resolve_roundtrip(self, tmp_path, + monkeypatch): server = make_server(tmp_path) client = FakeClientConnection() client.authenticated = True @@ -491,11 +514,12 @@ class TestBridge: assert req["choices"] == ["once", "session", "always", "deny"] assert req["description"] == "Destructive command" - # Resolve: with no real Hermes queue the resolve call fails-open to - # accepted=False but the pending entry must be consumed either way. from pheby import hermes_bridge as hb + import tools.approval + monkeypatch.setattr(tools.approval, "resolve_gateway_approval", + lambda *args, **kwargs: 1) ok = await hb.resolve_approval(req["approval_id"], "deny", None) - assert ok is True # pending entry existed; resolution attempted + assert ok is True # Double resolve → not found ok2 = await hb.resolve_approval(req["approval_id"], "once", None) assert ok2 is False @@ -552,6 +576,117 @@ class TestBridge: assert events[0]["type"] == proto.S_RUN_ACCEPTED assert events[1]["type"] == proto.S_MESSAGE_START + @pytest.mark.asyncio + async def test_chat_handler_passes_server_to_bridge(self, tmp_path): + server = make_server(tmp_path) + client = FakeClientConnection() + client.authenticated = True + cid = "a" * 32 + + await server._handle_chat_send(client, { + "type": proto.C_CHAT_SEND, + "conversation_id": cid, + "text": "hello through WebSocket", + }, "request-handler") + + assert server.adapter.handled[-1].text == "hello through WebSocket" + assert server.adapter.handled[-1].source.chat_id == cid + assert client.ws.events()[0]["type"] == proto.S_RUN_ACCEPTED + + @pytest.mark.asyncio + async def test_chat_dispatch_failure_closes_run(self, tmp_path, + monkeypatch): + server = make_server(tmp_path) + client = FakeClientConnection() + client.authenticated = True + server._clients["t"] = client + cid = "b" * 32 + + async def reject(_event): + raise RuntimeError("gateway unavailable") + + monkeypatch.setattr(server.adapter, "handle_message", reject) + from pheby import hermes_bridge as hb + with pytest.raises(RuntimeError, match="gateway unavailable"): + await hb.send_chat(server, cid, "hello", client, "request") + + assert hb.active_run(cid) is None + assert client.ws.events()[-1]["type"] == proto.S_RUN_FINISHED + assert client.ws.events()[-1]["status"] == "failed" + + @pytest.mark.asyncio + async def test_delete_removes_route_instead_of_resetting(self, tmp_path, + monkeypatch): + server = make_server(tmp_path) + cid = await server.router.new_conversation("Delete me") + session_key = f"agent:main:pheby:dm:{cid}" + + class Entry: + display_name = "Delete me" + + class Store: + def __init__(self): + self._lock = threading.Lock() + self._entries = {session_key: Entry()} + self.saved = False + + def _ensure_loaded_locked(self): + return None + + def _save(self): + self.saved = True + + def peek_session_id(self, key): + return "session-delete" if key == session_key else None + + class DB: + deleted = None + + def delete_session(self, session_id): + self.deleted = session_id + return True + + store, db = Store(), DB() + from pheby import hermes_bridge as hb + monkeypatch.setattr(hb, "_session_store", lambda: store) + monkeypatch.setattr(hb, "_session_db", lambda: db) + + assert await hb.delete_conversation(cid) is True + assert db.deleted == "session-delete" + assert session_key not in store._entries + assert store.saved is True + assert await server.router.get_name(cid) is None + + @pytest.mark.asyncio + async def test_models_snapshot_uses_current_hermes_signature( + self, tmp_path, monkeypatch): + make_server(tmp_path) + captured = {} + + def fake_list_picker_providers(**kwargs): + captured.update(kwargs) + return [{"slug": "test", "models": ["m1"]}] + + from pheby import hermes_bridge as hb + import hermes_cli.config as hermes_config + import hermes_cli.model_switch as model_switch + monkeypatch.setattr(hb, "_load_cfg", lambda: { + "model": {"default": "m1", "provider": "test"}, + "providers": {"test": {"base_url": "http://example"}}, + "model_catalog": {"excluded_providers": ["hidden"]}, + }) + monkeypatch.setattr( + hermes_config, "get_compatible_custom_providers", + lambda _cfg: [{"name": "test", "base_url": "http://example"}]) + monkeypatch.setattr( + model_switch, "list_picker_providers", fake_list_picker_providers) + + snapshot = await hb.models_snapshot() + assert snapshot["providers"][0]["slug"] == "test" + assert captured["current_model"] == "m1" + assert captured["excluded_providers"] == ["hidden"] + assert captured["custom_providers"][0]["name"] == "test" + @pytest.mark.asyncio async def test_cancel_run_interrupts_agent(self, tmp_path): server = make_server(tmp_path) @@ -565,6 +700,7 @@ class TestBridge: server.adapter.gateway_runner = runner from pheby import hermes_bridge as hb + hb._ACTIVE_RUNS["convX"] = {"run_id": "run-x", "started": 0} ok = await hb.cancel_run("convX", None) assert ok is True assert agent.interrupts # agent.interrupt called, not thread-kill @@ -612,6 +748,47 @@ class TestRecovery: assert h1["messages"] == h2["messages"] assert h1["conversation_id"] == h2["conversation_id"] == cid + @pytest.mark.asyncio + async def test_open_recovers_runtime_and_attachment_state(self, tmp_path): + server = make_server(tmp_path) + cid = await server.router.new_conversation("Recover") + source = Path(tmp_path) / "recovery.png" + source.write_bytes(b"\x89PNG recovery") + attachment = await server.store.register_file( + str(source), conversation_id=cid) + + from pheby import hermes_bridge as hb + hb._ACTIVE_RUNS[cid] = {"run_id": "run-recover", "started": 1} + hb.record_tool_event(cid, { + "type": proto.S_TOOL_EVENT, + "conversation_id": cid, + "run_id": "run-recover", + "tool_call_id": "tool-recover", + "tool_name": "terminal", + "status": "running", + }) + await hb.push_approval( + {"command": "echo hi", "description": "Run command"}, + f"agent:main:pheby:dm:{cid}") + await hb.push_clarify( + "clarify-recover", f"agent:main:pheby:dm:{cid}", + "Continue?", ["yes", "no"]) + + client = FakeClientConnection() + client.authenticated = True + await server._handle_conversation_open(client, { + "type": proto.C_CONVERSATION_OPEN, + "conversation_id": cid, + }, "recover") + snapshot = client.ws.events()[-1] + assert snapshot["attachments"][0]["attachment_id"] == \ + attachment["attachment_id"] + assert snapshot["run"]["run_id"] == "run-recover" + assert snapshot["tools"][0]["tool_call_id"] == "tool-recover" + assert snapshot["approvals"][0]["conversation_id"] == cid + assert snapshot["clarifications"][0]["clarify_id"] == \ + "clarify-recover" + @pytest.mark.asyncio async def test_broadcast_reaches_multiple_clients(self, tmp_path): server = make_server(tmp_path) @@ -756,6 +933,45 @@ class TestLiveServer: # Adapter unit checks # ═══════════════════════════════════════════════════════════════════════════ class TestAdapterUnits: + @pytest.mark.asyncio + async def test_stream_preview_keeps_run_open_until_finalize(self, tmp_path): + from pheby.adapter import PhebyAdapter + from gateway.config import PlatformConfig + from pheby import hermes_bridge as hb + + server = make_server(tmp_path) + client = FakeClientConnection() + client.authenticated = True + server._clients["t"] = client + adapter = PhebyAdapter(PlatformConfig( + enabled=True, extra={"secret": "test-secret-abc123"})) + adapter._server = server + adapter._loop = asyncio.get_running_loop() + adapter._drafts["conv"] = {"message_id": "draft-run", "text": ""} + hb.set_adapter(adapter) + hb.set_server(server) + hb._ACTIVE_RUNS["conv"] = {"run_id": "run-1", "started": 0} + + first = await adapter.send("conv", "hel", metadata={"expect_edits": True}) + assert first.message_id == "draft-run" + assert hb.active_run_id("conv") == "run-1" + assert client.ws.events()[-1]["type"] == proto.S_MESSAGE_DELTA + + final = await adapter.edit_message( + "conv", "draft-run", "hello", finalize=True) + await asyncio.sleep(0) + assert final.message_id == "draft-run" + assert hb.active_run_id("conv") is None + assert [e["type"] for e in client.ws.events()][-2:] == [ + proto.S_MESSAGE_COMPLETE, proto.S_RUN_FINISHED] + + def test_transport_auth_is_gateway_authorization(self): + from pheby.adapter import PhebyAdapter + from gateway.config import PlatformConfig + adapter = PhebyAdapter(PlatformConfig( + enabled=True, extra={"secret": "test-secret-abc123"})) + assert adapter.authorization_is_upstream is True + def test_redact_args(self): from pheby.adapter import _redact_args out = _redact_args({"query": "cats", "api_key": "sk-123", @@ -765,6 +981,16 @@ class TestAdapterUnits: assert out["query"] == "cats" assert out["long"].endswith("…") + def test_redact_args_recursively(self): + from pheby.adapter import _redact_args + out = _redact_args({ + "headers": {"Authorization": "Bearer secret"}, + "steps": [{"password": "hunter2", "value": "safe"}], + }) + assert out["headers"]["Authorization"] == "[redacted]" + assert out["steps"][0]["password"] == "[redacted]" + assert out["steps"][0]["value"] == "safe" + def test_sanitize_filename(self): from pheby.attachments import AttachmentStore assert AttachmentStore._sanitize_filename("../../etc/passwd") == "passwd"