Skip to content
Merged
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
107 changes: 107 additions & 0 deletions autobot-backend/chat_workflow/peer_messages_dispatch_16948_test.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,107 @@
# Copyright 2025-2026 mrveiss
# SPDX-License-Identifier: Apache-2.0
# AutoBot - AI-Powered Automation Platform
# Author: mrveiss
"""AI_STACK peer messages are drained at the real dispatch seam (#16948).

`_dispatch_tool_call` polls every 10ms and ran a peer's handler immediately,
interrupting whatever the recipient was doing mid-task
(`protocols/agent_communication.py:520`). `enforce_peer_messages` is called
at the top of `_dispatch_tool_call` (`chat_workflow/tool_handler.py`), the
one seam common to the graph path, the legacy fallback and delegated
subagents (#16948 seam map, confirmed). Drives the real function against a
real `LLMIterationContext`, not a fake standing in for both.
"""

from protocols.agent_communication import AgentIdentity
from protocols.agent_kind import AgentKind
from protocols.agent_presence import AgentPresenceRegistry
from protocols.peer_inbox import PeerInboxDirectory


def _ctx(**overrides):
from chat_workflow.models import LLMIterationContext

defaults = dict(
ollama_endpoint="http://localhost:11434",
selected_model="test-model",
session_id="sess-1",
terminal_session_id="term-1",
used_knowledge=False,
rag_citations=[],
workflow_messages=[],
)
defaults.update(overrides)
return LLMIterationContext(**defaults)


def _sender(name: str = "rag") -> AgentIdentity:
return AgentIdentity(agent_id=name, agent_type=name, kind=AgentKind.AI_STACK, name=name, tenant_id=None)


def _directory_with_chat_registered() -> PeerInboxDirectory:
presence = AgentPresenceRegistry(ttl_seconds=60)
presence.report(kind=AgentKind.AI_STACK, tenant_id=None, name="chat", instance_id="chat", busy=False)
return PeerInboxDirectory(presence)


class TestEnforcePeerMessages:
def test_a_queued_message_is_drained_into_ctx_context(self, monkeypatch):
from chat_workflow.tool_dispatch_guards import enforce_peer_messages

directory = _directory_with_chat_registered()
directory.send(kind=AgentKind.AI_STACK, name="chat", sender=_sender(), content="hello", message_id="m1")
monkeypatch.setattr("protocols.peer_inbox.get_peer_inbox_directory", lambda: directory)

ctx = _ctx()
enforce_peer_messages(ctx)

assert ctx.context["peer_messages"][0]["content"] == "hello"
assert ctx.context["peer_messages"][0]["sender_name"] == "rag"

def test_never_blocks_the_tool_call(self, monkeypatch):
"""The one property that distinguishes this from every other guard: it
always returns None, since a peer message is context, not an instruction
that could refuse a tool call (#16946 owner ruling)."""
from chat_workflow.tool_dispatch_guards import enforce_peer_messages

directory = _directory_with_chat_registered()
directory.send(kind=AgentKind.AI_STACK, name="chat", sender=_sender(), content="hello", message_id="m1")
monkeypatch.setattr("protocols.peer_inbox.get_peer_inbox_directory", lambda: directory)

assert enforce_peer_messages(_ctx()) is None

def test_no_ctx_is_a_no_op(self, monkeypatch):
directory = _directory_with_chat_registered()
directory.send(kind=AgentKind.AI_STACK, name="chat", sender=_sender(), content="hello", message_id="m1")
monkeypatch.setattr("protocols.peer_inbox.get_peer_inbox_directory", lambda: directory)

from chat_workflow.tool_dispatch_guards import enforce_peer_messages

assert enforce_peer_messages(None) is None
# The message is still queued -- a missing ctx drops the delivery
# opportunity for this call, not the message itself.
assert directory.inbox_for(kind=AgentKind.AI_STACK, tenant_id=None, name="chat").drain() != []

def test_a_message_queued_after_this_calls_own_drain_is_not_included(self, monkeypatch):
"""The mid-task-vs-next-turn boundary itself (#16948 AC4): a message
that arrives after THIS dispatch call already drained is held for the
next one, not retroactively spliced into a context already built."""
from chat_workflow.tool_dispatch_guards import enforce_peer_messages

directory = _directory_with_chat_registered()
monkeypatch.setattr("protocols.peer_inbox.get_peer_inbox_directory", lambda: directory)

ctx = _ctx()
enforce_peer_messages(ctx) # drains an empty inbox -- nothing queued yet
assert "peer_messages" not in ctx.context

# A message arrives mid-task, after this call's own drain already ran.
directory.send(kind=AgentKind.AI_STACK, name="chat", sender=_sender(), content="late", message_id="m2")
assert "peer_messages" not in ctx.context # not retroactively added

