diff --git a/.github/workflows/lint.yml b/.github/workflows/lint.yml index 9056185d27..6530567435 100644 --- a/.github/workflows/lint.yml +++ b/.github/workflows/lint.yml @@ -182,6 +182,11 @@ jobs: - name: Forbid literal /tmp paths outside the baseline run: python scripts/check_no_tmp_literals.py + # config.yaml is hand-edited and commented; a PyYAML dump of it destroys every comment + # (#92554, regressed repeatedly). All writers go through hermes_cli.config.atomic_config_write. + - name: Forbid config.yaml writers that bypass the comment-preserving writer + run: python scripts/check_config_yaml_writers.py + # The OS lanes import only files carrying the matching marker, so a test that fakes # macOS (is_macos -> True, sys.platform -> "darwin") without `platforms("macos")` is green on # Linux over a faked branch and never runs on macOS (#111866, AGENTS.md § Don't fake the host OS). diff --git a/.github/workflows/skills-index.yml b/.github/workflows/skills-index.yml index 3c36fc773b..7d6097e18f 100644 --- a/.github/workflows/skills-index.yml +++ b/.github/workflows/skills-index.yml @@ -54,9 +54,20 @@ jobs: # Plugin-catalog star counts follow the same rule as the index: GitHub is consulted # only here, on the schedule (one GraphQL request for every catalog repo), and docs # deploys reuse the artifact / live copy without touching the API. + # + # Fresh token: the App installation token minted at job start lives 1 h, and the + # index walk above runs 60-85 min, so the probe used an expired token and 401'd on + # every scheduled run from Sep 18 (#118113). Mint again right before the probe. + - name: Get GitHub App token for the stars probe + id: probe-token + uses: ./.github/actions/get-app-token + with: + client-id: ${{ vars.APP_CLIENT_ID }} + private-key: ${{ secrets.APP_PRIVATE_KEY }} + - name: Probe plugin catalog stars env: - GITHUB_TOKEN: ${{ steps.app-token.outputs.token }} + GITHUB_TOKEN: ${{ steps.probe-token.outputs.token }} run: python website/scripts/fetch-plugin-stars.py --probe - name: Upload index artifact diff --git a/agent/AGENTS.md b/agent/AGENTS.md index a230612e61..9bfd3d135f 100644 --- a/agent/AGENTS.md +++ b/agent/AGENTS.md @@ -98,6 +98,10 @@ cache break — keep it the only one. Full detail: - **Auxiliary (side-LLM) work** — curator, vision, embedding, title generation, session_search, compression — resolves through `agent/auxiliary_client.py::_resolve_auto_route`; each task can pin its own `provider/model/base_url/reasoning_effort` under `auxiliary:` in config.yaml. + Every physical attempt funnels through `_relay_sync_completion` / `_relay_async_completion` / + `_relay_sync_stream`, where `agent/auxiliary_hooks.py` emits `pre_auxiliary_call` / + `post_auxiliary_call` (observer-only, fail-open, `aux_task` set); the main-loop + `pre/post_api_request` events must NOT fire for aux calls (#79733). - Fallback models and credential pools are resolution-chain code: E2E them with real imports against a temp `HERMES_HOME`, not mocks (root rubric). diff --git a/agent/agent_init.py b/agent/agent_init.py index 9d40cbadb4..9bb42c8f7f 100644 --- a/agent/agent_init.py +++ b/agent/agent_init.py @@ -26,6 +26,7 @@ from agent.context_compressor import ContextCompressor from agent.agent_runtime_helpers import _ra from agent.iteration_budget import IterationBudget, normalize_budget_warning_ratio from agent.memory_manager import StreamingContextScrubber +from agent.memory_provider import is_core_memory_provider from agent.session_activity import ActivityProvenance from agent.model_metadata import ( MINIMUM_CONTEXT_LENGTH, fetch_model_metadata, is_local_endpoint, query_ollama_num_ctx @@ -1077,6 +1078,10 @@ def _load_tools(agent, enabled_toolsets, disabled_toolsets): # A finite -q run has no later session to learn for: no skill authoring tool (agent/oneshot_footprint.py). from agent.oneshot_footprint import prune_oneshot_tools agent.tools = prune_oneshot_tools(agent.tools or []) + from tools.connectors.turn import side_agent_tool_drops + drops = side_agent_tool_drops(agent) + if drops: + agent.tools = [t for t in agent.tools if t["function"]["name"] not in drops] agent.valid_tool_names = {tool["function"]["name"] for tool in agent.tools} if agent.tools else set() # Kanban guidance is session-static for the dispatcher-owned worker only. Profiles may @@ -1293,7 +1298,7 @@ def _init_memory(agent, _agent_cfg, skip_memory, platform): if not skip_memory: try: _mem_provider_name = mem_config.get("provider", "") if mem_config else "" - if _mem_provider_name and _mem_provider_name.strip(): + if not is_core_memory_provider(_mem_provider_name): from agent.memory_manager import MemoryManager as _MemoryManager from plugins.memory import load_memory_provider as _load_mem agent._memory_manager = _MemoryManager() @@ -1867,22 +1872,20 @@ def _select_context_engine(_agent_cfg): except Exception: _candidate = None if _candidate is not None and _candidate.name == _engine_name: - # Deep-copy the shared singleton so a child's update_model() can't mutate the - # parent's. Uncopyable state (locks, DB conns) → built-in with an ACCURATE message. - import copy + # The plugin system holds ONE shared instance; each agent gets its own so a child's + # update_model() can't mutate the parent's (#42449). clone_for_agent() defaults to + # deepcopy; engines with uncopyable state (locks, DB conns) override it. A failure + # falls back to the built-in compressor with an ACCURATE message, not "not found". try: - # Copy can fail for engines holding uncopyable state (locks, DB connections, clients); in - # that case fall back to the built-in compressor with an ACCURATE message rather than - # silently mislabelling it "not found". See #42449. - _selected_engine = copy.deepcopy(_candidate) + _selected_engine = _candidate.clone_for_agent() except Exception as _copy_err: _copy_failed = True _ra().logger.warning( "Context engine '%s' could not be safely copied for this " "agent (%s) — falling back to built-in compressor. Plugin " "engines that hold uncopyable state (locks, DB connections) " - "should implement __deepcopy__ to copy only mutable budget " - "state.", + "should override clone_for_agent() (or __deepcopy__) to copy " + "only mutable budget state.", _engine_name, _copy_err, ) @@ -2285,6 +2288,7 @@ _PASSTHROUGH_PARAMS = ( "enabled_toolsets", "disabled_toolsets", # Model response configuration (None = provider/model default) "max_tokens", "reasoning_config", "service_tier", + "side_agent", ) # Gateway identity params stored as ``agent._``. gateway_session_key is the stable # per-chat key (e.g. agent:main:telegram:dm:123). @@ -2338,22 +2342,8 @@ def init_agent( checkpoint_max_snapshots: int = 20, checkpoint_max_total_size_mb: int = 500, checkpoint_max_file_size_mb: int = 10, pass_session_id: bool = False, requested_provider: str = None, capabilities: Optional[Dict[str, bool]] = None, cwd: Optional[str] = None, + side_agent: bool = False, ): - """Initialize the AI Agent (body of :meth:`AIAgent.__init__`). - - Non-obvious parameters: - max_iterations: default unlimited (sys.maxsize); the budget is shared with subagents. - requested_provider: provider identity before runtime canonicalization. - cwd: logical session workspace, available to memory providers during construction; - None or empty leaves the runtime cwd resolver unpinned. - openrouter_min_coding_score: coding-score floor for ``openrouter/pareto-code`` only. - clarify_callback: ``(question, choices) -> str``; None → the clarify tool errors. - reasoning_config: None → ``{"enabled": True, "effort": "medium"}`` on OpenRouter. - prefill_messages: priming history. Anthropic Sonnet/Opus 4.6+ 400 on a trailing - assistant message — use structured outputs there instead. - skip_context_files: skip SOUL.md/.hermes.md/AGENTS.md/CLAUDE.md/.cursorrules injection; - load_soul_identity keeps ~/.hermes/SOUL.md as identity regardless. - """ _install_safe_stdio() _params = locals() diff --git a/agent/agent_runtime_helpers.py b/agent/agent_runtime_helpers.py index 74986d6a9d..5fb2dee9d6 100644 --- a/agent/agent_runtime_helpers.py +++ b/agent/agent_runtime_helpers.py @@ -2370,7 +2370,8 @@ def invoke_tool(agent, function_name: str, function_args: dict, effective_task_i no display logic. Used by the concurrent path; the sequential path keeps its own inline invocation for display.""" from agent.inline_tool_executors import ( - InlineToolContext, emit_terminal_post_tool_call, resolve_invoke_tool_executor, tool_hook_ids + InlineToolContext, apply_transform_tool_result, emit_terminal_post_tool_call, + resolve_invoke_tool_executor, tool_hook_ids ) if not isinstance(function_args, dict): function_args = {} @@ -2407,14 +2408,17 @@ def invoke_tool(agent, function_name: str, function_args: dict, effective_task_i def _execute(next_args: dict) -> Any: result = inline_executor(agent, next_args, inline_ctx) + call_args = next_args if isinstance(next_args, dict) else function_args + duration_ms = int((time.monotonic() - tool_start_time) * 1000) emit_terminal_post_tool_call( - agent, function_name=function_name, - function_args=next_args if isinstance(next_args, dict) else function_args, + agent, function_name=function_name, function_args=call_args, result=result, effective_task_id=effective_task_id, tool_call_id=tool_call_id, - duration_ms=int((time.monotonic() - tool_start_time) * 1000), - middleware_trace=_tool_middleware_trace, + duration_ms=duration_ms, middleware_trace=_tool_middleware_trace, + ) + return apply_transform_tool_result( + agent, function_name=function_name, function_args=call_args, result=result, + effective_task_id=effective_task_id, tool_call_id=tool_call_id, duration_ms=duration_ms, ) - return result else: def _execute(next_args: dict) -> Any: dispatch_kwargs = dict( diff --git a/agent/auxiliary_client.py b/agent/auxiliary_client.py index 4864eeb8e4..dc573ab9c2 100644 --- a/agent/auxiliary_client.py +++ b/agent/auxiliary_client.py @@ -2606,10 +2606,16 @@ def _relay_sync_completion( return _run_protected_sync_provider_call(callback, kwargs) provider_name, fallback_model, metadata = route from agent import relay_llm - return relay_llm.execute_current( - kwargs, lambda request: _run_protected_sync_provider_call(callback, request), - name=provider_name, model_name=str(kwargs.get("model") or fallback_model), - metadata=metadata, defer_logical_completion=True, + from agent.auxiliary_hooks import run_with_aux_hooks + model_name = str(kwargs.get("model") or fallback_model) + return run_with_aux_hooks( + lambda: relay_llm.execute_current( + kwargs, lambda request: _run_protected_sync_provider_call(callback, request), + name=provider_name, model_name=model_name, metadata=metadata, + defer_logical_completion=True, + ), + aux_task=str(metadata.get("auxiliary_task") or ""), metadata=metadata, client=client, kwargs=kwargs, + provider=provider_name, model=model_name, api_mode=str(metadata.get("api_mode") or ""), ) @@ -2627,9 +2633,15 @@ async def _relay_async_completion( return await callback(kwargs) provider_name, fallback_model, metadata = route from agent import relay_llm - return await relay_llm.execute_current_async( - kwargs, callback, name=provider_name, model_name=str(kwargs.get("model") or fallback_model), - metadata=metadata, defer_logical_completion=True, + from agent.auxiliary_hooks import arun_with_aux_hooks + model_name = str(kwargs.get("model") or fallback_model) + return await arun_with_aux_hooks( + lambda: relay_llm.execute_current_async( + kwargs, callback, name=provider_name, model_name=model_name, + metadata=metadata, defer_logical_completion=True, + ), + aux_task=str(metadata.get("auxiliary_task") or ""), metadata=metadata, client=client, kwargs=kwargs, + provider=provider_name, model=model_name, api_mode=str(metadata.get("api_mode") or ""), ) @@ -2647,10 +2659,15 @@ def _relay_sync_stream( return create(kwargs) provider_name, fallback_model, metadata = route from agent import relay_llm - return relay_llm.stream_current( - kwargs, create, name=provider_name, - model_name=str(kwargs.get("model") or fallback_model), finalizer=dict, metadata=metadata, - completed_response_predicate=lambda value: hasattr(value, "choices"), + from agent.auxiliary_hooks import run_with_aux_hooks + model_name = str(kwargs.get("model") or fallback_model) + return run_with_aux_hooks( + lambda: relay_llm.stream_current( + kwargs, create, name=provider_name, model_name=model_name, finalizer=dict, + metadata=metadata, completed_response_predicate=lambda value: hasattr(value, "choices"), + ), + aux_task=str(metadata.get("auxiliary_task") or ""), metadata=metadata, client=client, kwargs=kwargs, + provider=provider_name, model=model_name, api_mode=str(metadata.get("api_mode") or ""), streaming=True, ) diff --git a/agent/auxiliary_hooks.py b/agent/auxiliary_hooks.py new file mode 100644 index 0000000000..613d794ca2 --- /dev/null +++ b/agent/auxiliary_hooks.py @@ -0,0 +1,252 @@ +"""Plugin events for auxiliary LLM calls (#79733). + +``pre_auxiliary_call`` / ``post_auxiliary_call`` fire once per physical provider attempt at the +relay boundary of ``agent.auxiliary_client`` — the funnel every auxiliary task (titling, +compression, MoA advisors/aggregator, vision, approval, ...) shares, retries and fallbacks +included — carrying the ``pre_api_request`` / ``post_api_request`` payload shape plus +``aux_task``. They are deliberately DISTINCT events: the main-loop ``*_api_request`` events stay +turn-scoped, so observability plugins keyed on turn identity never see auxiliary traffic unless +they subscribe to these. Observer-only (returns ignored) and fail-open: a raising or hung +callback is logged and the auxiliary call proceeds untouched. +""" + +from __future__ import annotations + +import logging +import time +from typing import Any, Awaitable, Callable, Dict, Optional + +logger = logging.getLogger(__name__) + +PRE_AUXILIARY_CALL = "pre_auxiliary_call" +POST_AUXILIARY_CALL = "post_auxiliary_call" + + +def _parent_turn_identity() -> Dict[str, str]: + """``session_id`` / ``task_id`` / ``turn_id`` / ``platform`` of the main turn this auxiliary + call runs under, or empty strings for turn-less callers (cron, gateway idle work).""" + ident = {"session_id": "", "task_id": "", "turn_id": "", "platform": ""} + try: + from agent.relay_runtime import current_turn + + turn = current_turn() + except Exception: + return ident + if turn is None: + return ident + lease = getattr(turn, "lease", None) + ident["session_id"] = str(getattr(lease, "session_id", "") or "") + ident["platform"] = str(getattr(lease, "platform", "") or "") + ident["task_id"] = str(getattr(turn, "task_id", "") or "") + ident["turn_id"] = str(getattr(turn, "turn_id", "") or "") + return ident + + +def _system_prompt(messages: Any, kwargs: Dict[str, Any]) -> str: + if isinstance(kwargs.get("system"), str): # Anthropic + return kwargs["system"] + if isinstance(kwargs.get("instructions"), str): # Responses + return kwargs["instructions"] + if isinstance(messages, list) and messages and isinstance(messages[0], dict): + first = messages[0] + if first.get("role") == "system" and isinstance(first.get("content"), str): + return first["content"] + return "" + + +def _usage_summary(response: Any, *, provider: str, api_mode: str) -> Optional[Dict[str, Any]]: + raw_usage = getattr(response, "usage", None) + if response is None or not raw_usage: + return None + from dataclasses import asdict + + from agent.usage_pricing import normalize_usage + + cu = normalize_usage(raw_usage, provider=provider, api_mode=api_mode) + summary = asdict(cu) + summary.pop("raw_usage", None) + summary["prompt_tokens"] = cu.prompt_tokens + summary["total_tokens"] = cu.total_tokens + return summary + + +def _first_choice_message(response: Any) -> Any: + choices = getattr(response, "choices", None) + if isinstance(response, dict): + choices = response.get("choices") + if not choices: + return None, None + choice = choices[0] + if isinstance(choice, dict): + return choice.get("message"), choice.get("finish_reason") + return getattr(choice, "message", None), getattr(choice, "finish_reason", None) + + +def _field(obj: Any, name: str) -> Any: + return obj.get(name) if isinstance(obj, dict) else getattr(obj, name, None) + + +class _AuxCallHooks: + """Fires the pre/post pair for one provider attempt; the base payload is built once.""" + + def __init__( + self, *, aux_task: str, metadata: Dict[str, Any], client: Any, kwargs: Dict[str, Any], + provider: str, model: str, api_mode: str, streaming: bool, + ) -> None: + self.provider = provider + self.api_mode = api_mode + self.streaming = streaming + self.kwargs = kwargs + self.started_at = time.time() + self.base: Dict[str, Any] = dict(_parent_turn_identity()) + self.base.update( + aux_task=aux_task, + api_request_id=str(metadata.get("api_request_id") or ""), + retry_count=int(metadata.get("retry_count") or 0), + api_call_count=int(metadata.get("retry_count") or 0) + 1, + model=model, + provider=provider, + base_url=str(getattr(client, "base_url", "") or ""), + api_mode=api_mode, + streaming=streaming, + started_at=self.started_at, + message_count=len(kwargs.get("messages") or kwargs.get("input") or []), + ) + + def pre(self) -> None: + if not _has_hook(PRE_AUXILIARY_CALL): + return + from agent.api_request_hooks import ApiRequestHooksMixin as _Sanitize + + kwargs = self.kwargs + messages = kwargs.get("messages") + if not isinstance(messages, list): + messages = kwargs.get("input") # Responses API + if not isinstance(messages, list): + messages = [] + body = {k: v for k, v in kwargs.items() if k not in {"timeout", "http_client"}} + total_chars = sum(len(str(_field(m, "content") or "")) for m in messages) + _fire( + PRE_AUXILIARY_CALL, **self.base, + request_messages=list(messages), + system_prompt=_system_prompt(messages, kwargs), + tool_count=len(kwargs.get("tools") or []), + approx_input_tokens=total_chars // 4, + request_char_count=total_chars, + max_tokens=kwargs.get("max_tokens") or kwargs.get("max_completion_tokens"), + request=_Sanitize._sanitize_hook_payload({"method": "POST", "body": body}), + ) + + def post(self, response: Any = None, error: Optional[BaseException] = None) -> None: + if not _has_hook(POST_AUXILIARY_CALL): + return + from agent.api_request_hooks import ApiRequestHooksMixin as _Sanitize + + ended_at = time.time() + payload: Dict[str, Any] = dict( + self.base, ended_at=ended_at, api_duration=max(0.0, ended_at - self.started_at), + error=None if error is None else f"{type(error).__name__}: {error}"[:2000], + error_type=None if error is None else type(error).__name__, + ) + # A streamed response is handed back unconsumed (the MoA facade owns reassembly), so + # there is no usage/finish_reason to report yet; ``streaming`` tells the observer why. + if error is not None or self.streaming: + payload.update(finish_reason=None, response_model=None, usage=None, response=None, + assistant_content_chars=0, assistant_tool_call_count=0) + else: + message, finish_reason = _first_choice_message(response) + content = _field(message, "content") if message is not None else None + tool_calls = (_field(message, "tool_calls") if message is not None else None) or [] + payload.update( + finish_reason=finish_reason, + response_model=_field(response, "model"), + usage=_usage_summary(response, provider=self.provider, api_mode=self.api_mode), + response=_Sanitize._sanitize_hook_payload({ + "model": _field(response, "model"), + "finish_reason": finish_reason, + "assistant_message": { + "role": (_field(message, "role") if message is not None else None) or "assistant", + "content": content, + "tool_calls": tool_calls, + }, + "usage": payload.get("usage"), + }), + assistant_content_chars=len(content) if isinstance(content, str) else 0, + assistant_tool_call_count=len(tool_calls), + ) + _fire(POST_AUXILIARY_CALL, **payload) + + +def _has_hook(name: str) -> bool: + try: + from hermes_cli.lifecycle import has_hook + + return has_hook(name) + except Exception: + return False + + +def _fire(name: str, **payload: Any) -> None: + """Dispatch one event; a failing subscriber is logged, never propagated (the aux task's + result must not depend on an observer).""" + try: + from hermes_cli.lifecycle import invoke_hook + + invoke_hook(name, **payload) + except Exception: + logger.warning("%s plugin hook failed for aux_task=%s; continuing", + name, payload.get("aux_task"), exc_info=True) + + +def _hooks_or_none(**kw: Any) -> Optional[_AuxCallHooks]: + if not (_has_hook(PRE_AUXILIARY_CALL) or _has_hook(POST_AUXILIARY_CALL)): + return None + try: + hooks = _AuxCallHooks(**kw) + hooks.pre() + except Exception: + logger.warning("pre_auxiliary_call payload build failed; continuing", exc_info=True) + return None + return hooks + + +def _post_safely(hooks: Optional[_AuxCallHooks], response: Any = None, error: Any = None) -> None: + if hooks is None: + return + try: + hooks.post(response, error) + except Exception: + logger.warning("post_auxiliary_call payload build failed; continuing", exc_info=True) + + +def run_with_aux_hooks( + call: Callable[[], Any], *, aux_task: str, metadata: Dict[str, Any], client: Any, + kwargs: Dict[str, Any], provider: str, model: str, api_mode: str, streaming: bool = False, +) -> Any: + """Run one synchronous provider attempt between ``pre_auxiliary_call`` and + ``post_auxiliary_call``; the exception (if any) is reported in ``post`` and re-raised.""" + hooks = _hooks_or_none(aux_task=aux_task, metadata=metadata, client=client, kwargs=kwargs, + provider=provider, model=model, api_mode=api_mode, streaming=streaming) + try: + response = call() + except BaseException as exc: + _post_safely(hooks, error=exc) + raise + _post_safely(hooks, response=response) + return response + + +async def arun_with_aux_hooks( + call: Callable[[], Awaitable[Any]], *, aux_task: str, metadata: Dict[str, Any], client: Any, + kwargs: Dict[str, Any], provider: str, model: str, api_mode: str, +) -> Any: + """Async twin of :func:`run_with_aux_hooks`.""" + hooks = _hooks_or_none(aux_task=aux_task, metadata=metadata, client=client, kwargs=kwargs, + provider=provider, model=model, api_mode=api_mode, streaming=False) + try: + response = await call() + except BaseException as exc: + _post_safely(hooks, error=exc) + raise + _post_safely(hooks, response=response) + return response diff --git a/agent/context_compressor.py b/agent/context_compressor.py index f337eb65d8..707d571e39 100644 --- a/agent/context_compressor.py +++ b/agent/context_compressor.py @@ -12,12 +12,14 @@ import re import time import uuid from dataclasses import dataclass -from typing import Any, Dict, List, Optional, Tuple +from typing import Any, Dict, List, Optional, Sequence, Tuple from agent.image_eviction_policy import outbound_image_retire_count from agent.auxiliary_client import ( AuxiliaryExplicitCancellation, + _coerce_llm_message, _is_connection_error, + _message_field, aux_interrupt_protection, call_llm, extract_content_or_reasoning, @@ -158,6 +160,56 @@ def _response_finish_reason(response: Any) -> str: # into every subsequent iterative-update prompt. (Ported from earendil-works/pi#7048 / commit 97fa14e39.) _TRUNCATED_SUMMARY_MARKER = "finish_reason=length" +# A provider can return a natural-language refusal with finish_reason="stop". It is +# non-empty, so the usual response validation accepts it, but it contains none of +# the checkpoint needed to safely replace the compacted turns. Keep this narrow: +# a real summary may mention a refusal in a recorded turn, while a refusal as the +# whole response begins with one of these phrases and refers to the requested +# summary/checkpoint. +_SUMMARY_REFUSAL_PREFIX_RE = re.compile( + r"^\s*(?:(?:sorry|i(?:['’]m| am)\s+sorry|i\s+apologi[sz]e|as\s+an\s+ai)" + r"\s*[,;:]?\s*(?:but\s+)?)?(?:i|we)\s+" + r"(?:can(?:\s*not|['’]t)|could\s*not|couldn['’]t|won['’]t|will\s+not|must\s+decline|" + r"refuse\s+to|am\s+unable\s+to|am\s+not\s+able\s+to)\b" + r"|^\s*(?:i['’]?m|i\s+am)\s+(?:unable|not\s+able)\b", + re.IGNORECASE, +) + + +def _is_summary_refusal(content: str) -> bool: + """Return whether a complete response is a refusal instead of a summary.""" + normalized = " ".join(content.split()) + if not _SUMMARY_REFUSAL_PREFIX_RE.match(normalized): + return False + # A refusal-only body never carries the template's "## " section headings; a real summary + # that merely opens with a hedging preamble ("I cannot see earlier turns, but here is...") does. + if re.search(r"(?m)^##\s", content): + return False + # Limit the search to the opener so a structured checkpoint that records a + # historical refusal elsewhere is not rejected. Stems catch summary/summarize/summarise. + return any(term in normalized[:400].casefold() for term in ("summar", "checkpoint")) + + +def _response_refusal_text(response: Any) -> str: + """Explicit provider ``choices[0].message.refusal`` (str, or dict with message/reason/text); ``""`` when absent. + + OpenAI-style structured-output refusals put the refusal here and leave ``content`` as filler or + empty, so the prose detector never sees it. + """ + refusal = _message_field(_coerce_llm_message(response), "refusal") + if isinstance(refusal, dict): + refusal = refusal.get("message") or refusal.get("reason") or refusal.get("text") + return refusal.strip() if isinstance(refusal, str) else "" + + +def _is_refusal_response(response: Any, content: str) -> bool: + """Single refusal predicate for both summarizer paths. + + An explicit provider ``message.refusal`` wins even when ``content`` looks like a + summary; otherwise fall back to the prose detector on the extracted content. + """ + return bool(_response_refusal_text(response)) or _is_summary_refusal(content) + def _is_summary_access_or_quota_error(exc: Exception) -> bool: """Return True for non-retryable summary auth, permission, or quota errors.""" @@ -639,7 +691,12 @@ class _SummaryFailureKind: def _classify_summary_failure(e: Exception) -> _SummaryFailureKind: - """Classify a summary-call exception by status code / message shape.""" + """Classify a summary-call exception by status code / message shape. + + A "refusal content" RuntimeError (prose or provider ``refusal`` field) deliberately rides the + ``empty_content`` class — cooldown + main-model fallback + abort — so the "returned empty content" + fallback log line is expected for refusals. + """ status = _exc_status_code(e) err = str(e).lower() return _SummaryFailureKind( @@ -655,7 +712,9 @@ def _classify_summary_failure(e: Exception) -> _SummaryFailureKind: # HTTP 200 with empty body from a degraded provider, plus the sibling "no usable response" # shapes from _validate_llm_response. empty_content=isinstance(e, RuntimeError) and any( - m in err for m in ("empty content", "llm returned none response", "llm returned invalid response") + m in err for m in ( + "empty content", "refusal content", "llm returned none response", "llm returned invalid response", + ) ), # Truncated summary: one main-model retry, then ABORT preserving the session. truncated=isinstance(e, RuntimeError) and _TRUNCATED_SUMMARY_MARKER in err, @@ -932,7 +991,7 @@ def _build_recovery_footer(session_id: str, region_len: int) -> str: # identifier-preserving session log is produced by the SAME single summary request as the narrative summary # (one auxiliary LLM call per compaction attempt, total — #96603: the earlier per-chunk digest loop made up # to 28 extra aux calls and pushed compactions to 7-11 minutes on slow aux routes). Coverage over oversized -# regions comes from even input sampling (see ``_sample_summary_input``), and exact-needle defense comes +# regions comes from even record sampling (see ``_sample_summary_records``), and exact-needle defense comes # from the LLM-free anchor index below. _LEAN_SESSION_LOG_HEADING = "## Detailed Session Log (oldest first)" # Extra output-token guidance for the session-log section (single response). @@ -1957,6 +2016,10 @@ class ContextCompressor(SummaryDispatchMixin, MicroCompactionMixin, ContextEngin "total_duration_ms": None, "aux_call_duration_ms": None, "queue_wait_ms": None, "prompt_build_ms": None, "time_to_first_progress_ms": None, "summary_generation_ms": None, "commit_ms": None, "fallback_used": False, "commit_status": "unknown", "split_status": "unknown", "failure_class": None, + # Lean-sampling coverage (filled by _record_summary_input_coverage; None on the legacy path). + "summary_input_chars": None, "summary_input_sampled_chars": None, "summary_input_omitted_chars": None, + "summary_input_record_count": None, "summary_input_sampled_record_count": None, + "summary_input_elided_record_count": None, } self._active_compression_telemetry = self._last_compression_telemetry = telemetry return telemetry @@ -3243,8 +3306,8 @@ class ContextCompressor(SummaryDispatchMixin, MicroCompactionMixin, ContextEngin args = args[:self._TOOL_ARGS_HEAD] + "..." return f" {fn.get('name', '?')}({args})" - def _serialize_for_summary(self, turns: List[Dict[str, Any]]) -> str: - """Serialize turns into labeled, redacted text for the summarizer.""" + def _serialize_records_for_summary(self, turns: List[Dict[str, Any]]) -> List[str]: + """Serialize turns into a list of labeled, redacted records for the summarizer.""" # Lazy import: agent_runtime_helpers pulls heavy transitive imports. from agent.agent_runtime_helpers import strip_think_blocks parts = [] @@ -3266,7 +3329,11 @@ class ContextCompressor(SummaryDispatchMixin, MicroCompactionMixin, ContextEngin if role == "assistant" and msg.get("tool_calls", []): content += "\n[Tool calls:\n" + "\n".join(map(self._render_tool_call_for_summary, msg["tool_calls"])) + "\n]" parts.append(f"[{role.upper()}]: {content}") - return "\n\n".join(parts) + return parts + + def _serialize_for_summary(self, turns: List[Dict[str, Any]]) -> str: + """Serialize turns into labeled, redacted text for the summarizer.""" + return "\n\n".join(self._serialize_records_for_summary(turns)) def _fallback_anchors(self, turns_to_summarize: List[Dict[str, Any]]) -> Dict[str, list[str]]: """Locally extractable anchors: user asks, actions, files, blockers, last dropped turns.""" @@ -3465,31 +3532,157 @@ Summary generation was unavailable, so this is a best-effort deterministic fallb _SAMPLED_INPUT_SLICES = 8 @classmethod - def _sample_summary_input(cls, content: str) -> str: - """Cap summarizer input by EVEN SAMPLING across the whole region (lean mode). - The single request also produces the session log, so coverage must be uniform: head+tail - truncation would hide the entire middle from it.""" - if len(content) <= cls._SUMMARY_INPUT_MAX_CHARS: - return content - n = max(2, cls._SAMPLED_INPUT_SLICES) - marker_template = "\n\n...[{elided:,} chars elided — recover via session_search]...\n\n" - marker_reserve = len(marker_template.format(elided=len(content))) * (n - 1) - budget = max(cls._SUMMARY_INPUT_MAX_CHARS - marker_reserve, n) - slice_len = budget // n - stride = len(content) / n - parts: list[str] = [] - prev_end = 0 - for i in range(n): - start = int(i * stride) - if i == n - 1: - # Last slice anchors to the END: newest turns carry the most state. - start = max(start, len(content) - slice_len) - end = min(start + slice_len, len(content)) - if start > prev_end: - parts.append(marker_template.format(elided=start - prev_end)) - parts.append(content[start:end]) - prev_end = end - return "".join(parts) + def _bound_oversized_record(cls, record: str, limit: int) -> str: + """Bound an oversized record with an explicit intra-record truncation marker.""" + if len(record) <= limit: + return record + marker_template = "\n...[record truncated: {elided:,} chars elided — recover via session_search]...\n" + marker_reserve = len(marker_template.format(elided=len(record))) + if limit <= marker_reserve: + return record[:limit] + remaining = limit - marker_reserve + head_len = remaining // 2 + tail_len = remaining - head_len + head = record[:head_len].rstrip("\n") + tail = record[-tail_len:].lstrip("\n") + elided = len(record) - len(head) - len(tail) + return head + marker_template.format(elided=elided) + tail + + def _record_summary_input_coverage(self, coverage: Dict[str, int]) -> None: + """Expose lean sampling coverage without including transcript content in telemetry.""" + telemetry = getattr(self, "_active_compression_telemetry", None) + if not isinstance(telemetry, dict): + return + telemetry.update({ + "summary_input_chars": coverage["input_chars"], + "summary_input_sampled_chars": coverage["sampled_chars"], + "summary_input_omitted_chars": coverage["omitted_chars"], + "summary_input_record_count": coverage["record_count"], + "summary_input_sampled_record_count": coverage["sampled_record_count"], + "summary_input_elided_record_count": coverage["elided_record_count"], + }) + + @classmethod + def _sample_summary_records(cls, records: Sequence[str]) -> Tuple[str, Dict[str, int]]: + """Sample complete serialized records while retaining the character bound. + + Returns the bounded transcript and record-level coverage counters for compression + telemetry. `input_chars` counts raw serialized record content; `sampled_chars` counts the + *display* chars of retained records (after intra-record truncation by + `_bound_oversized_record`); neither includes separators or elision markers, so + `omitted_chars = input_chars - sampled_chars` also covers truncated-away bytes. + """ + input_chars = sum(len(r) for r in records) + + def _coverage(sampled_chars: int, sampled_record_count: int) -> Dict[str, int]: + return { + "input_chars": input_chars, "sampled_chars": sampled_chars, + "omitted_chars": input_chars - sampled_chars, "record_count": len(records), + "sampled_record_count": sampled_record_count, + "elided_record_count": len(records) - sampled_record_count, + } + + if not records: + return "", _coverage(0, 0) + + separator = "\n\n" + total_len = input_chars + len(separator) * (len(records) - 1) + if total_len <= cls._SUMMARY_INPUT_MAX_CHARS: + return separator.join(records), _coverage(input_chars, len(records)) + + n = max(1, min(cls._SAMPLED_INPUT_SLICES, len(records))) + marker_template = ( + "\n\n...[records {first:,}-{last:,}: {elided:,} chars elided — recover via session_search]...\n\n" + ) + marker_len = len(marker_template.format(first=len(records), last=len(records), elided=total_len)) + budget = max(cls._SUMMARY_INPUT_MAX_CHARS - marker_len * (n - 1), 1) + target = max(1, budget // n) + + # Oversized records are bounded to slice target with explicit intra-record truncation markers + # so they cannot consume other regions' budget or evict the newest record. + display_records = [cls._bound_oversized_record(r, target) for r in records] + + def _merged(slices: list[tuple[int, int]]) -> list[tuple[int, int]]: + out: list[tuple[int, int]] = [] + for s, e in slices: + if out and s <= out[-1][1]: + out[-1] = (out[-1][0], max(out[-1][1], e)) + else: + out.append((s, e)) + return out + + starts = [round(i * len(records) / n) for i in range(n)] + selected: list[tuple[int, int]] = [] + for index, start in enumerate(starts): + if index == len(starts) - 1: + # Anchor the last slice to the newest record at the end of the history. + end = len(records) + start = end - 1 + size = len(display_records[start]) + while start > 0 and size + len(separator) + len(display_records[start - 1]) <= target: + start -= 1 + size += len(separator) + len(display_records[start]) + else: + end = start + size = 0 + while end < len(records) and (size == 0 or size + len(display_records[end]) + len(separator) <= target): + size += len(display_records[end]) + (len(separator) if end > start else 0) + end += 1 + if end > start: + selected.append((start, end)) + selected = _merged(selected) + + def _render(slices: list[tuple[int, int]]) -> str: + parts: list[str] = [] + cursor = 0 + for s, e in slices: + if s > cursor: + sep_count = (s - cursor) if cursor == 0 else (s - cursor + 1) + elided = sum(len(records[i]) for i in range(cursor, s)) + len(separator) * sep_count + parts.append(marker_template.format(first=cursor + 1, last=s, elided=elided)) + parts.append(separator.join(display_records[s:e])) + cursor = e + return "".join(parts) + + # Budget extension: the greedy fill leaves each slice short of `target` by up to one record + # (5-43% of the cap unused for 8-20K records). Spend the headroom on whole neighbouring + # records, round-robin one record per slice per round so every region keeps an even share + # (the newest slice grows backward, older slices grow forward) — never past cap. + cap = cls._SUMMARY_INPUT_MAX_CHARS + rendered_len = len(_render(selected)) + grew = True + while grew: + grew = False + for idx in range(len(selected) - 1, -1, -1): + s, e = selected[idx] + if idx == len(selected) - 1: + nxt, grown = s - 1, (s - 1, e) + if nxt < (selected[idx - 1][1] if idx else 0): + continue + else: + nxt, grown = e, (s, e + 1) + if nxt >= selected[idx + 1][0]: + continue + if rendered_len + len(separator) + len(display_records[nxt]) > cap: + continue + selected[idx] = grown + new_len = len(_render(_merged(selected))) + # The pre-check above bounds the added record; the exact re-render catches the + # one thing it cannot see — a gap's first index gaining a digit or comma in the + # marker (e.g. 999 -> 1,000) when the render is already at cap. + if new_len > cap: + selected[idx] = (s, e) + continue + rendered_len = new_len + grew = True + selected = _merged(selected) + + # No overflow trim is needed: every slice holds <= `target` display chars (records are + # pre-bounded to `target`), there are <= n-1 markers each <= `marker_len` (widths computed + # at their maxima), and n*target + (n-1)*marker_len <= _SUMMARY_INPUT_MAX_CHARS by + # construction; the extension pass above only adds a record when the result stays <= cap. + shown = [i for s, e in selected for i in range(s, e)] + return _render(selected), _coverage(sum(len(display_records[i]) for i in shown), len(shown)) def _fallback_to_main_for_compression( self, e: Exception, reason: str, failed_model: Optional[str] = None @@ -3583,6 +3776,11 @@ Summary generation was unavailable, so this is a best-effort deterministic fallb # error, rather than replacing real context with an empty summary. if not content.strip(): raise RuntimeError(f"Context compression LLM returned empty content {where}") + if _is_refusal_response(response, content): + # Treat a refusal as unusable content. This deliberately reuses the + # established fallback/cooldown/abort path for an empty body, so it + # can never be committed as `_previous_summary`. + raise RuntimeError(f"Context compression LLM returned refusal content {where}") # A finish_reason of "length" means the summarizer hit its output token cap mid-generation: the text # present is PARTIAL. Persisting a partial summary as the compaction checkpoint silently truncates # the conversation's memory — the cut-off text replaces the real middle turns AND is fed back into @@ -3628,8 +3826,12 @@ Summary generation was unavailable, so this is a best-effort deterministic fallb _collect_ghosted_skill_names(turns_to_summarize) + _extract_pruned_skill_names(self._previous_summary or "") ))[:_MAX_PRUNED_SKILL_MARKERS] # Lean mode even-samples oversized input (one bounded request, never a second). - bound = self._sample_summary_input if getattr(self, "tail_mode", "lean") == "lean" else self._bound_summary_input - content_to_summarize = bound(self._serialize_for_summary(turns_to_summarize)) + if getattr(self, "tail_mode", "lean") == "lean": + records = self._serialize_records_for_summary(turns_to_summarize) + content_to_summarize, coverage = self._sample_summary_records(records) + self._record_summary_input_coverage(coverage) + else: + content_to_summarize = self._bound_summary_input(self._serialize_for_summary(turns_to_summarize)) has_user_turn = getattr(self, "_summary_has_user_turn", None) if has_user_turn is None: has_user_turn = self._transcript_has_real_user_turn(turns_to_summarize) diff --git a/agent/context_engine.py b/agent/context_engine.py index a052021b27..3fb272a50e 100644 --- a/agent/context_engine.py +++ b/agent/context_engine.py @@ -8,6 +8,7 @@ should_compress() / compress() -> on_session_end() at real session boundaries on (CLI exit, /reset, gateway expiry), never per-turn. """ +import copy import json from abc import ABC, abstractmethod from typing import Any, Dict, List, Optional @@ -222,6 +223,13 @@ class ContextEngine(ABC): "compression_count": self.compression_count, } + def clone_for_agent(self) -> "ContextEngine": + """Per-agent instance of a plugin-registered engine (the plugin system holds ONE shared + instance; every AIAgent gets its own so a child's update_model() cannot mutate the parent's). + Override when the engine holds uncopyable state (locks, DB connections): return a fresh + engine sharing the durable backend and copying only mutable budget state.""" + return copy.deepcopy(self) + def update_model( self, model: str, context_length: int, base_url: str = "", api_key: str = "", provider: str = "", api_mode: str = "", diff --git a/agent/curator.py b/agent/curator.py index 3585b44a0a..f27b64e3ab 100644 --- a/agent/curator.py +++ b/agent/curator.py @@ -287,7 +287,10 @@ CURATOR_REVIEW_PROMPT = ( "(imperative + one clause of why), the same lesson stated twice becomes " "one rule, and incident narration, PR/issue numbers, dates and quoted " "chatter are dropped — the rule must stand without the story. Moving a " - "file unchanged under references/ is filing, not consolidating.\n\n" + "file unchanged under references/ is filing, not consolidating. A SKILL.md " + "body over ~24k chars is a consolidation target on its own: skill_view loads " + "all of it into context for the rest of the session, so distill it to the " + "always-on rules and push topic depth into references/.\n\n" "Hard rules — do not violate:\n" "1. DO NOT touch bundled, hub-installed, or external-dir skills " "(`skills.external_dirs`). The candidate list below is already filtered " diff --git a/agent/display.py b/agent/display.py index 267c56aae6..ceb8ba4313 100644 --- a/agent/display.py +++ b/agent/display.py @@ -9,7 +9,7 @@ import re import sys import threading import time -from dataclasses import dataclass, field +from dataclasses import dataclass, field, replace from difflib import unified_diff from pathlib import Path from typing import Any, Iterator @@ -436,6 +436,16 @@ def _preview_skill_view(args: dict, max_len: int) -> str | None: return _tail_trunc(label, max_len) or None +def _preview_bridge_call(tool_name: str): + def _build(args: dict, max_len: int) -> str | None: + labels = bridge_tool_labels(tool_name, args) + if not labels: + return _primary_arg_preview(tool_name, args, max_len) + extra = f" +{len(labels) - 1}" if len(labels) > 1 else "" + return _tail_trunc(f"{labels[0].text}{extra}", max_len) or None + return _build + + # Tool-specific preview builders: f(args, max_len) -> preview. Tools not listed # fall through to the primary-argument lookup in build_tool_preview. _PREVIEW_BUILDERS = { @@ -445,6 +455,9 @@ _PREVIEW_BUILDERS = { "read_file": _preview_read_file, "memory": _preview_memory, "send_message": _preview_send_message, "skill_view": _preview_skill_view, "session_search": lambda args, _m: f"recall: \"{_clip(_oneline(args.get('query', '')), 25)}\"", + "tool_call": _preview_bridge_call("tool_call"), + "tool_search": _preview_bridge_call("tool_search"), + "tool_describe": _preview_bridge_call("tool_describe"), } @@ -461,6 +474,10 @@ def build_tool_preview(tool_name: str, args: dict, max_len: int | None = None) - builder = _PREVIEW_BUILDERS.get(tool_name) if builder is not None: return builder(args, max_len) + return _primary_arg_preview(tool_name, args, max_len) + + +def _primary_arg_preview(tool_name: str, args: dict, max_len: int) -> str | None: key = _PRIMARY_ARGS.get(tool_name) or next((k for k in _FALLBACK_PREVIEW_KEYS if k in args), None) if not key or key not in args: return None @@ -500,6 +517,37 @@ _TOOL_VERBS_NO_PREVIEW: frozenset[str] = frozenset({"skills_list", "session_sear # Verbs joined to the preview with " for " (search-style phrasing). _TOOL_VERBS_FOR_CONNECTOR: frozenset[str] = frozenset({"web_search", "search_files"}) +_BRIDGE_GENERATING = { + "tool_call": "a tool call", "tool_search": "a tool search", "tool_describe": "tool details", +} + + +def bridge_generating_phrase(tool_name: str) -> str | None: + return _BRIDGE_GENERATING.get(tool_name) if _friendly_tool_labels else None + + +def tool_labels_for_call(tool_name: str, args: dict | None) -> list: + try: + from tools.tool_labels import labels_for_call + labels = labels_for_call(tool_name, args or {}) + except Exception as exc: # noqa: BLE001 — display must never abort a turn + logger.debug("bridge labels failed for %s: %s", tool_name, exc) + return [] + skin = _get_skin() + overrides = (skin.tool_emojis if skin and skin.tool_emojis else None) or {} + return [replace(label, emoji=overrides[label.name]) if label.name in overrides else label + for label in labels] + + +def bridge_tool_labels(tool_name: str, args: dict | None) -> list: + return tool_labels_for_call(tool_name, args) if _friendly_tool_labels else [] + + +def tool_row_emoji(tool_name: str, args: dict | None = None, default: str = "⚡") -> str: + labels = bridge_tool_labels(tool_name, args) if args else [] + return labels[0].emoji if labels else get_tool_emoji(tool_name, default) + + def get_tool_verb(tool_name: str) -> str | None: """Friendly verb for a built-in tool, or None (labels disabled / no curated verb); callers compose ``f"{verb}{tool_verb_connector(tool)}{preview}"`` themselves.""" @@ -535,8 +583,12 @@ def build_status_phrase(tool_name: str, args: dict | None, max_len: int = 49) -> def build_tool_label(tool_name: str, args: dict, max_len: int | None = None) -> str | None: - """Human-phrased label ("Searching the web for ...") for curated built-ins; other - tools (or labels disabled) get the raw preview, so it is a drop-in for build_tool_preview.""" + labels = bridge_tool_labels(tool_name, args) + if labels: + label = labels[0] + preview = _tail_trunc(label.preview, max_len if max_len is not None else _tool_preview_max_len) + extra = f" +{len(labels) - 1}" if len(labels) > 1 else "" + return f"{label.text}{f' {preview}' if preview else ''}{extra}" verb = get_tool_verb(tool_name) if verb and tool_name in _TOOL_VERBS_NO_PREVIEW: return verb @@ -1085,15 +1137,59 @@ _CUTE_LINES = { } +def _cute_bridge_rows(tool_name: str, args: dict) -> list[str] | None: + labels = bridge_tool_labels(tool_name, args) + if not labels: + return None + return [f"┊ {label.emoji} {_cute_trunc(label.text)}" + f"{f' {_cute_trunc(label.preview)}' if label.preview else ''}" for label in labels] + + +_BRIDGE_CALL_INDEX_RE = re.compile(r"calls\[(\d+)\]") + + +def _entry_failure_suffix(entry: Any) -> str: + if not isinstance(entry, dict): + return "" + error = entry.get("error") + if isinstance(error, dict): + return f" [{_trim_error(str(error.get('message') or error.get('code') or 'error'))}]" + if not (error or entry.get("success") is False): + return "" + return _detect_tool_failure(str(entry.get("name") or ""), entry)[1] + + +def _bridge_row_suffixes(rows: int, result: Any, call_suffix: str) -> list[str]: + suffixes = [""] * rows + data = result if isinstance(result, dict) else safe_json_loads(result) + entries = data.get("results") if isinstance(data, dict) else None + if isinstance(entries, list): + for position, entry in enumerate(entries[:rows]): + suffix = _entry_failure_suffix(entry) + if not suffix: + continue + index = entry.get("index") + suffixes[index if isinstance(index, int) and 0 <= index < rows else position] = suffix + if call_suffix and not any(suffixes): + match = _BRIDGE_CALL_INDEX_RE.search(call_suffix) + named = int(match.group(1)) if match else rows - 1 + suffixes[named if 0 <= named < rows else rows - 1] = call_suffix + return suffixes + + def _get_cute_tool_message(tool_name: str, args: dict, duration: float, result: str | None = None) -> str: - """Tool completion line for CLI quiet mode: ``| {emoji} {verb:9} {detail} {duration}``, plus a - failure suffix from :func:`_detect_tool_failure`; the leading ``┊`` becomes the skin's tool prefix.""" args = redact_tool_args_for_display(tool_name, args) or args is_failure, failure_suffix = _detect_tool_failure(tool_name, result) render = _CUTE_LINES.get(tool_name) - body = render(args, result) if render else f"┊ ⚡ {tool_name[:9]:9} {_cute_trunc(build_tool_preview(tool_name, args) or '')}" - line = f"{body} {duration:.1f}s".replace("┊", get_skin_tool_prefix(), 1) - return f"{line}{failure_suffix}" if is_failure else line + rows = _cute_bridge_rows(tool_name, args) + if rows is None: + body = render(args, result) if render else f"┊ ⚡ {tool_name[:9]:9} {_cute_trunc(build_tool_preview(tool_name, args) or '')}" + rows, suffixes = [body], [failure_suffix if is_failure else ""] + else: + suffixes = _bridge_row_suffixes(len(rows), result, failure_suffix if is_failure else "") + rows = [*rows[:-1], f"{rows[-1]} {duration:.1f}s"] + prefix = get_skin_tool_prefix() + return "\n ".join(f"{row}{suffix}".replace("┊", prefix, 1) for row, suffix in zip(rows, suffixes)) def get_cute_tool_message(tool_name: str, args: dict, duration: float, result: str | None = None) -> str: diff --git a/agent/inline_tool_executors.py b/agent/inline_tool_executors.py index 10d22c7d9a..76a675e8f5 100644 --- a/agent/inline_tool_executors.py +++ b/agent/inline_tool_executors.py @@ -58,6 +58,31 @@ def emit_terminal_post_tool_call( pass +def apply_transform_tool_result( + agent, + *, + function_name: str, + function_args: dict, + result: Any, + effective_task_id: str, + tool_call_id: Optional[str], + duration_ms: int = 0, +) -> Any: + """Apply ``transform_tool_result`` to an inline-dispatched tool's result. + + Registry tools get this inside ``handle_function_call``; inline executors never + reach it, so the agent paths call the same helper (after the terminal + ``post_tool_call``) to keep the hook's "every tool" contract. Fail-open.""" + try: + from model_tools import _CallIds, _apply_transform_tool_result_hook + return _apply_transform_tool_result_hook( + function_name, function_args, result, duration_ms, + _CallIds(**tool_hook_ids(agent, effective_task_id, tool_call_id)), + ) + except Exception: + return result + + @dataclass class InlineToolContext: """Per-call state an inline executor may need beyond its arguments.""" diff --git a/agent/memory_provider.py b/agent/memory_provider.py index 938fd0b003..f9f4bcd4f6 100644 --- a/agent/memory_provider.py +++ b/agent/memory_provider.py @@ -39,6 +39,15 @@ PRE_COMPRESS_CHECKPOINT_API_VERSION = 2 # Default glyph for recall indicators; providers may use their own brand mark. INDICATOR_GLYPH = "🧠" +# ``memory.provider`` values that mean "the built-in store, no external plugin". The built-in +# store is core: doctor, migration and dependency refresh must never look these up as plugins. +CORE_MEMORY_PROVIDER_SENTINELS = frozenset({"", "default", "builtin", "built-in", "none"}) + + +def is_core_memory_provider(name: Optional[str]) -> bool: + """True when ``memory.provider`` selects the built-in store rather than an external plugin.""" + return str(name or "").strip().lower() in CORE_MEMORY_PROVIDER_SENTINELS + @dataclass(frozen=True) class RecallStatus: diff --git a/agent/micro_compaction.py b/agent/micro_compaction.py index 34ff76340a..47456f9ee7 100644 --- a/agent/micro_compaction.py +++ b/agent/micro_compaction.py @@ -148,12 +148,16 @@ class MicroCompactionMixin: message = response.choices[0].message content = message.get("content") if isinstance(message, dict) else getattr(message, "content", message) content = (content if isinstance(content, str) else str(content) if content else "").strip() + + from agent.agent_runtime_helpers import strip_think_blocks + content = strip_think_blocks(None, content).strip() if not content: logger.info("micro-summarization returned empty content") return None - - from agent.agent_runtime_helpers import strip_think_blocks - return strip_think_blocks(None, content).strip() or None + if _cc()._is_refusal_response(response, content): + logger.warning("micro-summarization returned refusal content — discarding unusable summary") + return None + return content def _needs_defrag(self) -> bool: """Return True when the rolling summary is large enough to defrag.""" diff --git a/agent/shell_hooks.py b/agent/shell_hooks.py index 08776486a0..402b740b78 100644 --- a/agent/shell_hooks.py +++ b/agent/shell_hooks.py @@ -445,6 +445,15 @@ def _parse_pre_tool_call(data: Dict[str, Any]) -> Optional[Dict[str, Any]]: for verb, _, _, payload in _PRE_TOOL_DIALECTS: if data.get(verb) == "modify" and isinstance(data.get(payload), dict): return {"action": "modify", "args": data[payload]} + # Hermes-only escalation to the human-approval gate (#92553). Claude-Code's ``decision: + # approve`` means auto-ALLOW, so it is deliberately not mapped onto this. + if data.get("action") == "approve": + directive: Dict[str, Any] = {"action": "approve"} + for key in ("message", "rule_key"): + value = data.get(key) + if isinstance(value, str) and value.strip(): + directive[key] = value.strip() + return directive return None diff --git a/agent/tool_executor.py b/agent/tool_executor.py index f6f84dcf11..b877a637a4 100644 --- a/agent/tool_executor.py +++ b/agent/tool_executor.py @@ -26,6 +26,7 @@ from agent.display import ( build_tool_label as _build_tool_label, get_cute_tool_message as _get_cute_tool_message_impl, get_tool_emoji as _get_tool_emoji, + tool_row_emoji as _tool_row_emoji, redact_tool_args_for_display as _redact_tool_args_for_display, _detect_tool_failure, ) @@ -33,6 +34,7 @@ from agent.message_sanitization import coalesce_tool_call_id from agent.inline_tool_executors import ( INLINE_TOOL_EXECUTORS, InlineToolContext, + apply_transform_tool_result, emit_terminal_post_tool_call, tool_hook_ids, ) @@ -1555,7 +1557,7 @@ def _start_quiet_tool_spinner(agent, function_name: str, function_args: dict, *, face = random.choice(KawaiiSpinner.get_waiting_faces()) if label is None: display_args = _redact_tool_args_for_display(function_name, function_args) or function_args - label = f"{_get_tool_emoji(function_name)} {_build_tool_label(function_name, display_args) or function_name}" + label = f"{_tool_row_emoji(function_name, display_args)} {_build_tool_label(function_name, display_args) or function_name}" spinner = KawaiiSpinner(f"{face} {label}", spinner_type='dots', print_fn=agent._print_fn) spinner.start() return spinner @@ -1592,6 +1594,7 @@ class _SequentialDispatch: is_delegate: bool = False finish_spinner: bool = True finish_in_finally: bool = True # inline tools print their completion line only on success + transform_applied: bool = False # True when execute already fired transform_tool_result def _resolve_sequential_dispatch(agent, ref: _ToolCallRef, messages: list) -> _SequentialDispatch: @@ -1656,6 +1659,7 @@ def _resolve_sequential_dispatch(agent, ref: _ToolCallRef, messages: list) -> _S error_log="handle_function_call raised for %s: %s", handles_keyboard_interrupt=True, finish_spinner=bool(agent.quiet_mode), + transform_applied=True, # handle_function_call fires transform_tool_result itself ) @@ -1726,19 +1730,28 @@ def _run_sequential_call( return managed, tool_duration -def _publish_sequential_result(agent, messages: list, ref: _ToolCallRef, managed: _ManagedToolResult, *, tool_duration: float, index: int, budget: BudgetConfig) -> bool: +def _publish_sequential_result(agent, messages: list, ref: _ToolCallRef, managed: _ManagedToolResult, *, tool_duration: float, index: int, budget: BudgetConfig, transform_applied: bool) -> bool: """Terminal hook → observe → commit → completion callbacks/print for one sequential result; False when the incremental flush failed (the caller must stop the batch).""" ref.args, ref.trace, function_result = managed.args, managed.middleware_trace, managed.result _execution_timed_out = isinstance(function_result, (_ToolTimeoutResult, _ToolCancelledResult)) - # Multimodal dict results (_multimodal=True) are not sliceable as strings. - _result_len = len(function_result) if isinstance(function_result, str) else len(str(function_result)) - _is_error_result, _ = _detect_tool_failure(ref.name, function_result) # Inline-dispatched runtime tools never reach handle_function_call, so the # executor owns the one terminal post_tool_call per tool_call_id (the inner # observer is suppressed); also stops an abandoned timeout worker reporting late. + # transform_tool_result follows the observer, unless the dispatch already fired it. if not managed.blocked and not _execution_timed_out: ref.emit_post(agent, function_result, duration_ms=int(tool_duration * 1000)) + if not transform_applied: + function_result = apply_transform_tool_result( + agent, function_name=ref.name, function_args=ref.args, result=function_result, + effective_task_id=ref.task_id, tool_call_id=ref.call_id, + duration_ms=int(tool_duration * 1000), + ) + # Classify the result the model will actually see, i.e. after any transform; the + # registry and concurrent paths both classify post-transform. + # Multimodal dict results (_multimodal=True) are not sliceable as strings. + _result_len = len(function_result) if isinstance(function_result, str) else len(str(function_result)) + _is_error_result, _ = _detect_tool_failure(ref.name, function_result) committed = _commit_tool_result( agent, messages, ref, function_result, budget=budget, tool_duration=tool_duration, is_error=_is_error_result, blocked=managed.blocked, @@ -1809,7 +1822,8 @@ def _execute_tool_calls_sequential(agent, assistant_message, messages: list, eff display_index=i, tool_start_time=tool_start_time, ) - if not _publish_sequential_result(agent, messages, ref, managed, tool_duration=tool_duration, index=i, budget=_tool_budget): + if not _publish_sequential_result(agent, messages, ref, managed, tool_duration=tool_duration, index=i, + budget=_tool_budget, transform_applied=dispatch.transform_applied): return if agent._interrupt_requested and i < len(tool_calls): diff --git a/agent/turn_context.py b/agent/turn_context.py index c655e1d11c..2fc545ffd3 100644 --- a/agent/turn_context.py +++ b/agent/turn_context.py @@ -21,6 +21,7 @@ from agent.conversation_compression import recover_rotated_compression_session from agent.iteration_budget import IterationBudget from agent.memory_manager import build_memory_context_block from agent.memory_provider import is_trivial_prompt +from agent.message_content import flatten_message_text from agent.message_metadata import append_message, stamp_message_timestamp from agent.model_metadata import estimate_messages_tokens_rough, estimate_request_tokens_rough from agent.image_token_cost import bind_image_token_cost @@ -78,20 +79,29 @@ def _agent_stale_thinking_on_wire(agent: Any) -> bool: return True +def compose_multimodal_context_part( + ext_prefetch_cache: str, plugin_user_context: str, +) -> Optional[str]: + """The ephemeral context of one turn (memory prefetch + ``pre_llm_call``) as one text + block; ``None`` when nothing is injected. The string sidecar appends it to ``content``; + a multimodal (list) turn carries it as a durable text part (#71998).""" + fenced = build_memory_context_block(ext_prefetch_cache) if ext_prefetch_cache else "" + injections = [part for part in (fenced, plugin_user_context) if part] + return "\n\n".join(injections) if injections else None + + def compose_user_api_content( content: Any, ext_prefetch_cache: str, plugin_user_context: str ) -> Optional[str]: - """Compose the API-bound content of the current turn's user message. + """Compose the API-bound content of the current turn's string user message. Single source for the ``api_content`` sidecar and the wire bytes so they never drift - (what turn N sends is what turn N+1 replays). ``None`` when nothing is injected.""" + (what turn N sends is what turn N+1 replays). ``None`` when nothing is injected or the + content is not a string (list content takes the text-part path).""" if not isinstance(content, str): return None - fenced = build_memory_context_block(ext_prefetch_cache) if ext_prefetch_cache else "" - injections = [part for part in (fenced, plugin_user_context) if part] - if not injections: - return None - return content + "\n\n" + "\n\n".join(injections) + injection = compose_multimodal_context_part(ext_prefetch_cache, plugin_user_context) + return None if injection is None else content + "\n\n" + injection def substitute_api_content(api_msg: Dict[str, Any]) -> Optional[str]: @@ -136,7 +146,7 @@ def consume_surface_switch_note(agent: Any) -> str: return _pop_turn_note(agent, "_surface_switch_note") -def append_notes_to_multimodal_content(content: Any, notes: str) -> bool: +def append_notes_to_multimodal_content(content: Any, notes: Optional[str]) -> bool: """Append must-deliver notes as a durable text part on a multimodal (list) user message (the sidecar path returns ``None`` for non-string content).""" if not notes or not isinstance(content, list): @@ -831,6 +841,15 @@ def _bind_interrupt_scope(agent: Any, ra) -> None: agent._interrupt_thread_signal_pending = False +def _memory_query_text(original_user_message: Any) -> str: + """Semantic text of the turn for memory queries: a multimodal (list) turn carries its text + in parts, so keying off ``isinstance(str)`` collapsed it to ``""`` and ``is_trivial_prompt`` + skipped prefetch entirely. An image-only turn still flattens to ``""`` (trivial).""" + if isinstance(original_user_message, (str, list)): + return flatten_message_text(original_user_message) + return "" + + def _memory_turn_start_and_prefetch( agent: Any, original_user_message: Any, turn_author: Optional[Dict[str, Any]] = None, ) -> str: @@ -839,7 +858,7 @@ def _memory_turn_start_and_prefetch( Returns the prefetch text (``""`` when nothing was injected).""" if not agent._memory_manager: return "" - _query = original_user_message if isinstance(original_user_message, str) else "" + _query = _memory_query_text(original_user_message) # The author rides along so a provider can attribute THIS turn, not whoever opened the session. _author = turn_author if isinstance(turn_author, dict) else {} with suppress(Exception): @@ -909,6 +928,38 @@ def _stamp_api_content_sidecar( logger.warning("api_content backfill failed for session=%s", agent.session_id or "none", exc_info=True) +def _append_multimodal_context( + agent: Any, turn_user_msg: Dict[str, Any], ext_prefetch_cache: str, plugin_user_context: str, + *, preflight_compressed: bool, +) -> None: + """Multimodal (list) content takes no string sidecar: the turn's context becomes a durable + text part on the current turn's live list (the gateway must-deliver-note channel, #71998), + so wire, persisted row, compaction and replay all carry the same parts. Runs once per turn, + before the first request; historical rows are never touched. + + A user row another writer materialized BEFORE the prologue (in-place preflight compaction, + a close/early flush that raced it) is updated in place: the crash persist marker-skips that + message, so without this a resumed session replays a view the model never saw. Same + ``_row_id``-under-lock protocol as the string sidecar backfill; the row keeps its writer's + shape (compaction inserted the raw parts, a flush the text projection).""" + _mm_ctx = compose_multimodal_context_part(ext_prefetch_cache, plugin_user_context) + if not append_notes_to_multimodal_content(turn_user_msg.get("content"), _mm_ctx): + return + from agent.session_persistence import _durable_content, _persist_lock + + with _persist_lock(agent): + _row_id = turn_user_msg.get("_row_id") + _db = getattr(agent, "_session_db", None) + if _db is None or not isinstance(_row_id, int): + return + _in_place_compacted = preflight_compressed and bool(getattr(agent, "_last_compaction_in_place", False)) + content = turn_user_msg["content"] if _in_place_compacted else _durable_content(turn_user_msg["content"]) + try: + _db.set_user_message_content(agent.session_id, _row_id, content) + except Exception: + logger.warning("multimodal context backfill failed for session=%s", agent.session_id or "none", exc_info=True) + + def _persist_turn_start( agent: Any, messages: List[Any], conversation_history: Optional[List[Any]], pending_cli_message: Any, @@ -1072,24 +1123,26 @@ def build_turn_context( _bind_interrupt_scope(agent, ra) ext_prefetch_cache = _memory_turn_start_and_prefetch(agent, original_user_message, turn_author) - # Sidecar skipped for codex_app_server/MoA. - if ( - not moa_active - and getattr(agent, "api_mode", None) != "codex_app_server" - and 0 <= current_turn_user_idx < len(messages) - and messages[current_turn_user_idx].get("role") == "user" - ): - _stamp_api_content_sidecar( - agent, messages, current_turn_user_idx, ext_prefetch_cache, - plugin_user_context, preflight_compressed=compaction.compressed, - ) + # Title the session now: titling depends only on the user's ask (before any injected + # context lands on list content), so it runs concurrently with the turn. Daemon thread, + # no-op once titled; it ensures the session row itself. + _maybe_title_session_at_turn_start(agent, messages) + + # Sidecar skipped for codex_app_server/MoA; list content carries its context as a part in every mode. + if 0 <= current_turn_user_idx < len(messages) and messages[current_turn_user_idx].get("role") == "user": + if isinstance(messages[current_turn_user_idx].get("content"), list): + _append_multimodal_context( + agent, messages[current_turn_user_idx], ext_prefetch_cache, plugin_user_context, + preflight_compressed=compaction.compressed, + ) + elif not moa_active and getattr(agent, "api_mode", None) != "codex_app_server": + _stamp_api_content_sidecar( + agent, messages, current_turn_user_idx, ext_prefetch_cache, + plugin_user_context, preflight_compressed=compaction.compressed, + ) _persist_turn_start(agent, messages, conversation_history, pending_cli_message) - # Title the session now: the row exists and titling depends only on the user's ask, - # so it runs concurrently with the turn. Daemon thread, no-op once titled. - _maybe_title_session_at_turn_start(agent, messages) - return TurnContext( user_message=user_message, original_user_message=original_user_message, messages=messages, conversation_history=conversation_history, active_system_prompt=active_system_prompt, diff --git a/agent/turn_final_response.py b/agent/turn_final_response.py index af36bc717f..779a473003 100644 --- a/agent/turn_final_response.py +++ b/agent/turn_final_response.py @@ -301,6 +301,22 @@ def finish_text_response( final_response = None return _verdict("continue") + # Plugins rewrite the reply BEFORE it is appended and flushed: SQLite treats a non-blank + # assistant row as settled, so a transform after this write would reach the user but never + # the stored/replayed transcript (#44239). finalize_turn reads the recorded outcome; like + # there, an interrupted turn keeps the raw text. + from agent.turn_finalizer import apply_llm_output_transform + _transformed = False + if not getattr(agent, "_interrupt_requested", False): + final_response, _transformed, _ = apply_llm_output_transform( + agent, final_response, turn_id=getattr(agent, "_current_turn_id", "") or "", logger=logger, + ) + if _transformed: + if _promoted: + final_msg["api_content"] = final_response + else: + final_msg["content"] = final_response + append_message(messages, final_msg) # Make the answer durable before leaving the loop (_DB_PERSISTED_MARKER keeps # _persist_session idempotent). Failure must NOT abort the turn: finalize retries. diff --git a/agent/turn_finalizer.py b/agent/turn_finalizer.py index e1c2de569e..3db49dd22f 100644 --- a/agent/turn_finalizer.py +++ b/agent/turn_finalizer.py @@ -413,21 +413,16 @@ def _apply_output_hooks( agent, final_response, logger, *, platform, effective_task_id, turn_id, original_user_message, messages, ) -> Tuple[Any, bool, Optional[Any]]: - """Fire ``transform_llm_output`` then ``post_llm_call`` once per turn after the tool loop. - Returns ``(final_response, transformed, pre_transform_response)``.""" - transformed, pre_transform = False, None - # First hook to return a string wins; None/empty leaves the text unchanged. - for _hook_result in _invoke_hook_safely( - "transform_llm_output", logger, - response_text=final_response, - session_id=agent.session_id or "", - model=agent.model, - platform=platform, - turn_id=turn_id, # per-turn identity for the hook callback gate - ): - if isinstance(_hook_result, str) and _hook_result: - pre_transform, final_response, transformed = final_response, _hook_result, True - break + """Resolve the turn's ``transform_llm_output`` outcome, then fire ``post_llm_call`` once per + turn after the tool loop. Returns ``(final_response, transformed, pre_transform_response)``. + + The transform itself normally already ran before the assistant row was first persisted + (``apply_llm_output_transform`` from ``finish_text_response`` / ``_persist_step``); this + call returns that recorded outcome, and only fires the hook here when no earlier seam saw a + response (e.g. text that only appeared through ``_explain_abnormal_exit``).""" + final_response, transformed, pre_transform = apply_llm_output_transform( + agent, final_response, turn_id=turn_id, platform=platform, logger=logger, + ) # Detached forks are internal work and must not publish turns under the parent's session ID. if not getattr(agent, "_persist_disabled", False): _invoke_hook_safely( @@ -444,6 +439,48 @@ def _apply_output_hooks( return final_response, transformed, pre_transform +def apply_llm_output_transform( + agent, final_response, *, turn_id, platform=None, logger=None, +) -> Tuple[Any, bool, Optional[Any]]: + """Fire ``transform_llm_output`` once per turn and return + ``(final_response, transformed, pre_transform_response)``. + + Called BEFORE the final assistant row is first persisted — from ``finish_text_response`` + ahead of its durable flush, and from ``finalize_turn._persist_step`` ahead of the + recovery-path tail close — so the text the user sees is the text stored in SQLite/JSON and + replayed next turn (#44239). SQLite treats a non-blank assistant row as settled (a re-flush + adopts the stored content rather than overwriting it), so transforming after that first + write can never reach the durable store. Idempotent per ``turn_id``: later callers in the + same turn get the recorded outcome instead of a second hook firing. Only the current + turn's not-yet-written text is touched — earlier turns and the system prompt are never + rewritten (prompt-cache invariant).""" + if logger is None: + from agent.conversation_loop import logger + recorded = getattr(agent, "_llm_output_transform", None) + if isinstance(recorded, tuple) and len(recorded) == 3 and recorded[0] == turn_id: + _, transformed, pre_transform = recorded + return final_response, transformed, pre_transform + if not final_response: + return final_response, False, None + if platform is None: + platform = getattr(agent, "platform", None) or "" + transformed, pre_transform = False, None + # First hook to return a string wins; None/empty leaves the text unchanged. + for _hook_result in _invoke_hook_safely( + "transform_llm_output", logger, + response_text=final_response, + session_id=agent.session_id or "", + model=agent.model, + platform=platform, + turn_id=turn_id, # per-turn identity for the hook callback gate + ): + if isinstance(_hook_result, str) and _hook_result: + pre_transform, final_response, transformed = final_response, _hook_result, True + break + agent._llm_output_transform = (turn_id, transformed, pre_transform) + return final_response, transformed, pre_transform + + def finalize_turn( agent, *, final_response, api_call_count, interrupted, failed, messages, conversation_history, effective_task_id, turn_id, user_message, original_user_message, _should_review_memory, @@ -508,6 +545,12 @@ def finalize_turn( final_response, _recovered_from_stream = _recover_final_from_stream( agent, final_response, interrupted, failed ) + # Recovery paths (stream-recovered / prior-turn text) reach here with a response no + # earlier seam transformed; the normal text turn already did this before its flush and + # gets the recorded outcome back. Either way the tail close below writes the text the + # user will see, never the raw model text (#44239). + if final_response and not interrupted: + final_response, _, _ = apply_llm_output_transform(agent, final_response, turn_id=turn_id, logger=logger) _close_transcript_tail(agent, messages, final_response, interrupted, _recovered_from_stream) if not interrupted and not failed: _micro_compact_after_turn(agent, messages, final_response, logger) diff --git a/apps/desktop/e2e/completed-reply-refresh.spec.ts b/apps/desktop/e2e/completed-reply-refresh.spec.ts new file mode 100644 index 0000000000..e8bd69645b --- /dev/null +++ b/apps/desktop/e2e/completed-reply-refresh.spec.ts @@ -0,0 +1,147 @@ +import { expect, test } from '@playwright/test' + +import { MOCK_REPLY } from '../../../tests-js/scripts/mock-server' + +import { setupMockBackend, waitForAppReady } from './fixtures' + +// Real Electron, backend, SQLite and stream. Only delivery of one real history +// request is delayed: admit while idle, read while the next turn is held, then +// deliver after completion. The mock provider intentionally repeats its answer. +test('a delayed history read cannot remove the latest completed reply', async () => { + test.setTimeout(150_000) + const nextPrompt = 'Complete the second request for the refresh regression' + const fixture = await setupMockBackend({ mockServer: { holdFirstStreamForPrompt: nextPrompt } }) + const { app, page, mock } = fixture + const events: string[] = [] + page.on('websocket', socket => { + socket.on('framereceived', frame => { + try { + const message = JSON.parse(String(frame.payload)) + + if (message.method === 'event') { + events.push(message.params?.type) + } + } catch { + /* Non-JSON transport frames carry no turn event. */ + } + }) + }) + + try { + await waitForAppReady(fixture) + + const composer = page + .locator('[data-slot="composer-root"] [contenteditable="true"]') + .filter({ visible: true }) + .first() + + await composer.fill('First request for the refresh regression') + await composer.press('Enter') + const answers = page.locator('[data-slot="aui_assistant-message-content"]').filter({ hasText: MOCK_REPLY }) + await expect(answers).toHaveCount(1, { timeout: 60_000 }) + await expect.poll(() => events.filter(type => type === 'message.complete').length).toBeGreaterThan(0) + const stored = await page.evaluate(() => decodeURIComponent(location.hash.slice(2).split('?')[0])) + expect(stored).not.toBe('') + + // The private IPC handler table is used only by this test. Preserve the + // original handler (including backend routing and actual SQLite reads). + await app.evaluate(({ ipcMain }, sessionId) => { + const handlers = (ipcMain as unknown as { _invokeHandlers: Map Promise> }) + ._invokeHandlers + + const original = handlers.get('hermes:api')! + + const control = { + admitted: false, + sampled: false, + releaseRead: () => {}, + releaseDelivery: () => {}, + snapshot: null as any + } + + const readGate = new Promise(resolve => { + control.releaseRead = resolve + }) + + const deliveryGate = new Promise(resolve => { + control.releaseDelivery = resolve + }) + + ;(globalThis as any).__historyRead = control + ipcMain.removeHandler('hermes:api') + ipcMain.handle('hermes:api', async (event: unknown, request: { path: string }) => { + if (!control.admitted && request.path.startsWith(`/api/sessions/${sessionId}/messages?`)) { + control.admitted = true + await readGate + const snapshot = await original(event, request) + control.snapshot = snapshot + control.sampled = true + await deliveryGate + + return snapshot + } + + return original(event, request) + }) + }, stored) + + // An actual backend metadata write produces the production change tick. + await page.evaluate(async sessionId => { + await (window as any).hermesDesktop.api({ + path: `/api/sessions/${sessionId}`, + method: 'PATCH', + body: { title: 'Completed reply regression' } + }) + }, stored) + await expect + .poll(() => app.evaluate(() => (globalThis as any).__historyRead.admitted), { timeout: 20_000 }) + .toBe(true) + await composer.fill(nextPrompt) + await composer.press('Enter') + await mock.waitForHeldStream() + await app.evaluate(() => { + ;(globalThis as any).__historyRead.releaseRead() + }) + await expect.poll(() => app.evaluate(() => (globalThis as any).__historyRead.sampled)).toBe(true) + const sampled = await app.evaluate(() => (globalThis as any).__historyRead.snapshot) + expect(sampled.messages.some((message: any) => message.role === 'user' && message.content === nextPrompt)).toBe( + true + ) + expect( + sampled.messages.filter((message: any) => message.role === 'assistant' && message.content === MOCK_REPLY) + ).toHaveLength(1) + const completeCount = events.filter(type => type === 'message.complete').length + mock.releaseHeldStream() + await expect + .poll(() => events.filter(type => type === 'message.complete').length, { timeout: 60_000 }) + .toBeGreaterThan(completeCount) + await expect(answers).toHaveCount(2) + await page.screenshot({ path: test.info().outputPath('completed-before-history.png') }) + await app.evaluate(() => { + ;(globalThis as any).__historyRead.releaseDelivery() + }) + // A bounded quiet window checks that the answer stays painted after the + // delayed IPC promise and the renderer effects have been processed. + await page.waitForTimeout(1500) + await expect(answers).toHaveCount(2) + await expect(page.locator('[data-slot="aui_thread-viewport"]')).toContainText(nextPrompt) + await page.screenshot({ path: test.info().outputPath('completed-after-history.png') }) + await test.info().attach('transport-receipt', { + body: JSON.stringify( + { stored, events, sampledRows: sampled.messages.length, visibleAnswers: await answers.count() }, + null, + 2 + ), + contentType: 'application/json' + }) + } finally { + await app + .evaluate(() => { + ;(globalThis as any).__historyRead?.releaseRead() + ;(globalThis as any).__historyRead?.releaseDelivery() + }) + .catch(() => undefined) + mock.releaseHeldStream() + await fixture.cleanup() + } +}) diff --git a/apps/desktop/electron/backend-exit-recovery.test.ts b/apps/desktop/electron/backend-exit-recovery.test.ts index f7fc379518..512d7cf4f6 100644 --- a/apps/desktop/electron/backend-exit-recovery.test.ts +++ b/apps/desktop/electron/backend-exit-recovery.test.ts @@ -89,3 +89,66 @@ test('a backend that dies after every ready is respawned at most maxRespawns tim assert.equal(latch.claim(empty), true) assert.equal(latch.isCrashLooping(), false) }) + +test('a claimed recovery that fails before ready can retry within the same crash-loop budget', () => { + let clock = 1_000 + const latch = createBackendExitRecoveryLatch({ maxRespawns: 3, windowMs: 120_000, now: () => clock }) + const empty = { hasCurrentOwner: false, hasPendingStart: false, intentionalTeardown: false } + + assert.equal(latch.claim(empty), true, 'ready backend death grants the first recovery') + clock += 1_000 + assert.equal(latch.retryAfterFailedStart(empty), true, 'pre-ready failure grants a bounded retry') + clock += 1_000 + assert.equal(latch.retryAfterFailedStart(empty), true, 'the final budgeted retry is still admitted') + clock += 1_000 + assert.equal(latch.retryAfterFailedStart(empty), false, 'a fourth recovery attempt is refused') + assert.equal(latch.isCrashLooping(), true) +}) + +test('a failed recovery does not release its claim while another owner/start or teardown is present', () => { + const latch = createBackendExitRecoveryLatch() + const empty = { hasCurrentOwner: false, hasPendingStart: false, intentionalTeardown: false } + + assert.equal(latch.claim(empty), true) + assert.equal(latch.retryAfterFailedStart({ ...empty, hasPendingStart: true }), false) + assert.equal(latch.claim(empty), false, 'the pending-start refusal preserves the original claim') + + latch.reset() + assert.equal(latch.claim(empty), true) + assert.equal(latch.retryAfterFailedStart({ ...empty, hasCurrentOwner: true }), false) + assert.equal(latch.claim(empty), false, 'the current-owner refusal preserves the original claim') + + latch.reset() + assert.equal(latch.claim(empty), true) + assert.equal(latch.retryAfterFailedStart({ ...empty, intentionalTeardown: true }), false) + assert.equal(latch.claim(empty), false, 'intentional teardown does not re-arm recovery') +}) + +test('a failed start that owns no recovery claim is not re-armed and spends no budget', () => { + let clock = 1_000 + const latch = createBackendExitRecoveryLatch({ maxRespawns: 3, windowMs: 120_000, now: () => clock }) + const empty = { hasCurrentOwner: false, hasPendingStart: false, intentionalTeardown: false } + + // Nothing claimed yet (fresh latch): a pre-ready failure of a user-driven + // start is not the supervisor's retry to take. + assert.equal(latch.retryAfterFailedStart(empty), false, 'fresh latch has no claim to re-arm') + assert.equal(latch.isCrashLooping(), false) + + // A ready backend released the claim: a later failed start still owns nothing. + assert.equal(latch.claim(empty), true) + latch.reset() + clock += 1_000 + assert.equal(latch.retryAfterFailedStart(empty), false, 'reset() leaves nothing to re-arm') + assert.equal(latch.isCrashLooping(), false) + + // Neither refusal consumed the window: the remaining two grants are intact. + clock += 1_000 + assert.equal(latch.claim(empty), true, 'second respawn') + latch.reset() + clock += 1_000 + assert.equal(latch.claim(empty), true, 'third respawn') + latch.reset() + clock += 1_000 + assert.equal(latch.claim(empty), false, 'fourth is the real budget exhaustion') + assert.equal(latch.isCrashLooping(), true) +}) diff --git a/apps/desktop/electron/backend-exit-recovery.ts b/apps/desktop/electron/backend-exit-recovery.ts index e7a5223ae7..a97a217e16 100644 --- a/apps/desktop/electron/backend-exit-recovery.ts +++ b/apps/desktop/electron/backend-exit-recovery.ts @@ -39,6 +39,28 @@ export function createBackendExitRecoveryLatch({ let respawnedAt: number[] = [] let crashLooping = false + const blocked = (state: BackendExitRecoveryState) => + state.hasCurrentOwner || state.hasPendingStart || state.intentionalTeardown + + const claim = (state: BackendExitRecoveryState): boolean => { + if (claimed || blocked(state)) { + return false + } + + const at = now() + respawnedAt = respawnedAt.filter(t => at - t < windowMs) + crashLooping = respawnedAt.length >= maxRespawns + + if (crashLooping) { + return false + } + + respawnedAt.push(at) + claimed = true + + return true + } + return { /** * True exactly once per empty slot; `reset()` when a backend becomes ready @@ -47,25 +69,24 @@ export function createBackendExitRecoveryLatch({ * within `windowMs` is a crash loop, and the supervisor stops respawning * (`isCrashLooping()`) instead of cycling child + error toast forever. */ - claim(state: BackendExitRecoveryState): boolean { - if (claimed || state.hasCurrentOwner || state.hasPendingStart || state.intentionalTeardown) { + claim, + /** + * Re-arm only the recovery attempt that already owns this latch and failed + * before reaching ready. A concurrent owner/start or intentional teardown + * keeps the existing claim intact; its eventual ready transition owns the + * normal reset. A real retry consumes the same crash-loop budget as every + * other supervisor respawn. + */ + retryAfterFailedStart(state: BackendExitRecoveryState): boolean { + if (!claimed || blocked(state)) { return false } - const at = now() - respawnedAt = respawnedAt.filter(t => at - t < windowMs) - crashLooping = respawnedAt.length >= maxRespawns + claimed = false - if (crashLooping) { - return false - } - - respawnedAt.push(at) - claimed = true - - return true + return claim(state) }, - /** True when the last `claim` was refused because the respawn budget for the window is spent. */ + /** True when the last claim attempt was refused because the respawn budget for the window is spent. */ isCrashLooping(): boolean { return crashLooping }, diff --git a/apps/desktop/electron/backend-start-failure.test.ts b/apps/desktop/electron/backend-start-failure.test.ts index ef488e7968..b3c1247787 100644 --- a/apps/desktop/electron/backend-start-failure.test.ts +++ b/apps/desktop/electron/backend-start-failure.test.ts @@ -23,6 +23,13 @@ test('never latches a REMOTE failure so recovery stays retryable without a resta assert.equal(shouldLatchBackendStartFailure({ attemptedRemote: true }), false) }) +test('never latches a supervisor-owned respawn failure (it has its own bounded crash-loop budget)', () => { + // A pre-ready child exit during a supervisor respawn must be able to spend + // the remaining crash-loop slots instead of becoming a permanent local latch. + assert.equal(shouldLatchBackendStartFailure({ attemptedRemote: false, supervisorRecovery: true }), false) + assert.equal(shouldLatchBackendStartFailure({ attemptedRemote: false, supervisorRecovery: false }), true) +}) + test('the two branches are mutually exclusive (a failure either latches or stays retryable)', () => { for (const attemptedRemote of [true, false]) { const latched = shouldLatchBackendStartFailure({ attemptedRemote }) diff --git a/apps/desktop/electron/backend-start-failure.ts b/apps/desktop/electron/backend-start-failure.ts index 66ea47a0a3..5134290d59 100644 --- a/apps/desktop/electron/backend-start-failure.ts +++ b/apps/desktop/electron/backend-start-failure.ts @@ -28,6 +28,11 @@ export interface BackendStartFailureContext { * cloud) primary backend rather than spawning a local child. */ attemptedRemote: boolean + /** + * True when the boot that just failed was a supervisor-owned respawn after + * an unexpected primary exit, not an initial or user-driven start. + */ + supervisorRecovery?: boolean } /** @@ -35,9 +40,14 @@ export interface BackendStartFailureContext { * Latch local failures (prevent install-restart loops); never latch remote * failures (they are transient and must stay retryable so recovery paths work * without an app restart). + * + * A supervisor-owned respawn never latches either: it already has its own + * bounded crash-loop budget, so a pre-ready child exit must not become a + * permanent local boot latch before that budget can run. Initial and + * user-driven starts keep the fail-closed latch. */ export function shouldLatchBackendStartFailure(context: BackendStartFailureContext): boolean { - return !context.attemptedRemote + return !context.attemptedRemote && !context.supervisorRecovery } export interface RemoteReauthFailureContext { diff --git a/apps/desktop/electron/main.ts b/apps/desktop/electron/main.ts index 4ec40c4f6b..41e03aafbf 100644 --- a/apps/desktop/electron/main.ts +++ b/apps/desktop/electron/main.ts @@ -12140,22 +12140,81 @@ function releaseHostSpawnReservation() { hostSpawnReservation = null } -function startHermes(): Promise>> { +function startHermes({ supervisorRecovery = false }: { supervisorRecovery?: boolean } = {}): Promise>> { primaryRecoverySuppressed = false primaryStartsInFlight += 1 const start: Promise>> = - localBackendLifecycle.start(runHermesStart) + localBackendLifecycle.start(() => runHermesStart({ supervisorRecovery })) const releaseStart = (): void => { primaryStartsInFlight -= 1 } + // Ordering contract: this reaction is registered on the SAME promise the + // caller receives, before any caller `.catch`, so releaseStart has already + // run (primaryStartsInFlight back to 0) when runPrimaryRecoverySpawn's + // `.catch` evaluates primaryRecoveryState(). Returning a derived promise + // (start.then(...)) or wrapping `start` would invert that order: every + // pre-ready retry would see hasPendingStart:true, be refused, and leave the + // recovery claim stuck with no retry and no UI. void start.then(releaseStart, releaseStart) return start } +function primaryRecoveryState() { + return { + hasCurrentOwner: backendConnectionState.getProcess() !== null || backendConnectionState.getPromise() !== null, + hasPendingStart: primaryStartsInFlight > 0, + intentionalTeardown: primaryRecoverySuppressed || isQuittingForHandoff || backendShutdown.hasStarted() + } +} + +function reportPrimaryRecoveryCrashLoop(code: number | null, signal: string | null): boolean { + if (!primaryExitRecovery.isCrashLooping()) { + return false + } + + const message = + 'Hermes backend keeps crashing right after it restarts; not restarting it again. Relaunch Hermes Desktop.' + + rememberLog(`[supervisor] ${message}`) + sendBackendExit({ code, signal, error: message }) + + return true +} + +const firstLine = (text: string): string => (text || '').split('\n').find(Boolean) || '' + +function runPrimaryRecoverySpawn(code: number | null, signal: string | null) { + startHermes({ supervisorRecovery: true }).catch(respawnError => { + rememberLog(`[supervisor] backend respawn failed: ${firstLine(respawnError.message)}`) + + // Terminal boot failures still own their existing recovery UI. Only a + // supervisor-owned respawn that failed transiently before ready may spend + // another bounded recovery slot. + const latched = latchedBootFailure() + + if (latched) { + rememberLog(`[supervisor] respawn refused: boot failure latched: ${firstLine(latched.message)}`) + + return + } + + // releaseStart (startHermes) already ran: same-promise reaction order, so + // hasPendingStart is false here. See the ordering contract in startHermes. + if (primaryExitRecovery.retryAfterFailedStart(primaryRecoveryState())) { + rememberLog('[supervisor] backend respawn failed before ready; retrying within crash-loop budget') + runPrimaryRecoverySpawn(code, signal) + + return + } + + reportPrimaryRecoveryCrashLoop(code, signal) + }) +} + // A ready primary child died. When its exit leaves the primary slot with no // owner and no start in flight (outside an intentional teardown), the // supervisor owns the respawn (#112344): the stale-classified exit used to @@ -12172,34 +12231,31 @@ function scheduleUnexpectedPrimaryRecovery({ return false } - const claimed: boolean = primaryExitRecovery.claim({ - hasCurrentOwner: backendConnectionState.getProcess() !== null || backendConnectionState.getPromise() !== null, - hasPendingStart: primaryStartsInFlight > 0, - intentionalTeardown: primaryRecoverySuppressed || isQuittingForHandoff || backendShutdown.hasStarted() - }) + const claimed = primaryExitRecovery.claim(primaryRecoveryState()) if (!claimed) { - if (primaryExitRecovery.isCrashLooping()) { - const message: string = - 'Hermes backend keeps crashing right after it restarts; not restarting it again. Relaunch Hermes Desktop.' - - rememberLog(`[supervisor] ${message}`) - sendBackendExit({ code, signal, error: message }) - - return true - } - - return false + return reportPrimaryRecoveryCrashLoop(code, signal) } rememberLog('[supervisor] backend exit left no primary owner and no start in flight; respawning') sendBackendExit({ code, signal, ...(error ? { error } : {}) }) - startHermes().catch(respawnError => rememberLog(`[supervisor] backend respawn failed: ${respawnError.message}`)) + runPrimaryRecoverySpawn(code, signal) return true } -async function runHermesStart(): Promise>> { +/** + * The terminal boot failure currently latched in this process, if any. These + * latches are cleared only by an explicit recovery path (reset, repair, + * apply-config, confirmed sign-in, or the child 'exit' handler), never by a + * retry, so both the per-request short-circuit in runHermesStart and the + * supervisor's respawn refusal must consult the same trio in the same order. + */ +function latchedBootFailure(): Error | null { + return bootstrapFailure ?? backendStartFailure ?? remoteReauthFailure ?? null +} + +async function runHermesStart({ supervisorRecovery = false }: { supervisorRecovery?: boolean } = {}): Promise>> { // Only the single-instance lock holder may reap/spawn/claim the desktop // backend. A lock-losing instance must stay inert even if some path reaches // here (e.g. the deferred-quit window before `ready`): its reapOrphans() @@ -12218,19 +12274,19 @@ async function runHermesStart(): Promise startHermes), so a log line here would flood the + // bounded rememberLog ring and evict the lines that explain the original + // failure. The supervisor logs the refusal once in runPrimaryRecoverySpawn. + const latched = latchedBootFailure() - if (backendStartFailure) { - throw backendStartFailure - } - - // A confirmed remote reauth rejection is terminal until the user signs in. - // Short-circuiting here keeps the boot-failure overlay latched and its - // "Sign in" button clickable, instead of re-driving boot on every retry. - if (remoteReauthFailure) { - throw remoteReauthFailure + if (latched) { + throw latched } // E2E: simulate a boot failure without breaking the real backend. The boot @@ -12680,7 +12736,8 @@ async function runHermesStart(): Promise { + afterEach(cleanup) + + it('renders one control per schema type from the table, secrets masked with no value echoed', () => { + render( true)} />) + + expect((screen.getByLabelText('API URL') as HTMLInputElement).value).toBe('https://a') + expect((screen.getByLabelText(/^Retries/) as HTMLInputElement).type).toBe('number') + expect(screen.getByRole('switch', { name: 'Verbose' }).getAttribute('aria-checked')).toBe('false') + expect(screen.getByRole('combobox', { name: 'Mode' })).toBeTruthy() + const secret = screen.getByLabelText(/^API key/) as HTMLInputElement + expect(secret.type).toBe('password') + expect(secret.value).toBe('') + expect(secret.placeholder).toContain('set') + expect((screen.getByLabelText(/^Extra/) as HTMLTextAreaElement).value).toBe(JSON.stringify({ a: 1 }, null, 2)) + // Nothing changed yet → nothing to save. + expect((screen.getByRole('button', { name: 'Save settings' }) as HTMLButtonElement).disabled).toBe(true) + }) + + it('submits only what changed, coerced to wire types, secrets routed by env name; clears secrets after save', async () => { + const onSave = vi.fn(async () => true) + render() + + fireEvent.change(screen.getByLabelText(/^Retries/), { target: { value: '9' } }) + fireEvent.click(screen.getByRole('switch', { name: 'Verbose' })) + fireEvent.change(screen.getByLabelText(/^API key/), { target: { value: 'sk-x' } }) + fireEvent.submit(screen.getByTestId('p-settings-form')) + + await vi.waitFor(() => + expect(onSave).toHaveBeenCalledWith({ secrets: { DEMO_API_KEY: 'sk-x' }, values: { retries: 9, verbose: true } }) + ) + await vi.waitFor(() => expect((screen.getByLabelText(/^API key/) as HTMLInputElement).value).toBe('')) + + // A value the plugin cannot accept never reaches the backend. + expect(() => collectChanges(FIELDS, { ...initialDraft(FIELDS), mode: 'reckless' })).toThrow(/Mode/) + expect(() => collectChanges(FIELDS, { ...initialDraft(FIELDS), retries: 'five' })).toThrow(/number/) + }) +}) diff --git a/apps/desktop/src/app/capabilities/plugins/plugin-settings-form.tsx b/apps/desktop/src/app/capabilities/plugins/plugin-settings-form.tsx new file mode 100644 index 0000000000..55a4bf96b7 --- /dev/null +++ b/apps/desktop/src/app/capabilities/plugins/plugin-settings-form.tsx @@ -0,0 +1,299 @@ +import { type ReactNode, useMemo, useState } from 'react' + +import { Button } from '@/components/ui/button' +import { Field, FieldHint } from '@/components/ui/field' +import { Input } from '@/components/ui/input' +import { Select, SelectContent, SelectItem, SelectTrigger, SelectValue } from '@/components/ui/select' +import { Switch } from '@/components/ui/switch' +import { Textarea } from '@/components/ui/textarea' +import { useI18n } from '@/i18n' +import type { PluginSettingField, PluginSettingFieldType } from '@/store/agent-plugins' + +// A plugin manifest's `config_schema` rendered as a form. Every non-secret +// value is edited as TEXT (the draft) and coerced to its wire type on save, so a +// half-typed number never fights the input; secrets are drafted separately and +// only ever sent to the `.env` credential route by the caller. + +export type PluginSettingsDraft = Record + +export interface PluginSettingsSave { + values: Record + secrets: Record +} + +interface ControlProps { + field: PluginSettingField + id: string + raw: string + disabled: boolean + onChange: (raw: string) => void + /** Localised hint for a secret that already has a stored value. */ + secretSetHint: string +} + +/** What the input shows before the user touches it. */ +const INITIAL_TEXT: Record string> = { + boolean: field => (field.value === true ? 'true' : 'false'), + enum: field => String(field.value ?? field.choices?.[0] ?? ''), + json: field => (field.value === undefined || field.value === null ? '' : JSON.stringify(field.value, null, 2)), + number: field => (field.value === undefined || field.value === null ? '' : String(field.value)), + secret: () => '', + string: field => String(field.value ?? '') +} + +/** Text → wire value; a thrown Error is the field's validation message. */ +const COERCE: Record, (raw: string, field: PluginSettingField) => unknown> = { + boolean: raw => raw === 'true', + enum: (raw, field) => { + if (!field.choices?.includes(raw)) { + throw new Error(`${field.label}: not one of ${field.choices?.join(', ') ?? ''}`) + } + + return raw + }, + json: (raw, field) => { + const parsed: unknown = JSON.parse(raw || 'null') + + if (parsed === null || typeof parsed !== 'object') { + throw new Error(`${field.label}: expected a JSON list or object`) + } + + return parsed + }, + number: (raw, field) => { + const n = Number(raw) + + if (raw.trim() === '' || Number.isNaN(n)) { + throw new Error(`${field.label}: expected a number`) + } + + return n + }, + string: raw => raw +} + +const TEXT_CLASS = 'h-7 text-xs' + +function StringControl({ disabled, id, onChange, raw }: ControlProps) { + return ( + onChange(e.currentTarget.value)} + value={raw} + /> + ) +} + +function NumberControl({ disabled, id, onChange, raw }: ControlProps) { + return ( + onChange(e.currentTarget.value)} + type="number" + value={raw} + /> + ) +} + +function BooleanControl({ disabled, field, id, onChange, raw }: ControlProps) { + return ( + onChange(on ? 'true' : 'false')} + /> + ) +} + +function EnumControl({ disabled, field, id, onChange, raw }: ControlProps) { + return ( + + ) +} + +function SecretControl({ disabled, field, id, onChange, raw, secretSetHint }: ControlProps) { + return ( + onChange(e.currentTarget.value)} + placeholder={field.has_value ? secretSetHint : ''} + type="password" + value={raw} + /> + ) +} + +function JsonControl({ disabled, id, onChange, raw }: ControlProps) { + return ( +