feat(pheby): accept inbound chat attachments
This commit is contained in:
@@ -179,6 +179,7 @@ class AttachmentStore:
|
||||
"kind": kind,
|
||||
"conversation_id": conversation_id,
|
||||
"message_id": message_id,
|
||||
"direction": "outbound",
|
||||
"created_at": protocol.now_iso(),
|
||||
"created_epoch": time.time(),
|
||||
"retention_days": self._retention_days,
|
||||
@@ -197,6 +198,83 @@ class AttachmentStore:
|
||||
attachment_id, kind, size, conversation_id)
|
||||
return self.describe(attachment_id)
|
||||
|
||||
async def register_bytes(
|
||||
self,
|
||||
data: bytes,
|
||||
*,
|
||||
conversation_id: str,
|
||||
filename: str,
|
||||
mime_type: Optional[str] = None,
|
||||
message_id: Optional[str] = None,
|
||||
) -> Optional[Dict[str, Any]]:
|
||||
"""Register an *inbound* upload: raw client bytes → adapter storage.
|
||||
|
||||
Returns the attachment descriptor dict, or ``None`` on failure.
|
||||
"""
|
||||
fname = self._sanitize_filename(filename)
|
||||
attachment_id = protocol.new_id()
|
||||
mime = (mime_type or "").strip()
|
||||
if mime in ("", "application/octet-stream"):
|
||||
mime = guess_mime(fname)
|
||||
if not isinstance(mime, str) or "/" not in mime or len(mime) > 120:
|
||||
mime = "application/octet-stream"
|
||||
is_image = mime.startswith(IMAGE_MIME_PREFIXES) or (
|
||||
Path(fname).suffix.lower() in IMAGE_EXTS)
|
||||
if Path(fname).suffix.lower() in VIDEO_EXTS:
|
||||
kind = "video"
|
||||
elif Path(fname).suffix.lower() in AUDIO_EXTS:
|
||||
kind = "audio"
|
||||
elif is_image:
|
||||
kind = "image"
|
||||
else:
|
||||
kind = "document"
|
||||
size = len(data)
|
||||
|
||||
try:
|
||||
async with self._lock:
|
||||
self._load_index()
|
||||
blob = self._blob_path(attachment_id, fname)
|
||||
blob.parent.mkdir(parents=True, exist_ok=True)
|
||||
await asyncio.to_thread(blob.write_bytes, data)
|
||||
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,
|
||||
"direction": "inbound",
|
||||
"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] inbound attachment registration failed",
|
||||
exc_info=True)
|
||||
return None
|
||||
|
||||
logger.info(
|
||||
"[pheby] inbound attachment registered: id=%s kind=%s size=%d "
|
||||
"conv=%s", attachment_id, kind, size, conversation_id)
|
||||
return self.describe(attachment_id)
|
||||
|
||||
def anchor_message(self, attachment_id: str, message_id: str) -> None:
|
||||
"""Attach *message_id* to a registered attachment (thread anchoring)."""
|
||||
self._load_index()
|
||||
meta = self._meta.get(attachment_id)
|
||||
if meta is not None:
|
||||
meta["message_id"] = str(message_id)
|
||||
self._save_index()
|
||||
|
||||
def get_blob_path(self, attachment_id: str) -> Optional[Path]:
|
||||
"""Filesystem path for an attachment (server-side use only)."""
|
||||
return self.resolve_blob(attachment_id)
|
||||
|
||||
# ── lookup / download ────────────────────────────────────────────────
|
||||
def describe(self, attachment_id: str) -> Optional[Dict[str, Any]]:
|
||||
"""Public descriptor for an attachment (no server paths)."""
|
||||
@@ -210,6 +288,7 @@ class AttachmentStore:
|
||||
"size": int(meta.get("size", 0)),
|
||||
"kind": meta.get("kind", "document"),
|
||||
"inline_image": meta.get("kind") == "image",
|
||||
"direction": meta.get("direction", "outbound"),
|
||||
"conversation_id": meta.get("conversation_id"),
|
||||
"message_id": meta.get("message_id"),
|
||||
"created_at": meta.get("created_at"),
|
||||
@@ -261,7 +340,8 @@ class AttachmentStore:
|
||||
out = []
|
||||
for aid in list(self._meta):
|
||||
desc = self.describe(aid)
|
||||
if desc and desc.get("conversation_id") == conversation_id:
|
||||
if desc and desc.get("conversation_id") == conversation_id and (
|
||||
desc.get("direction") != "inbound" or desc.get("message_id")):
|
||||
out.append(desc)
|
||||
return out
|
||||
|
||||
|
||||
@@ -225,6 +225,16 @@ def _iso(value: Any) -> Optional[str]:
|
||||
return None
|
||||
|
||||
|
||||
def _display_user_text(text: str) -> str:
|
||||
"""Hide agent-only inline file context from the user-facing transcript."""
|
||||
start = "[Pheby user message]\n"
|
||||
end = "\n[/Pheby user message]"
|
||||
if start not in text:
|
||||
return text
|
||||
body = text.split(start, 1)[1]
|
||||
return body.split(end, 1)[0] if end in body else text
|
||||
|
||||
|
||||
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.
|
||||
@@ -259,6 +269,8 @@ async def conversation_history(conversation_id: str, limit: int
|
||||
continue
|
||||
content = row.get("content")
|
||||
text = content if isinstance(content, str) else str(content or "")
|
||||
if role == "user":
|
||||
text = _display_user_text(text)
|
||||
# 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":
|
||||
@@ -374,12 +386,18 @@ async def delete_conversation(conversation_id: str) -> bool:
|
||||
# Chat runs
|
||||
# ═══════════════════════════════════════════════════════════════════════════
|
||||
async def send_chat(server: Any, conversation_id: str, text: str,
|
||||
client: Any, request_id: Optional[str]) -> None:
|
||||
client: Any, request_id: Optional[str],
|
||||
attachment_ids: Optional[List[str]] = None) -> 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.
|
||||
|
||||
*attachment_ids* (optional) reference inbound uploads already registered
|
||||
in ``server.store``; they are anchored to the user message, handed to
|
||||
Hermes as ``media_urls``/``media_types`` (tool-accessible local files),
|
||||
and small text files are additionally inlined into the message text.
|
||||
"""
|
||||
adapter = _current_adapter()
|
||||
if adapter is None or not hasattr(adapter, "handle_message"):
|
||||
@@ -392,13 +410,68 @@ async def send_chat(server: Any, conversation_id: str, text: str,
|
||||
|
||||
run_id = uuid.uuid4().hex[:16]
|
||||
source = _source_for(conversation_id)
|
||||
message_id = uuid.uuid4().hex[:12]
|
||||
|
||||
# ── inbound attachments ────────────────────────────────────────────────
|
||||
media_urls: List[str] = []
|
||||
media_types: List[str] = []
|
||||
media_text_inlined: List[bool] = []
|
||||
inline_blocks: List[str] = []
|
||||
attachment_descs: List[Dict[str, Any]] = []
|
||||
for aid in attachment_ids or []:
|
||||
desc = server.store.describe(aid)
|
||||
if desc is None or desc.get("conversation_id") != conversation_id or \
|
||||
desc.get("direction") != "inbound" or desc.get("message_id"):
|
||||
await client.send_json(proto.error_event(
|
||||
proto.ERR_NOT_FOUND,
|
||||
"Unknown or already sent attachment for this conversation", request_id))
|
||||
return
|
||||
attachment_descs.append(desc)
|
||||
for desc in attachment_descs:
|
||||
aid = desc["attachment_id"]
|
||||
blob = server.store.get_blob_path(aid)
|
||||
if blob is None:
|
||||
await client.send_json(proto.error_event(
|
||||
proto.ERR_NOT_FOUND, "Attachment unavailable", request_id))
|
||||
return
|
||||
media_urls.append(str(blob))
|
||||
mime = str(desc.get("mime_type", "application/octet-stream"))
|
||||
media_types.append(mime)
|
||||
# Small text files become part of the message text so the model
|
||||
# reads them directly without a tool round-trip.
|
||||
fname = str(desc.get("filename", "file"))
|
||||
inlined = False
|
||||
if mime.startswith("text/") and int(desc.get("size", 0)) <= \
|
||||
proto.INLINE_TEXT_BYTES:
|
||||
try:
|
||||
content = blob.read_text(encoding="utf-8", errors="replace")
|
||||
inline_blocks.append(
|
||||
f"Attached file: {fname}\n```{fname.rsplit('.', 1)[-1]}\n"
|
||||
f"{content}\n```")
|
||||
inlined = True
|
||||
except OSError:
|
||||
pass
|
||||
media_text_inlined.append(inlined)
|
||||
|
||||
effective_text = text
|
||||
if attachment_descs:
|
||||
effective_text = f"[Pheby user message]\n{text}\n[/Pheby user message]"
|
||||
if inline_blocks:
|
||||
effective_text += "\n\n" + "\n\n".join(inline_blocks)
|
||||
|
||||
from gateway.platforms.base import MessageEvent, MessageType
|
||||
event = MessageEvent(
|
||||
text=text,
|
||||
message_type=MessageType.TEXT,
|
||||
text=effective_text,
|
||||
message_type=(MessageType.PHOTO if media_types and
|
||||
all(t.startswith("image/") for t in media_types)
|
||||
else (MessageType.DOCUMENT if media_types
|
||||
else MessageType.TEXT)),
|
||||
source=source,
|
||||
message_id=uuid.uuid4().hex[:12],
|
||||
message_id=message_id,
|
||||
metadata={"pheby_run_id": run_id},
|
||||
media_urls=media_urls,
|
||||
media_types=media_types,
|
||||
media_text_inlined=media_text_inlined,
|
||||
)
|
||||
|
||||
with _RUN_LOCK:
|
||||
@@ -438,6 +511,14 @@ async def send_chat(server: Any, conversation_id: str, text: str,
|
||||
# returns quickly; the eventual reply arrives through adapter.send().
|
||||
try:
|
||||
await adapter.handle_message(event)
|
||||
for desc in attachment_descs:
|
||||
aid = desc["attachment_id"]
|
||||
server.store.anchor_message(aid, message_id)
|
||||
await server.broadcast({
|
||||
"type": proto.S_ATTACHMENT_ADDED,
|
||||
"conversation_id": conversation_id,
|
||||
"attachment": server.store.describe(aid),
|
||||
})
|
||||
except Exception:
|
||||
with _RUN_LOCK:
|
||||
current = _ACTIVE_RUNS.get(conversation_id)
|
||||
|
||||
@@ -36,6 +36,8 @@ 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
|
||||
MAX_UPLOAD_BYTES = 64 * 1024 * 1024 # inbound attachment upload cap
|
||||
INLINE_TEXT_BYTES = 100 * 1024 # text/* files below this are inlined
|
||||
|
||||
# ── Error codes (machine-readable) ───────────────────────────────────────────
|
||||
ERR_UNAUTHORIZED = "unauthorized"
|
||||
|
||||
+82
-2
@@ -61,11 +61,13 @@ class PhebyServer:
|
||||
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 = _web.Application(client_max_size=max(proto.MAX_UPLOAD_BYTES,
|
||||
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)
|
||||
app.router.add_post("/attachments", self._handle_attachment_upload)
|
||||
self._app = app
|
||||
self._runner = web.AppRunner(app, access_log=None)
|
||||
await self._runner.setup()
|
||||
@@ -171,6 +173,57 @@ class PhebyServer:
|
||||
},
|
||||
)
|
||||
|
||||
async def _handle_attachment_upload(
|
||||
self, request: web.Request) -> web.Response:
|
||||
"""Inbound attachment upload: raw body → adapter storage.
|
||||
|
||||
The descriptor is returned only to the uploader. Other clients see
|
||||
the attachment after chat.send successfully claims it.
|
||||
"""
|
||||
if not self._check_http_secret(request):
|
||||
return web.json_response(
|
||||
{"error": {"code": proto.ERR_UNAUTHORIZED,
|
||||
"message": "Authentication required"}},
|
||||
status=401)
|
||||
conversation_id = request.query.get("conversation_id", "")
|
||||
if not ConversationRouter.is_valid_conversation_id(conversation_id):
|
||||
return web.json_response(
|
||||
{"error": {"code": proto.ERR_BAD_REQUEST,
|
||||
"message": "Invalid conversation_id"}},
|
||||
status=400)
|
||||
declared = request.content_length
|
||||
if declared is not None and declared > proto.MAX_UPLOAD_BYTES:
|
||||
return web.json_response(
|
||||
{"error": {"code": proto.ERR_TOO_LARGE,
|
||||
"message": f"Upload exceeds "
|
||||
f"{proto.MAX_UPLOAD_BYTES} bytes"}},
|
||||
status=413)
|
||||
data = await request.content.read(proto.MAX_UPLOAD_BYTES + 1)
|
||||
if len(data) > proto.MAX_UPLOAD_BYTES:
|
||||
return web.json_response(
|
||||
{"error": {"code": proto.ERR_TOO_LARGE,
|
||||
"message": f"Upload exceeds "
|
||||
f"{proto.MAX_UPLOAD_BYTES} bytes"}},
|
||||
status=413)
|
||||
if not data:
|
||||
return web.json_response(
|
||||
{"error": {"code": proto.ERR_BAD_REQUEST,
|
||||
"message": "Empty upload body"}},
|
||||
status=400)
|
||||
filename = request.query.get("filename") or "file.bin"
|
||||
mime = request.headers.get("Content-Type", "").split(";")[0].strip()
|
||||
desc = await self.store.register_bytes(
|
||||
data, conversation_id=conversation_id, filename=filename,
|
||||
mime_type=mime or None)
|
||||
if desc is None:
|
||||
return web.json_response(
|
||||
{"error": {"code": proto.ERR_INTERNAL,
|
||||
"message": "Attachment registration failed"}},
|
||||
status=500)
|
||||
logger.info("[pheby] attachment upload: id=%s bytes=%d conv=%s",
|
||||
desc["attachment_id"], desc["size"], conversation_id)
|
||||
return web.json_response({"attachment": desc}, status=201)
|
||||
|
||||
# ── WebSocket handler ────────────────────────────────────────────────
|
||||
async def _handle_ws(self, request: web.Request) -> web.WebSocketResponse:
|
||||
# Bind server/adapter identity to THIS task's context. ContextVars set
|
||||
@@ -468,8 +521,35 @@ class PhebyServer:
|
||||
proto.ERR_TOO_LARGE,
|
||||
f"text exceeds {proto.MAX_TEXT_CHARS} chars", request_id))
|
||||
return
|
||||
raw_ids = message.get("attachment_ids")
|
||||
attachment_ids: List[str] = []
|
||||
if raw_ids is not None:
|
||||
if not isinstance(raw_ids, list) or \
|
||||
not all(isinstance(x, str) for x in raw_ids) or \
|
||||
len(raw_ids) > 10 or len(raw_ids) != len(set(raw_ids)):
|
||||
await client.send_json(proto.error_event(
|
||||
proto.ERR_BAD_REQUEST,
|
||||
"attachment_ids must be at most 10 unique ids",
|
||||
request_id))
|
||||
return
|
||||
for aid in raw_ids:
|
||||
desc = self.store.describe(aid)
|
||||
if desc is None or \
|
||||
desc.get("conversation_id") != conversation_id or \
|
||||
desc.get("direction") != "inbound":
|
||||
await client.send_json(proto.error_event(
|
||||
proto.ERR_NOT_FOUND,
|
||||
"Unknown attachment for this conversation", request_id))
|
||||
return
|
||||
if desc.get("message_id"):
|
||||
await client.send_json(proto.error_event(
|
||||
proto.ERR_BAD_REQUEST,
|
||||
"Attachment already belongs to a message", request_id))
|
||||
return
|
||||
attachment_ids.append(aid)
|
||||
await self.bridge.send_chat(
|
||||
self, conversation_id, text, client, request_id)
|
||||
self, conversation_id, text, client, request_id,
|
||||
attachment_ids=attachment_ids)
|
||||
|
||||
async def _handle_run_cancel(self, client, message, request_id):
|
||||
conversation_id = str(message.get("conversation_id", ""))
|
||||
|
||||
Reference in New Issue
Block a user