# It surfaces only on the NEXT dispatch call's own drain -- a fresh ctx,
# matching a fresh LLM turn.
next_ctx = _ctx()
enforce_peer_messages(next_ctx)
assert next_ctx.context["peer_messages"][0]["content"] == "late"
19 changes: 19 additions & 0 deletions autobot-backend/chat_workflow/tool_dispatch_guards.py
Original file line number Diff line number Diff line change
Expand Up @@ -334,3 +334,22 @@ def enforce_work_item_approval(
"work_item_id": work_item_id,
},
)


def enforce_peer_messages(ctx: "LLMIterationContext | None") -> None:
"""Drain the "chat" role's peer inbox into `ctx.context` at this seam (#16948).

Never blocks: a peer message is context, never an instruction (#16946
owner ruling) -- any tool call it prompts still goes through every gate
above, exactly as it would for a human-typed message. A no-op without a
`ctx` to carry the drained messages onto.
"""
if ctx is None:
return None
from protocols.agent_kind import AgentKind
from protocols.peer_inbox import get_peer_inbox_directory

drained = get_peer_inbox_directory().inbox_for(kind=AgentKind.AI_STACK, tenant_id=None, name="chat").drain()
if drained:
ctx.context.setdefault("peer_messages", []).extend(e.to_dict() for e in drained)
return None
16 changes: 8 additions & 8 deletions autobot-backend/chat_workflow/tool_handler.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@
enforce_config_protection,
enforce_fact_forcing,
enforce_forbidden_work,
enforce_peer_messages,
enforce_pre_action_verifier,
enforce_repetition,
enforce_work_item_approval,
Expand Down Expand Up @@ -1465,9 +1466,7 @@ async def _execute_terminal_command(
if not self.terminal_tool:
return {"status": "error", "error": "Terminal tool not available"}

# Ensure terminal session exists for this conversation
if not self.terminal_tool.active_sessions.get(session_id):
# Create session
session_result = await self.terminal_tool.create_session(
agent_id=f"chat_agent_{session_id}",
conversation_id=session_id,
Expand All @@ -1478,7 +1477,6 @@ async def _execute_terminal_command(
if session_result.get("status") != "success":
return session_result

# Execute command
result = await self.terminal_tool.execute_command(
conversation_id=session_id, command=command, description=description
)
Expand Down Expand Up @@ -2214,7 +2212,6 @@ async def _handle_command_error(
command,
repairable_error.message,
)
# Emit REPAIRABLE_ERROR hook
await _emit_repairable_error(
Exception(repairable_error.message),
session_id,
Expand All @@ -2232,7 +2229,6 @@ async def _handle_command_error(
},
)
else:
# Emit CRITICAL_ERROR hook for non-repairable errors
await _emit_critical_error(Exception(error), session_id, {"command": command})
additional_response_parts.append(f"\n\n❌ Command execution failed: {error}")
yield WorkflowMessage(
Expand Down Expand Up @@ -2264,17 +2260,14 @@ def _classify_command_error(self, command: str, error: str, stderr: str) -> Repa
# AttributeError out of the tool-call generator.
combined = f"{str(error or '').lower()} {str(stderr or '').lower()}"

# Check for critical (non-repairable) errors first
if any(p in combined for p in _CRITICAL_ERROR_PATTERNS):
logger.warning("[Issue #655] Critical error (out of memory): %s", error)
return None

# Check against repairable error patterns
result = _match_repairable_error(combined, command, error)
if result:
return result

# Default: treat as repairable with generic suggestion
return RepairableException(
message=f"Command failed: {error}",
suggestion="Check the error details and try an alternative approach",
Expand Down Expand Up @@ -3274,6 +3267,10 @@ def _enforce_work_item_approval(
"""
return enforce_work_item_approval(tool_call, ctx, execution_results)

def _enforce_peer_messages(self, ctx: "LLMIterationContext" | None) -> None:
"""Drain queued peer messages into ctx.context (#16948) -- never blocks."""
return enforce_peer_messages(ctx)

async def _dispatch_tool_call(
self,
tool_call: dict[str, Any],
Expand All @@ -3295,6 +3292,9 @@ async def _dispatch_tool_call(
"""
tool_name = tool_call["name"]

# #16948: peer messages are context, drained at this same live seam.
self._enforce_peer_messages(ctx)

# GH#11145: enforce the acting agent's forbidden_work manifest at the single
# production dispatch seam — before any tool-specific branch. Every tool call
# funnels through here, so this is the one place the capability boundary is
Expand Down
4 changes: 4 additions & 0 deletions autobot-backend/protocols/agent_presence.py
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@
from autobot_shared.logging_manager import get_logger
from autobot_shared.singleton_factory import lazy_singleton
from protocols.agent_kind import AgentKind
from protocols.idle_notice import notify_agent_idle

logger = get_logger(__name__)

Expand Down Expand Up @@ -160,7 +161,10 @@ def report(
f"{name!r} ({kind.value}, tenant={tenant_id!r}) is already live under instance "
f"{existing.instance_id!r}; refusing instance {instance_id!r}"
)
was_busy = existing.busy if existing is not None else None
self._entries[key] = _Record(instance_id=instance_id, busy=busy, detail=detail, last_seen=now)
if was_busy is True and busy is False:
notify_agent_idle(kind=kind, tenant_id=tenant_id, name=name)

def deregister(self, *, kind: AgentKind, tenant_id: str | None, name: str, instance_id: str) -> None:
"""Explicit removal (clean shutdown) -- only the reporting instance may do this."""
Expand Down
108 changes: 108 additions & 0 deletions autobot-backend/protocols/idle_notice.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
#!/usr/bin/env python3
# Copyright 2025-2026 mrveiss
# SPDX-License-Identifier: Apache-2.0
# AutoBot - AI-Powered Automation Platform
# Author: mrveiss
"""Subscribe once to be told when a peer next goes idle (#16949).

Absent before this: presence (#16947) can be *polled* for busy/idle, but
nothing pushed a notice on the transition -- an agent waiting on a peer had
to poll, which is what this layer exists to avoid.

Built on the existing event transport (`events/bus.py`'s `EventBus`), not a
new one: `AgentPresenceRegistry.report()` publishes `EVT_AGENT_IDLE`
exactly on a busy=True -> busy=False transition (never on a first report
that starts idle -- that is "already free", handled below by checking
current state before waiting, not a transition to notify about).
`wait_for_idle()` is the one-shot subscriber: it resolves immediately if
the target is already idle, otherwise waits for exactly one matching event
or the bounded timeout, whichever comes first, and tells the caller which.
"""

from __future__ import annotations

import asyncio
from typing import TYPE_CHECKING

from autobot_shared.env_utils import env_float
from events.bus import EventBus, PersistStrategy, get_event_bus
from protocols.agent_kind import AgentKind

if TYPE_CHECKING:
from protocols.agent_presence import AgentPresenceRegistry

EVT_AGENT_IDLE = "agent_idle"

#: Bounded wait for #16949's one-shot subscription -- a peer that never goes
#: idle within this must not leave the subscriber waiting forever.
IDLE_NOTICE_TIMEOUT_ENV = "AUTOBOT_IDLE_NOTICE_TIMEOUT_SECONDS"
DEFAULT_IDLE_NOTICE_TIMEOUT_SECONDS = 300.0


def idle_notice_timeout_seconds() -> float:
return env_float(IDLE_NOTICE_TIMEOUT_ENV, DEFAULT_IDLE_NOTICE_TIMEOUT_SECONDS)


def notify_agent_idle(*, kind: AgentKind, tenant_id: str | None, name: str) -> None:
"""Publish the idle transition. Fire-and-forget -- no running loop, no notice.

Called from `AgentPresenceRegistry.report()`, a synchronous method every
presence feed adapter calls from an `async def` (a real running loop in
every production case); a sync unit test calling `report()` directly
simply gets no notification scheduled, which is correct for a test that
never awaits anything.
"""
try:
loop = asyncio.get_running_loop()
except RuntimeError:
return
payload = {"kind": kind.value, "tenant_id": tenant_id, "name": name}
channel = f"agent:{name}"
loop.create_task(get_event_bus().publish(channel, EVT_AGENT_IDLE, payload, persist=PersistStrategy.NONE))


class IdleWaitExpiredError(Exception):
"""Raised by `wait_for_idle()` when the peer never went idle in time."""


async def wait_for_idle(
presence: "AgentPresenceRegistry",
*,
kind: AgentKind,
tenant_id: str | None,
name: str,
bus: EventBus | None = None,
timeout_seconds: float | None = None,
) -> None:
"""Return once, the moment *name* is idle; raise `IdleWaitExpiredError` on timeout.

Resolves immediately without subscribing at all if the target is
already idle -- "the next time" a caller finds an already-free peer is
now, not a transition to wait for.
"""
bus = bus if bus is not None else get_event_bus()
timeout = timeout_seconds if timeout_seconds is not None else idle_notice_timeout_seconds()

current = next((e for e in presence.list_live(tenant_id) if (e.kind, e.name) == (kind, name)), None)
if current is not None and not current.busy:
return

loop = asyncio.get_event_loop()
resolved: asyncio.Future[None] = loop.create_future()

async def _listener(event_data: dict) -> None:
# EventManager wraps every in-process delivery as {"type", "payload"}
# (event_manager.py::publish) -- the raw notify_agent_idle() payload
# is one level down, not at this dict's top level.
fields = event_data.get("payload", {})
if (fields.get("kind"), fields.get("tenant_id"), fields.get("name")) == (kind.value, tenant_id, name):
if not resolved.done():
resolved.set_result(None)

bus.subscribe(EVT_AGENT_IDLE, _listener)
try:
await asyncio.wait_for(resolved, timeout=timeout)
except asyncio.TimeoutError:
raise IdleWaitExpiredError(f"{name!r} ({kind.value}) did not go idle within {timeout}s") from None
finally:
bus.unsubscribe(EVT_AGENT_IDLE, _listener)
Loading
Loading