Split adapter.py monolith into focused modules; restore Ruff complexity defaults (issue #12)
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.
This commit is contained in:
1 parent
7faaf2aa1c
commit
b8e756c3dd
26 files changed
+3101
-2775
No files matched your search
@@ -0,0 +1,255 @@
|
||||
"""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()
|
||||
Reference in new issue
Block a user