fix(gateway): busy-mode steer also reaches the parent's active subagents
A parent blocked inside delegate_task only drains its steer queue after the tool returns, i.e. after the child finishes. On Telegram (busy_input_mode=steer, the /steer command, and the PRIORITY path) the gateway queued the text on that parent alone and acked "Steered into current run", so a looping delegated child never saw it (#112095). Fan the text out to running_agent._active_children — the same identity-scoped snapshot interrupt() uses — and make the ack say where it went.
This commit is contained in:
@@ -221,6 +221,40 @@ class GatewayBusySessionMixin:
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
@staticmethod
|
||||
def _steer_active_subagents(running_agent: Any, text: str) -> int:
|
||||
"""Queue *text* into every live child of *running_agent*; returns how many accepted it.
|
||||
|
||||
A parent blocked inside ``delegate_task`` only drains its own steer queue after the tool
|
||||
returns, i.e. after the child finishes — so a steer aimed at a looping child would sit
|
||||
unread for the whole delegation (#112095, Telegram). Children are the parent's own
|
||||
``_active_children`` (identity-scoped, same snapshot ``interrupt()`` fans out to)."""
|
||||
children = getattr(running_agent, "_active_children", None)
|
||||
if not isinstance(children, (list, tuple, set)) or not children:
|
||||
return 0
|
||||
lock = getattr(running_agent, "_active_children_lock", None)
|
||||
try:
|
||||
with lock if lock is not None else contextlib.nullcontext():
|
||||
snapshot = list(children)
|
||||
except Exception:
|
||||
return 0
|
||||
accepted = 0
|
||||
for child in snapshot:
|
||||
steer = getattr(child, "steer", None)
|
||||
if not callable(steer):
|
||||
continue
|
||||
try:
|
||||
accepted += bool(steer(text))
|
||||
except Exception as exc:
|
||||
logger.warning("Steer into subagent %r failed: %s", getattr(child, "_delegate_id", child), exc)
|
||||
return accepted
|
||||
|
||||
def _steer_running_agent(self, running_agent: Any, text: str) -> bool:
|
||||
"""``running_agent.steer(text)`` plus fan-out to its active subagents (see
|
||||
:meth:`_steer_active_subagents`); True when the parent or any child queued it."""
|
||||
accepted = bool(running_agent.steer(text))
|
||||
return bool(self._steer_active_subagents(running_agent, text)) or accepted
|
||||
|
||||
async def _session_has_compression_in_flight(self, session_key: str) -> bool:
|
||||
"""True when a compression lock is held for this session's id (callers demote interrupt →
|
||||
queue, else a follow-up against the pre-rotation parent orphans compression siblings).
|
||||
@@ -535,6 +569,8 @@ class GatewayBusySessionMixin:
|
||||
"""Call ``running_agent.<verb>(text)`` (steer/redirect); False + warning on failure."""
|
||||
try:
|
||||
call_text = self._steer_text_with_origin(text, event) if event else text
|
||||
if verb == "steer":
|
||||
return self._steer_running_agent(running_agent, call_text)
|
||||
return bool(getattr(running_agent, verb)(call_text))
|
||||
except Exception as exc:
|
||||
logger.warning("Gateway %s failed for session %s: %s", verb, session_key, exc)
|
||||
@@ -615,7 +651,10 @@ class GatewayBusySessionMixin:
|
||||
except Exception:
|
||||
pass
|
||||
status_detail = f" ({', '.join(status_parts)})" if status_parts else ""
|
||||
if is_steer_mode:
|
||||
if is_steer_mode and self._agent_has_active_subagents(running_agent):
|
||||
head = "⏩ Steered into current run and its active subagent(s)"
|
||||
tail = ". Your message arrives after their next tool call."
|
||||
elif is_steer_mode:
|
||||
head, tail = "⏩ Steered into current run", ". Your message arrives after the next tool call."
|
||||
elif is_redirect_mode:
|
||||
head, tail = "↪ Redirected current run", ". I'll adjust using your correction."
|
||||
@@ -939,14 +978,15 @@ class GatewayBusySessionMixin:
|
||||
if not running_agent or not hasattr(running_agent, "steer"):
|
||||
return _queue_fallback("No active agent — /steer queued for the next turn.")
|
||||
try:
|
||||
accepted = running_agent.steer(self._steer_text_with_origin(steer_text, event))
|
||||
accepted = self._steer_running_agent(running_agent, self._steer_text_with_origin(steer_text, event))
|
||||
except Exception as exc:
|
||||
logger.warning("Steer failed for session %s: %s", quick_key, exc)
|
||||
return f"⚠️ Steer failed: {exc}"
|
||||
if not accepted:
|
||||
return "Steer rejected (empty payload)."
|
||||
preview = steer_text[:60] + ("..." if len(steer_text) > 60 else "")
|
||||
return f"⏩ Steer queued — arrives after the next tool call: '{preview}'"
|
||||
target = "run and its active subagent(s)" if self._agent_has_active_subagents(running_agent) else "run"
|
||||
return f"⏩ Steer queued into current {target} — arrives after the next tool call: '{preview}'"
|
||||
|
||||
async def _busy_goal_command(self, event: MessageEvent, quick_key: str, source):
|
||||
# Control verbs are safe mid-run (state only); setting new goal text is rejected so we don't
|
||||
|
||||
@@ -625,7 +625,7 @@ class GatewayInboundMixin:
|
||||
steered = False
|
||||
if self._hm_text_only(event) and steer_text and hasattr(running_agent, "steer"):
|
||||
try:
|
||||
steered = bool(running_agent.steer(self._steer_text_with_origin(steer_text, event)))
|
||||
steered = self._steer_running_agent(running_agent, self._steer_text_with_origin(steer_text, event))
|
||||
except Exception as exc:
|
||||
logger.warning("PRIORITY steer failed for session %s: %s", _quick_key, exc)
|
||||
if steered:
|
||||
|
||||
69
tests/gateway/test_busy_steer_subagent_fanout.py
Normal file
69
tests/gateway/test_busy_steer_subagent_fanout.py
Normal file
@@ -0,0 +1,69 @@
|
||||
"""Busy-mode steer reaches the subagent the parent is blocked on (#112095, Telegram surface).
|
||||
|
||||
A parent inside ``delegate_task`` drains its own steer queue only after the tool returns — i.e.
|
||||
after the child finishes — so a steer that lands on the parent alone is never read by a looping
|
||||
child. Every gateway steer entry point must fan the text out to ``_active_children`` and the ack
|
||||
must say so.
|
||||
"""
|
||||
|
||||
import threading
|
||||
|
||||
import pytest
|
||||
|
||||
from gateway.config import GatewayConfig, Platform
|
||||
from gateway.platforms.event import MessageEvent, MessageType
|
||||
from gateway.run import GatewayRunner
|
||||
from gateway.session import SessionSource
|
||||
|
||||
|
||||
class _Agent:
|
||||
def __init__(self, children=()):
|
||||
self.payload = None
|
||||
self._active_children = list(children)
|
||||
self._active_children_lock = threading.Lock()
|
||||
|
||||
def steer(self, text):
|
||||
self.payload = text
|
||||
return True
|
||||
|
||||
|
||||
def _event(text="focus on rows 10-20"):
|
||||
source = SessionSource(platform=Platform.TELEGRAM, chat_id="c", user_id="u", chat_type="dm")
|
||||
return MessageEvent(text=text, message_type=MessageType.TEXT, source=source, message_id="m")
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@pytest.mark.parametrize("route", ["busy_steer_mode", "priority", "explicit_command"])
|
||||
async def test_busy_steer_fans_out_to_active_subagents(route, tmp_path, monkeypatch):
|
||||
monkeypatch.setattr("gateway.run._hermes_home", tmp_path)
|
||||
runner = GatewayRunner(config=GatewayConfig())
|
||||
child_a, child_b = _Agent(), _Agent()
|
||||
parent = _Agent(children=[child_a, child_b])
|
||||
runner._session_state("key").turn.agent = parent
|
||||
event = _event("/steer focus on rows 10-20" if route == "explicit_command" else "focus on rows 10-20")
|
||||
|
||||
if route == "busy_steer_mode":
|
||||
outcome = await runner._resolve_busy_steer_or_redirect(event, "key", "steer", parent)
|
||||
assert outcome.steered is True
|
||||
elif route == "priority":
|
||||
runner._hm_busy_steer(event, parent, "key")
|
||||
else:
|
||||
reply = await runner._busy_steer_command(event, "key", event.source)
|
||||
assert "subagent" in reply
|
||||
|
||||
assert parent.payload and parent.payload.endswith("focus on rows 10-20")
|
||||
# The looping child — the one actually doing the work — gets the same text, not just the parent.
|
||||
assert child_a.payload == parent.payload
|
||||
assert child_b.payload == parent.payload
|
||||
|
||||
|
||||
def test_busy_steer_ack_names_subagents(tmp_path, monkeypatch):
|
||||
monkeypatch.setattr("gateway.run._hermes_home", tmp_path)
|
||||
runner = GatewayRunner(config=GatewayConfig())
|
||||
parent = _Agent(children=[_Agent()])
|
||||
kwargs = dict(is_steer_mode=True, is_queue_mode=False, is_redirect_mode=False,
|
||||
demoted_for_subagents=False, demoted_for_compression=False)
|
||||
with_children = runner._compose_busy_ack_message(_event(), 0.0, None, parent, **kwargs)
|
||||
without = runner._compose_busy_ack_message(_event(), 0.0, None, _Agent(), **kwargs)
|
||||
assert "subagent" in with_children
|
||||
assert "subagent" not in without
|
||||
Reference in New Issue
Block a user