fix: join delegated work before finite chat exits

This commit is contained in:
Victor Kyriazakos
2026-09-10 10:33:47 +00:00
committed by Teknium
parent fa75692211
commit 49ef015ca3
5 changed files with 451 additions and 3 deletions

View File

@@ -0,0 +1,189 @@
"""Finite chat must consume parallel delegated results before the CLI exits."""
import http.server
import json
import os
from pathlib import Path
import subprocess
import sys
import threading
from types import SimpleNamespace
from typing import Any
from unittest.mock import MagicMock
import pytest
REPO_ROOT = Path(__file__).resolve().parents[2]
@pytest.mark.parametrize("mode", ["quiet", "oneshot", "redirected"])
def test_finite_chat_joins_parallel_children_before_final_response(tmp_path, mode):
"""Real parser, CLI, agent loops and delegate_task. Only inference is synthetic.
Both child HTTP requests must reach the barrier before either can finish.
The provider synthesizes a final answer only from returned tool results,
not from child requests or transcripts that the parent has not consumed.
"""
home = tmp_path / "profile"
home.mkdir()
barrier = threading.Barrier(2)
lock = threading.Lock()
children, joined_results, errors = [], [], []
workers = ("WORKER_ALPHA", "WORKER_BETA")
class Provider(http.server.BaseHTTPRequestHandler):
def do_GET(self):
self.send_error(404)
def do_POST(self):
try:
request = json.loads(self.rfile.read(int(self.headers["Content-Length"])))
messages = request.get("messages", [])
if not messages:
self.send_error(404)
return
users = [str(m.get("content", "")) for m in messages if m["role"] == "user"]
worker = next((w for w in workers if w in users[-1]), None)
results = [json.loads(m["content"]) for m in messages if m["role"] == "tool"]
message: dict[str, Any]
if worker:
with lock:
children.append(worker)
barrier.wait(timeout=20)
message = {"role": "assistant", "content": worker + "_COMPLETE"}
elif results:
with lock:
joined_results.extend(results)
summaries = json.dumps(results)
joined = all(w + "_COMPLETE" in summaries for w in workers)
message = {"role": "assistant", "content": (
"FANOUT_JOINED_AND_SYNTHESIZED" if joined else "PARENT_ENDED_BEFORE_JOIN"
)}
else:
message = {"role": "assistant", "content": None, "tool_calls": [{
"id": "call_fanout", "type": "function", "function": {
"name": "delegate_task", "arguments": json.dumps({
"tasks": [{"goal": f"Complete {w} and return its completion token."}
for w in workers],
}),
},
}]}
response = {
"id": "chatcmpl-local", "object": "chat.completion", "created": 1,
"model": "test-model", "choices": [{
"index": 0, "message": message,
"finish_reason": "tool_calls" if "tool_calls" in message else "stop",
}], "usage": {"prompt_tokens": 20, "completion_tokens": 10, "total_tokens": 30},
}
content_type = "application/json"
if request.get("stream"):
response["object"] = "chat.completion.chunk"
response["choices"][0]["delta"] = response["choices"][0].pop("message")
for index, tool in enumerate(message.get("tool_calls", [])):
tool["index"] = index
raw = ("data: " + json.dumps(response) + "\n\ndata: [DONE]\n\n").encode()
content_type = "text/event-stream"
else:
raw = json.dumps(response).encode()
self.send_response(200)
self.send_header("Content-Type", content_type)
self.send_header("Content-Length", str(len(raw)))
self.end_headers()
self.wfile.write(raw)
except Exception as exc:
with lock:
errors.append(repr(exc))
self.send_error(500)
def log_message(self, format, *args):
pass
server = http.server.ThreadingHTTPServer(("127.0.0.1", 0), Provider)
thread = threading.Thread(target=server.serve_forever, daemon=True)
thread.start()
url = f"http://127.0.0.1:{server.server_port}/v1"
(home / "config.yaml").write_text(
f"model:\n provider: custom\n base_url: {url}\n api_mode: chat_completions\n"
"memory:\n memory_enabled: false\n user_profile_enabled: false\n"
"terminal:\n env: local\n oneshot_completion_wait_seconds: 1\n"
f"delegation:\n max_concurrent_children: 2\n base_url: {url}\n"
" api_key: local-test-only\n model: test-model\n",
encoding="utf-8",
)
query = tmp_path / "query.txt"
query.write_text("Delegate two independent tasks, then synthesize their results.", encoding="utf-8")
# Retain native Windows location variables, but never inherited credentials
# or another session's finite/approval/runtime markers.
env = {key: os.environ[key] for key in (
"PATH", "SYSTEMROOT", "WINDIR", "COMSPEC", "TEMP", "TMP", "LOCALAPPDATA", "APPDATA",
) if key in os.environ}
env.update(HOME=str(tmp_path), USERPROFILE=str(tmp_path), HERMES_HOME=str(home),
HERMES_MANAGED_DIR=str(tmp_path / "managed"), TERMINAL_CWD=str(tmp_path),
OPENAI_BASE_URL=url, OPENAI_API_KEY="local-test-only", PYTHONPATH=str(REPO_ROOT),
PYTHONDONTWRITEBYTECODE="1", LANG="C.UTF-8")
mode_flags = {"quiet": ["-Q"], "oneshot": ["--oneshot"], "redirected": []}
command = [
sys.executable, "-c", "from hermes_cli.main import main; main()", "chat",
*mode_flags[mode], "--provider", "custom", "--model", "test-model",
"--toolsets", "delegation", "--ignore-rules", "--query-file", str(query),
"--reasoning", "high", "--max-turns", "10", "--run-budget", "60",
]
try:
result = subprocess.run(command, cwd=tmp_path, env=env, stdin=subprocess.DEVNULL,
capture_output=True, text=True, encoding="utf-8", timeout=75)
finally:
server.shutdown()
server.server_close()
thread.join(timeout=5)
diagnostic = (result.stdout, result.stderr, joined_results, children, errors)
assert result.returncode == 0, diagnostic
assert "FANOUT_JOINED_AND_SYNTHESIZED" in result.stdout, diagnostic
assert "PARENT_ENDED_BEFORE_JOIN" not in result.stdout, diagnostic
assert sorted(children) == sorted(workers), diagnostic
assert not errors, diagnostic
assert len(joined_results) == 1, diagnostic
assert [r["status"] for r in joined_results[0]["results"]] == ["completed", "completed"]
assert [r["summary"] for r in joined_results[0]["results"]] == [w + "_COMPLETE" for w in workers]
@pytest.mark.parametrize("query,image", [("Delegate a task", None), (None, "image.png")])
def test_tty_seeded_chat_keeps_background_delegation(monkeypatch, query, image):
"""A TTY -q (or image seed) still owns a later-turn completion consumer."""
import cli
import tools.delegate_tool as dt
from gateway.session_context import reset_session_vars
from run_agent import AIAgent
from tests.tools.test_delegate import _make_mock_parent
reset_session_vars()
monkeypatch.delenv("HERMES_SINGLE_QUERY_SESSION", raising=False)
monkeypatch.delenv("HERMES_KANBAN_TASK", raising=False)
monkeypatch.setattr(cli.sys.stdin, "isatty", lambda: True)
monkeypatch.setattr(cli.sys.stdout, "isatty", lambda: True)
monkeypatch.setattr(cli, "_collect_query_images", lambda q, i: (q, [i] if i else []))
parent = _make_mock_parent()
parent.session_id = "interactive-parent"
monkeypatch.setattr(dt, "_build_child_agent", lambda **kw: MagicMock())
monkeypatch.setattr(dt, "_resolve_delegation_credentials", lambda *a, **kw: {
"model": "test-model", "provider": "custom", "base_url": None, "api_key": None,
"api_mode": None, "command": None, "args": None,
})
dispatched = []
def dispatch(unit, unit_id, slot_key, routing):
dispatched.append(unit)
return {"status": "dispatched", "delegation_id": "interactive-delegation"}
monkeypatch.setattr("tools.delegate_tool_dispatch._dispatch_unit", dispatch)
seeded = SimpleNamespace(run=lambda: AIAgent._dispatch_delegate_task(
parent, {"tasks": [{"goal": "independent task"}]},
))
try:
result = cli._run_single_query_mode(seeded, query, image, False, False)
assert isinstance(result, str)
payload = json.loads(result)
assert payload["status"] == "dispatched"
assert len(dispatched) == 1
assert "HERMES_SINGLE_QUERY_SESSION" not in os.environ
finally:
reset_session_vars()

