Fix Hermes adapter integration and recovery
This commit is contained in:
+183
-147
@@ -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:<conversation_id>
|
||||
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"]
|
||||
|
||||
Reference in New Issue
Block a user