feat: Pheby platform adapter/plugin for Hermes (protocol v1)

Third-party Hermes platform plugin serving the Pheby WebSocket+HTTPS
protocol for the native Android client, behind Caddy:

- Platform adapter (BasePlatformAdapter subclass) registered via the
  documented plugin system (plugins.platforms pattern, kind: platform)
- Multiple conversations mapped to Hermes sessions (authoritative state)
- Streaming chat via GatewayStreamConsumer edit path (message.delta)
- Structured tool events (never fake tool text in chat), incl.
  post_tool_call hook relay with Hermes tool_call_ids
- Native approval + clarification round-trips via tools.approval /
  tools.clarify_gateway primitives
- Attachment delivery: adapter-owned copies, opaque IDs, persistent
  metadata, 7-day retention + safe cleanup, authenticated HTTPS download
- Model + reasoning-effort query/change via Hermes picker data and
  session/global overrides
- Shared-secret auth (constant-time, lockout), size limits, path-
  traversal-proof attachment resolution, minimal health endpoint
- Protocol v1 spec with JSON examples (docs/PROTOCOL.md)
- Install/config/Caddy/security docs + discovered Hermes limitations
- 37 passing tests (auth, conversations, protocol, attachments incl.
  expiry/traversal, tool events, approvals, clarifications, cancel,
  reconnect re-sync, live WS smoke tests) — gateway-free fakes