View File

@@ -0,0 +1,228 @@
"""Finite dispatch returns sibling outcomes, not an orphaned async handle.
Only child construction/conversation is synthetic: dispatch, parallel executor,
child timeout/status conversion, interrupt propagation and cleanup are real.
"""
from __future__ import annotations
import json
import threading
from types import SimpleNamespace
import pytest
from agent.interrupt_control import InterruptControlMixin
from gateway import session_context as sc
from tools import async_delegation, delegate_tool as dt
from tools.delegate_tool_child_run import _attach_child
from tools.process_registry import process_registry
class _Parent(InterruptControlMixin, SimpleNamespace):
pass
class _Child:
def __init__(self, outcome="completed"):
self.outcome = outcome
self.tool_progress_callback = None
self._credential_pool = None
self._delegate_saved_tool_names = []
self._delegate_role = "leaf"
self._delegate_depth = 1
self._subagent_id = None
self.model = "test-model"
self.session_prompt_tokens = self.session_completion_tokens = 0
self.session_estimated_cost_usd = 0.0
self.session_cost_status = "unknown"
self.started = threading.Event()
self.interrupted = threading.Event()
self.unwinding = threading.Event()
self.allow_finish = threading.Event()
self.finished = threading.Event()
self.closed = threading.Event()
self.close_while_running = False
self.worker = None
def run_conversation(self, **_kwargs):
self.worker = threading.current_thread()
self.started.set()
try:
if self.outcome == "slow":
assert self.interrupted.wait(5), "child never received stop"
self.unwinding.set()
assert self.allow_finish.wait(5), "child teardown was not released"
return {"final_response": "partial", "completed": False,
"interrupted": True, "api_calls": 1, "messages": []}
if self.outcome == "error":
raise RuntimeError("synthetic child exception")
if self.outcome == "failed":
return {"final_response": "provider rejection", "completed": False,
"failed": True, "error": "synthetic provider rejection",
"api_calls": 1, "messages": []}
return {"final_response": "sibling evidence", "completed": True,
"api_calls": 1, "messages": []}
finally:
self.finished.set()
def hard_interrupt(self, _reason=None, **_kwargs):
self.interrupted.set()
def get_activity_summary(self):
return {"api_call_count": 1}
def close(self):
self.close_while_running |= not self.finished.is_set()
self.closed.set()
@pytest.fixture
def harness(monkeypatch, tmp_path):
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
monkeypatch.setenv("HERMES_SINGLE_QUERY_SESSION", "1")
tokens = [(v, v.set(sc._UNSET)) for v in sc._VAR_MAP.values()]
tokens += [(v, v.set(sc._UNSET)) for v in
(sc._SESSION_ASYNC_DELIVERY, sc._SESSION_HISTORY_DELIVERY)]
async_delegation._reset_for_tests()
parent = _Parent(session_id="finite-outcomes", _delegate_depth=0,
_current_task_id=None, _interrupt_requested=False,
_execution_thread_id=None, quiet_mode=True,
_active_children=[], _active_children_lock=threading.Lock())
children = []
def build_child(task_index, parent_agent, **_kwargs):
child = children[task_index]
_attach_child(parent_agent, child)
return child
monkeypatch.setattr(dt, "_build_child_agent", build_child)
monkeypatch.setattr(dt, "_load_config", lambda: {})
monkeypatch.setattr(dt, "_get_max_concurrent_children", lambda: 2)
monkeypatch.setattr(dt, "_get_worktree_isolation", lambda: False)
monkeypatch.setattr(dt, "_get_child_timeout", lambda: 4)
monkeypatch.setattr(dt, "_resolve_delegation_credentials", lambda *_a, **_k: {
"model": "test-model", "provider": None, "base_url": None,
"api_key": None, "api_mode": None, "command": None, "args": None,
})
def dispatch(*outcomes):
children.extend(_Child(outcome) for outcome in outcomes)
return json.loads(dt.delegate_task(
tasks=[{"goal": f"Exercise child outcome {i}"} for i in range(len(children))],
background=True, parent_agent=parent,
))
yield parent, children, dispatch
# Release on assertion failure too. Never leave conversation workers behind.
for child in children:
child.interrupted.set()
child.allow_finish.set()
async_delegation._reset_for_tests()
for child in children:
if child.started.is_set():
assert child.finished.wait(5)
assert child.closed.wait(5)
child.worker.join(5)
assert not child.worker.is_alive()
assert not child.close_while_running
assert not parent._active_children
while not process_registry.completion_queue.empty():
process_registry.completion_queue.get_nowait()
for var, token in reversed(tokens):
var.reset(token)
def _joined(result, statuses):
assert result.get("status") != "dispatched", result
assert [entry["task_index"] for entry in result["results"]] == [0, 1]
assert [entry["status"] for entry in result["results"]] == statuses
assert result["results"][0]["summary"] == "sibling evidence"
assert process_registry.completion_queue.empty()
@pytest.mark.parametrize("outcome", ["error", "failed"])
def test_finite_batch_returns_success_and_child_failure(harness, outcome):
_parent, children, dispatch = harness
result = dispatch("completed", outcome)
_joined(result, ["completed", outcome])
assert "synthetic" in result["results"][1]["error"]
assert result["results"][1]["exit_reason"] == "error"
assert all(child.closed.is_set() for child in children)
def test_finite_batch_returns_timeout_while_child_unwinds(harness, monkeypatch):
_parent, children, dispatch = harness
monkeypatch.setattr(dt, "_get_child_timeout", lambda: 0.5)
result = dispatch("completed", "slow")
_joined(result, ["completed", "timeout"])
assert "timed out" in result["results"][1]["error"]
assert result["results"][1]["exit_reason"] == "timeout"
slow = children[1]
assert slow.unwinding.wait(1)
assert not slow.finished.is_set()
assert not slow.closed.is_set(), "timeout must not close an unwinding child"
slow.allow_finish.set()
assert slow.closed.wait(2)
def test_finite_batch_returns_parent_interruption(harness):
parent, children, dispatch = harness
children.extend([_Child(), _Child("slow")])
slow = children[1]
stop_errors = []
def stop_parent():
try:
assert slow.started.wait(3)
assert children[0].closed.wait(3)
parent.hard_interrupt("test parent stop")
except Exception as exc:
stop_errors.append(exc)
finally:
slow.allow_finish.set()
stopper = threading.Thread(target=stop_parent)
stopper.start()
try:
result = dispatch()
_joined(result, ["completed", "interrupted"])
assert parent._interrupt_requested is True
assert slow.interrupted.is_set()
finally:
slow.allow_finish.set()
stopper.join(5)
assert not stopper.is_alive()
assert not stop_errors
def test_finite_marker_overrides_api_history_continuation(harness):
_parent, _children, dispatch = harness
sc.set_session_vars(platform="api_server", chat_id="history-session",
session_key="history-session", session_id="history-session",
session_history_delivery="1", async_delivery=False)
assert sc.session_history_delivery_supported()
result = dispatch("completed", "error")
_joined(result, ["completed", "error"])
@pytest.mark.parametrize("marker", [None, "0", "false"])
@pytest.mark.parametrize("api_history", [False, True])
def test_nonfinite_marker_preserves_background_dispatch(harness, monkeypatch, marker, api_history):
_parent, children, dispatch = harness
if marker is None:
monkeypatch.delenv("HERMES_SINGLE_QUERY_SESSION")
else:
monkeypatch.setenv("HERMES_SINGLE_QUERY_SESSION", marker)
if api_history:
sc.set_session_vars(platform="api_server", chat_id="history-session",
session_key="history-session", session_id="history-session",
session_history_delivery="1", async_delivery=False)
result = dispatch("completed", "completed")
assert result["status"] == "dispatched", result
assert result["mode"] == "background"
assert "results" not in result
event = process_registry.completion_queue.get(timeout=5)
assert event["type"] == "async_delegation"
if api_history:
assert event["origin_session_id"] == "history-session"
assert all(child.closed.wait(2) for child in children)

