refactor(agent/transports): compact ACP/SSL/stream helper modules (copilot_acp_client, acp_openai_bridge, ssl_*, stream_*, jiter_preload)

This commit is contained in:
Teknium
2026-09-02 13:18:33 -07:00
parent d2b3e545cb
commit a53d886584
7 changed files with 235 additions and 615 deletions

View File

@@ -1,33 +1,12 @@
"""OpenAI-shape bridge shared by Hermes' ACP clients.
An ACP agent (``copilot --acp``, and the ACP CLIs that reach Hermes as
providers) speaks the Agent Client Protocol, which has no OpenAI-style
``tools``/``tool_calls`` channel: a prompt is text, and a response is text plus
the agent's *own* tool notifications. Hermes' agentic surface — ``memory``,
``todo``, ``skill_manage`` and friends — is dispatched from OpenAI-shaped
``tool_calls``, so on an ACP provider it can only work if the schemas travel
*into* the prompt as text and the calls are parsed back *out* of the response
text.
``agent/copilot_acp_client.py`` already carried a private copy of that bridge.
This module is that code, lifted verbatim into one place so every ACP client
shares it instead of re-deriving the wire contract:
* :func:`render_tool_bridge_sections` — prompt sections describing the
forwarded tools and the ``<tool_call>{...}</tool_call>`` contract.
* :func:`extract_tool_calls_from_text` — parse those blocks back into
``ChatCompletionMessageToolCall`` objects and return the response text with
the blocks stripped.
* :func:`completion_to_stream_chunks` — re-shape a one-shot ACP response as
OpenAI stream chunks for callers that asked for ``stream=True`` (an ACP turn
is inherently one-shot from Hermes' perspective).
The one axis clients differ on is *which* tools they forward, so
``render_tool_bridge_sections`` takes an optional allowlist. A CLI with no tools
of its own (Copilot) forwards everything Hermes offers; a CLI that is an
autonomous agent with its own read/edit/execute tools must forward only Hermes'
agent-level tools, because re-offering the overlapping ones makes Hermes re-run
work the agent already finished.
ACP has no OpenAI-style ``tools``/``tool_calls`` channel, so Hermes' tool
schemas travel INTO the prompt as text (:func:`render_tool_bridge_sections`) and
calls are parsed back OUT of the response text (:func:`extract_tool_calls_from_text`).
Clients differ only in WHICH tools they forward (``allowlist``): a CLI with no
tools of its own forwards everything; an autonomous agent with its own
read/edit/execute tools forwards only Hermes' agent-level tools, since
re-offering overlapping ones makes Hermes redo finished work.
"""
from __future__ import annotations
@@ -48,7 +27,6 @@ TOOL_CALL_JSON_RE = re.compile(
re.DOTALL,
)
# The contract sentence shared by every ACP client: how to emit a call.
TOOL_CALL_CONTRACT = (
"Available tools (OpenAI function schema). "
"When using a tool, emit ONLY <tool_call>{...}</tool_call> with one JSON object "
@@ -69,41 +47,31 @@ __all__ = [
class StreamChunks(list):
"""Stream chunks that can still carry response-level attributes.
Hermes reads provider-level extras off the object returned by
``chat.completions.create`` (e.g. ``hermes_projected_messages``, consumed by
``agent/provider_projection.py``). A plain list of chunks would silently drop
them on the ``stream=True`` path, so ACP clients return this instead and copy
the extras onto it.
"""
"""Chunk list that also carries response-level attributes (e.g. ``hermes_projected_messages``)
Hermes reads off the ``create`` result; a plain list would drop them on the stream path."""
def completion_to_stream_chunks(completion: SimpleNamespace) -> StreamChunks:
"""Convert a one-shot ACP response into OpenAI-style stream chunks.
"""Re-shape a one-shot ACP response as OpenAI stream chunks (data chunk + usage chunk).
Response-level attributes other than ``choices``/``usage``/``model`` are
copied onto the returned object so nothing a caller reads off the completion
is lost when it asked to stream.
Response-level attributes other than choices/usage/model are copied onto the result.
"""
choice = completion.choices[0]
message = choice.message
tool_call_deltas = None
if message.tool_calls:
tool_call_deltas = []
for index, tool_call in enumerate(message.tool_calls):
tool_call_deltas.append(
SimpleNamespace(
index=index,
id=getattr(tool_call, "id", None),
type=getattr(tool_call, "type", "function"),
function=SimpleNamespace(
name=getattr(tool_call.function, "name", None),
arguments=getattr(tool_call.function, "arguments", None),
),
)
tool_call_deltas = [
SimpleNamespace(
index=index,
id=getattr(tool_call, "id", None),
type=getattr(tool_call, "type", "function"),
function=SimpleNamespace(
name=getattr(tool_call.function, "name", None),
arguments=getattr(tool_call.function, "arguments", None),
),
)
for index, tool_call in enumerate(message.tool_calls)
]
delta = SimpleNamespace(
role="assistant",
content=message.content or None,
@@ -112,21 +80,11 @@ def completion_to_stream_chunks(completion: SimpleNamespace) -> StreamChunks:
reasoning=getattr(message, "reasoning", None),
)
data_chunk = SimpleNamespace(
choices=[
SimpleNamespace(
index=0,
delta=delta,
finish_reason=choice.finish_reason,
)
],
choices=[SimpleNamespace(index=0, delta=delta, finish_reason=choice.finish_reason)],
model=completion.model,
usage=None,
)
usage_chunk = SimpleNamespace(
choices=[],
model=completion.model,
usage=completion.usage,
)
usage_chunk = SimpleNamespace(choices=[], model=completion.model, usage=completion.usage)
chunks = StreamChunks([data_chunk, usage_chunk])
for key, value in vars(completion).items():
if key not in ("choices", "usage", "model"):
@@ -134,12 +92,7 @@ def completion_to_stream_chunks(completion: SimpleNamespace) -> StreamChunks:
return chunks
def build_openai_tool_call(
*,
call_id: str,
name: str,
arguments: str,
) -> ChatCompletionMessageToolCall:
def build_openai_tool_call(*, call_id: str, name: str, arguments: str) -> ChatCompletionMessageToolCall:
"""Build an OpenAI-compatible tool-call object for downstream handling."""
return ChatCompletionMessageToolCall(
id=call_id,
@@ -155,18 +108,11 @@ def tool_specs_from_openai_tools(
*,
allowlist: Iterable[str] | None = None,
) -> list[dict[str, Any]]:
"""Flatten OpenAI ``tools`` into ``{name, description, parameters}`` specs.
Malformed entries are skipped. When ``allowlist`` is given, only tools whose
name is in it survive — that is how a client forwards just Hermes'
agent-level tools instead of the whole toolset.
"""
"""Flatten OpenAI ``tools`` into ``{name, description, parameters}`` specs; malformed entries are skipped."""
allowed = {str(n).strip() for n in allowlist} if allowlist is not None else None
specs: list[dict[str, Any]] = []
for t in tools or []:
if not isinstance(t, dict):
continue
fn = t.get("function") or {}
fn = t.get("function") or {} if isinstance(t, dict) else None
if not isinstance(fn, dict):
continue
name = fn.get("name")
@@ -175,13 +121,7 @@ def tool_specs_from_openai_tools(
name = name.strip()
if allowed is not None and name not in allowed:
continue
specs.append(
{
"name": name,
"description": fn.get("description", ""),
"parameters": fn.get("parameters", {}),
}
)
specs.append({"name": name, "description": fn.get("description", ""), "parameters": fn.get("parameters", {})})
return specs
@@ -191,31 +131,22 @@ def render_tool_bridge_sections(
*,
allowlist: Iterable[str] | None = None,
) -> list[str]:
"""Prompt sections that carry the forwarded tool schemas + choice hint.
Returns an empty list when no tool survives filtering and no choice hint was
requested, so callers can splice the result into their section list
unconditionally.
"""
"""Prompt sections carrying the forwarded tool schemas + choice hint (empty list when neither applies)."""
specs = tool_specs_from_openai_tools(tools, allowlist=allowlist)
sections: list[str] = []
if specs:
sections.append(
TOOL_CALL_CONTRACT + "\n" + json.dumps(specs, ensure_ascii=False)
)
sections.append(TOOL_CALL_CONTRACT + "\n" + json.dumps(specs, ensure_ascii=False))
if tool_choice is not None:
sections.append(f"Tool choice hint: {json.dumps(tool_choice, ensure_ascii=False)}")
return sections
def extract_tool_calls_from_text(
text: str,
) -> tuple[list[ChatCompletionMessageToolCall], str]:
def extract_tool_calls_from_text(text: str) -> tuple[list[ChatCompletionMessageToolCall], str]:
"""Pull ``<tool_call>`` blocks out of an ACP response.
Returns ``(tool_calls, cleaned_text)`` where ``cleaned_text`` is the
response with the consumed blocks removed, so the assistant message doesn't
show raw JSON to the user.
Returns ``(tool_calls, cleaned_text)`` with the consumed blocks removed so the
assistant message doesn't show raw JSON. Bare-JSON fallback runs only when no
XML block parsed.
"""
if not isinstance(text, str) or not text.strip():
return [], ""
@@ -228,9 +159,7 @@ def extract_tool_calls_from_text(
obj = json.loads(raw_json)
except Exception:
return
if not isinstance(obj, dict):
return
fn = obj.get("function")
fn = obj.get("function") if isinstance(obj, dict) else None
if not isinstance(fn, dict):
return
fn_name = fn.get("name")
@@ -242,27 +171,15 @@ def extract_tool_calls_from_text(
call_id = obj.get("id")
if not isinstance(call_id, str) or not call_id.strip():
call_id = f"acp_call_{len(extracted)+1}"
extracted.append(
build_openai_tool_call(
call_id=call_id,
name=fn_name.strip(),
arguments=fn_args,
)
)
extracted.append(build_openai_tool_call(call_id=call_id, name=fn_name.strip(), arguments=fn_args))
for m in TOOL_CALL_BLOCK_RE.finditer(text):
raw = m.group(1)
_try_add_tool_call(raw)
_try_add_tool_call(m.group(1))
consumed_spans.append((m.start(), m.end()))
# Only try bare-JSON fallback when no XML blocks were found.
if not extracted:
for m in TOOL_CALL_JSON_RE.finditer(text):
raw = m.group(0)
_try_add_tool_call(raw)
_try_add_tool_call(m.group(0))
consumed_spans.append((m.start(), m.end()))
if not consumed_spans:
return extracted, text.strip()
@@ -282,6 +199,5 @@ def extract_tool_calls_from_text(
cursor = max(cursor, end)
if cursor < len(text):
parts.append(text[cursor:])
cleaned = "\n".join(p.strip() for p in parts if p and p.strip()).strip()
return extracted, cleaned

View File

@@ -1,9 +1,7 @@
"""OpenAI-compatible shim that forwards Hermes requests to `copilot --acp`.
This adapter lets Hermes treat the GitHub Copilot ACP server as a chat-style
backend. Each request starts a short-lived ACP session, sends the formatted
conversation as a single prompt, collects text chunks, and converts the result
back into the minimal shape Hermes expects from an OpenAI client.
Each request starts a short-lived ACP session, sends the formatted conversation
as one prompt, collects text chunks, and returns the minimal OpenAI-client shape.
"""
from __future__ import annotations
@@ -35,27 +33,24 @@ ACP_MARKER_BASE_URL = "acp://copilot"
logger = logging.getLogger(__name__)
_DEFAULT_TIMEOUT_SECONDS = 900.0
# Stderr fingerprint of the deprecated `gh copilot` CLI extension
# (https://github.blog/changelog/2025-09-25-upcoming-deprecation-of-gh-copilot-cli-extension).
# We require BOTH the literal product name ("gh-copilot") AND a deprecation
# marker, so generic stderr from the NEW `@github/copilot` CLI — whose repo
# is github.com/github/copilot-cli and which legitimately mentions "copilot-cli"
# in its own banners and error messages — doesn't get misclassified as the
# deprecated extension.
# Stderr fingerprint of the deprecated `gh copilot` extension. Require BOTH the
# product name AND a deprecation marker: the NEW `@github/copilot` CLI (repo
# github/copilot-cli) legitimately mentions "copilot-cli" in its own banners.
_DEPRECATION_REQUIRED = ("gh-copilot",)
_DEPRECATION_MARKERS = (
"has been deprecated",
"no commands will be executed",
)
_DEPRECATION_MARKERS = ("has been deprecated", "no commands will be executed")
_ROLE_LABELS = {"system": "System", "user": "User", "assistant": "Assistant", "tool": "Tool", "context": "Context"}
# Probe verdicts per binary path so the ~50ms --help cost is paid once per
# process. Only definitive True/False is cached; an inconclusive probe is not,
# so a CLI installed mid-session is picked up.
_ACP_PROBE_CACHE: dict[str, bool] = {}
def _is_gh_copilot_deprecation_message(stderr_text: str) -> bool:
"""True iff stderr looks like the deprecated gh-copilot extension's banner."""
lower = stderr_text.lower()
if not any(req in lower for req in _DEPRECATION_REQUIRED):
return False
return any(marker in lower for marker in _DEPRECATION_MARKERS)
return any(req in lower for req in _DEPRECATION_REQUIRED) and any(m in lower for m in _DEPRECATION_MARKERS)
def _resolve_command() -> str:
@@ -68,42 +63,18 @@ def _resolve_command() -> str:
def _resolve_args() -> list[str]:
raw = os.getenv("HERMES_COPILOT_ACP_ARGS", "").strip()
if not raw:
return ["--acp", "--stdio"]
return shlex.split(raw)
# Probe verdicts cached per binary path so repeated prompts against a
# CLI that supports --acp pay the ~50ms --help cost exactly once per
# process. Only definitive verdicts (True/False) are cached; an
# inconclusive probe (binary missing, --help crashed or timed out) is
# not cached so a CLI installed mid-session is picked up.
_ACP_PROBE_CACHE: dict[str, bool] = {}
return shlex.split(raw) if raw else ["--acp", "--stdio"]
def _acp_supported(command: str, args: list[str]) -> bool | None:
"""Tri-state probe: does ``command`` accept the ACP args we'd pass?
"""Tri-state probe: does ``command`` accept ``--acp``?
Different CLI versions support different transports. The GitHub
Copilot CLI (`@github/copilot`, late 2025+) ships with ``--acp``;
older releases (and Claude Code v2.x as of Aug 2026) do not.
Spawning a CLI that doesn't recognize the flag silently exits
with code 1 and ``error: unknown option '--acp'`` on stderr,
after which every delegate_task call hangs the parent for
``child_timeout_seconds`` (default 600s) waiting for stdout
that never arrives.
Returns:
- ``True`` — help text advertises ``--acp``; safe to spawn.
- ``False`` — help ran cleanly but ``--acp`` is absent; spawning
would hang, so the caller should fast-fail with a clear error.
- ``None`` — inconclusive (binary missing, --help failed or
timed out). The caller must fall through to the normal spawn
path, which surfaces the existing "Could not start Copilot ACP
command" error with full context.
Only probes when ``--acp`` is actually among ``args``: a custom
HERMES_COPILOT_ACP_ARGS transport is the operator's business.
A CLI without the flag (older releases, Claude Code v2.x) exits 1 with
``error: unknown option '--acp'`` and the parent then waits the full
child timeout for stdout that never arrives. True = help advertises --acp;
False = help ran cleanly without it (caller fast-fails); None = inconclusive
(binary missing / --help failed), caller falls through to the normal spawn error.
Only probes when ``--acp`` is among ``args`` — a custom transport is the operator's business.
"""
if "--acp" not in args:
return True
@@ -111,32 +82,25 @@ def _acp_supported(command: str, args: list[str]) -> bool | None:
if cached is not None:
return cached
try:
probe = subprocess.run(
[command, "--help"],
capture_output=True, text=True, timeout=5,
)
probe = subprocess.run([command, "--help"], capture_output=True, text=True, timeout=5)
except (FileNotFoundError, subprocess.TimeoutExpired, OSError):
return None
if probe.returncode != 0:
# --help itself failed; can't tell anything about --acp.
return None
# Match ``--acp`` as a flag in the help text; tolerate spacing and
# variants like ``[--acp]``.
# ``--acp`` as a flag token; tolerate spacing and ``[--acp]`` variants.
verdict = bool(re.search(r"(?:^|[\s\[])--acp(?:[\s=\],]|$)", probe.stdout, re.MULTILINE))
_ACP_PROBE_CACHE[command] = verdict
return verdict
def _resolve_home_dir() -> str:
"""Return a stable HOME for child ACP processes."""
"""Return a stable HOME for child ACP processes; /tmp as a last resort so the child never starts HOME-less."""
home = os.environ.get("HOME", "").strip()
if home:
return home
expanded = os.path.expanduser("~")
if expanded and expanded != "~":
return expanded
try:
import pwd
@@ -145,46 +109,29 @@ def _resolve_home_dir() -> str:
return resolved
except Exception:
pass
# Last resort: /tmp (writable on any POSIX system). Avoids crashing the
# subprocess with no HOME; callers can set HERMES_HOME explicitly if they
# need a different writable dir.
return "/tmp"
def _build_subprocess_env() -> dict[str, str]:
# Copilot ACP is a model-driving CLI executor: it legitimately needs LLM
# provider credentials. Route through the central helper so Tier-1 secrets
# (gateway bot tokens, GitHub auth, infra) are still stripped (#29157).
# Copilot ACP drives a model and legitimately needs LLM provider credentials;
# the central helper still strips Tier-1 secrets (bot tokens, GitHub auth, infra).
env = hermes_subprocess_env(inherit_credentials=True)
home = _resolve_home_dir()
env["HOME"] = home
env["HOME"] = _resolve_home_dir()
from hermes_constants import apply_subprocess_home_env
apply_subprocess_home_env(env)
return env
def _jsonrpc_result(message_id: Any, result: Any) -> dict[str, Any]:
return {"jsonrpc": "2.0", "id": message_id, "result": result}
def _jsonrpc_error(message_id: Any, code: int, message: str) -> dict[str, Any]:
return {
"jsonrpc": "2.0",
"id": message_id,
"error": {
"code": code,
"message": message,
},
}
return {"jsonrpc": "2.0", "id": message_id, "error": {"code": code, "message": message}}
def _permission_denied(message_id: Any) -> dict[str, Any]:
return {
"jsonrpc": "2.0",
"id": message_id,
"result": {
"outcome": {
"outcome": "cancelled",
}
},
}
return _jsonrpc_result(message_id, {"outcome": {"outcome": "cancelled"}})
def _model_selection_request(
@@ -264,14 +211,10 @@ def _format_messages_as_prompt(
"IMPORTANT: If you take an action with a tool, you MUST output tool calls using <tool_call>{...}</tool_call> blocks with JSON exactly in OpenAI function-call shape.",
"If no tool is needed, answer normally.",
]
# Deliberately no "requested model" line in the prompt: the model is
# applied for real via ACP session/set_model, and when the backend can't
# honor it (org-policy-disabled id) a prompt-text mention makes the
# serving model FALSELY self-identify as the requested one. Identity
# must come from the backend, not from prompt suggestion.
# Copilot has no tools of its own that would collide with Hermes', so it
# forwards the whole toolset (no allowlist).
# Deliberately no "requested model" line: the model is applied for real via ACP
# session/set_model; a prompt-text mention makes a substituted backend model
# FALSELY self-identify as the requested one. Identity comes from the backend.
# Copilot has no tools of its own that collide with Hermes', so forward the whole toolset.
sections.extend(_render_tool_bridge_sections(tools, tool_choice))
transcript: list[str] = []
@@ -279,28 +222,13 @@ def _format_messages_as_prompt(
if not isinstance(message, dict):
continue
role = str(message.get("role") or "unknown").strip().lower()
if role == "tool":
role = "tool"
elif role not in {"system", "user", "assistant"}:
if role not in _ROLE_LABELS:
role = "context"
content = message.get("content")
rendered = _render_message_content(content)
if not rendered:
continue
label = {
"system": "System",
"user": "User",
"assistant": "Assistant",
"tool": "Tool",
"context": "Context",
}.get(role, role.title())
transcript.append(f"{label}:\n{rendered}")
rendered = _render_message_content(message.get("content"))
if rendered:
transcript.append(f"{_ROLE_LABELS[role]}:\n{rendered}")
if transcript:
sections.append("Conversation transcript:\n\n" + "\n\n".join(transcript))
sections.append("Continue the conversation from the latest user request.")
return "\n\n".join(section.strip() for section in sections if section and section.strip())
@@ -313,7 +241,7 @@ def _render_message_content(content: Any) -> str:
if isinstance(content, dict):
if "text" in content:
return str(content.get("text") or "").strip()
if "content" in content and isinstance(content.get("content"), str):
if isinstance(content.get("content"), str):
return str(content.get("content") or "").strip()
return json.dumps(content, ensure_ascii=True)
if isinstance(content, list):
@@ -342,6 +270,58 @@ def _ensure_path_within_cwd(path_text: str, cwd: str) -> Path:
return resolved
def _effective_timeout(timeout: Any) -> float:
"""Normalise a float or httpx.Timeout-like object to wall-clock seconds (largest component wins)."""
if timeout is None:
return _DEFAULT_TIMEOUT_SECONDS
if isinstance(timeout, (int, float)):
return float(timeout)
_candidates = [getattr(timeout, attr, None) for attr in ("read", "write", "connect", "pool", "timeout")]
_numeric = [float(v) for v in _candidates if isinstance(v, (int, float))]
return max(_numeric) if _numeric else _DEFAULT_TIMEOUT_SECONDS
def _fs_read_text_file(params: dict[str, Any], cwd: str) -> Any:
path = _ensure_path_within_cwd(str(params.get("path") or ""), cwd)
block_error = get_read_block_error(str(path))
if block_error:
raise PermissionError(block_error)
try:
content = path.read_text(encoding="utf-8")
except FileNotFoundError:
content = ""
line = params.get("line")
limit = params.get("limit")
if isinstance(line, int) and line > 1:
lines = content.splitlines(keepends=True)
start = line - 1
end = start + limit if isinstance(limit, int) and limit > 0 else None
content = "".join(lines[start:end])
if content:
content = redact_sensitive_text(content, force=True)
return {"content": content}
def _fs_write_text_file(params: dict[str, Any], cwd: str) -> Any:
path = _ensure_path_within_cwd(str(params.get("path") or ""), cwd)
denied = get_write_denied_error(str(path))
if denied:
raise PermissionError(denied)
# Approval-gated paths (e.g. ~/.ssh/config) are only soft-gated for interactive
# tools, but the ACP shim has no human channel to confirm — fail closed.
if is_write_approval_required(str(path)):
raise PermissionError(
f"Write denied: '{path}' requires interactive approval "
"and cannot be written through the ACP file bridge."
)
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(str(params.get("content") or ""), encoding="utf-8")
return None
_FS_HANDLERS = {"fs/read_text_file": _fs_read_text_file, "fs/write_text_file": _fs_write_text_file}
class _ACPChatCompletions:
def __init__(self, client: "CopilotACPClient"):
self._client = client
@@ -389,7 +369,6 @@ class CopilotACPClient:
self._active_process_lock = threading.Lock()
def close(self) -> None:
proc: subprocess.Popen[str] | None
with self._active_process_lock:
proc = self._active_process
self._active_process = None
@@ -416,34 +395,10 @@ class CopilotACPClient:
stream: bool = False,
**_: Any,
) -> Any:
prompt_text = _format_messages_as_prompt(
messages or [],
model=model,
tools=tools,
tool_choice=tool_choice,
)
# Normalise timeout: run_agent.py may pass an httpx.Timeout object
# (used natively by the OpenAI SDK) rather than a plain float.
if timeout is None:
_effective_timeout = _DEFAULT_TIMEOUT_SECONDS
elif isinstance(timeout, (int, float)):
_effective_timeout = float(timeout)
else:
# httpx.Timeout or similar — pick the largest component so the
# subprocess has enough wall-clock time for the full response.
_candidates = [
getattr(timeout, attr, None)
for attr in ("read", "write", "connect", "pool", "timeout")
]
_numeric = [float(v) for v in _candidates if isinstance(v, (int, float))]
_effective_timeout = max(_numeric) if _numeric else _DEFAULT_TIMEOUT_SECONDS
prompt_text = _format_messages_as_prompt(messages or [], model=model, tools=tools, tool_choice=tool_choice)
response_text, reasoning_text = self._run_prompt(
prompt_text,
timeout_seconds=_effective_timeout,
model=model,
prompt_text, timeout_seconds=_effective_timeout(timeout), model=model
)
tool_calls, cleaned_text = _extract_tool_calls_from_text(response_text)
usage = SimpleNamespace(
@@ -459,33 +414,14 @@ class CopilotACPClient:
reasoning_content=reasoning_text or None,
reasoning_details=None,
)
finish_reason = "tool_calls" if tool_calls else "stop"
choice = SimpleNamespace(message=assistant_message, finish_reason=finish_reason)
completion = SimpleNamespace(
choices=[choice],
usage=usage,
model=model or "copilot-acp",
)
if stream:
return _completion_to_stream_chunks(completion)
return completion
choice = SimpleNamespace(message=assistant_message, finish_reason="tool_calls" if tool_calls else "stop")
completion = SimpleNamespace(choices=[choice], usage=usage, model=model or "copilot-acp")
return _completion_to_stream_chunks(completion) if stream else completion
def _run_prompt(
self,
prompt_text: str,
*,
timeout_seconds: float,
model: str | None = None,
) -> tuple[str, str]:
# Fast-fail when the CLI doesn't support the ACP args we'd pass.
# Without this guard, a CLI like Claude Code v2.x exits with
# ``error: unknown option '--acp'`` immediately, then the parent
# ACP loop waits the full ``child_timeout_seconds`` (default 600s)
# for stdout that never arrives. The probe costs ~50ms and turns
# a 600s silent hang into a 280ms clear error.
# ``None`` (inconclusive probe — e.g. binary missing) falls
# through to the spawn below, which raises the established
# "Could not start Copilot ACP command" error.
def _spawn(self) -> subprocess.Popen[str]:
# Fast-fail when the CLI rejects --acp: without the probe the parent waits
# the full child timeout for stdout that never arrives. ``None`` falls
# through to the spawn, which raises the established start error.
if _acp_supported(self._acp_command, self._acp_args) is False:
preview = " ".join(self._acp_args[:3]) if self._acp_args else "(none)"
raise RuntimeError(
@@ -498,17 +434,8 @@ class CopilotACPClient:
f"HERMES_COPILOT_ACP_COMMAND / HERMES_COPILOT_ACP_ARGS "
f"to a working pair."
)
# Note the model Hermes selected; it is applied after session/new via
# the ACP-native `session/set_model` call. The CLI's `--model` spawn
# flag is deliberately NOT used here: `copilot --acp` validates it
# (an unknown id aborts the spawn) but then ignores it for the actual
# session, so it adds a failure mode without selecting anything.
requested_model = str(model or "").strip()
try:
# Hide the console the CLI child would otherwise flash on Windows
# (#56747). Hide-only — stdio pipes stay intact for the ACP wire.
# Hide the console the child would flash on Windows; stdio pipes stay intact.
from hermes_cli._subprocess_compat import windows_hide_flags
proc = subprocess.Popen(
@@ -527,21 +454,24 @@ class CopilotACPClient:
f"Could not start Copilot ACP command '{self._acp_command}'. "
"Install GitHub Copilot CLI or set HERMES_COPILOT_ACP_COMMAND/COPILOT_CLI_PATH."
) from exc
if proc.stdin is None or proc.stdout is None:
proc.kill()
raise RuntimeError("Copilot ACP process did not expose stdin/stdout pipes.")
self.is_closed = False
with self._active_process_lock:
self._active_process = proc
return proc
def _run_prompt(self, prompt_text: str, *, timeout_seconds: float, model: str | None = None) -> tuple[str, str]:
# The CLI's `--model` spawn flag is deliberately NOT used: `copilot --acp`
# validates it (unknown id aborts the spawn) but ignores it for the session.
# The model is applied after session/new via ACP-native model selection.
requested_model = str(model or "").strip()
proc = self._spawn()
inbox: queue.Queue[dict[str, Any]] = queue.Queue()
stderr_tail: deque[str] = deque(maxlen=40)
def _stdout_reader() -> None:
if proc.stdout is None:
return
for line in proc.stdout:
try:
inbox.put(json.loads(line))
@@ -554,24 +484,15 @@ class CopilotACPClient:
for line in proc.stderr:
stderr_tail.append(line.rstrip("\n"))
out_thread = threading.Thread(target=_stdout_reader, daemon=True)
err_thread = threading.Thread(target=_stderr_reader, daemon=True)
out_thread.start()
err_thread.start()
threading.Thread(target=_stdout_reader, daemon=True).start()
threading.Thread(target=_stderr_reader, daemon=True).start()
next_id = 0
def _request(method: str, params: dict[str, Any], *, text_parts: list[str] | None = None, reasoning_parts: list[str] | None = None) -> Any:
nonlocal next_id
next_id += 1
request_id = next_id
payload = {
"jsonrpc": "2.0",
"id": request_id,
"method": method,
"params": params,
}
proc.stdin.write(json.dumps(payload) + "\n")
proc.stdin.write(json.dumps({"jsonrpc": "2.0", "id": request_id, "method": method, "params": params}) + "\n")
proc.stdin.flush()
deadline = time.monotonic() + timeout_seconds
@@ -582,23 +503,15 @@ class CopilotACPClient:
msg = inbox.get(timeout=0.1)
except queue.Empty:
continue
if self._handle_server_message(
msg,
process=proc,
cwd=self._acp_cwd,
text_parts=text_parts,
reasoning_parts=reasoning_parts,
msg, process=proc, cwd=self._acp_cwd, text_parts=text_parts, reasoning_parts=reasoning_parts
):
continue
if msg.get("id") != request_id:
continue
if "error" in msg:
err = msg.get("error") or {}
raise RuntimeError(
f"Copilot ACP {method} failed: {err.get('message') or err}"
)
raise RuntimeError(f"Copilot ACP {method} failed: {err.get('message') or err}")
return msg.get("result")
stderr_text = "\n".join(stderr_tail).strip()
@@ -626,35 +539,17 @@ class CopilotACPClient:
"initialize",
{
"protocolVersion": 1,
"clientCapabilities": {
"fs": {
"readTextFile": True,
"writeTextFile": True,
}
},
"clientInfo": {
"name": "hermes-agent",
"title": "Hermes Agent",
"version": "0.0.0",
},
"clientCapabilities": {"fs": {"readTextFile": True, "writeTextFile": True}},
"clientInfo": {"name": "hermes-agent", "title": "Hermes Agent", "version": "0.0.0"},
},
)
session = _request(
"session/new",
{
"cwd": self._acp_cwd,
"mcpServers": [],
},
) or {}
session = _request("session/new", {"cwd": self._acp_cwd, "mcpServers": []}) or {}
session_id = str(session.get("sessionId") or "").strip()
if not session_id:
raise RuntimeError("Copilot ACP did not return a sessionId.")
# Select the model Hermes asked for. Prefer the stable ACP v1
# session-config API: session/new advertises a category="model"
# select option and session/set_config_option updates it. Copilot
# still exposes the older models/session/set_model extension too,
# so retain that only as compatibility fallback for older agents.
# Prefer the stable ACP v1 session-config API (category="model" select
# option + session/set_config_option); session/set_model is the fallback.
if requested_model and requested_model != "copilot-acp":
try:
selection = _model_selection_request(session, requested_model)
@@ -679,15 +574,7 @@ class CopilotACPClient:
reasoning_parts: list[str] = []
_request(
"session/prompt",
{
"sessionId": session_id,
"prompt": [
{
"type": "text",
"text": prompt_text,
}
],
},
{"sessionId": session_id, "prompt": [{"type": "text", "text": prompt_text}]},
text_parts=text_parts,
reasoning_parts=reasoning_parts,
)
@@ -704,18 +591,16 @@ class CopilotACPClient:
text_parts: list[str] | None,
reasoning_parts: list[str] | None,
) -> bool:
"""Consume a server->client message; True when handled (notification or request answered)."""
method = msg.get("method")
if not isinstance(method, str):
return False
if method == "session/update":
params = msg.get("params") or {}
update = params.get("update") or {}
update = (msg.get("params") or {}).get("update") or {}
kind = str(update.get("sessionUpdate") or "").strip()
content = update.get("content") or {}
chunk_text = ""
if isinstance(content, dict):
chunk_text = str(content.get("text") or "")
chunk_text = str(content.get("text") or "") if isinstance(content, dict) else ""
if kind == "agent_message_chunk" and chunk_text and text_parts is not None:
text_parts.append(chunk_text)
elif kind == "agent_thought_chunk" and chunk_text and reasoning_parts is not None:
@@ -727,66 +612,15 @@ class CopilotACPClient:
message_id = msg.get("id")
params = msg.get("params") or {}
if method == "session/request_permission":
response = _permission_denied(message_id)
elif method == "fs/read_text_file":
elif method in _FS_HANDLERS:
try:
path = _ensure_path_within_cwd(str(params.get("path") or ""), cwd)
block_error = get_read_block_error(str(path))
if block_error:
raise PermissionError(block_error)
try:
content = path.read_text(encoding="utf-8")
except FileNotFoundError:
content = ""
line = params.get("line")
limit = params.get("limit")
if isinstance(line, int) and line > 1:
lines = content.splitlines(keepends=True)
start = line - 1
end = start + limit if isinstance(limit, int) and limit > 0 else None
content = "".join(lines[start:end])
if content:
content = redact_sensitive_text(content, force=True)
response = {
"jsonrpc": "2.0",
"id": message_id,
"result": {
"content": content,
},
}
except Exception as exc:
response = _jsonrpc_error(message_id, -32602, str(exc))
elif method == "fs/write_text_file":
try:
path = _ensure_path_within_cwd(str(params.get("path") or ""), cwd)
denied = get_write_denied_error(str(path))
if denied:
raise PermissionError(denied)
# Approval-gated paths (e.g. ~/.ssh/config) are not hard-denied
# for interactive tools, but the ACP shim has no human channel
# to confirm the write — fail closed here.
if is_write_approval_required(str(path)):
raise PermissionError(
f"Write denied: '{path}' requires interactive approval "
"and cannot be written through the ACP file bridge."
)
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(str(params.get("content") or ""), encoding="utf-8")
response = {
"jsonrpc": "2.0",
"id": message_id,
"result": None,
}
response = _jsonrpc_result(message_id, _FS_HANDLERS[method](params, cwd))
except Exception as exc:
response = _jsonrpc_error(message_id, -32602, str(exc))
else:
response = _jsonrpc_error(
message_id,
-32601,
f"ACP client method '{method}' is not supported by Hermes yet.",
)
response = _jsonrpc_error(message_id, -32601, f"ACP client method '{method}' is not supported by Hermes yet.")
process.stdin.write(json.dumps(response) + "\n")
process.stdin.flush()

View File

@@ -1,11 +1,9 @@
"""Best-effort early import for the OpenAI SDK's native streaming parser.
"""Best-effort early import of the OpenAI SDK's native streaming parser.
The OpenAI SDK imports ``jiter`` while constructing streaming chat-completion
responses. On some Windows installs the native extension can be imported
directly from the Hermes venv, but the first import fails when it happens later
inside the threaded streaming request path. Loading it once during agent
package import avoids that import-order failure while preserving the normal
SDK error path for genuinely missing or broken installs.
On some Windows installs ``jiter``'s native extension imports fine from the venv
but fails when first imported later inside the threaded streaming path. Loading
it once at agent-package import avoids that while preserving the SDK's normal
error path for genuinely broken installs.
"""
from __future__ import annotations
@@ -18,19 +16,15 @@ _JITER_PRELOAD_ERROR: Exception | None = None
def preload_jiter_native_extension() -> bool:
"""Import jiter's native extension early if it is available."""
global _JITER_PRELOADED, _JITER_PRELOAD_ERROR
if _JITER_PRELOADED:
return True
try:
importlib.import_module("jiter.jiter")
from jiter import from_json as _from_json # noqa: F401
except Exception as exc:
_JITER_PRELOAD_ERROR = exc
return False
_JITER_PRELOADED = True
_JITER_PRELOAD_ERROR = None
return True

View File

@@ -1,8 +1,5 @@
"""Preventive SSL CA certificate checks for Hermes Agent.
This module catches broken CA bundle paths before OpenAI/httpx turns them into
opaque ``FileNotFoundError: [Errno 2] No such file or directory`` failures.
"""
"""Preventive SSL CA certificate checks — catch broken CA bundle paths before
OpenAI/httpx turns them into an opaque ``FileNotFoundError``."""
from __future__ import annotations
@@ -15,32 +12,23 @@ from agent.errors import SSLConfigurationError
logger = logging.getLogger(__name__)
_CA_BUNDLE_ENV_VARS = (
"HERMES_CA_BUNDLE",
"SSL_CERT_FILE",
"REQUESTS_CA_BUNDLE",
"CURL_CA_BUNDLE",
)
_CA_BUNDLE_ENV_VARS = ("HERMES_CA_BUNDLE", "SSL_CERT_FILE", "REQUESTS_CA_BUNDLE", "CURL_CA_BUNDLE")
_SKIP_VALUES = {"1", "true", "yes", "on"}
_REPAIR_HINT = (
"Repair: run `hermes doctor --fix` (auto-reinstalls certifi), or "
"manually: python -m pip install --force-reinstall certifi openai httpx\n"
"If you configured a custom corporate CA bundle, fix or unset the "
"broken CA bundle environment variable."
)
def _skip_ssl_guard_enabled() -> bool:
return os.getenv("HERMES_SKIP_SSL_GUARD", "").strip().lower() in _SKIP_VALUES
def _repair_hint() -> str:
return (
"Repair: run `hermes doctor --fix` (auto-reinstalls certifi), or "
"manually: python -m pip install --force-reinstall certifi openai httpx\n"
"If you configured a custom corporate CA bundle, fix or unset the "
"broken CA bundle environment variable."
)
def _ssl_err(message: str) -> SSLConfigurationError:
"""Create a consistent, user-actionable SSL configuration error."""
return SSLConfigurationError(f"{message}\n{_repair_hint()}")
return SSLConfigurationError(f"{message}\n{_REPAIR_HINT}")
def _validate_bundle_path(label: str, value: str, *, require_substantial: bool = False) -> None:
@@ -58,8 +46,7 @@ def _validate_bundle_path(label: str, value: str, *, require_substantial: bool =
try:
loaded_certs = ctx.get_ca_certs()
except NotImplementedError:
# truststore-backed SSLContext (Windows OS trust store) doesn't
# implement get_ca_certs(); bundle was already validated above.
# truststore-backed SSLContext (Windows OS trust store) lacks get_ca_certs(); loading above already validated it.
return
if not loaded_certs:
raise _ssl_err(f"{label} CA bundle at {value} did not load any certificates")
@@ -68,34 +55,23 @@ def _validate_bundle_path(label: str, value: str, *, require_substantial: bool =
def verify_ca_bundle() -> None:
"""Verify configured and bundled CA certificates are present and loadable.
Raises:
SSLConfigurationError: If an explicit CA-bundle environment variable
points at a bad path, or if certifi's bundled ``cacert.pem`` is
missing/corrupt.
Raises SSLConfigurationError when an explicit CA-bundle env var points at a
bad path or certifi's bundled ``cacert.pem`` is missing/corrupt.
"""
if _skip_ssl_guard_enabled():
logger.debug("SSL CA bundle guard skipped via HERMES_SKIP_SSL_GUARD")
return
for env_var in _CA_BUNDLE_ENV_VARS:
value = os.getenv(env_var)
if value:
_validate_bundle_path(env_var, value)
try:
import certifi
except Exception as exc:
raise _ssl_err(f"certifi is not importable: {exc}") from exc
ca_bundle = str(certifi.where())
_validate_bundle_path("certifi", ca_bundle, require_substantial=True)
_validate_bundle_path("certifi", str(certifi.where()), require_substantial=True)
def verify_ca_bundle_with_fallback() -> None:
"""Backward-compatible wrapper for older call sites.
The old PR name mentioned a platform fallback, but allowing startup with a
broken certifi bundle still leaves httpx/OpenAI and requests call sites
failing later. Keep the wrapper name but enforce the same check.
"""
"""Backward-compatible name for older call sites; a broken certifi bundle fails later anyway, so enforce the same check."""
verify_ca_bundle()

View File

@@ -10,13 +10,13 @@ from typing import Any, Optional
logger = logging.getLogger(__name__)
_CA_BUNDLE_ENV_VARS = ("HERMES_CA_BUNDLE", "SSL_CERT_FILE", "REQUESTS_CA_BUNDLE", "CURL_CA_BUNDLE")
def _coerce_insecure(ssl_verify: Any) -> bool:
if ssl_verify is False:
return True
if isinstance(ssl_verify, str) and ssl_verify.strip().lower() in {"false", "0", "no", "off"}:
return True
return False
return isinstance(ssl_verify, str) and ssl_verify.strip().lower() in {"false", "0", "no", "off"}
def resolve_httpx_verify(
@@ -25,17 +25,8 @@ def resolve_httpx_verify(
ssl_verify: Any = None,
base_url: str = "",
) -> bool | ssl.SSLContext:
"""Resolve httpx ``verify`` for provider HTTP clients.
Priority:
1. ``ssl_verify: false`` — disable verification (local dev only)
2. explicit ``ca_bundle`` (per-provider ``ssl_ca_cert`` config field)
3. ``HERMES_CA_BUNDLE``, ``SSL_CERT_FILE``, ``REQUESTS_CA_BUNDLE``,
``CURL_CA_BUNDLE`` env vars
4. ``True`` (httpx/certifi default)
``base_url`` is used only for the insecure-mode warning message.
"""
"""Resolve httpx ``verify``: ``ssl_verify: false`` > explicit ``ca_bundle`` >
CA-bundle env vars > ``True`` (certifi default). ``base_url`` only feeds the warning."""
if _coerce_insecure(ssl_verify):
logger.warning(
"TLS certificate verification DISABLED (ssl_verify: false) for %s — "
@@ -45,13 +36,11 @@ def resolve_httpx_verify(
)
return False
effective_ca = (
(ca_bundle or "").strip()
or os.getenv("HERMES_CA_BUNDLE", "").strip()
or os.getenv("SSL_CERT_FILE", "").strip()
or os.getenv("REQUESTS_CA_BUNDLE", "").strip()
or os.getenv("CURL_CA_BUNDLE", "").strip()
)
effective_ca = (ca_bundle or "").strip()
for env_var in _CA_BUNDLE_ENV_VARS:
if effective_ca:
break
effective_ca = os.getenv(env_var, "").strip()
if effective_ca:
ca_path = str(Path(effective_ca).expanduser())
if os.path.isfile(ca_path):

View File

@@ -1,15 +1,10 @@
"""Stream diagnostics — per-attempt counters, exception chains, retry logging.
When a streaming chat-completions request dies mid-response, we want to
know why: which Cloudflare edge served the request, which OpenRouter
downstream provider answered, how many bytes/chunks we got before the
drop, the HTTP status, the underlying httpx error class. These helpers
collect that info and emit it both to ``agent.log`` (full detail) and to
the user-facing status line (compact).
All helpers are extracted from :class:`AIAgent` for cleanliness.
``run_agent`` keeps thin forwarder methods so existing call sites and
tests that patch ``run_agent.<helper>`` keep working.
When a streaming request dies mid-response these helpers record WHY (which CF
edge / OpenRouter downstream served it, bytes+chunks before the drop, HTTP
status, underlying httpx error class) to ``agent.log`` in full and to the
user-facing status line compactly. ``run_agent`` keeps thin forwarders so
existing call sites and tests patching ``run_agent.<helper>`` keep working.
"""
from __future__ import annotations
@@ -20,10 +15,7 @@ from typing import Any, Dict, List, Optional
logger = logging.getLogger(__name__)
# Per-attempt stream diagnostic headers. Lowercased; httpx returns
# CIMultiDict so case-insensitive lookups already work, but we read .get()
# on the dict from agent.log for free-form post-hoc analysis.
# Lowercased upstream headers captured per attempt for post-hoc analysis.
STREAM_DIAG_HEADERS = (
"cf-ray",
"cf-cache-status",
@@ -39,12 +31,7 @@ STREAM_DIAG_HEADERS = (
def stream_diag_init() -> Dict[str, Any]:
"""Return a fresh per-attempt diagnostic dict.
Mutated in-place by the streaming functions and read from the retry
block when a stream dies. Lives on ``request_client_holder`` so it
survives across the closure boundary.
"""
"""Fresh per-attempt diagnostic dict; mutated in place by the streaming functions and read by the retry block."""
return {
"started_at": time.time(),
"first_chunk_at": None,
@@ -56,12 +43,7 @@ def stream_diag_init() -> Dict[str, Any]:
def stream_diag_capture_response(agent: Any, diag: Dict[str, Any], http_response: Any) -> None:
"""Snapshot interesting headers + HTTP status from the live stream.
Called once at stream open (before iterating chunks) so the metadata
survives even if the stream dies before any chunk arrives. Failures
are swallowed — diag is best-effort.
"""
"""Snapshot headers + HTTP status at stream open (so they survive a drop before the first chunk). Best-effort."""
if http_response is None or not isinstance(diag, dict):
return
try:
@@ -71,14 +53,12 @@ def stream_diag_capture_response(agent: Any, diag: Dict[str, Any], http_response
try:
headers = getattr(http_response, "headers", None) or {}
captured: Dict[str, str] = {}
# Allow per-agent override of the headers list (back-compat).
target_headers = getattr(agent, "_STREAM_DIAG_HEADERS", STREAM_DIAG_HEADERS)
for name in target_headers:
# Per-agent override of the headers list (back-compat).
for name in getattr(agent, "_STREAM_DIAG_HEADERS", STREAM_DIAG_HEADERS):
try:
val = headers.get(name)
if val:
# Truncate single-value to keep log lines bounded.
captured[name] = str(val)[:120]
captured[name] = str(val)[:120] # keep log lines bounded
except Exception:
continue
diag["headers"] = captured
@@ -87,14 +67,11 @@ def stream_diag_capture_response(agent: Any, diag: Dict[str, Any], http_response
def flatten_exception_chain(error: BaseException) -> str:
"""Return a compact ``Outer(msg) <- Inner(msg) <- ...`` rendering.
"""Compact ``Outer(msg) <- Inner(msg) <- ...`` rendering.
OpenAI SDK wraps httpx errors as ``APIConnectionError`` /
``APIError`` and only the wrapper's class is visible at the catch
site — but the underlying ``RemoteProtocolError`` /
``ConnectError`` / ``ReadError`` is what tells us WHY the stream
died. Walks ``__cause__`` then ``__context__`` (deduped, max 4
deep) to surface the chain in one line.
The OpenAI SDK wraps httpx errors so only the wrapper class is visible at the
catch site; the inner RemoteProtocolError/ConnectError/ReadError says WHY the
stream died. Walks ``__cause__`` then ``__context__`` (deduped, max 4 deep).
"""
seen: List[BaseException] = []
link: Optional[BaseException] = error
@@ -102,9 +79,7 @@ def flatten_exception_chain(error: BaseException) -> str:
if link in seen:
break
seen.append(link)
nxt = getattr(link, "__cause__", None) or getattr(
link, "__context__", None
)
nxt = getattr(link, "__cause__", None) or getattr(link, "__context__", None)
if nxt is None or nxt is link:
break
link = nxt
@@ -127,19 +102,12 @@ def log_stream_retry(
mid_tool_call: bool,
diag: Optional[Dict[str, Any]] = None,
) -> None:
"""Record a transient stream-drop and retry to ``agent.log``.
"""Structured WARNING to ``agent.log`` for a transient stream drop + retry.
Always logs a structured WARNING so users have a breadcrumb regardless
of UI verbosity. Subagents in particular benefit because their
retries no longer spam the parent's terminal — but the file log keeps
full detail (provider, error class, attempt, base_url, subagent_id).
When *diag* is provided (the per-attempt stream-diagnostic dict from
:func:`stream_diag_init`), the WARNING also captures upstream headers
(cf-ray, x-openrouter-provider, x-openrouter-id), HTTP status, bytes
streamed before the drop, and elapsed time on the dying attempt.
These are the breadcrumbs needed to answer "is one CF edge / one
downstream provider responsible, or is it random across runs?"
Always logged regardless of UI verbosity (subagent retries no longer spam the
parent's terminal but keep full detail here). With *diag*, also records upstream
headers, HTTP status, bytes/chunks streamed, elapsed and TTFB on the dying attempt —
enough to tell "one CF edge / downstream provider" from "random across runs".
"""
try:
try:
@@ -148,14 +116,11 @@ def log_stream_retry(
_summary = str(error)
if _summary and len(_summary) > 240:
_summary = _summary[:240] + "…"
# Inner-cause chain (httpx errors hide under openai.APIError).
try:
_chain = flatten_exception_chain(error)
except Exception:
_chain = type(error).__name__
# Per-attempt counters and upstream headers.
_now = time.time()
_bytes = 0
_chunks = 0
@@ -174,9 +139,7 @@ def log_stream_retry(
_ttfb = max(0.0, float(_first) - _started)
headers = diag.get("headers") or {}
if isinstance(headers, dict) and headers:
_headers_repr = " ".join(
f"{k}={v}" for k, v in headers.items()
)
_headers_repr = " ".join(f"{k}={v}" for k, v in headers.items())
if diag.get("http_status") is not None:
_http_status = str(diag.get("http_status"))
except Exception:
@@ -220,35 +183,16 @@ def emit_stream_drop(
mid_tool_call: bool,
diag: Optional[Dict[str, Any]] = None,
) -> None:
"""Emit a single user-visible line for a stream drop+retry.
"""One compact user-visible status line for a stream drop+retry, plus the full WARNING via log_stream_retry.
Both top-level agents and subagents announce drops in the UI — the
parent prefixes subagent lines with ``[subagent-N]`` via ``log_prefix``
so they're easy to attribute. All cases also write a structured
WARNING to ``agent.log`` via :func:`log_stream_retry` with the full
diagnostic detail (subagent_id, provider, base_url, error_type,
cf-ray, x-openrouter-provider, bytes/chunks, elapsed) for post-hoc
analysis.
The user-visible status line is intentionally compact: provider,
error class, attempt N/M, plus ``after Xs`` when the stream dropped
mid-flight. Full diagnostic detail goes to ``agent.log`` only —
``hermes logs --level WARNING | grep "Stream drop"`` to inspect.
Subagent lines get a ``[subagent-N]`` prefix from ``log_prefix``. ``after Xs``
distinguishes "couldn't connect" (0s) from "died mid-stream" (idle-kill / proxy timeout).
"""
kind = "drop mid tool-call" if mid_tool_call else "drop"
log_stream_retry(
agent,
kind=kind,
error=error,
attempt=attempt,
max_attempts=max_attempts,
mid_tool_call=mid_tool_call,
diag=diag,
agent, kind=kind, error=error, attempt=attempt, max_attempts=max_attempts, mid_tool_call=mid_tool_call, diag=diag
)
provider = agent.provider or "provider"
# Compose a brief "after Xs" suffix when we have timing data — helps
# the user distinguish "couldn't connect" (0s) from "died after 30s
# of streaming" (likely upstream idle-kill or proxy timeout).
_suffix = ""
if isinstance(diag, dict):
try:
@@ -262,10 +206,7 @@ def emit_stream_drop(
f"⚠️ {provider} stream {kind} ({type(error).__name__}){_suffix} "
f"— reconnecting, retry {attempt}/{max_attempts}"
)
agent._touch_activity(
f"stream retry {attempt}/{max_attempts} "
f"after {type(error).__name__}"
)
agent._touch_activity(f"stream retry {attempt}/{max_attempts} after {type(error).__name__}")
except Exception:
pass

View File

@@ -1,23 +1,11 @@
"""Best-effort accessors for the single-writer stream fence (#65991).
"""Best-effort accessors for the single-writer stream fence.
The fence itself lives on ``AIAgent`` (``_claim_stream_writer`` /
``_stream_writer_is_current`` in ``run_agent.py``), but the streaming code paths
that use it live in *other* modules — ``chat_completion_helpers`` (chat /
anthropic / bedrock) and ``codex_runtime`` (codex responses). Calling the fence
directly as ``agent._claim_stream_writer()`` from those modules makes them
hard-depend on the method being present on whatever object is passed in as
``agent``.
That coupling is a latent crash: a partially-updated checkout (the streaming
helper module newer than ``run_agent``), a hot-reloaded gateway, a duck-typed
agent, or a test double without the method turns an *additive* safety net into a
fatal ``AttributeError`` that aborts the whole turn. A cron job died exactly
this way with ``'AIAgent' object has no attribute '_claim_stream_writer'``.
The fence is only ever allowed to drop a *provably* superseded stream — never
the sole legitimate writer. So when the guard is unavailable (or raises), the
correct degradation is "no fence": keep streaming. These helpers make the
claim/check best-effort to guarantee that.
The fence lives on ``AIAgent`` (``_claim_stream_writer`` / ``_stream_writer_is_current``)
but is used from other streaming modules. Calling it directly would turn an
*additive* safety net into a fatal AttributeError on a partially-updated checkout,
hot-reloaded gateway, duck-typed agent, or test double (a cron job died this way).
The fence may only drop a *provably* superseded stream, never the sole writer, so
when it is unavailable or raises the correct degradation is "no fence": keep streaming.
"""
from __future__ import annotations
@@ -29,33 +17,18 @@ logger = logging.getLogger(__name__)
def claim_stream_writer(agent: Any) -> int:
"""Claim the delta sink for the calling stream attempt, best-effort.
Returns the agent's monotonic writer token when the fence is available, or
``0`` when the agent doesn't expose it (or the claim raised). A ``0`` token
pairs with :func:`stream_writer_is_current` always returning ``True``, so a
guard-less agent is simply never fenced instead of crashing the turn.
"""
"""Claim the delta sink for this stream attempt; ``0`` (never fenced) when the agent lacks the fence or the claim raised."""
claim = getattr(agent, "_claim_stream_writer", None)
if callable(claim):
try:
return int(claim())
except Exception:
logger.debug(
"stream single-writer: claim failed; proceeding unfenced",
exc_info=True,
)
logger.debug("stream single-writer: claim failed; proceeding unfenced", exc_info=True)
return 0
def stream_writer_is_current(agent: Any, token: int) -> bool:
"""True when ``token`` is still the active writer, best-effort.
A falsy token (from a claim that no-oped) or an agent without the fence
means we cannot prove supersession, so the stream is treated as current and
never fenced. This preserves the single-writer invariant's one-way promise:
only a demonstrably stale writer is ever stopped.
"""
"""True when ``token`` is still the active writer; a falsy token or a fence-less agent cannot prove supersession, so True."""
if not token:
return True
is_current = getattr(agent, "_stream_writer_is_current", None)
@@ -63,8 +36,5 @@ def stream_writer_is_current(agent: Any, token: int) -> bool:
try:
return bool(is_current(token))
except Exception:
logger.debug(
"stream single-writer: is_current check failed; treating as current",
exc_info=True,
)
logger.debug("stream single-writer: is_current check failed; treating as current", exc_info=True)
return True