6be1601ccd
Deliverables now inherit the active assistant draft message ID so clients can render them beside the reply that produced them.
529 lines
22 KiB
Python
529 lines
22 KiB
Python
"""Pheby platform adapter — the Hermes gateway ↔ Pheby protocol bridge.
|
|
|
|
Extends ``BasePlatformAdapter`` like every other platform (Telegram,
|
|
Discord, ntfy, …) so the full gateway pipeline — sessions, tool approval,
|
|
clarify, deliverables, streaming — works unchanged on the Pheby platform.
|
|
|
|
Outbound mapping:
|
|
* ``send`` / ``edit_message`` → chat draft events (S_MESSAGE_DELTA etc.)
|
|
* 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
|
|
|
|
Inbound mapping (reversed): Pheby WS messages are turned into MessageEvents
|
|
delivered through ``handle_message`` so the gateway treats them identically
|
|
to any other platform's messages.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import importlib.util
|
|
import logging
|
|
from typing import Any, Dict, Optional
|
|
|
|
AIOHTTP_AVAILABLE = importlib.util.find_spec("aiohttp") is not None
|
|
|
|
from gateway.config import Platform, PlatformConfig
|
|
from gateway.platforms.base import (
|
|
BasePlatformAdapter,
|
|
SendResult,
|
|
)
|
|
|
|
from . import protocol as proto
|
|
from . import hermes_bridge
|
|
from .config import PhebyConfig, load_config
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
class PhebyAdapter(BasePlatformAdapter):
|
|
"""Serve the Pheby WebSocket/HTTP protocol and map it onto the gateway."""
|
|
|
|
# Async tools (background terminal tasks, delegate_task) may wake a later
|
|
# turn on this platform — the WS is a persistent push channel.
|
|
supports_async_delivery: bool = True
|
|
|
|
# Pheby clients re-render every delta; we accumulate full text and the
|
|
# client truncates nothing — no platform length limit.
|
|
MAX_MESSAGE_LENGTH = 0
|
|
|
|
def __init__(self, config: PlatformConfig):
|
|
# Platform("pheby") resolves via the enum's _missing_() hook once the
|
|
# plugin registry knows the name; in bare unit tests (registry not
|
|
# populated) fall back to a synthetic enum member so the adapter can
|
|
# still be constructed and tested.
|
|
try:
|
|
platform = Platform("pheby")
|
|
except ValueError:
|
|
platform = object.__new__(Platform)
|
|
platform._value_ = "pheby"
|
|
platform._name_ = "PHEBY"
|
|
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
|
|
|
|
@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:
|
|
if not AIOHTTP_AVAILABLE:
|
|
logger.warning("[pheby] aiohttp not installed — cannot serve")
|
|
return False
|
|
if not self._pcfg.enabled:
|
|
self._set_fatal_error(
|
|
"pheby_no_secret",
|
|
"PHEBY_SECRET is not set — refusing to start the Pheby "
|
|
"server without a credential. Set it in ~/.hermes/.env.",
|
|
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)",
|
|
proto.PROTOCOL_VERSION)
|
|
return True
|
|
|
|
async def disconnect(self) -> None:
|
|
self._running = False
|
|
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")
|
|
|
|
# ── outbound: chat text ──────────────────────────────────────────────
|
|
def _conv_from_chat_id(self, chat_id: str) -> str:
|
|
return str(chat_id)
|
|
|
|
async def send(
|
|
self,
|
|
chat_id: str,
|
|
content: str,
|
|
reply_to: Optional[str] = None,
|
|
metadata: Optional[Dict[str, Any]] = None,
|
|
**kwargs,
|
|
) -> SendResult:
|
|
"""Deliver assistant text (final response, commentary, or notices).
|
|
|
|
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)
|
|
metadata = metadata or {}
|
|
|
|
# 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,
|
|
"run_id": hermes_bridge.active_run_id(conversation_id),
|
|
"message_id": message_id,
|
|
"text": content,
|
|
"ts": proto.now_iso(),
|
|
}
|
|
await self._server.broadcast(event)
|
|
hermes_bridge.note_run_finished(conversation_id, "completed")
|
|
return SendResult(success=True, message_id=message_id)
|
|
|
|
async def edit_message(
|
|
self,
|
|
chat_id: str,
|
|
message_id: str,
|
|
content: str,
|
|
finalize: bool = False,
|
|
metadata: Optional[Dict[str, Any]] = None,
|
|
**kwargs,
|
|
) -> SendResult:
|
|
"""Streaming path: GatewayStreamConsumer edits the in-place draft.
|
|
|
|
The stream-consumer contract requires concrete adapters to accept
|
|
``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")
|
|
conversation_id = self._conv_from_chat_id(chat_id)
|
|
draft = self._drafts.setdefault(conversation_id, {
|
|
"message_id": message_id or f"draft-{proto.new_id()[:12]}",
|
|
"text": "",
|
|
})
|
|
draft["text"] = content # consumer sends cumulative text
|
|
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 ─────────────────────────────────────────
|
|
def format_tool_event(self, event: Any, *, mode: str = "all",
|
|
preview_max_len: int = 40) -> Optional[str]:
|
|
"""Emit tool activity as structured JSON — never as fake chat text.
|
|
|
|
Returning a truthy marker would put prose in chat; instead we push an
|
|
S_TOOL_EVENT broadcast and return None so the gateway's text queue
|
|
stays clean. (The dispatcher treats None as "adapter ate the event".)
|
|
"""
|
|
# 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 ``on_pre_tool_call``.
|
|
"""
|
|
try:
|
|
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")
|
|
status = str(kwargs.get("status") or "")
|
|
duration_ms = kwargs.get("duration_ms") or 0
|
|
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 f"t-{proto.new_id()[:12]}"),
|
|
"tool_name": tool_name,
|
|
"status": "completed" if status in ("ok", "success", "")
|
|
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
|
|
# gets outcome status; verbose content stays in Hermes.
|
|
"ts": proto.now_iso(),
|
|
}
|
|
error_message = kwargs.get("error_message")
|
|
if error_message and event["status"] == "failed":
|
|
event["error"] = proto.safe_str(error_message, 200)
|
|
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)
|
|
|
|
# ── 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_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({
|
|
"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)
|
|
return SendResult(success=False, error=str(exc))
|
|
|
|
# ── clarification ────────────────────────────────────────────────────
|
|
async def send_clarify(
|
|
self,
|
|
chat_id: str,
|
|
question: str,
|
|
choices: Optional[list],
|
|
clarify_id: str,
|
|
session_key: str,
|
|
metadata: Optional[Dict[str, Any]] = None,
|
|
) -> SendResult:
|
|
"""Native structured clarify prompt (buttons on the client)."""
|
|
if self._server is None:
|
|
return SendResult(success=False, error="server not running")
|
|
await hermes_bridge.push_clarify(clarify_id, session_key, question,
|
|
choices)
|
|
# Text capture is unnecessary: the client responds through
|
|
# clarify.respond, which resolves the entry directly.
|
|
return SendResult(success=True, message_id=clarify_id)
|
|
|
|
# ── deliverables (attachments) ───────────────────────────────────────
|
|
def _active_assistant_message_id(self, conversation_id: str) -> Optional[str]:
|
|
"""Return the live assistant draft that a deliverable belongs beneath."""
|
|
draft = self._drafts.get(conversation_id)
|
|
message_id = draft.get("message_id") if isinstance(draft, dict) else None
|
|
return str(message_id) if message_id else None
|
|
|
|
async def _register_and_broadcast(
|
|
self,
|
|
file_path: str,
|
|
conversation_id: str,
|
|
*,
|
|
message_id: Optional[str] = None,
|
|
kind_hint: Optional[str] = None,
|
|
filename: Optional[str] = None,
|
|
) -> Optional[Dict[str, Any]]:
|
|
if self._server is None:
|
|
return None
|
|
desc = await self._server.store.register_file(
|
|
file_path,
|
|
conversation_id=conversation_id,
|
|
message_id=message_id or self._active_assistant_message_id(conversation_id),
|
|
filename=filename,
|
|
kind_hint=kind_hint,
|
|
)
|
|
if desc is None:
|
|
return None
|
|
await self._server.broadcast({
|
|
"type": proto.S_ATTACHMENT_ADDED,
|
|
"conversation_id": conversation_id,
|
|
"attachment": desc,
|
|
"ts": proto.now_iso(),
|
|
})
|
|
return desc
|
|
|
|
async def send_document(
|
|
self,
|
|
chat_id: str,
|
|
file_path: str,
|
|
caption: Optional[str] = None,
|
|
file_name: Optional[str] = None,
|
|
reply_to: Optional[str] = None,
|
|
metadata: Optional[Dict[str, Any]] = None,
|
|
**kwargs,
|
|
) -> SendResult:
|
|
conversation_id = self._conv_from_chat_id(chat_id)
|
|
desc = await self._register_and_broadcast(
|
|
file_path, conversation_id, filename=file_name)
|
|
if desc is None:
|
|
return SendResult(success=False, error="attachment failed")
|
|
if caption:
|
|
await self.send(chat_id, caption)
|
|
return SendResult(success=True, message_id=desc["attachment_id"])
|
|
|
|
async def send_image_file(
|
|
self,
|
|
chat_id: str,
|
|
image_path: str,
|
|
caption: Optional[str] = None,
|
|
reply_to: Optional[str] = None,
|
|
metadata: Optional[Dict[str, Any]] = None,
|
|
**kwargs,
|
|
) -> SendResult:
|
|
conversation_id = self._conv_from_chat_id(chat_id)
|
|
desc = await self._register_and_broadcast(
|
|
image_path, conversation_id, kind_hint="image")
|
|
if desc is None:
|
|
return SendResult(success=False, error="attachment failed")
|
|
if caption:
|
|
await self.send(chat_id, caption)
|
|
return SendResult(success=True, message_id=desc["attachment_id"])
|
|
|
|
async def send_voice(
|
|
self,
|
|
chat_id: str,
|
|
audio_path: str,
|
|
caption: Optional[str] = None,
|
|
reply_to: Optional[str] = None,
|
|
metadata: Optional[Dict[str, Any]] = None,
|
|
**kwargs,
|
|
) -> SendResult:
|
|
conversation_id = self._conv_from_chat_id(chat_id)
|
|
desc = await self._register_and_broadcast(
|
|
audio_path, conversation_id, kind_hint="voice")
|
|
if desc is None:
|
|
return SendResult(success=False, error="attachment failed")
|
|
return SendResult(success=True, message_id=desc["attachment_id"])
|
|
|
|
async def send_video(
|
|
self,
|
|
chat_id: str,
|
|
video_path: str,
|
|
caption: Optional[str] = None,
|
|
reply_to: Optional[str] = None,
|
|
metadata: Optional[Dict[str, Any]] = None,
|
|
**kwargs,
|
|
) -> SendResult:
|
|
conversation_id = self._conv_from_chat_id(chat_id)
|
|
desc = await self._register_and_broadcast(
|
|
video_path, conversation_id, kind_hint="video")
|
|
if desc is None:
|
|
return SendResult(success=False, error="attachment failed")
|
|
return SendResult(success=True, message_id=desc["attachment_id"])
|
|
|
|
# ── misc contract ────────────────────────────────────────────────────
|
|
async def get_chat_info(self, chat_id: str) -> Dict[str, Any]:
|
|
return {"name": str(chat_id), "type": "dm"}
|
|
|
|
# Standalone cron/send_message delivery (out-of-process).
|
|
def standalone_send(self):
|
|
async def _send(pconfig, chat_id: str, message: str, **kwargs):
|
|
# Out-of-process there is no WS server; deliver via a transient
|
|
# connection to our own HTTP/WS endpoint is overkill — instead
|
|
# cron jobs targeting Pheby run in-process with the gateway.
|
|
return {"error": "pheby standalone delivery requires the gateway "
|
|
"(deliver within gateway-managed processes)"}
|
|
return _send
|
|
|
|
|
|
def _redact_args(args: Optional[Dict[str, Any]]) -> Optional[Dict[str, Any]]:
|
|
"""Recursively strip likely-secret values before sending args to client."""
|
|
if not isinstance(args, dict):
|
|
return None
|
|
sensitive = ("key", "token", "secret", "password", "credential", "auth")
|
|
|
|
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"]
|