Skip to content

Commit 45c6d28

Browse files
authored
feat(hitl): add agent-driven human-in-the-loop via request_human_input (#51)
An agent can call request_human_input() to block synchronously on a threading.Event until a human resolves its question, or raise TimeoutError after a configured number of seconds -- there is no silent auto-approve on timeout. inform_user() sends a non-blocking progress update to whoever is watching the session. - By default a request blocks indefinitely. A timeout is set per call (the tool's own `timeout` argument) or as a manifest-wide default (a `{kind: system, name: request_human_input, params: {timeout: N}}` entry in spec.tools); the call's own value always wins. - Non-interactive runs (lab/benchmark runs, and batch CLI runs) auto-resolve instead of blocking, since there is no human present to answer; real interactive sessions always block for a genuine response. The auto-resolve decision defaults to "approve" and is configurable the same two ways as the timeout (spec.tools params, or the MAS_HITL_AUTO_RESOLVE_DECISION env var; the manifest value wins if both are set). - HitlResolverRegistry is the side-channel between a blocked tool call and whatever surfaces the request to a human (a Webex card, a CLI prompt, etc.) and later resolves it back. - Every delegated agent's own SessionController receives the same trace/trace_timestamps/trace_engine/trace_summary/trace_color configuration as the entry agent. - mas-ctl's CLI trace and this HITL bridge are both ExchangePlugin subscribers (subscribe_exchange/unsubscribe_exchange over a list of subscribers): tool-call errors are always surfaced in RED and forwarded to an external exchange_listener, regardless of --trace/--verbose. - request_human_input/inform_user are now available to every agent including delegates (needed so a delegated sub-agent can ask for human input too); mock_llm.py's schema-driven tool-selection fallback (used by MockModelAccess for offline/CI runs) is updated to skip these two system tools when picking a "nothing else matched" default, so a mocked delegate still exercises its own real tool instead of getting its call routed to request_human_input and auto-resolved. Regression tests cover: both exchange plugins staying subscribed exactly once per session (with the bridge still forwarding events when --trace/--verbose are off); the timeout/auto-resolve precedence rules above; and mock_llm.py's system-tool-skipping fallback. Two golden-run fixtures (extensions, lifecycle-control) are regenerated to reflect request_human_input/inform_user now being present on every agent's tool list. Full `task ci` (unit, controller, functional/golden-run/tutorial, dry-run, chat smoke, HITL gate) passes clean. Signed-off-by: Jordan Augé <augjorda@cisco.com>
1 parent 94cee03 commit 45c6d28

51 files changed

Lines changed: 3571 additions & 1264 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎ctl/src/mas/ctl/benchmark/runner.py‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -267,9 +267,17 @@ def run(
267267
flavour: Any = None,
268268
**kwargs: Any,
269269
) -> RunResult:
270+
import os
271+
270272
from mas.lab.benchmark.runners.fixtures import write_tool_fixtures_sidecar
271273
from mas.lab.inputs import RunInput
272274

275+
# Bench/lab runs are non-interactive by nature — no human is present to
276+
# resolve agent-initiated HITL requests, so auto-resolve them instead of
277+
# blocking for up to 60s per request (see execute_run_mas for the CLI
278+
# equivalent).
279+
os.environ.setdefault("MAS_HITL_AUTO_RESOLVE", "1")
280+
273281
ri: RunInput | None = run_input if isinstance(run_input, RunInput) else None
274282

275283
queries = ri.scripted_queries() if ri else ([prompt] if prompt else [])

‎ctl/src/mas/ctl/executor/mas_session.py‎

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -148,6 +148,11 @@ def prepare_delegation_entry_session(
148148
display: Any = None,
149149
verbose: int = 0,
150150
session_id: str = "",
151+
trace: bool = False,
152+
trace_timestamps: bool = False,
153+
trace_engine: bool = False,
154+
trace_summary: bool = False,
155+
trace_color: bool = False,
151156
) -> PreparedEntrySession:
152157
"""Wire dynamic-delegation entry agent (same path as ``execute_run_mas``).
153158
@@ -190,6 +195,11 @@ def prepare_delegation_entry_session(
190195
verbose=verbose,
191196
from_agent=entry_id,
192197
session_id=resolved_session_id,
198+
trace=trace,
199+
trace_timestamps=trace_timestamps,
200+
trace_engine=trace_engine,
201+
trace_summary=trace_summary,
202+
trace_color=trace_color,
193203
),
194204
entry_agent_id=entry_id,
195205
mas_config=compose.mas_config,
@@ -212,6 +222,11 @@ def wire_peer_delegation(
212222
verbose: int = 0,
213223
already_wired: "set[str] | None" = None,
214224
session_id: str = "",
225+
trace: bool = False,
226+
trace_timestamps: bool = False,
227+
trace_engine: bool = False,
228+
trace_summary: bool = False,
229+
trace_color: bool = False,
215230
) -> list[str]:
216231
"""Wire delegation onto every agent that declares its own ``delegates_to``
217232
peers in the MAS workflow topology — not just the entry agent.
@@ -247,6 +262,11 @@ def wire_peer_delegation(
247262
verbose=verbose,
248263
from_agent=entry_id,
249264
session_id=session_id,
265+
trace=trace,
266+
trace_timestamps=trace_timestamps,
267+
trace_engine=trace_engine,
268+
trace_summary=trace_summary,
269+
trace_color=trace_color,
250270
)
251271
newly_wired: list[str] = []
252272
for agent in compose.bind.agents:
@@ -288,6 +308,11 @@ def make_workflow_send(
288308
verbose: int,
289309
from_agent: str = "",
290310
session_id: str = "",
311+
trace: bool = False,
312+
trace_timestamps: bool = False,
313+
trace_engine: bool = False,
314+
trace_summary: bool = False,
315+
trace_color: bool = False,
291316
) -> RunTurnFn:
292317
"""Run one agent turn inside a multi-agent workflow (sequential or delegation).
293318
@@ -393,6 +418,11 @@ def send(
393418
config=ConversationConfig(single_turn=True),
394419
session_id=state["session_id"],
395420
working_memory_key=memory_key,
421+
trace=trace,
422+
trace_timestamps=trace_timestamps,
423+
trace_engine=trace_engine,
424+
trace_summary=trace_summary,
425+
trace_color=trace_color,
396426
)
397427
result = controller.run_turn(prompt, turn_id=turn_id, parent_call_id=parent_call_id)
398428
# Do NOT close observability after a delegated sub-turn: in a multi-agent
@@ -408,6 +438,18 @@ def send(
408438
state["prev_agent"] = agent_id
409439
if turn_failed(result):
410440
raise RuntimeError(f"agent {agent_id!r} turn failed")
441+
442+
# Propagate awaiting_hitl state via side channel (not return value)
443+
# This allows external systems (Webex bot) to detect and resolve HITL
444+
# from delegated agents without breaking the delegation contract.
445+
if result.awaiting_hitl:
446+
from mas.runtime.boundary.hitl.registry import get_hitl_resolver_registry
447+
registry = get_hitl_resolver_registry()
448+
# Check if there are any pending HITL requests for this agent
449+
if registry.has_pending(state["session_id"], agent_id):
450+
# Mark in state that this agent has pending HITL
451+
state.setdefault("pending_hitl_agents", set()).add(agent_id)
452+
411453
return result.text
412454

413455
return send

‎ctl/src/mas/ctl/executor/run_mas.py‎

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -60,11 +60,22 @@ def execute_run_mas(
6060
trace_color: bool = False,
6161
) -> int:
6262
"""Compose → materialize → SessionController on entry agent."""
63+
import os
64+
6365
from mas.ctl.session.controller import ConversationConfig, SessionController, close_observability
6466
from mas.ctl.session.hitl_config import resolve_hitl_from_manifest
6567
from mas.ctl.session.controller import run_session_loop
6668
from mas.ctl.ui.stdout import StdoutConversationDisplay
6769

70+
# Batch/CLI runs with auto-hitl (the default) have no external resolver
71+
# (Webex bot, operator console) listening for agent-initiated
72+
# request_human_input() calls, so the synchronous HITL wait in
73+
# manifest_tool_provider would otherwise always time out. Signal batch
74+
# mode via env var (mirrors the existing MAS_MANIFEST_RESOLVE_REFS
75+
# pattern) so it auto-resolves instead of blocking. Interactive sessions
76+
# never set this, so real HITL resolution still blocks as intended.
77+
os.environ["MAS_HITL_AUTO_RESOLVE"] = "1" if (auto_hitl and not interactive) else "0"
78+
6879
scripted = list(queries or [])
6980
if prompt:
7081
scripted.insert(0, prompt)
@@ -115,6 +126,11 @@ def execute_run_mas(
115126
entry_id=entry,
116127
display=display,
117128
verbose=verbose,
129+
trace=trace,
130+
trace_timestamps=trace_timestamps,
131+
trace_engine=trace_engine,
132+
trace_summary=trace_summary,
133+
trace_color=trace_color,
118134
)
119135
except KeyError as exc:
120136
logger.error("%s", exc)
@@ -136,6 +152,11 @@ def execute_run_mas(
136152
verbose=verbose,
137153
already_wired={entry},
138154
session_id=prepared.session_id,
155+
trace=trace,
156+
trace_timestamps=trace_timestamps,
157+
trace_engine=trace_engine,
158+
trace_summary=trace_summary,
159+
trace_color=trace_color,
139160
)
140161

141162
if runtime_params:

‎ctl/src/mas/ctl/session/bootstrap.py‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -62,6 +62,8 @@ class InstantiationOptions:
6262
enable_observability: bool = True
6363
enable_governance: bool = True
6464
enable_coordination: bool = True
65+
hitl_contract: object | None = None
66+
user_io_contract: object | None = None
6567

6668

6769
def instantiate_runtime(
@@ -209,6 +211,8 @@ def instantiate_runtime(
209211
options.manifest_dir or Path.cwd(),
210212
app_root=options.app_root or options.manifest_dir,
211213
workspace_root=ws.root if ws.found else None,
214+
hitl_contract=options.hitl_contract,
215+
user_io_contract=options.user_io_contract,
212216
)
213217
return instance, store
214218

‎ctl/src/mas/ctl/session/controller.py‎

Lines changed: 87 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44

55
from __future__ import annotations
66

7+
import logging
78
import sys
89
import time
910
import uuid
@@ -15,16 +16,59 @@
1516
sync_working_memory_in,
1617
sync_working_memory_out,
1718
)
18-
from mas.runtime.driver.driver import DriverTrace
19+
from mas.runtime.boundary.obs.exchange_plugin import ExchangePlugin
20+
from mas.runtime.driver.driver import DriverTrace, ExchangeRecord
1921
from mas.runtime.driver.instance import RuntimeInstance
2022
from mas.runtime.schema.egress import EmitClientResponse
2123

2224
from mas.ctl.session.exchange_log import (
25+
CliTraceExchangePlugin,
2326
TraceFormatOptions,
24-
print_exchange,
2527
)
2628
from mas.ctl.ui.display import ConversationDisplay
2729

30+
_logger = logging.getLogger("mas.runtime")
31+
_RED = "\033[1;31m"
32+
_RESET = "\033[0m"
33+
34+
35+
class _ToolErrorAndListenerBridge(ExchangePlugin):
36+
"""Second, independent ExchangePlugin subscriber alongside CliTraceExchangePlugin.
37+
38+
Two behaviors that aren't display formatting, so they don't belong in
39+
CliTraceExchangePlugin (a generic, reusable trace renderer):
40+
- Always-on tool-call error surfacing: a TOOL->AGENT exchange whose text
41+
is a JSON error result is printed in RED to stderr and logged at ERROR
42+
level, regardless of --trace/--verbose, so a failing tool call (e.g.
43+
run_skill_script returning {"error": ...}) can never silently
44+
disappear from view.
45+
- Forwards every exchange record to an optional external listener (e.g.
46+
a chat-UI plugin) via SessionController.exchange_listener, if set.
47+
Subscribing this as its own plugin (rather than folding it into
48+
CliTraceExchangePlugin) keeps each subscriber single-purpose, and proves
49+
the additive subscribe_exchange() interface actually supports more than
50+
one concurrent consumer.
51+
"""
52+
53+
def __init__(self) -> None:
54+
self.agent_id = "n/a"
55+
self.exchange_listener: Any | None = None
56+
57+
def configure(self, *, agent_id: str, exchange_listener: Any | None) -> None:
58+
self.agent_id = agent_id
59+
self.exchange_listener = exchange_listener
60+
61+
def on_exchange(self, record: ExchangeRecord) -> None:
62+
if record.tag == "TOOL->AGENT" and '"error"' in record.text:
63+
message = f"[{self.agent_id}] TOOL ERROR: {record.text.strip()}"
64+
print(f"{_RED}{message}{_RESET}", file=sys.stderr, flush=True)
65+
_logger.error(message)
66+
if callable(self.exchange_listener):
67+
try:
68+
self.exchange_listener(record)
69+
except Exception:
70+
_logger.exception("session exchange listener failed")
71+
2872

2973
@dataclass
3074
class ConversationConfig:
@@ -63,6 +107,7 @@ class SessionController:
63107
agent_id: str = "n/a"
64108
llm_id: str = "gpt-4o-mini"
65109
obs_recorder: Any | None = None
110+
exchange_listener: Any | None = None
66111
# One id for the whole MAS run — every turn this controller ever runs
67112
# shares it (see _run_user_turn). Empty here means "mint a fresh one";
68113
# explicitly pass an existing value (e.g. from an entry agent's own
@@ -79,6 +124,8 @@ class SessionController:
79124
working_memory_key: str = ""
80125
_turn: int = 0
81126
_trace_turn_start: float = 0.0
127+
_trace_plugin: CliTraceExchangePlugin | None = field(default=None, repr=False)
128+
_bridge_plugin: _ToolErrorAndListenerBridge | None = field(default=None, repr=False)
82129

83130
def __post_init__(self) -> None:
84131
if not self.session_id:
@@ -99,35 +146,50 @@ def _trace_format_options(self) -> TraceFormatOptions:
99146
)
100147

101148
def _setup_exchange_tracing(self) -> None:
102-
"""Setup realtime exchange tracing via driver callback (if trace or verbose)."""
103-
# Skip if neither trace nor verbose logging requested
149+
"""Configure the CLI trace + error/listener-bridge ExchangePlugins.
150+
151+
Subscribes two persistent, single-purpose plugins to the driver's
152+
exchange_plugins list the first time this runs, then only updates
153+
their config on every subsequent turn — the driver's subscriber
154+
list is additive (see KernelDriver.subscribe_exchange), so
155+
re-running this every turn no longer risks discarding any other
156+
plugin (e.g. an external chat-UI plugin) registered on the same
157+
driver:
158+
- CliTraceExchangePlugin: normal exchange display (AGENT/LLM/TOOL
159+
lines), gated by --trace or --verbose as before.
160+
- _ToolErrorAndListenerBridge: tool-call *errors* are always
161+
surfaced — unconditionally printed in RED to stderr and logged
162+
at ERROR level, regardless of trace/verbose settings, so a
163+
failing tool call (e.g. run_skill_script returning
164+
{"error": ...}) can never silently disappear from view. Also
165+
forwards every exchange record to self.exchange_listener, if set.
166+
"""
167+
if self._trace_plugin is None:
168+
self._trace_plugin = CliTraceExchangePlugin()
169+
self.instance.driver.subscribe_exchange(self._trace_plugin)
170+
if self._bridge_plugin is None:
171+
self._bridge_plugin = _ToolErrorAndListenerBridge()
172+
self.instance.driver.subscribe_exchange(self._bridge_plugin)
173+
self._bridge_plugin.configure(agent_id=self.agent_id, exchange_listener=self.exchange_listener)
174+
175+
# Skip normal display if neither trace nor verbose logging requested
176+
# (the error/listener bridge stays active regardless).
104177
if not self.trace and self.verbose < 1:
105-
self.instance.driver.on_exchange = None
178+
self._trace_plugin.configure(
179+
trace=False, verbose=0, agent_id=self.agent_id, fmt=self._trace_format_options()
180+
)
106181
self.instance.driver.capture_engine_io = False
107182
return
108183

109184
self._trace_turn_start = time.perf_counter()
110185
self.instance.driver.capture_engine_io = self.trace_engine
111-
fmt = self._trace_format_options()
112-
logger = __import__("logging").getLogger("mas.runtime")
113-
114-
def on_exchange(ex: object) -> None:
115-
from mas.runtime.driver.driver import ExchangeRecord
116-
from mas.ctl.session.exchange_log import format_exchange
117-
118-
if not isinstance(ex, ExchangeRecord):
119-
return
120-
121-
# Realtime stderr output (if --trace) — primary display path
122-
if self.trace:
123-
print_exchange(ex, err=sys.stderr, agent_id=self.agent_id, fmt=fmt)
124-
# Verbose logging only if NOT using --trace (alternative logging path)
125-
elif self.verbose >= 1:
126-
formatted = format_exchange(self.agent_id, ex, fmt=fmt).strip()
127-
for line in formatted.splitlines():
128-
logger.info("[%s] %s", self.agent_id, line)
129-
130-
self.instance.driver.on_exchange = on_exchange
186+
self._trace_plugin.configure(
187+
trace=self.trace,
188+
verbose=self.verbose,
189+
agent_id=self.agent_id,
190+
fmt=self._trace_format_options(),
191+
)
192+
131193

132194
def _handle_list_skills(self) -> TurnResult:
133195
"""Handle /skills command — list all available skills."""

0 commit comments

Comments
 (0)