Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 35 additions & 0 deletions gateway/platforms/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -1022,6 +1022,27 @@ async def handle_message(self, event: MessageEvent) -> None:

# Check if there's already an active handler for this session
if session_key in self._active_sessions:
# Gateway slash commands like /approve, /deny, /stop, /new, /status,
# etc. are control-plane messages. They must bypass adapter-level
# session serialization so the gateway can resolve them immediately
# even while an agent turn is active or blocked on approval.
if event.is_command():
try:
from hermes_cli.commands import resolve_command as _resolve_gateway_command
_cmd = event.get_command()
if _cmd and _resolve_gateway_command(_cmd):
logger.debug("[%s] Dispatching gateway command %s immediately while session %s is active", self.name, _cmd, session_key)
task = asyncio.create_task(self._process_priority_command(event))
try:
self._background_tasks.add(task)
except TypeError:
return
if hasattr(task, "add_done_callback"):
task.add_done_callback(self._background_tasks.discard)
return
except Exception:
pass

# Special case: photo bursts/albums frequently arrive as multiple near-
# simultaneous messages. Queue them without interrupting the active run,
# then process them immediately after the current task finishes.
Expand Down Expand Up @@ -1086,6 +1107,20 @@ def _get_human_delay() -> float:
min_ms, max_ms = 800, 2500
return random.uniform(min_ms / 1000.0, max_ms / 1000.0)

async def _process_priority_command(self, event: MessageEvent) -> None:
"""Process a gateway slash command immediately while another turn is active.

This bypasses adapter-level session serialization for control-plane
commands like /approve, /deny, /stop, /new, and /status so they can
resolve blocked or running turns instead of being queued behind them.
"""
if not self._message_handler:
return
response = await self._message_handler(event)
if response:
_thread_metadata = {"thread_id": event.source.thread_id} if event.source.thread_id else None
await self.send(event.source.chat_id, response, metadata=_thread_metadata)

async def _process_message_background(self, event: MessageEvent, session_key: str) -> None:
"""Background task that actually processes the message."""
# Track delivery outcomes for the processing-complete hook
Expand Down
18 changes: 18 additions & 0 deletions tests/e2e/test_telegram_commands.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
import pytest

from gateway.platforms.base import SendResult
from gateway.session import build_session_key
from tests.e2e.conftest import (
make_adapter,
make_event,
Expand Down Expand Up @@ -86,6 +87,23 @@ async def test_stop_when_no_agent_running(self, adapter):
response_lower = response_text.lower()
assert "no" in response_lower or "stop" in response_lower or "not running" in response_lower

@pytest.mark.asyncio
async def test_approve_bypasses_active_session_serialization(self, adapter):
session_key = build_session_key(make_source())
adapter._active_sessions[session_key] = asyncio.Event()
adapter._message_handler = AsyncMock(return_value="✅ Command approved. The agent is resuming...")
adapter.set_message_handler(adapter._message_handler)

event = make_event("/approve always")
adapter.send.reset_mock()
await adapter.handle_message(event)
await asyncio.sleep(0.3)

adapter._message_handler.assert_awaited()
adapter.send.assert_called_once()
response_text = adapter.send.call_args[1].get("content") or adapter.send.call_args[0][1]
assert "approved" in response_text.lower()

@pytest.mark.asyncio
async def test_commands_shows_listing(self, adapter):
send = await send_and_capture(adapter, "/commands")
Expand Down