No Hermes core modifications required.
This commit is contained in:
2026-09-02 19:31:17 +00:00
parent 31c16de665
commit 53ffdb8aa2
19 changed files with 4688 additions and 0 deletions
+81
View File
@@ -0,0 +1,81 @@
"""Pheby — Hermes Agent platform adapter/plugin for the Pheby Android client.
A third-party Hermes platform plugin. Install this package directory as
``~/.hermes/plugins/pheby/`` (HERMES_HOME/plugins/pheby), enable it in
config.yaml (``plugins.enabled: [pheby]``, ``platforms.pheby.enabled: true``),
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"
def register(ctx: Any) -> None:
"""Plugin entry point — called by the Hermes plugin system at startup."""
from .adapter import PhebyAdapter
from .config import (
check_requirements,
env_enablement,
is_connected,
validate_config,
)
adapter_holder: dict = {"adapter": None}
def _factory(cfg: Any) -> PhebyAdapter:
adapter = PhebyAdapter(cfg)
adapter_holder["adapter"] = adapter
return adapter
ctx.register_platform(
name="pheby",
label="Pheby",
adapter_factory=_factory,
check_fn=check_requirements,
validate_config=validate_config,
is_connected=is_connected,
required_env=["PHEBY_SECRET"],
install_hint="pip install aiohttp # already a Hermes dependency; "
"set PHEBY_SECRET in ~/.hermes/.env",
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,
platform_hint=(
"You are communicating with the user via Pheby, a private "
"native Android client over WebSocket. Respond in normal "
"markdown; the client renders it natively. Attachments you "
"produce via MEDIA: tags are delivered as downloadable files "
"and inline image previews."
),
)
# 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).
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("post_tool_call", _post_tool_call)
except Exception:
logger.debug("[pheby] post_tool_call hook registration failed",
exc_info=True)
logger.info("[pheby] plugin registered (platform 'pheby')")
__all__ = ["register", "__version__"]
+485
View File
@@ -0,0 +1,485 @@
"""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.)
* ``format_tool_event`` → 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 logging
import mimetypes
import os
import time
from pathlib import Path
from typing import Any, Dict, List, Optional, Tuple
try:
import aiohttp
from aiohttp import web
AIOHTTP_AVAILABLE = True
except ImportError: # pragma: no cover
AIOHTTP_AVAILABLE = False
from gateway.config import Platform, PlatformConfig
from gateway.platforms.base import (
BasePlatformAdapter,
MessageEvent,
MessageType,
SendResult,
)
from . import protocol as proto
from . import hermes_bridge
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."""
# 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._drafts: Dict[str, Dict[str, Any]] = {} # conv → draft state
self._typing: Dict[str, float] = {}
# ── 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
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:
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
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).
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``.
"""
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]}"
# A draft exists while the turn streams; the final text supersedes
# the draft and closes it out.
draft = self._drafts.pop(conversation_id, None)
event = {
"type": proto.S_MESSAGE_COMPLETE,
"conversation_id": conversation_id,
"message_id": (draft or {}).get("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)
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=`` even when ignored (it's False during progressive
edits; the final content always arrives via ``send()``).
"""
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
await self._server.broadcast({
"type": proto.S_MESSAGE_DELTA,
"conversation_id": conversation_id,
"message_id": draft["message_id"],
"text": content,
"ts": proto.now_iso(),
})
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".)
"""
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]}")
# -- Hermes plugin hooks (registered in __init__.py register()) --------
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``.
"""
try:
conversation_id = self._active_conversation_id()
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,
"tool_call_id": str(kwargs.get("tool_call_id")
or self._lookup_tool(tool_name,
conversation_id)),
"tool_name": tool_name,
"status": "completed" if status in ("ok", "success", "")
else "failed" if status == "error" 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)
asyncio.ensure_future(self._server.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)."""
try:
await hermes_bridge.push_approval(approval_data, session_key)
except Exception:
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)
# ── 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) ───────────────────────────────────────
async def _register_and_broadcast(
self,
file_path: str,
conversation_id: str,
*,
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=None,
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
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."""
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
__all__ = ["PhebyAdapter", "AIOHTTP_AVAILABLE"]
+357
View File
@@ -0,0 +1,357 @@
"""Pheby attachment store — adapter-owned copies of Hermes deliverables.
Design (spec: "Generated attachments / Deliverable Mode"):
* The agent produces files via Hermes's normal deliverable pipeline. The
adapter intercepts ``send_document`` / ``send_image_file`` / ``send_voice``
/ ``send_video`` and *copies* the source file into adapter-owned storage,
registering an opaque 32-hex attachment ID.
* The client only ever sees attachment IDs — never server paths. Downloads
resolve ID → registered file inside the storage root; path traversal and
arbitrary filesystem reads are impossible by construction.
* Metadata (JSON, one file per attachment) persists across restarts so valid
attachments survive a Hermes/Pheby restart.
* Cleanup deletes only files this store owns (inside its own storage root,
matched by registered IDs) after the retention window (default 7 days).
Hermes-owned originals elsewhere on disk are never touched.
"""
from __future__ import annotations
import asyncio
import hashlib
import hmac
import json
import logging
import mimetypes
import os
import shutil
import time
from pathlib import Path
from typing import Any, Dict, List, Optional
from . import protocol
logger = logging.getLogger(__name__)
IMAGE_MIME_PREFIXES = ("image/",)
IMAGE_EXTS = {".png", ".jpg", ".jpeg", ".gif", ".webp", ".bmp", ".heic"}
# Extensions Hermes's deliverable system treats as audio/video (used to pick
# a sensible kind for inline preview decisions).
AUDIO_EXTS = {".ogg", ".opus", ".mp3", ".wav", ".m4a", ".flac"}
VIDEO_EXTS = {".mp4", ".mov", ".avi", ".mkv", ".webm"}
def guess_mime(filename: str, fallback: str = "application/octet-stream") -> str:
"""Best-effort MIME type for a filename (stdlib mimetypes + extras)."""
ext = Path(filename).suffix.lower()
explicit = {
".md": "text/markdown", ".yml": "application/yaml",
".yaml": "application/yaml", ".toml": "application/toml",
".log": "text/plain", ".apk": "application/vnd.android.package-archive",
".ogg": "audio/ogg", ".opus": "audio/opus",
}
if ext in explicit:
return explicit[ext]
guessed, _ = mimetypes.guess_type(filename)
return guessed or fallback
class AttachmentStore:
"""Owns adapter-managed attachment copies and their metadata."""
def __init__(self, root: Path, retention_days: int = 7,
index_path: Optional[Path] = None):
self._root = Path(root).resolve()
self._retention_days = max(0, int(retention_days))
self._index_path = (
Path(index_path) if index_path else self._root / "attachments.json"
)
self._lock = asyncio.Lock()
self._meta: Dict[str, Dict[str, Any]] = {}
self._loaded = False
# ── paths ────────────────────────────────────────────────────────────
@property
def root(self) -> Path:
return self._root
def _blob_path(self, attachment_id: str, filename: str) -> Path:
"""Blob location: ``blobs/<aa>/<id>__<sanitized-filename>``."""
safe_name = self._sanitize_filename(filename)
return self._root / "blobs" / attachment_id[:2] / f"{attachment_id}__{safe_name}"
def _meta_path(self, attachment_id: str) -> Path:
return self._root / "meta" / f"{attachment_id}.json"
@staticmethod
def _sanitize_filename(filename: str) -> str:
"""Strip path separators/control chars from a stored filename."""
name = os.path.basename(str(filename or "").replace("\\", "/")).strip()
name = "".join(c for c in name if c.isprintable() and c not in '/\\')
return name[:120] or "file.bin"
# ── persistence ──────────────────────────────────────────────────────
def _load_index(self) -> None:
if self._loaded:
return
self._loaded = True
try:
if self._index_path.exists():
data = json.loads(self._index_path.read_text(encoding="utf-8"))
if isinstance(data, dict):
self._meta = {
k: v for k, v in data.items()
if isinstance(k, str) and isinstance(v, dict)
and protocol.is_valid_attachment_id(k)
}
except Exception:
logger.warning("[pheby] attachment index unreadable; starting empty",
exc_info=True)
def _save_index(self) -> None:
try:
self._index_path.parent.mkdir(parents=True, exist_ok=True)
tmp = self._index_path.with_suffix(".tmp")
tmp.write_text(
json.dumps(self._meta, ensure_ascii=False, indent=1),
encoding="utf-8")
os.replace(tmp, self._index_path)
except Exception:
logger.error("[pheby] failed to persist attachment index",
exc_info=True)
# ── registration ─────────────────────────────────────────────────────
async def register_file(
self,
source_path: str,
*,
conversation_id: str,
message_id: Optional[str] = None,
filename: Optional[str] = None,
kind_hint: Optional[str] = None,
) -> Optional[Dict[str, Any]]:
"""Copy *source_path* into adapter storage and register metadata.
Returns the attachment descriptor dict, or ``None`` when the source
is missing/unsafe. The original file is never modified or deleted.
"""
try:
src = Path(source_path).expanduser().resolve(strict=True)
except (OSError, RuntimeError, ValueError):
logger.warning("[pheby] deliverable not found: %s",
protocol.safe_str(source_path, 120))
return None
if not src.is_file():
return None
fname = self._sanitize_filename(filename or src.name)
attachment_id = protocol.new_id()
mime = guess_mime(fname)
is_image = mime.startswith(IMAGE_MIME_PREFIXES) or (
Path(fname).suffix.lower() in IMAGE_EXTS)
if kind_hint == "voice":
kind = "voice"
elif kind_hint == "video" or Path(fname).suffix.lower() in VIDEO_EXTS:
kind = "video"
elif kind_hint == "audio" or Path(fname).suffix.lower() in AUDIO_EXTS:
kind = "audio"
elif is_image:
kind = "image"
else:
kind = "document"
try:
size = src.stat().st_size
async with self._lock:
self._load_index()
blob = self._blob_path(attachment_id, fname)
blob.parent.mkdir(parents=True, exist_ok=True)
# Copy under the lock so cleanup can never race a half-written
# blob (cleanup only deletes registered+expired entries).
await asyncio.to_thread(shutil.copy2, str(src), str(blob))
meta: Dict[str, Any] = {
"attachment_id": attachment_id,
"filename": fname,
"mime_type": mime,
"size": size,
"kind": kind,
"conversation_id": conversation_id,
"message_id": message_id,
"created_at": protocol.now_iso(),
"created_epoch": time.time(),
"retention_days": self._retention_days,
"blob": blob.name,
"blob_subdir": blob.parent.name,
}
self._meta[attachment_id] = meta
self._save_index()
except Exception:
logger.error("[pheby] attachment registration failed for %s",
protocol.safe_str(source_path, 120), exc_info=True)
return None
logger.info(
"[pheby] attachment registered: id=%s kind=%s size=%d conv=%s",
attachment_id, kind, size, conversation_id)
return self.describe(attachment_id)
# ── lookup / download ────────────────────────────────────────────────
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:
return None
return {
"attachment_id": attachment_id,
"filename": meta.get("filename", "file.bin"),
"mime_type": meta.get("mime_type", "application/octet-stream"),
"size": int(meta.get("size", 0)),
"kind": meta.get("kind", "document"),
"inline_image": meta.get("kind") == "image",
"conversation_id": meta.get("conversation_id"),
"message_id": meta.get("message_id"),
"created_at": meta.get("created_at"),
"expires_at": self._expires_at_iso(meta),
"download_path": f"/attachments/{attachment_id}",
}
def _expires_at_iso(self, meta: Dict[str, Any]) -> Optional[str]:
retention = int(meta.get("retention_days", self._retention_days))
if retention <= 0:
return None
created = float(meta.get("created_epoch", 0) or 0)
if not created:
return None
import datetime as _dt
return _dt.datetime.fromtimestamp(
created + retention * 86400, tz=_dt.timezone.utc).isoformat()
def resolve_blob(self, attachment_id: str) -> Optional[Path]:
"""Resolve an ID to its blob path — only for registered IDs.
Returns ``None`` for unknown, expired, or malformed IDs. The
returned path is always inside the storage root (the blob filename
comes from sanitized metadata, never client input).
"""
if not protocol.is_valid_attachment_id(attachment_id):
return None
meta = self._meta.get(attachment_id)
if not meta:
return None
if self._is_expired(meta):
return None
blob = (self._root / "blobs" / str(meta.get("blob_subdir", "")) /
str(meta.get("blob", "")))
try:
resolved = blob.resolve(strict=True)
except (OSError, RuntimeError, ValueError):
return None
# Defense in depth: blob must live inside our storage root.
try:
resolved.relative_to(self._root)
except ValueError:
return None
if not resolved.is_file():
return None
return resolved
def list_for_conversation(self, conversation_id: str) -> List[Dict[str, Any]]:
out = []
for aid in list(self._meta):
desc = self.describe(aid)
if desc and desc.get("conversation_id") == conversation_id:
out.append(desc)
return out
# ── cleanup ──────────────────────────────────────────────────────────
def _is_expired(self, meta: Dict[str, Any]) -> bool:
retention = int(meta.get("retention_days", self._retention_days))
if retention <= 0:
return False
created = float(meta.get("created_epoch", 0) or 0)
return created > 0 and (time.time() - created) > retention * 86400
async def cleanup_expired(self) -> int:
"""Delete expired adapter-owned blobs + metadata. Returns count.
Only deletes blobs this store registered (inside its own root, keyed
by ID). Never touches anything outside the storage root.
"""
async with self._lock:
self._load_index()
expired = [aid for aid, m in self._meta.items() if self._is_expired(m)]
removed = 0
for aid in expired:
meta = self._meta.pop(aid, None)
if not meta:
continue
blob = (self._root / "blobs" / str(meta.get("blob_subdir", "")) /
str(meta.get("blob", "")))
try:
resolved = blob.resolve(strict=False)
resolved.relative_to(self._root) # containment check
if resolved.is_file():
await asyncio.to_thread(resolved.unlink)
removed += 1
except (OSError, RuntimeError, ValueError):
logger.warning("[pheby] cleanup skipped blob for %s", aid)
try:
mp = self._meta_path(aid)
if mp.exists():
await asyncio.to_thread(mp.unlink)
except OSError:
pass
if expired:
self._save_index()
logger.info("[pheby] cleaned %d expired attachment(s)", removed)
return removed
# ── legacy per-id meta files (restart durability helper) ─────────────
def hydrate_legacy_meta(self) -> None:
"""Read any per-ID ``meta/*.json`` files from older versions."""
self._load_index()
try:
meta_dir = self._root / "meta"
if not meta_dir.is_dir():
return
for mp in meta_dir.glob("*.json"):
aid = mp.stem
if aid in self._meta or not protocol.is_valid_attachment_id(aid):
continue
try:
data = json.loads(mp.read_text(encoding="utf-8"))
if isinstance(data, dict):
self._meta[aid] = data
except Exception:
continue
except Exception:
logger.debug("[pheby] legacy meta hydration skipped", exc_info=True)
def wipe_all(self) -> None:
"""Test helper: remove everything this store owns."""
self._meta = {}
self._loaded = True
if self._root.exists():
shutil.rmtree(self._root, ignore_errors=True)
def constant_time_equals(a: str, b: str) -> bool:
"""Length-safe constant-time string comparison for secrets."""
a_b = a.encode("utf-8")
b_b = b.encode("utf-8")
return len(a_b) == len(b_b) and hmac.compare_digest(a_b, b_b)
def hash_secret_for_log(secret: str) -> str:
"""Short non-reversible fingerprint for log lines (never the secret)."""
if not secret:
return "unset"
return hashlib.sha256(secret.encode("utf-8")).hexdigest()[:8]
__all__ = [
"AttachmentStore", "constant_time_equals", "hash_secret_for_log",
"guess_mime", "IMAGE_EXTS", "AUDIO_EXTS", "VIDEO_EXTS",
]
+179
View File
@@ -0,0 +1,179 @@
"""Pheby plugin configuration.
Secrets come from environment variables (Hermes convention: ``~/.hermes/.env``
is loaded by Hermes itself before plugins load). Non-secret behavior lives in
the platform's ``extra`` block in ``config.yaml`` under ``platforms.pheby``.
Resolution precedence for every key (highest wins):
1. environment variable (secrets *must* come from env)
2. ``platforms.pheby.extra.<key>`` in config.yaml
3. built-in default
"""
from __future__ import annotations
import os
from dataclasses import dataclass, field
from typing import Any, Dict, Optional
# Environment variable names
ENV_SECRET = "PHEBY_SECRET" # shared credential (required to serve)
ENV_BIND_HOST = "PHEBY_BIND_HOST"
ENV_PORT = "PHEBY_PORT"
ENV_DEBUG = "PHEBY_DEBUG"
ENV_LOG_CHAT_CONTENT = "PHEBY_LOG_CHAT_CONTENT"
# Defaults
DEFAULT_BIND_HOST = "127.0.0.1" # safe behind a local Caddy reverse proxy
DEFAULT_PORT = 8620
DEFAULT_RETENTION_DAYS = 7
DEFAULT_MAX_SECRET_LEN = 1024
@dataclass
class PhebyConfig:
"""Resolved runtime configuration for the Pheby server + adapter."""
# Shared secret for WS/HTTP auth. Empty disables the adapter entirely.
secret: str = ""
bind_host: str = DEFAULT_BIND_HOST
port: int = DEFAULT_PORT
# Attachment storage directory; a per-instance subdir is created inside.
storage_dir: str = ""
# Attachment retention in days (0 = keep forever — not recommended).
retention_days: int = DEFAULT_RETENTION_DAYS
# Verbose protocol logging (still never logs secrets).
debug: bool = False
# When True, debug logs MAY include chat text and tool previews.
log_chat_content: bool = False
# Path to a persistent JSON index of registered attachments. Empty =
# derive from storage_dir.
index_path: str = ""
# Extra config passthrough (whole ``extra`` dict) for future keys.
extra: Dict[str, Any] = field(default_factory=dict)
# ------------------------------------------------------------------
@property
def enabled(self) -> bool:
"""The adapter only serves when a secret is configured."""
return bool(self.secret)
@property
def attachments_root(self) -> str:
if self.storage_dir:
return self.storage_dir
# Resolved lazily by the attachment store (needs get_hermes_home).
return ""
def _env_bool(name: str) -> bool:
return os.getenv(name, "").strip().lower() in ("1", "true", "yes", "on")
def _coerce_int(value: Any, default: int) -> int:
try:
return int(value)
except (TypeError, ValueError):
return default
def load_config(extra: Optional[Dict[str, Any]] = None) -> PhebyConfig:
"""Build a :class:`PhebyConfig` from env + ``platforms.pheby.extra``.
Environment variables win over YAML ``extra`` keys. Never raises.
"""
extra = dict(extra or {})
def _pick(env_name: str, key: str, default: Any = "") -> Any:
env = os.getenv(env_name, "")
if env.strip():
return env.strip()
val = extra.get(key)
if val is None or (isinstance(val, str) and not val.strip()):
return default
return val
secret = os.getenv(ENV_SECRET, "").strip()
if not secret:
# A YAML secret is allowed for local testing but strongly discouraged;
# env always wins and docs recommend env-only.
secret = str(extra.get("secret", "") or "").strip()
if len(secret) > DEFAULT_MAX_SECRET_LEN:
secret = secret[:DEFAULT_MAX_SECRET_LEN]
port = _coerce_int(
_pick(ENV_PORT, "port", DEFAULT_PORT), DEFAULT_PORT)
if not (0 < port < 65536):
port = DEFAULT_PORT
retention = _coerce_int(
_pick("", "attachment_retention_days", DEFAULT_RETENTION_DAYS),
DEFAULT_RETENTION_DAYS)
if retention < 0:
retention = DEFAULT_RETENTION_DAYS
debug = _env_bool(ENV_DEBUG) or bool(extra.get("debug", False))
log_chat = _env_bool(ENV_LOG_CHAT_CONTENT) or bool(
extra.get("log_chat_content", False))
cfg = PhebyConfig(
secret=secret,
bind_host=str(_pick(ENV_BIND_HOST, "bind_host", DEFAULT_BIND_HOST)),
port=port,
storage_dir=str(_pick("", "attachment_storage_dir", "") or ""),
retention_days=retention,
debug=bool(debug),
log_chat_content=bool(log_chat),
index_path=str(_pick("", "attachment_index_path", "") or ""),
extra=extra,
)
return cfg
def check_requirements() -> bool:
"""Platform-entry dependency check: aiohttp available + secret set."""
try:
import aiohttp # noqa: F401
except ImportError:
return False
return bool(os.getenv(ENV_SECRET, "").strip())
def validate_config(config: Any) -> bool:
"""Gateway config validation: at minimum a secret must be resolvable."""
extra = getattr(config, "extra", {}) or {}
if os.getenv(ENV_SECRET, "").strip():
return True
return bool(str(extra.get("secret", "") or "").strip())
def is_connected(config: Any) -> bool:
"""True when Pheby is configured (env or config.yaml)."""
return validate_config(config)
def env_enablement() -> Optional[Dict[str, Any]]:
"""Seed ``PlatformConfig.extra`` from env for env-only setups.
Mirrors the ntfy adapter pattern so ``hermes gateway status`` reflects a
PHEBY_SECRET-only deployment without instantiating the server.
"""
secret = os.getenv(ENV_SECRET, "").strip()
if not secret:
return None
seed: Dict[str, Any] = {}
host = os.getenv(ENV_BIND_HOST, "").strip()
if host:
seed["bind_host"] = host
port = os.getenv(ENV_PORT, "").strip()
if port:
seed["port"] = port
return seed
__all__ = [
"ENV_SECRET", "ENV_BIND_HOST", "ENV_PORT", "ENV_DEBUG",
"ENV_LOG_CHAT_CONTENT", "DEFAULT_BIND_HOST", "DEFAULT_PORT",
"DEFAULT_RETENTION_DAYS", "PhebyConfig", "load_config",
"check_requirements", "validate_config", "is_connected", "env_enablement",
]
+148
View File
@@ -0,0 +1,148 @@
"""Pheby conversation/session routing — maps Pheby conversations to Hermes sessions.
Hermes remains the authoritative source of conversation state. Each Pheby
conversation is one Hermes session reached through a stable ``SessionSource``
keyed ``pheby:dm:<conversation_id>``. Pheby keeps only a thin, rebuildable
mapping (conversation ID → display name) in ``pheby_conversations.json`` under
HERMES_HOME; everything else (history, titles, tokens) is read from the
SessionStore / SessionDB on demand.
"""
from __future__ import annotations
import asyncio
import json
import logging
from pathlib import Path
from typing import Any, Dict, List, Optional
from hermes_constants import get_hermes_home
from . import protocol
logger = logging.getLogger(__name__)
PLATFORM_NAME = "pheby"
_INDEX_FILE = "pheby_conversations.json"
_INDEX_MAX_ENTRIES = 500
class ConversationRouter:
"""Maps stable Pheby conversation IDs to Hermes session sources."""
def __init__(self) -> None:
self._lock = asyncio.Lock()
self._names: Dict[str, str] = {}
self._loaded = False
# ── persistence ──────────────────────────────────────────────────────
@property
def _index_path(self) -> Path:
return Path(get_hermes_home()) / _INDEX_FILE
def _load(self) -> None:
if self._loaded:
return
self._loaded = True
try:
if self._index_path.exists():
data = json.loads(self._index_path.read_text(encoding="utf-8"))
if isinstance(data, dict):
self._names = {
str(k): str(v)
for k, v in data.items()
if isinstance(k, str) and isinstance(v, (str, int))
}
except Exception:
logger.warning("[pheby] conversation index unreadable", exc_info=True)
def _save(self) -> None:
try:
# Bound the index; oldest-written entries lose (dict order).
if len(self._names) > _INDEX_MAX_ENTRIES:
keep = list(self._names.items())[-_INDEX_MAX_ENTRIES:]
self._names = dict(keep)
tmp = self._index_path.with_suffix(".tmp")
tmp.write_text(json.dumps(self._names, ensure_ascii=False, indent=1),
encoding="utf-8")
import os
os.replace(tmp, self._index_path)
except Exception:
logger.error("[pheby] failed to persist conversation index",
exc_info=True)
# ── ID management ────────────────────────────────────────────────────
async def ensure_conversation(self, conversation_id: str,
name: Optional[str] = None) -> str:
"""Register a conversation ID (client-generated or server-new)."""
async with self._lock:
self._load()
cid = str(conversation_id or protocol.new_id())
if not cid or len(cid) > 128:
cid = protocol.new_id()
if cid not in self._names:
self._names[cid] = (name or "").strip() or "New chat"
self._save()
elif name:
self._names[cid] = name.strip()
self._save()
return cid
async def new_conversation(self, name: Optional[str] = None) -> str:
return await self.ensure_conversation(protocol.new_id(), name)
async def rename(self, conversation_id: str, name: str) -> bool:
async with self._lock:
self._load()
cid = str(conversation_id)
if cid not in self._names:
return False
self._names[cid] = (name or "").strip() or self._names[cid]
self._save()
return True
async def forget(self, conversation_id: str) -> bool:
"""Remove the local index entry (session deletion is handled via DB)."""
async with self._lock:
self._load()
existed = str(conversation_id) in self._names
self._names.pop(str(conversation_id), None)
if existed:
self._save()
return existed
async def get_name(self, conversation_id: str) -> Optional[str]:
async with self._lock:
self._load()
return self._names.get(str(conversation_id))
async def known_ids(self) -> List[str]:
async with self._lock:
self._load()
return list(self._names.keys())
# ── session key ──────────────────────────────────────────────────────
@staticmethod
def session_key_for(conversation_id: str) -> str:
"""The Hermes gateway session key for a Pheby conversation.
Session keys are built by ``gateway.session.build_session_key`` for
DM sources as ``agent:main:pheby:dm:<chat_id>``. Conversation IDs are
server/opaque-controlled (32-hex or client GUIDs validated below), so
the chat_id component is safe to embed in the key.
"""
return f"agent:main:{PLATFORM_NAME}:dm:{conversation_id}"
@staticmethod
def is_valid_conversation_id(value: Any) -> bool:
"""Accept opaque IDs up to 128 chars from a safe alphabet."""
if not isinstance(value, str) or not value or len(value) > 128:
return False
allowed = set(
"abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789-_"
)
return all(c in allowed for c in value)
__all__ = ["ConversationRouter", "PLATFORM_NAME"]
+715
View File
@@ -0,0 +1,715 @@
"""Hermes gateway integration — runs, approvals, clarifications, models.
This module is the ONLY place that touches Hermes internals, so every Hermes
API dependency is documented and defensive (getattr + try/except) to survive
normal Hermes upgrades. All Hermes imports are deferred (inside functions)
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 uuid
from typing import Any, Dict, List, Optional, Tuple
from . import protocol as proto
from .conversations import ConversationRouter
logger = logging.getLogger(__name__)
# Pending interactive requests (approval/clarify) keyed by opaque ID → context.
# Single-user app, but a dict keeps the protocol multi-client friendly.
_PENDING_APPROVALS: Dict[str, Dict[str, Any]] = {}
_PENDING_CLARIFIES: Dict[str, Dict[str, Any]] = {}
_RUN_LOCK = threading.Lock()
_ACTIVE_RUNS: Dict[str, Dict[str, Any]] = {} # conversation_id → run info
def _runner() -> Any:
"""The GatewayRunner back-reference injected into the adapter."""
adapter = _current_adapter()
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()
def set_adapter(adapter: Any) -> None:
_ADAPTER_CTX.set(adapter)
def _session_store() -> Any:
runner = _runner()
return getattr(runner, "session_store", None) if runner else None
def _session_db() -> Any:
runner = _runner()
db = getattr(runner, "_session_db", None) if runner else None
return getattr(db, "_db", db) if db else None
def _source_for(conversation_id: str, user_name: str = "Chris"):
"""Build the SessionSource for a Pheby conversation (deferred import)."""
adapter = _current_adapter()
if adapter is not None:
return adapter.build_source(
chat_id=conversation_id,
chat_name=conversation_id,
chat_type="dm",
user_id="pheby-client",
user_name=user_name,
)
# Fallback (tests / standalone): construct directly.
from gateway.config import Platform
from gateway.session import SessionSource
return SessionSource(
platform=Platform("pheby"),
chat_id=str(conversation_id),
chat_name=str(conversation_id),
chat_type="dm",
user_id="pheby-client",
user_name=user_name,
)
def _session_key_for(conversation_id: str) -> str:
"""Compute the gateway session key for a conversation.
Prefers the SessionStore's own key builder (authoritative); falls back to
the documented deterministic shape used by ``build_session_key`` for DM
sources (``agent:main:<platform>:dm:<chat_id>``).
"""
store = _session_store()
if store is not None:
try:
source = _source_for(conversation_id)
return store._generate_session_key(source)
except Exception:
logger.debug("[pheby] session key via store failed", exc_info=True)
return ConversationRouter.session_key_for(conversation_id)
# ═══════════════════════════════════════════════════════════════════════════
# Conversations
# ═══════════════════════════════════════════════════════════════════════════
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()
# 1. Sessions Hermes already tracks for the pheby platform.
store = _session_store()
if store is not None:
try:
entries = await asyncio.to_thread(store.list_sessions)
for entry in entries:
origin = getattr(entry, "origin", None)
platform = getattr(getattr(origin, "platform", None),
"value", "")
if platform != "pheby":
continue
cid = str(getattr(origin, "chat_id", "") or "")
if not cid or cid in seen:
continue
seen.add(cid)
out.append({
"conversation_id": cid,
"name": (getattr(entry, "display_name", None)
or _router_name(router, cid) or cid),
"session_id": getattr(entry, "session_id", None),
"last_active": _iso(getattr(entry, "updated_at", None)),
"source": "hermes",
})
except Exception:
logger.debug("[pheby] session store listing failed", exc_info=True)
# 2. Router-known conversations (incl. freshly created, no messages yet).
if router is not None:
for cid in await router.known_ids():
if cid in seen:
continue
seen.add(cid)
out.append({
"conversation_id": cid,
"name": await router.get_name(cid) or cid,
"session_id": None,
"last_active": None,
"source": "pheby",
})
out.sort(key=lambda c: (c.get("last_active") is None,
c.get("last_active") or ""), reverse=False)
out.sort(key=lambda c: c.get("last_active") or "", reverse=True)
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
except AttributeError:
return None
async def conversation_history(conversation_id: str, limit: int
) -> Tuple[List[Dict[str, Any]], bool]:
"""Load transcript rows for a conversation from Hermes state.db.
Returns ``(messages, found)``. ``found`` is False when neither the
session store nor the session DB knows the conversation.
"""
messages: List[Dict[str, Any]] = []
found = False
store = _session_store()
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:
session_id = str(entry)
found = True
except Exception:
logger.debug("[pheby] peek_session_id failed", exc_info=True)
db = _session_db()
if db is not None and session_id:
try:
rows = await asyncio.to_thread(
db.get_messages_as_conversation, session_id)
for row in rows[-limit:]:
role = row.get("role")
if role not in ("user", "assistant"):
continue
content = row.get("content")
text = content if isinstance(content, str) else str(content or "")
# Tool-call rows can surface as assistant rows with empty
# content; skip empties so the client transcript stays clean.
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,
"role": role,
"text": text,
"ts": row.get("timestamp") if isinstance(
row.get("timestamp"), str) else None,
})
found = True
except Exception:
logger.debug("[pheby] transcript load failed", exc_info=True)
# A router-known conversation with no messages yet is still "found" so a
# fresh client can open it as an empty chat.
if not found:
server = _current_server()
if server is not None:
name = await server.router.get_name(conversation_id)
if name is not None:
found = True
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.
"""
server = _current_server()
router = server.router if server else None
if router is None:
return False
if not await router.forget(conversation_id):
return False
store = _session_store()
db = _session_db()
session_id = None
if store is not None:
try:
session_id = await asyncio.to_thread(
store.peek_session_id, _session_key_for(conversation_id))
except Exception:
session_id = None
if session_id and db is not None:
try:
await asyncio.to_thread(db.delete_session, session_id)
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)
_ACTIVE_RUNS.pop(conversation_id, None)
return True
# ═══════════════════════════════════════════════════════════════════════════
# Chat runs
# ═══════════════════════════════════════════════════════════════════════════
async def send_chat(server: Any, conversation_id: str, text: str,
client: Any, request_id: Optional[str]) -> None:
"""Deliver a user message into the Hermes gateway for this conversation.
The gateway's full pipeline (auth, sessions, tools, approvals, clarify,
deliverables, streaming) runs on the adapter's message handler. Pheby
adds nothing to the agent loop.
"""
adapter = _current_adapter()
if adapter is None or not hasattr(adapter, "handle_message"):
await client.send_json(proto.error_event(
proto.ERR_INTERNAL, "Gateway not connected yet", request_id))
return
# Register the conversation so it survives restarts.
await server.router.ensure_conversation(conversation_id)
run_id = uuid.uuid4().hex[:16]
source = _source_for(conversation_id)
from gateway.platforms.base import MessageEvent, MessageType
event = MessageEvent(
text=text,
message_type=MessageType.TEXT,
source=source,
message_id=uuid.uuid4().hex[:12],
metadata={"pheby_run_id": run_id},
)
_ACTIVE_RUNS[conversation_id] = {
"run_id": run_id,
"started": asyncio.get_event_loop().time(),
}
await client.send_json({
"type": proto.S_RUN_ACCEPTED,
"conversation_id": conversation_id,
"run_id": run_id,
**({"request_id": request_id} if request_id else {}),
})
draft_message_id = f"draft-{run_id}"
await server.broadcast({
"type": proto.S_MESSAGE_START,
"conversation_id": conversation_id,
"run_id": run_id,
"message_id": draft_message_id,
})
# Track the draft in the adapter so send()/edit_message() associate the
# final text with the announced draft message id.
if getattr(adapter, "_drafts", None) is not None:
adapter._drafts.setdefault(conversation_id, {
"message_id": draft_message_id, "text": ""})
# The base adapter's handle_message() spawns background tasks and
# returns quickly; the eventual reply arrives through adapter.send().
await adapter.handle_message(event)
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
server = _current_server()
if server is None:
return
payload = {
"type": proto.S_RUN_FINISHED,
"conversation_id": conversation_id,
"status": status,
}
if run_id:
payload["run_id"] = run_id
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))
except RuntimeError:
pass
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:
return False # stale run id — nothing to cancel
session_key = _session_key_for(conversation_id)
runner = _runner()
interrupted = False
if runner is not None:
# 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()):
try:
agent.interrupt("Cancelled by Pheby client")
invalidate = getattr(
runner, "_invalidate_session_run_generation", None)
if callable(invalidate):
invalidate(session_key, reason="pheby_cancel")
interrupted = True
except Exception:
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
except Exception:
logger.debug("[pheby] adapter interrupt failed", exc_info=True)
note_run_finished(conversation_id,
"cancelled" if interrupted else "idle")
return interrupted
# ═══════════════════════════════════════════════════════════════════════════
# Approvals
# ═══════════════════════════════════════════════════════════════════════════
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", ""))
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 = {
"type": proto.S_APPROVAL_REQUEST,
"approval_id": approval_id,
"session_key": session_key,
"command": proto.safe_str(command, 2000),
"description": proto.safe_str(
approval_data.get("description", ""), 1000),
"choices": choices,
"ts": proto.now_iso(),
}
_PENDING_APPROVALS[approval_id] = {
"session_key": session_key,
"created": asyncio.get_event_loop().time(),
}
server = _current_server()
if server is not None:
await server.broadcast(event)
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
try:
from tools.approval import resolve_gateway_approval
count = await asyncio.to_thread(
resolve_gateway_approval,
pending["session_key"], choice, False, reason)
ok = count > 0
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
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()
for aid in [a for a, p in _PENDING_APPROVALS.items()
if now - p["created"] > max_age]:
_PENDING_APPROVALS.pop(aid, None)
# ═══════════════════════════════════════════════════════════════════════════
# Clarifications
# ═══════════════════════════════════════════════════════════════════════════
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 = {
"type": proto.S_CLARIFY_REQUEST,
"clarify_id": clarify_id,
"session_key": session_key,
"question": proto.safe_str(question, 2000),
"choices": [proto.safe_str(c, 300) for c in choices]
if choices else None,
"allow_free_text": True, # Hermes clarify always permits "Other"
"ts": proto.now_iso(),
}
_PENDING_CLARIFIES[clarify_id] = {
"session_key": session_key,
"created": asyncio.get_event_loop().time(),
}
server = _current_server()
if server is not None:
await server.broadcast(event)
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)
if pending is None:
return False
try:
from tools.clarify_gateway import resolve_gateway_clarify
ok = await asyncio.to_thread(
resolve_gateway_clarify, clarify_id, response)
if not ok:
# Might be an awaiting-text open-ended clarify: route via the
# session text path instead.
from tools.clarify_gateway import \
resolve_text_response_for_session
ok = await asyncio.to_thread(
resolve_text_response_for_session,
pending["session_key"], response)
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),
})
return bool(ok)
# ═══════════════════════════════════════════════════════════════════════════
# Models & reasoning
# ═══════════════════════════════════════════════════════════════════════════
async def models_snapshot() -> Dict[str, Any]:
"""Providers + models Hermes currently exposes (credential-aware)."""
def _collect() -> Dict[str, Any]:
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 "")
providers = list_picker_providers(
current_provider=current_provider,
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
)
return {"providers": providers, "current_model": current_model,
"current_provider": current_provider}
try:
data = await asyncio.to_thread(_collect)
except Exception:
logger.error("[pheby] model listing failed", exc_info=True)
data = {"providers": [], "current_model": "", "current_provider": "",
"error": "Model catalog unavailable"}
data["supported_reasoning_efforts"] = list(proto.REASONING_EFFORTS)
data["ts"] = proto.now_iso()
return data
async def current_model_snapshot() -> Dict[str, Any]:
def _collect() -> Dict[str, Any]:
cfg = _load_cfg()
model_cfg = (cfg.get("model") or {}) if isinstance(cfg, dict) else {}
return {"model": str(model_cfg.get("default", "") or ""),
"provider": str(model_cfg.get("provider", "") or "")}
try:
data = await asyncio.to_thread(_collect)
except Exception:
data = {"model": "", "provider": "", "error": "Config unavailable"}
data["ts"] = proto.now_iso()
return data
async def set_model(model: str, provider: Optional[str],
conversation_id: Optional[str]) -> Dict[str, Any]:
"""Change the active model via Hermes's session/global override path."""
if not model:
return {"ok": False, "code": proto.ERR_BAD_REQUEST,
"message": "model is required"}
try:
from hermes_cli.model_switch import switch_model
cfg = _load_cfg()
model_cfg = (cfg.get("model") or {}) if isinstance(cfg, dict) else {}
result = await asyncio.to_thread(
switch_model,
model,
str(model_cfg.get("provider", "openrouter") or "openrouter"),
str(model_cfg.get("default", "") or ""),
str(model_cfg.get("base_url", "") or ""),
"", # current_api_key — runtime resolution handles credentials
False, # is_global → session-scoped when conversation given
provider or "",
cfg.get("providers") if isinstance(cfg, dict) else None,
None,
)
except Exception as exc:
logger.error("[pheby] switch_model failed", exc_info=True)
return {"ok": False, "code": proto.ERR_BAD_REQUEST,
"message": proto.safe_str(exc, 200)}
ok = bool(getattr(result, "success", False))
if not ok:
return {"ok": False, "code": proto.ERR_BAD_REQUEST,
"message": proto.safe_str(getattr(result, "error", ""),
300)}
resolved_model = getattr(result, "model", model)
resolved_provider = getattr(result, "provider", provider or "")
override = {"model": resolved_model}
if resolved_provider:
override["provider"] = resolved_provider
store = _session_store()
if conversation_id and store is not None:
try:
await asyncio.to_thread(store.set_model_override,
_session_key_for(conversation_id),
override)
return {"ok": True, "model": resolved_model,
"provider": resolved_provider, "scope": "conversation"}
except Exception:
logger.debug("[pheby] session model override failed",
exc_info=True)
# Global fallback: persist via Hermes config save (same path /model
# --global uses).
try:
await asyncio.to_thread(_save_global_model, resolved_model,
resolved_provider)
return {"ok": True, "model": resolved_model,
"provider": resolved_provider, "scope": "global"}
except Exception as exc:
logger.error("[pheby] global model save failed", exc_info=True)
return {"ok": False, "code": proto.ERR_INTERNAL,
"message": proto.safe_str(exc, 200)}
def _save_global_model(model: str, provider: str) -> None:
from hermes_cli.config import load_config, save_config_value
save_config_value("model.default", model)
if provider:
save_config_value("model.provider", provider)
async def reasoning_snapshot() -> Dict[str, Any]:
def _collect() -> Dict[str, Any]:
from hermes_constants import resolve_reasoning_config
cfg = _load_cfg()
model_cfg = (cfg.get("model") or {}) if isinstance(cfg, dict) else {}
resolved = resolve_reasoning_config(
cfg, str(model_cfg.get("default", "") or ""))
if resolved is None:
return {"effort": None, "enabled": None}
if resolved.get("enabled") is False:
return {"effort": "none", "enabled": False}
return {"effort": resolved.get("effort"), "enabled": True}
try:
data = await asyncio.to_thread(_collect)
except Exception:
data = {"effort": None, "enabled": None,
"error": "Config unavailable"}
data["supported_efforts"] = ["none"] + list(proto.REASONING_EFFORTS)
data["ts"] = proto.now_iso()
return data
async def set_reasoning(effort: str,
conversation_id: Optional[str]) -> Dict[str, Any]:
"""Set reasoning effort (Hermes levels + 'none' to disable)."""
if effort not in ("none",) + proto.REASONING_EFFORTS:
return {"ok": False, "code": proto.ERR_BAD_REQUEST,
"message": f"effort must be one of: none, "
f"{', '.join(proto.REASONING_EFFORTS)}"}
parsed = {"enabled": False} if effort == "none" else {
"enabled": True, "effort": effort}
runner = _runner()
if runner is not None and conversation_id:
try:
await asyncio.to_thread(
runner._set_session_reasoning_override,
_session_key_for(conversation_id), parsed)
return {"ok": True, "effort": effort, "scope": "conversation"}
except Exception:
logger.debug("[pheby] session reasoning override failed",
exc_info=True)
try:
await asyncio.to_thread(_save_global_reasoning, effort)
return {"ok": True, "effort": effort, "scope": "global"}
except Exception as exc:
return {"ok": False, "code": proto.ERR_INTERNAL,
"message": proto.safe_str(exc, 200)}
def _save_global_reasoning(effort: str) -> None:
from hermes_cli.config import save_config_value
save_config_value("agent.reasoning_effort",
False if effort == "none" else effort)
def _load_cfg() -> Dict[str, Any]:
from hermes_cli.config import load_config
return load_config() or {}
# ═══════════════════════════════════════════════════════════════════════════
# Server context
# ═══════════════════════════════════════════════════════════════════════════
_SERVER_CTX: contextvars.ContextVar = contextvars.ContextVar(
"pheby_server", default=None)
def set_server(server: Any) -> None:
_SERVER_CTX.set(server)
def _current_server() -> Any:
return _SERVER_CTX.get()
__all__ = [
"set_adapter", "set_server", "list_conversations", "conversation_history",
"delete_conversation", "send_chat", "cancel_run", "note_run_finished",
"push_approval", "resolve_approval", "fail_stale_approvals",
"push_clarify", "resolve_clarify", "models_snapshot",
"current_model_snapshot", "set_model", "reasoning_snapshot",
"set_reasoning",
]
+47
View File
@@ -0,0 +1,47 @@
name: pheby
label: Pheby
kind: platform
version: 1.0.0
description: >
Pheby platform adapter for Hermes Agent — serves a WebSocket + HTTPS
protocol for the Pheby native Android client (Kotlin/Compose) behind a
Caddy reverse proxy. Exposes multiple conversations (Hermes sessions),
streamed chat, structured tool events, native approval/clarification
round-trips, model + reasoning-effort selection, and agent-generated
attachments with 7-day adapter-managed retention. Single-user, shared-secret
auth (PHEBY_SECRET). No Hermes core modifications required.
author: Pheby (for Chris)
requires_env:
- name: PHEBY_SECRET
description: "Shared credential every WS connection and attachment download must present (generate with: openssl rand -hex 32)"
prompt: "Pheby shared secret"
password: true
optional_env:
- name: PHEBY_BIND_HOST
description: "Bind address (default: 127.0.0.1 — safe behind a local Caddy)"
prompt: "Bind host (or empty for 127.0.0.1)"
password: false
- name: PHEBY_PORT
description: "Listen port (default: 8620)"
prompt: "Listen port (or empty for 8620)"
password: false
- name: PHEBY_DEBUG
description: "Verbose protocol logging (true/false, default false)"
prompt: "Enable debug protocol logging? (true/false)"
password: false
- name: PHEBY_LOG_CHAT_CONTENT
description: "Debug logs may include chat text and tool previews (default false)"
prompt: "Log chat content in debug mode? (true/false)"
password: false
- name: PHEBY_HOME_CHANNEL
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
+193
View File
@@ -0,0 +1,193 @@
"""Pheby protocol v1 — message types, error codes, and (de)serialization helpers.
The protocol is JSON-over-WebSocket. Every message (both directions) has a
``type`` field. Client → server requests may carry a ``request_id`` (any
string) which is echoed on the direct reply so the client can correlate
RPC-style calls. Server → client events are broadcast to all authenticated
connections and carry the IDs needed to associate them with a conversation,
message, run, tool call, approval, clarification, or attachment.
This module is intentionally dependency-free (stdlib only) so it can be unit
tested without Hermes/aiohttp installed.
"""
from __future__ import annotations
import json
import logging
import re
import uuid
from datetime import datetime, timezone
from typing import Any, Dict, Optional, Tuple
logger = logging.getLogger(__name__)
# ── Protocol version ─────────────────────────────────────────────────────────
PROTOCOL_VERSION = 1
# Bump when a wire-incompatible change lands. Clients negotiate via the
# ``hello`` handshake; the server refuses mismatches with a ``version_mismatch``
# error instead of guessing.
# ── Limits ───────────────────────────────────────────────────────────────────
MAX_TEXT_CHARS = 64_000 # client chat message body limit
MAX_WS_MESSAGE_BYTES = 2 * 1024 * 1024 # inbound WebSocket frame cap (aiohttp)
MAX_HISTORY_MESSAGES = 500 # per conversation.open fetch cap
AUTH_TIMEOUT_SECONDS = 10.0 # hello must arrive within this window
AUTH_FAILURE_LOCKOUT_SECONDS = 60.0 # repeated auth failures lock the source
AUTH_FAILURE_THRESHOLD = 5 # failures before lockout
# ── Error codes (machine-readable) ───────────────────────────────────────────
ERR_UNAUTHORIZED = "unauthorized"
ERR_AUTH_TIMEOUT = "auth_timeout"
ERR_VERSION_MISMATCH = "version_mismatch"
ERR_BAD_REQUEST = "bad_request"
ERR_INVALID_JSON = "invalid_json"
ERR_UNKNOWN_TYPE = "unknown_type"
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_TOO_LARGE = "too_large"
ERR_RATE_LIMITED = "rate_limited"
ERR_INTERNAL = "internal_error"
ERR_NOT_IMPLEMENTED = "not_implemented"
# ── Client → server message types ────────────────────────────────────────────
C_HELLO = "hello"
C_PING = "ping"
C_CONVERSATION_LIST = "conversation.list"
C_CONVERSATION_OPEN = "conversation.open"
C_CONVERSATION_CREATE = "conversation.create"
C_CONVERSATION_RENAME = "conversation.rename"
C_CONVERSATION_DELETE = "conversation.delete"
C_CHAT_SEND = "chat.send"
C_RUN_CANCEL = "run.cancel"
C_APPROVAL_RESPOND = "approval.respond"
C_CLARIFY_RESPOND = "clarify.respond"
C_MODELS_LIST = "models.list"
C_MODEL_SET = "model.set"
C_MODEL_CURRENT = "models.current"
C_REASONING_SET = "reasoning.set"
C_REASONING_CURRENT = "reasoning.current"
# ── Server → client message types ────────────────────────────────────────────
S_READY = "ready"
S_PONG = "pong"
S_ERROR = "error"
S_CONVERSATION_SNAPSHOT = "conversation.snapshot" # reply to conversation.list
S_CONVERSATION_CREATED = "conversation.created"
S_CONVERSATION_RENAMED = "conversation.renamed"
S_CONVERSATION_UPDATED = "conversation.updated" # auto-title etc.
S_CONVERSATION_DELETED = "conversation.deleted"
S_CONVERSATION_HISTORY = "conversation.history" # reply to conversation.open
S_RUN_ACCEPTED = "run.accepted"
S_RUN_FINISHED = "run.finished"
S_MESSAGE_START = "message.start" # streaming draft opened
S_MESSAGE_DELTA = "message.delta" # cumulative streamed draft text
S_MESSAGE_COMPLETE = "message.complete" # final assistant message
S_TOOL_EVENT = "tool.event" # structured tool activity
S_APPROVAL_REQUEST = "approval.request"
S_APPROVAL_RESOLVED = "approval.resolved"
S_CLARIFY_REQUEST = "clarify.request"
S_CLARIFY_RESOLVED = "clarify.resolved"
S_ATTACHMENT_ADDED = "attachment.added"
S_MODELS_SNAPSHOT = "models.snapshot" # reply to models.list
S_MODEL_CURRENT_SNAPSHOT = "model.current" # reply to models.current
S_MODEL_CHANGED = "model.changed" # after model.set accepted
S_REASONING_SNAPSHOT = "reasoning.snapshot" # reply to reasoning.current
S_REASONING_CHANGED = "reasoning.changed" # after reasoning.set accepted
# Reasoning effort levels supported by Hermes (hermes_constants).
REASONING_EFFORTS = ("minimal", "low", "medium", "high", "xhigh", "max", "ultra")
_ATTACHMENT_ID_RE = re.compile(r"^[0-9a-f]{32}$")
def now_iso() -> str:
"""UTC timestamp in ISO-8601 format."""
return datetime.now(timezone.utc).isoformat()
def new_id() -> str:
"""Generate an opaque 32-hex identifier (also the attachment ID shape)."""
return uuid.uuid4().hex
def is_valid_attachment_id(value: str) -> bool:
"""True when *value* looks like one of our opaque attachment IDs."""
return bool(isinstance(value, str) and _ATTACHMENT_ID_RE.match(value))
def encode_message(payload: Dict[str, Any]) -> str:
"""Serialize a protocol message to a JSON string (compact, UTF-8)."""
return json.dumps(payload, ensure_ascii=False, separators=(",", ":"))
def decode_message(raw: Any) -> Tuple[Optional[Dict[str, Any]], Optional[str]]:
"""Parse one inbound WebSocket text frame.
Returns ``(message, None)`` on success or ``(None, error_code)`` when the
frame is not a valid JSON object. Never raises.
"""
try:
data = json.loads(raw)
except (json.JSONDecodeError, UnicodeDecodeError, TypeError, ValueError):
return None, ERR_INVALID_JSON
if not isinstance(data, dict):
return None, ERR_INVALID_JSON
if not isinstance(data.get("type"), str) or not data["type"]:
return None, ERR_INVALID_JSON
return data, None
def error_event(code: str, message: str, request_id: Optional[str] = None,
**extra: Any) -> Dict[str, Any]:
"""Build a server → client error event with a machine-readable code."""
event: Dict[str, Any] = {
"type": S_ERROR,
"error": {"code": code, "message": str(message)[:500]},
"ts": now_iso(),
}
if request_id is not None:
event["request_id"] = request_id
if extra:
event.update(extra)
return event
def safe_str(value: Any, max_len: int = 500) -> str:
"""Coerce *value* to a bounded string for logs and summaries."""
if value is None:
return ""
text = value if isinstance(value, str) else json.dumps(
value, ensure_ascii=False, default=str)
return text[:max_len]
__all__ = [
"PROTOCOL_VERSION", "MAX_TEXT_CHARS", "MAX_WS_MESSAGE_BYTES",
"MAX_HISTORY_MESSAGES", "AUTH_TIMEOUT_SECONDS",
"AUTH_FAILURE_LOCKOUT_SECONDS", "AUTH_FAILURE_THRESHOLD",
"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_INTERNAL", "ERR_NOT_IMPLEMENTED",
"C_HELLO", "C_PING", "C_CONVERSATION_LIST", "C_CONVERSATION_OPEN",
"C_CONVERSATION_CREATE", "C_CONVERSATION_RENAME", "C_CONVERSATION_DELETE",
"C_CHAT_SEND", "C_RUN_CANCEL", "C_APPROVAL_RESPOND", "C_CLARIFY_RESPOND",
"C_MODELS_LIST", "C_MODEL_SET", "C_MODEL_CURRENT", "C_REASONING_SET",
"C_REASONING_CURRENT",
"S_READY", "S_PONG", "S_ERROR", "S_CONVERSATION_SNAPSHOT",
"S_CONVERSATION_CREATED", "S_CONVERSATION_RENAMED", "S_CONVERSATION_UPDATED",
"S_CONVERSATION_DELETED", "S_CONVERSATION_HISTORY", "S_RUN_ACCEPTED",
"S_RUN_FINISHED", "S_MESSAGE_START", "S_MESSAGE_DELTA",
"S_MESSAGE_COMPLETE", "S_TOOL_EVENT", "S_APPROVAL_REQUEST",
"S_APPROVAL_RESOLVED", "S_CLARIFY_REQUEST", "S_CLARIFY_RESOLVED",
"S_ATTACHMENT_ADDED", "S_MODELS_SNAPSHOT", "S_MODEL_CURRENT_SNAPSHOT",
"S_MODEL_CHANGED", "S_REASONING_SNAPSHOT", "S_REASONING_CHANGED",
"REASONING_EFFORTS",
"now_iso", "new_id", "is_valid_attachment_id", "encode_message",
"decode_message", "error_event", "safe_str",
]
+562
View File
@@ -0,0 +1,562 @@
"""Pheby HTTP + WebSocket server (aiohttp).
Listens on localhost/plain HTTP behind Caddy. Routes:
* ``GET /health`` — unauthenticated liveness (minimal info).
* ``GET /ws`` — WebSocket; first frame must be ``hello``.
* ``GET /attachments/{id}`` — authenticated attachment download (streamed).
All message routing lives in :meth:`PhebyServer.handle_client_message`;
Hermes integration (runs, approvals, clarifications, models) lives in
:mod:`.hermes_bridge` to keep this module focused on protocol + transport.
"""
from __future__ import annotations
import asyncio
import logging
import time
from pathlib import Path
from typing import Any, Dict, List, Optional
from aiohttp import web
from . import protocol as proto
from .attachments import AttachmentStore, constant_time_equals
from .config import PhebyConfig
from .conversations import ConversationRouter
from . import hermes_bridge
from .ws_client import ClientConnection
logger = logging.getLogger(__name__)
class PhebyServer:
"""Owns the aiohttp app, connected clients, and shared subsystems."""
def __init__(self, config: PhebyConfig, adapter: Any = None):
self.config = config
self.adapter = adapter # PhebyAdapter (may be None in tests)
self.router = ConversationRouter()
from hermes_constants import get_hermes_home
root = config.attachments_root or str(
Path(get_hermes_home()) / "pheby-attachments")
self.store = AttachmentStore(
root=Path(root),
retention_days=config.retention_days,
index_path=Path(config.index_path) if config.index_path else None,
)
self.bridge = hermes_bridge # bridge function module
self._clients: Dict[str, ClientConnection] = {}
self._auth_failures: Dict[str, List[float]] = {}
self._cleanup_task: Optional[asyncio.Task] = None
self._app: Optional[web.Application] = None
self._runner: Optional[web.AppRunner] = None
self._site: Optional[web.TCPSite] = None
self._conn_counter = 0
# ── lifecycle ────────────────────────────────────────────────────────
async def start(self) -> bool:
from aiohttp import web as _web # local import keeps import light
self.store.hydrate_legacy_meta()
app = _web.Application(client_max_size=proto.MAX_WS_MESSAGE_BYTES)
app.router.add_get("/health", self._handle_health)
app.router.add_get("/ws", self._handle_ws)
app.router.add_get("/attachments/{attachment_id}",
self._handle_attachment_download)
self._app = app
self._runner = web.AppRunner(app, access_log=None)
await self._runner.setup()
self._site = web.TCPSite(self._runner, self.config.bind_host,
self.config.port)
try:
await self._site.start()
except OSError as exc:
logger.error("[pheby] failed to bind %s:%s — %s",
self.config.bind_host, self.config.port, exc)
await self.stop()
return False
self._cleanup_task = asyncio.create_task(self._cleanup_loop())
logger.info("[pheby] serving on http://%s:%d (attachments: %s)",
self.config.bind_host, self.config.port,
self.store.root)
return True
async def stop(self) -> None:
if self._cleanup_task:
self._cleanup_task.cancel()
try:
await self._cleanup_task
except asyncio.CancelledError:
pass
self._cleanup_task = None
for client in list(self._clients.values()):
client.closed = True
try:
await client.ws.close()
except Exception:
pass
self._clients.clear()
if self._runner:
await self._runner.cleanup()
self._runner = None
self._site = None
self._app = None
logger.info("[pheby] server stopped")
# ── background cleanup ───────────────────────────────────────────────
async def _cleanup_loop(self) -> None:
"""Hourly expired-attachment sweep; first sweep after 5 minutes."""
try:
await asyncio.sleep(300)
while True:
try:
await self.store.cleanup_expired()
except Exception:
logger.error("[pheby] attachment cleanup failed",
exc_info=True)
await asyncio.sleep(3600)
except asyncio.CancelledError:
pass
# ── HTTP handlers ────────────────────────────────────────────────────
async def _handle_health(self, request: web.Request) -> web.Response:
"""Minimal unauthenticated liveness probe."""
return web.json_response({"status": "ok"})
def _check_http_secret(self, request: web.Request) -> bool:
header = request.headers.get("Authorization", "")
if header.startswith("Bearer "):
token = header[7:].strip()
elif header.startswith("ApiKey "):
token = header[7:].strip()
else:
token = request.headers.get("X-Pheby-Secret", "").strip()
if not token:
return False
return constant_time_equals(token, self.config.secret)
async def _handle_attachment_download(
self, request: web.Request) -> web.StreamResponse:
attachment_id = request.match_info.get("attachment_id", "")
if not self._check_http_secret(request):
return web.json_response(
{"error": {"code": proto.ERR_UNAUTHORIZED,
"message": "Authentication required"}},
status=401)
blob = self.store.resolve_blob(attachment_id)
if blob is None:
# Expired, unknown, or malformed — same minimal response so the
# endpoint leaks nothing about implementation details.
return web.json_response(
{"error": {"code": proto.ERR_NOT_FOUND,
"message": "Attachment unavailable"}},
status=404)
desc = self.store.describe(attachment_id) or {}
safe_name = desc.get("filename", "file.bin")
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-Type": desc.get("mime_type",
"application/octet-stream"),
},
)
# ── WebSocket handler ────────────────────────────────────────────────
async def _handle_ws(self, request: web.Request) -> web.WebSocketResponse:
ws = web.WebSocketResponse(max_msg_size=proto.MAX_WS_MESSAGE_BYTES,
heartbeat=30.0, autoping=True)
await ws.prepare(request)
self._conn_counter += 1
conn_id = f"c{self._conn_counter}"
peer = request.remote or "unknown"
if self._is_locked_out(peer):
logger.warning("[pheby] auth lockout active for %s — refusing",
peer)
await ws.close(code=4401, message=b"locked out")
return ws
client = ClientConnection(ws, conn_id)
self._clients[conn_id] = client
logger.info("[pheby] client %s connected from %s", conn_id, peer)
try:
# Auth phase: hello must arrive within the window.
try:
authed = await asyncio.wait_for(
self._authenticate(client, peer),
timeout=proto.AUTH_TIMEOUT_SECONDS)
except asyncio.TimeoutError:
await client.send_json(proto.error_event(
proto.ERR_AUTH_TIMEOUT, "hello not received in time"))
await ws.close()
return ws
if not authed:
await ws.close(code=4401, message=b"unauthorized")
return ws
await client.send_json({
"type": proto.S_READY,
"protocol_version": proto.PROTOCOL_VERSION,
"server": "pheby",
"ts": proto.now_iso(),
})
await client.read_loop(self)
finally:
self._clients.pop(conn_id, None)
logger.info("[pheby] client %s disconnected (authed=%s, %.0fs)",
conn_id, client.authenticated,
time.time() - client.connected_at)
return ws
def _is_locked_out(self, peer: str) -> bool:
fails = self._auth_failures.get(peer)
if not fails:
return False
cutoff = time.time() - proto.AUTH_FAILURE_LOCKOUT_SECONDS
recent = [t for t in fails if t > cutoff]
self._auth_failures[peer] = recent
return len(recent) >= proto.AUTH_FAILURE_THRESHOLD
def _record_auth_failure(self, peer: str) -> None:
self._auth_failures.setdefault(peer, []).append(time.time())
async def _authenticate(self, client: ClientConnection,
peer: str) -> bool:
"""Wait for the hello frame and validate the shared secret."""
msg = await client.ws.receive(timeout=proto.AUTH_TIMEOUT_SECONDS + 5)
if msg.type != "text" and not hasattr(msg, "data"):
return False
message, err = proto.decode_message(msg.data)
if err or message is None:
await client.send_json(
proto.error_event(err or proto.ERR_BAD_REQUEST,
"Expected hello message"))
return False
if message.get("type") != proto.C_HELLO:
await client.send_json(proto.error_event(
proto.ERR_UNAUTHORIZED, "First message must be hello"))
self._record_auth_failure(peer)
return False
supplied = str(message.get("secret", ""))
if not supplied or not constant_time_equals(supplied,
self.config.secret):
logger.warning("[pheby] auth failure from %s", peer)
self._record_auth_failure(peer)
# Small delay to slow brute force; constant-time compare already
# used for the secret itself.
await asyncio.sleep(0.5)
await client.send_json(proto.error_event(
proto.ERR_UNAUTHORIZED, "Invalid secret"))
return False
requested = message.get("protocol_version")
if requested is not None and int(requested) != proto.PROTOCOL_VERSION:
await client.send_json(proto.error_event(
proto.ERR_VERSION_MISMATCH,
f"Protocol version mismatch: server={proto.PROTOCOL_VERSION}, "
f"client={requested}"))
return False
client.authenticated = True
client.protocol_version = proto.PROTOCOL_VERSION
logger.info("[pheby] client %s authenticated", client.conn_id)
return True
# ── broadcast ────────────────────────────────────────────────────────
async def broadcast(self, payload: Dict[str, Any]) -> None:
"""Send an event to every authenticated client."""
for client in list(self._clients.values()):
if client.authenticated and not client.closed:
await client.send_json(payload)
def has_clients(self) -> bool:
return any(c.authenticated and not c.closed
for c in self._clients.values())
# ── inbound dispatch ─────────────────────────────────────────────────
async def handle_client_message(self, client: ClientConnection,
raw: str) -> None:
message, err = proto.decode_message(raw)
if err or message is None:
await client.send_json(proto.error_event(
err or proto.ERR_BAD_REQUEST, "Malformed message"))
return
mtype = message.get("type", "")
request_id = message.get("request_id")
if self.config.debug:
# Verbose protocol logging — never logs secrets; chat content
# only when explicitly configured (privacy default).
safe = {k: v for k, v in message.items()
if k not in ("secret",)}
if not self.config.log_chat_content and mtype == proto.C_CHAT_SEND:
safe = dict(safe)
safe["text"] = f"<{len(str(message.get('text', '')))} chars>"
logger.info("[pheby] << %s", proto.safe_str(safe, 400))
try:
handler = self._HANDLERS.get(mtype)
if handler is None:
await client.send_json(proto.error_event(
proto.ERR_UNKNOWN_TYPE, f"Unknown message type: {mtype}",
request_id))
return
await handler(self, client, message, request_id)
except Exception:
logger.error("[pheby] handler failed for %s", mtype, exc_info=True)
await client.send_json(proto.error_event(
proto.ERR_INTERNAL, "Internal server error", request_id))
# ── simple handlers ──────────────────────────────────────────────────
async def _handle_ping(self, client: ClientConnection, message: Dict,
request_id: Optional[str]) -> None:
await client.send_json({"type": proto.S_PONG,
"ts": proto.now_iso(),
**({"request_id": request_id}
if request_id else {})})
async def _handle_conversation_list(self, client, message, request_id):
conversations = await self.bridge.list_conversations()
await client.send_json({
"type": proto.S_CONVERSATION_SNAPSHOT,
"conversations": conversations,
**({"request_id": request_id} if request_id else {}),
})
async def _handle_conversation_open(self, client, message, request_id):
conversation_id = str(message.get("conversation_id", ""))
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
limit = message.get("limit", proto.MAX_HISTORY_MESSAGES)
try:
limit = max(1, min(int(limit), proto.MAX_HISTORY_MESSAGES))
except (TypeError, ValueError):
limit = proto.MAX_HISTORY_MESSAGES
history, found = await self.bridge.conversation_history(
conversation_id, limit)
if not found:
await client.send_json(proto.error_event(
proto.ERR_CONVERSATION_NOT_FOUND,
"Conversation not found", request_id))
return
await client.send_json({
"type": proto.S_CONVERSATION_HISTORY,
"conversation_id": conversation_id,
"messages": history,
**({"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)
await client.send_json({
"type": proto.S_CONVERSATION_CREATED,
"conversation_id": cid,
"name": await self.router.get_name(cid),
**({"request_id": request_id} if request_id else {}),
})
await self.broadcast({
"type": proto.S_CONVERSATION_UPDATED,
"conversation_id": cid,
"name": await self.router.get_name(cid),
})
async def _handle_conversation_rename(self, client, message, request_id):
conversation_id = str(message.get("conversation_id", ""))
name = str(message.get("name", "")).strip()
if not ConversationRouter.is_valid_conversation_id(conversation_id) \
or not name or len(name) > 200:
await client.send_json(proto.error_event(
proto.ERR_BAD_REQUEST,
"conversation_id and name (≤200 chars) required", request_id))
return
ok = await self.router.rename(conversation_id, name)
if not ok:
await client.send_json(proto.error_event(
proto.ERR_CONVERSATION_NOT_FOUND, "Conversation not found",
request_id))
return
event = {
"type": proto.S_CONVERSATION_RENAMED,
"conversation_id": conversation_id,
"name": name,
**({"request_id": request_id} if request_id else {}),
}
await client.send_json(event)
await self.broadcast({k: v for k, v in event.items()
if k != "request_id"})
async def _handle_conversation_delete(self, client, message, request_id):
conversation_id = str(message.get("conversation_id", ""))
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
ok = await self.bridge.delete_conversation(conversation_id)
if not ok:
await client.send_json(proto.error_event(
proto.ERR_CONVERSATION_NOT_FOUND, "Conversation not found",
request_id))
return
event = {
"type": proto.S_CONVERSATION_DELETED,
"conversation_id": conversation_id,
**({"request_id": request_id} if request_id else {}),
}
await client.send_json(event)
await self.broadcast({k: v for k, v in event.items()
if k != "request_id"})
async def _handle_chat_send(self, client, message, request_id):
conversation_id = str(message.get("conversation_id", ""))
text = message.get("text")
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
if not isinstance(text, str) or not text.strip():
await client.send_json(proto.error_event(
proto.ERR_BAD_REQUEST, "text is required", request_id))
return
if len(text) > proto.MAX_TEXT_CHARS:
await client.send_json(proto.error_event(
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)
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"}}),
**({"request_id": request_id} if request_id else {}),
})
async def _handle_approval_respond(self, client, message, request_id):
approval_id = str(message.get("approval_id", ""))
choice = str(message.get("choice", ""))
reason = message.get("reason")
resolved = await self.bridge.resolve_approval(
approval_id, choice, str(reason) if reason else None)
if not resolved:
await client.send_json(proto.error_event(
proto.ERR_APPROVAL_NOT_FOUND,
"Unknown or already-resolved approval", request_id))
return
await client.send_json({
"type": proto.S_APPROVAL_RESOLVED,
"approval_id": approval_id,
"choice": choice,
**({"request_id": request_id} if request_id else {}),
})
async def _handle_clarify_respond(self, client, message, request_id):
clarify_id = str(message.get("clarify_id", ""))
response = message.get("response")
resolved = await self.bridge.resolve_clarify(
clarify_id, str(response) if response is not None else "")
if not resolved:
await client.send_json(proto.error_event(
proto.ERR_CLARIFY_NOT_FOUND,
"Unknown or already-resolved clarification", request_id))
return
await client.send_json({
"type": proto.S_CLARIFY_RESOLVED,
"clarify_id": clarify_id,
**({"request_id": request_id} if request_id else {}),
})
async def _handle_models_list(self, client, message, request_id):
snapshot = await self.bridge.models_snapshot()
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()
snapshot["type"] = proto.S_MODEL_CURRENT_SNAPSHOT
if request_id:
snapshot["request_id"] = request_id
await client.send_json(snapshot)
async def _handle_model_set(self, client, message, request_id):
model = str(message.get("model", "")).strip()
provider = message.get("provider")
conversation_id = message.get("conversation_id")
result = await self.bridge.set_model(
model, str(provider) if provider else None,
str(conversation_id) if conversation_id else None)
if not result.get("ok"):
await client.send_json(proto.error_event(
result.get("code", proto.ERR_BAD_REQUEST),
result.get("message", "Model change failed"), request_id))
return
event = {
"type": proto.S_MODEL_CHANGED,
"model": result.get("model"),
"provider": result.get("provider"),
"scope": result.get("scope", "global"),
**({"request_id": request_id} if request_id else {}),
}
await client.send_json(event)
await self.broadcast({k: v for k, v in event.items()
if k != "request_id"})
async def _handle_reasoning_current(self, client, message, request_id):
snapshot = await self.bridge.reasoning_snapshot()
snapshot["type"] = proto.S_REASONING_SNAPSHOT
if request_id:
snapshot["request_id"] = request_id
await client.send_json(snapshot)
async def _handle_reasoning_set(self, client, message, request_id):
effort = str(message.get("effort", "")).strip().lower()
conversation_id = message.get("conversation_id")
result = await self.bridge.set_reasoning(
effort, str(conversation_id) if conversation_id else None)
if not result.get("ok"):
await client.send_json(proto.error_event(
result.get("code", proto.ERR_BAD_REQUEST),
result.get("message", "Reasoning change failed"), request_id))
return
event = {
"type": proto.S_REASONING_CHANGED,
"effort": result.get("effort"),
"scope": result.get("scope", "global"),
**({"request_id": request_id} if request_id else {}),
}
await client.send_json(event)
await self.broadcast({k: v for k, v in event.items()
if k != "request_id"})
_HANDLERS = {
proto.C_PING: _handle_ping,
proto.C_CONVERSATION_LIST: _handle_conversation_list,
proto.C_CONVERSATION_OPEN: _handle_conversation_open,
proto.C_CONVERSATION_CREATE: _handle_conversation_create,
proto.C_CONVERSATION_RENAME: _handle_conversation_rename,
proto.C_CONVERSATION_DELETE: _handle_conversation_delete,
proto.C_CHAT_SEND: _handle_chat_send,
proto.C_RUN_CANCEL: _handle_run_cancel,
proto.C_APPROVAL_RESPOND: _handle_approval_respond,
proto.C_CLARIFY_RESPOND: _handle_clarify_respond,
proto.C_MODELS_LIST: _handle_models_list,
proto.C_MODEL_SET: _handle_model_set,
proto.C_MODEL_CURRENT: _handle_model_current,
proto.C_REASONING_SET: _handle_reasoning_set,
proto.C_REASONING_CURRENT: _handle_reasoning_current,
}
__all__ = ["PhebyServer"]
+72
View File
@@ -0,0 +1,72 @@
"""Per-client WebSocket connection state and inbound dispatch."""
from __future__ import annotations
import asyncio
import logging
import time
from typing import Any, Dict, Optional
from aiohttp import web, WSMsgType
from . import protocol as proto
logger = logging.getLogger(__name__)
class ClientConnection:
"""One authenticated WebSocket client (single-user app: at most a few)."""
def __init__(self, ws: "web.WebSocketResponse", conn_id: str):
self.ws = ws
self.conn_id = conn_id
self.authenticated = False
self.protocol_version: Optional[int] = None
self.connected_at = time.time()
# Serialized outbound writes — aiohttp WS frames are not safe to
# compose concurrently from multiple tasks.
self._send_lock = asyncio.Lock()
self.closed = False
async def send_json(self, payload: Dict[str, Any]) -> bool:
"""Thread-safe send; returns False when the socket is going away."""
if self.closed:
return False
try:
async with self._send_lock:
await self.ws.send_str(proto.encode_message(payload))
return True
except (ConnectionError, RuntimeError, asyncio.CancelledError):
self.closed = True
return False
# ------------------------------------------------------------------
async def read_loop(self, server: Any) -> None:
"""Receive/dispatch loop; exits on close, error, or auth timeout."""
try:
async for msg in self.ws:
if msg.type == WSMsgType.TEXT:
if len(msg.data) > proto.MAX_WS_MESSAGE_BYTES:
await self.send_json(proto.error_event(
proto.ERR_TOO_LARGE,
"WebSocket message exceeds server limit"))
continue
await server.handle_client_message(self, msg.data)
elif msg.type == WSMsgType.BINARY:
await self.send_json(proto.error_event(
proto.ERR_BAD_REQUEST,
"Binary frames are not part of the Pheby protocol"))
elif msg.type in (WSMsgType.CLOSE, WSMsgType.CLOSING,
WSMsgType.CLOSED):
break
elif msg.type == WSMsgType.ERROR:
logger.warning("[pheby] ws error on %s: %s",
self.conn_id, msg.data)
break
except asyncio.CancelledError:
raise
except Exception:
logger.warning("[pheby] client %s read loop crashed",
self.conn_id, exc_info=True)
finally:
self.closed = True