adapter.py was a 3,493-line monolith. Split it into focused modules with clear separation of responsibilities, bringing it down to ~857 lines: - Module-level helpers: hooks, classify, pickers, commands, setup, defaults, secrets - Frame-handler mixins: inbound, tool_frames, push_frames, media_frames, picker_frames, channel_frames, query_frames - mixin_base: IrisAdapterBase (declaration-only base for shared attrs) - adapter.py now holds only IrisAdapter (the composition of the 7 mixins + BasePlatformAdapter), register(), and test-facing re-exports The mixins come before BasePlatformAdapter in the MRO so their methods override the base; super() calls (e.g. send_image) still resolve to BasePlatformAdapter. No circular imports; dispatch.py and http_server.py (instance-method callers) are unaffected. Ruff complexity ceilings (PLR0911/0912/0913/0915) restored to Ruff's built-in defaults (12/50/6/5) instead of "just above the current maxima", which ratchets the bar down as code grows. The existing genuinely-complex functions (frame builders mirroring the wire schema, the QR matrix builder, the dispatch table) carry an explicit `# noqa: PLR09xx` marking them as reviewed, frozen exceptions; new code is held to the default ceilings. All 125 tests green (94 test_android + 31 test_android_http); no new ruff errors introduced.
256 lines
10 KiB
Python
256 lines
10 KiB
Python
"""Inbound ``message.send`` handling (app -> agent).
|
|
|
|
Mixin for ``adapter.IrisAdapter``. Echoes the user message to all devices
|
|
(multi-device sync + ack), resolves ``media_refs`` to
|
|
``MessageEvent.media_urls``/``media_types``, applies auto-threading, and
|
|
hands the ``MessageEvent`` to ``handle_message()`` (the gateway's command
|
|
pipeline + agent turn).
|
|
"""
|
|
|
|
import asyncio
|
|
import logging
|
|
import threading
|
|
import time
|
|
import uuid
|
|
from typing import Any
|
|
|
|
from gateway.platforms.base import MessageEvent, MessageType
|
|
|
|
from . import protocol
|
|
from .classify import _derive_thread_name
|
|
from .mixin_base import IrisAdapterBase
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
class InboundHandlers(IrisAdapterBase):
|
|
"""Inbound message.send (see module docstring)."""
|
|
|
|
async def on_message_send(self, frame: protocol.Frame, device_id: str) -> None: # noqa: PLR0912,PLR0915
|
|
"""Handle an inbound ``message.send`` frame.
|
|
|
|
Echoes the user message to all devices (multi-device sync + ack),
|
|
then builds a ``MessageEvent`` and hands it to ``handle_message()``
|
|
(the gateway's command pipeline + agent turn).
|
|
|
|
M4: ``media_refs`` reference completed ``media.upload``s; they are
|
|
resolved to ``MessageEvent.media_urls``/``media_types`` (local paths
|
|
the agent's vision/audio tools can read) and echoed in the user
|
|
message's ``media[]`` so every device renders the attachments.
|
|
"""
|
|
payload = frame.payload
|
|
text = payload.get("text")
|
|
text = text if isinstance(text, str) else ""
|
|
|
|
refs_raw = payload.get("media_refs")
|
|
media_refs = (
|
|
[r for r in refs_raw if isinstance(r, str) and r] if isinstance(refs_raw, list) else []
|
|
)
|
|
|
|
if not text.strip() and not media_refs:
|
|
await self._reply(
|
|
device_id,
|
|
protocol.error(
|
|
protocol.ERR_UNSUPPORTED, "message.send requires non-empty text", id=frame.id
|
|
),
|
|
)
|
|
return
|
|
|
|
chat_id = frame.chat_id or payload.get("chat_id")
|
|
if not isinstance(chat_id, str) or not chat_id.strip():
|
|
chat_id = self.home_channel
|
|
chat_id = chat_id.strip()
|
|
|
|
# Automation channels are read-only for the user: they only receive
|
|
# gateway-originated output (cron jobs, webhooks). Reject direct sends
|
|
# (the app hides the composer for them, this is the server-side
|
|
# enforcement).
|
|
target = self._channels.get(chat_id)
|
|
if target is not None and target.get("automation"):
|
|
await self._reply(
|
|
device_id,
|
|
protocol.error(
|
|
protocol.ERR_UNSUPPORTED,
|
|
f"{target.get('name') or chat_id} is an automation channel "
|
|
"(read-only: cron/webhook output only)",
|
|
id=frame.id,
|
|
),
|
|
)
|
|
return
|
|
|
|
thread_id = frame.thread_id or payload.get("thread_id")
|
|
if not isinstance(thread_id, str) or not thread_id.strip():
|
|
thread_id = None
|
|
|
|
reply_to = payload.get("reply_to")
|
|
if not isinstance(reply_to, str) or not reply_to.strip():
|
|
reply_to = None
|
|
|
|
# Auto-threading (the app's Threads setting, docs/06 §6.3): a message
|
|
# in a channel's flat lane gets its own fresh thread, the way Telegram
|
|
# topic mode mints a topic per new conversation. The thread is named
|
|
# instantly from the user's opening message (derived title) and the
|
|
# LLM upgrades the name in the background. The user echo, the agent
|
|
# turn, and all streaming frames then carry the new thread_id.
|
|
# Skipped for slash commands (session-scoped, not conversation
|
|
# starters) and replies (they continue where the user is). Threading
|
|
# is only active on the default channel; other channels stay flat.
|
|
auto_thread = bool(payload.get("auto_thread"))
|
|
default_entry = self._channels.default()
|
|
if (
|
|
auto_thread
|
|
and thread_id is None
|
|
and text.strip()
|
|
and not text.lstrip().startswith("/")
|
|
and reply_to is None
|
|
and default_entry is not None
|
|
and chat_id == default_entry["chat_id"]
|
|
):
|
|
entry = self._channels.create(
|
|
name=_derive_thread_name(text),
|
|
kind="thread",
|
|
parent_chat_id=chat_id,
|
|
)
|
|
thread_id = entry["chat_id"]
|
|
# Bare broadcast (like channel.create): the directory is
|
|
# re-served on hello.ack, so no outbox entry is needed.
|
|
await self._broadcast_both(protocol.channel_created(entry, auto=True))
|
|
self._schedule_thread_title_upgrade(entry["chat_id"], text)
|
|
|
|
# M4: resolve media refs (single-use; unknown ref -> error).
|
|
media_urls: list[str] = []
|
|
media_types: list[str] = []
|
|
media_wire: list[dict[str, Any]] = []
|
|
for ref in media_refs:
|
|
entry = self._media.get_inbound(ref)
|
|
if entry is None:
|
|
await self._reply(
|
|
device_id,
|
|
protocol.error(
|
|
protocol.ERR_UNSUPPORTED, f"unknown media_ref {ref}", id=frame.id
|
|
),
|
|
)
|
|
return
|
|
media_urls.append(entry.path)
|
|
media_types.append(entry.mime)
|
|
media_wire.append(
|
|
{
|
|
"media_id": entry.media_id,
|
|
"kind": entry.kind,
|
|
"mime": entry.mime,
|
|
"size": entry.size,
|
|
"filename": entry.filename,
|
|
}
|
|
)
|
|
|
|
device = self._devices.get(device_id) or {}
|
|
user_name = device.get("name") or device_id
|
|
|
|
# Echo to all devices: the sender confirms (server-assigned id),
|
|
# other devices see the message too (single-user, multi-device).
|
|
# Routed through _broadcast_or_log (not a bare broadcast) so the echo
|
|
# is appended to the outbox: the app's ChatStore is in-memory only, so
|
|
# after a process death / activity recreation the only way the user's
|
|
# own message is restored is via the sync replay. Without this, user
|
|
# messages vanish on reconnect while bot messages (already parked)
|
|
# survive.
|
|
message_id = f"m_{uuid.uuid4().hex[:16]}"
|
|
echo = protocol.message(
|
|
chat_id=chat_id,
|
|
message_id=message_id,
|
|
role=protocol.ROLE_USER,
|
|
text=text,
|
|
thread_id=thread_id,
|
|
media=media_wire or None,
|
|
reply_to=reply_to,
|
|
ts=int(time.time() * 1000),
|
|
)
|
|
await self._broadcast_or_log(chat_id, echo)
|
|
# Refs are consumed by this message (no replay).
|
|
for ref in media_refs:
|
|
self._media.pop_inbound(ref)
|
|
|
|
# M4: a new user turn starts -- stale offer association is dropped.
|
|
self._last_message_id.pop(chat_id, None)
|
|
|
|
kind = media_wire[0]["kind"] if media_wire else None
|
|
if kind == "image":
|
|
message_type = MessageType.PHOTO
|
|
elif kind == "video":
|
|
message_type = MessageType.VIDEO
|
|
elif kind == "audio":
|
|
message_type = MessageType.AUDIO
|
|
elif kind == "voice":
|
|
message_type = MessageType.VOICE
|
|
elif kind == "document":
|
|
message_type = MessageType.DOCUMENT
|
|
else:
|
|
message_type = MessageType.TEXT
|
|
|
|
source = self.build_source(
|
|
chat_id=chat_id,
|
|
chat_name=self._channel_name(chat_id),
|
|
chat_type="dm",
|
|
user_id=device_id,
|
|
user_name=user_name,
|
|
thread_id=thread_id,
|
|
)
|
|
event = MessageEvent(
|
|
text=text,
|
|
message_type=message_type,
|
|
user_id=device_id,
|
|
user_name=user_name,
|
|
source=source,
|
|
message_id=message_id,
|
|
reply_to_message_id=reply_to,
|
|
media_urls=media_urls,
|
|
media_types=media_types,
|
|
)
|
|
await self.handle_message(event)
|
|
# M5: acknowledge the user message to the originating device (the
|
|
# app shows ✓✓) at the moment it is handed to the agent.
|
|
await self._reply(
|
|
device_id,
|
|
protocol.read_receipt(chat_id, message_id),
|
|
)
|
|
|
|
def _schedule_thread_title_upgrade(self, thread_id: str, text: str) -> None:
|
|
"""Upgrade an auto-created thread's name with the model's title.
|
|
|
|
Stage 2 of hermes' two-stage session titling (``agent/title_generator
|
|
.py``): the thread was created with an instant derived name; this
|
|
background call on the ``title_generation`` auxiliary task replaces it
|
|
with the model's title and broadcasts ``channel.renamed``. Best-effort
|
|
— any failure (config, model, network) leaves the derived name in
|
|
place, and a thread the user already renamed or archived is untouched.
|
|
"""
|
|
loop = asyncio.get_running_loop()
|
|
|
|
def _work() -> None:
|
|
try:
|
|
from agent.title_generator import generate_title
|
|
|
|
title = generate_title(text)
|
|
except Exception:
|
|
logger.debug("Thread title upgrade failed", exc_info=True)
|
|
return
|
|
if not title:
|
|
return
|
|
entry = self._channels.get(thread_id)
|
|
if entry is None or entry.get("archived"):
|
|
return
|
|
if (entry.get("name") or "") == title:
|
|
return
|
|
renamed = self._channels.rename(thread_id, title)
|
|
if renamed is None:
|
|
return
|
|
try:
|
|
asyncio.run_coroutine_threadsafe(
|
|
self._broadcast_both(protocol.channel_renamed(renamed)),
|
|
loop,
|
|
)
|
|
except Exception:
|
|
logger.debug("Thread title rename broadcast failed", exc_info=True)
|
|
|
|
threading.Thread(target=_work, daemon=True, name="iris-thread-title").start()
|