View File

@@ -529,8 +529,10 @@ _DESCRIPTION_HEAD = (
"Spawn subagents in isolated contexts; each gets its own conversation, terminal session, and toolset, and only its "
"final summary returns to you. Pass every task in `tasks` — one entry spawns one subagent, several run in parallel "
"(limit in the tasks description).\n\n"
"Runs in the background: dispatch returns immediately with live transcript paths, and the call's results re-enter "
"the conversation as a new message when its subagents finish ({delivery}). Results are delivered only "
"Sessions without a later-result consumer (including one-shot CLI and cron) join parallel children "
"and return results in this tool call. "
"Otherwise runs in the background: dispatch returns live transcript paths and results re-enter "
"as a new message when subagents finish ({delivery}). Background results are delivered only "
"BETWEEN your turns: finish whatever does not depend on them, then give a one-line status and END YOUR TURN. Never "
"wait or poll on transcripts, artifact files, or CI for a child. "
"While children run, `action` (list/steer/stop) controls them live — steer when a transcript shows a "

View File

@@ -193,7 +193,7 @@ _SYNC_FALLBACK_NOTES = {
"no_async": (
"background=true is not available in this session — it cannot "
"receive a detached subagent result after the turn ends (a "
"one-shot runner such as `hermes -z`, a cron job, a Kanban "
"finite chat using -Q, --oneshot, or non-TTY stdio, `hermes -z`, a cron job, a Kanban "
"worker, or a stateless HTTP endpoint). The subagent(s) ran SYNCHRONOUSLY and the result is included above."
),
"at_capacity": (
@@ -215,6 +215,14 @@ def _resolve_async_wake_sid(origin_wake_sid: str, origin_session_history_deliver
API completion only persists a row; this does not authorize a model wake. The
continuation must read that row, not an authoritative caller-owned snapshot."""
from gateway.session_context import get_session_env
# Finite chat owns no later turn to consume a detached result. Reuse its
# approval/lifecycle marker without disabling terminal notify completions:
# those have their own bounded exit linger and durable result receipts.
if get_session_env("HERMES_SINGLE_QUERY_SESSION") == "1":
return None
try:
# Finite sessions cannot route a detached subagent result back to the agent after their turn/process
# ends. This includes stateless HTTP requests (#10760) and one-shot Kanban workers (#63169). Fall

View File

@@ -146,6 +146,27 @@ hermes chat --ignore-user-config --ignore-rules -q "Repro without my personal se
hermes chat --safe-mode -q "Is this bug mine or Hermes'?"
```
#### Delegation in finite chat runs
When chat answers and exits (`-Q`, `chat --oneshot`, or a query with non-TTY
stdio), `delegate_task` waits for its children and returns their results to the
parent in the same turn. Batch children still run in parallel, subject to
`delegation.max_concurrent_children`. The parent can use those results in its
final response before the CLI exits.
- **Automatic joining:** no opt-in or background-mode override is needed.
Interactive TTY chat and messaging sessions keep background delegation.
- **Existing safeguards:** delegation limits, timeouts, cancellation, and
`approvals.single_query_mode` still apply. Joining does not auto-approve commands
or guarantee successful child outcomes. Inspect results and verify artifacts.
- **Terminal completions:** this does not change background terminal notification
behavior or the bounded `terminal.oneshot_completion_wait_seconds` exit wait.
That setting is not a delegation timeout.
Delegation remains process-local. Interrupting or terminating the parent can
cancel unfinished children. Use a durable scheduler for work that must survive
the initiating process.
### `hermes -z <prompt>` — scripted one-shot
For programmatic callers (shell scripts, CI, cron, parent processes piping in a prompt), `hermes -z` is the purest one-shot entry point: **single prompt in, final response text out, nothing else on stdout or stderr.** No banner, no spinner, no tool previews, no `Session:` line — just the agent's final reply as plain text.