Files
ARIA b8e756c3dd
CI / Gateway plugin tests (pull_request) Successful in 4m59s
CI / Kotlin tests (android host + desktop) (pull_request) Successful in 7m5s
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.
2026-08-24 20:55:00 +02:00

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()