diff --git a/.env.example b/.env.example index 02985b54a4..78c0108e1c 100644 --- a/.env.example +++ b/.env.example @@ -533,3 +533,12 @@ IMAGE_TOOLS_DEBUG=false # GOOGLE_CHAT_ALLOW_ALL_USERS=false # Set true to skip the allowlist # GOOGLE_CHAT_HOME_CHANNEL= # Default space (spaces/XXXX) for cron delivery # GOOGLE_CHAT_HOME_CHANNEL_NAME= # Display name for the home channel + +# ============================================================================= +# reddit-reading skill (optional) — app-only credentials, NOT a user login +# ============================================================================= +# The skill works with no credentials via Reddit's public feeds (~1 request/minute). +# For faster access with scores and nested comments, register a free "script" app +# at https://www.reddit.com/prefs/apps and paste its id and secret here. +# REDDIT_CLIENT_ID= +# REDDIT_CLIENT_SECRET= diff --git a/agent/anthropic_endpoints.py b/agent/anthropic_endpoints.py index 82b96c5af1..237a562176 100644 --- a/agent/anthropic_endpoints.py +++ b/agent/anthropic_endpoints.py @@ -72,6 +72,25 @@ def _is_kimi_family_endpoint(base_url: str | None, model: str | None = None) -> ) + +_DEEPSEEK_THINKING_MODEL_PREFIXES = ( + "deepseek-r", "deepseek-v4", "deepseek_v4", "deepseek-pro", + "deepseek_pro", "deepseek-flash", "deepseek_flash", +) + + +def _model_name_is_deepseek_thinking(model: str | None) -> bool: + """Known DeepSeek thinking families behind an Anthropic-compatible relay. + + Strip vendor namespaces, but do not treat arbitrary DeepSeek chat/distill + names as evidence of the thinking replay contract. + """ + if not isinstance(model, str): + return False + name = model.strip().lower().rsplit("/", 1)[-1] + return bool(name) and name.startswith(_DEEPSEEK_THINKING_MODEL_PREFIXES) + + def _is_deepseek_anthropic_endpoint(base_url: str | None) -> bool: """DeepSeek's ``/anthropic`` route. In thinking mode DeepSeek requires prior-turn ``thinking`` blocks to round-trip while the generic third-party path strips them; its blocks are unsigned, diff --git a/agent/anthropic_message_convert.py b/agent/anthropic_message_convert.py index 2335345ea7..5309ba7fc9 100644 --- a/agent/anthropic_message_convert.py +++ b/agent/anthropic_message_convert.py @@ -12,7 +12,7 @@ from typing import Any, Dict, List, Optional, Tuple from agent.anthropic_endpoints import ( _is_deepseek_anthropic_endpoint, _is_kimi_family_endpoint, _is_nous_portal_endpoint, - _is_third_party_anthropic_endpoint, + _is_third_party_anthropic_endpoint, _model_name_is_deepseek_thinking, ) logger = logging.getLogger(__name__) @@ -565,7 +565,9 @@ def _manage_thinking_signatures(result: List[Dict[str, Any]], base_url: str | No """ is_third_party = _is_third_party_anthropic_endpoint(base_url) and not _is_nous_portal_endpoint(base_url) is_kimi = _is_kimi_family_endpoint(base_url, model) - is_deepseek = _is_deepseek_anthropic_endpoint(base_url) + is_deepseek = _is_deepseek_anthropic_endpoint(base_url) or ( + is_third_party and _model_name_is_deepseek_thinking(model) + ) last_assistant_idx = next((i for i in range(len(result) - 1, -1, -1) if result[i].get("role") == "assistant"), None) for idx, m in _assistant_block_lists(result): if is_kimi: diff --git a/agent/auxiliary_client.py b/agent/auxiliary_client.py index ee492ef6fe..e9f5fe79b8 100644 --- a/agent/auxiliary_client.py +++ b/agent/auxiliary_client.py @@ -2957,7 +2957,8 @@ def _contains_any(text: str, needles: Tuple[str, ...]) -> bool: # Billing-body markers (credit exhaustion wrapped in 402/403/404/429 bodies), plus daily/weekly quota -# exhaustion (functionally credit exhaustion; "resource exhausted" is the Vertex/gRPC quota phrasing). +# exhaustion (functionally credit exhaustion; "resource exhausted" is the Vertex/gRPC quota phrasing — +# also serialized by SDK wrappers and NIM as RESOURCE_EXHAUSTED / ResourceExhausted / resource-exhausted). _PAYMENT_KEYWORDS = ( "credits", "insufficient funds", "can only afford", "billing", "payment required", "out of funds", "run out of funds", "balance_depleted", "no usable credits", @@ -2965,6 +2966,7 @@ _PAYMENT_KEYWORDS = ( "requires a subscription", "upgrade for access", "upgrade for higher limits", "reached your session usage limit", "quota exceeded", "quota_exceeded", "too many tokens per day", "daily limit", "tokens per day", "daily quota", "resource exhausted", + "resource_exhausted", "resource-exhausted", "resourceexhausted", "weekly usage limit", "weekly limit", ) diff --git a/agent/chat_completion_helpers.py b/agent/chat_completion_helpers.py index 49e005c934..7751836887 100644 --- a/agent/chat_completion_helpers.py +++ b/agent/chat_completion_helpers.py @@ -378,15 +378,24 @@ def _validated_openrouter_provider_sort(raw_sort: Any) -> Optional[str]: def _provider_preferences_for_agent(agent) -> Dict[str, Any]: - """Build the validated provider-routing object shared by request paths.""" - preferences: Dict[str, Any] = {} - for key, value in (("only", agent.providers_allowed), ("ignore", agent.providers_ignored), - ("order", agent.providers_order), ("sort", _validated_openrouter_provider_sort(agent.provider_sort)), - ("require_parameters", True if agent.provider_require_parameters else None), - ("data_collection", agent.provider_data_collection)): - if value: - preferences[key] = value - return preferences + """Build the validated provider-routing object shared by request paths. + + ``provider_routing.models.`` overlays the flat constructor values for the CURRENT + ``agent.model`` (so ``/model`` switches, fallbacks, and delegated children on another + model each get their own pins without any surface re-plumbing the kwargs).""" + flat = {"only": agent.providers_allowed, "ignore": agent.providers_ignored, "order": agent.providers_order, + "sort": agent.provider_sort, "require_parameters": agent.provider_require_parameters, + "data_collection": agent.provider_data_collection} + per_model = {} + with contextlib.suppress(Exception): + from hermes_cli.config import load_config_readonly + from hermes_constants import resolve_per_model_provider_routing + _pr = load_config_readonly().get("provider_routing") + per_model = resolve_per_model_provider_routing(agent.model, (_pr or {}).get("models") if isinstance(_pr, dict) else None) + merged = {**flat, **{k: v for k, v in per_model.items() if k in flat}} + merged["sort"] = _validated_openrouter_provider_sort(merged["sort"]) + merged["require_parameters"] = True if merged["require_parameters"] else None + return {key: value for key, value in merged.items() if value} def _prompt_cache_scope_for_agent(agent) -> "str | None": @@ -1720,6 +1729,19 @@ def build_assistant_message(agent, assistant_message, finish_reason: str) -> dic value = getattr(assistant_message, attr, None) if value: msg[attr] = value + if attr == "codex_reasoning_items": + from agent.codex_responses_adapter import ( + has_replayable_native_compaction_checkpoint, + ) + + note_checkpoint = getattr( + agent.context_compressor, "note_native_compaction_checkpoint", None + ) + if ( + callable(note_checkpoint) + and has_replayable_native_compaction_checkpoint(agent, [msg]) + ): + note_checkpoint() if assistant_tool_calls: msg["tool_calls"] = [_assistant_tool_call_dict(agent, tc, i) for i, tc in enumerate(assistant_tool_calls)] @@ -1797,8 +1819,10 @@ def _fallback_reason_text(reason: "FailoverReason | None") -> str: def _is_anthropic_wire_url(url: str) -> bool: - """Same host match as determine_api_mode() / _detect_api_mode_for_url().""" - return url.rstrip("/").lower().endswith("/anthropic") or base_url_hostname(url) == "api.anthropic.com" + """Same Messages-only host match as determine_api_mode() / _detect_api_mode_for_url(): api.anthropic.com, + a /anthropic suffix, or Kimi Code's api.kimi.com/coding (its /chat/completions 404s — #77256).""" + from hermes_cli.providers import host_mandated_api_mode + return host_mandated_api_mode(url) == "anthropic_messages" def _fallback_api_mode_hint(fb: dict, fb_provider: str, fb_base_url_hint: Optional[str]) -> tuple[bool, str]: @@ -2925,19 +2949,13 @@ class _StreamingCall: role = "assistant" _diag = self._new_diag() self._writer_token = self._attempt_request_client = self._attempt_stream_response = None + from agent.chat_completion_helpers_relay import RelayChatAccumulator + relay_response = RelayChatAccumulator() def _open_stream(next_api_kwargs: dict[str, Any]): timeout = _httpx.Timeout(connect=conn_cap, read=read_timeout, write=base_timeout, pool=conn_cap) return self._open_chat_stream({**next_api_kwargs, "stream": True, "timeout": timeout}) - def _relay_final_response() -> dict[str, Any]: - tool_calls.materialize() - message = {"role": role, "content": "".join(content_parts) or None, - "reasoning_content": "".join(reasoning_parts) or None, - "tool_calls": [tool_calls_acc[i] for i in sorted(tool_calls_acc)] or None} - return {"model": model_name, "usage": usage_obj, - "choices": [{"message": message, "finish_reason": finish_reason or "stop"}]} - def _flush_pending_stream_text(): pending_parts = list(pending_text_parts) pending_text_parts.clear() @@ -2946,8 +2964,8 @@ class _StreamingCall: from agent import relay_llm stream = self._set_managed_stream(relay_llm.stream(self.api_kwargs, _open_stream, - **_relay_stream_identity(self.agent, "provider"), finalizer=_relay_final_response, - on_stream_created=self._chat_stream_created, + **_relay_stream_identity(self.agent, "provider"), finalizer=relay_response.finalize, + on_stream_created=self._chat_stream_created, on_chunk=relay_response.observe, accept_chunk=lambda chunk: self._accept_chat_chunk(stream_attempt_id, chunk), completed_response_predicate=lambda value: hasattr(value, "choices"), metadata=_relay_stream_metadata(self.agent, "chat_completions"), defer_logical_completion=True)) diff --git a/agent/chat_completion_helpers_relay.py b/agent/chat_completion_helpers_relay.py new file mode 100644 index 0000000000..37dd8b850f --- /dev/null +++ b/agent/chat_completion_helpers_relay.py @@ -0,0 +1,74 @@ +"""Relay-side accumulator for the chat_completions streaming wire. + +Relay invokes its collector for every post-intercept chunk and then its finalizer as soon +as the provider stream ends — concurrently with Hermes' consumer thread, which may not have +read the last chunk yet. The finalizer therefore builds Relay's recorded response from +collector-observed state only, never from the consumer loop's closures. Sibling of +``relay_llm.AnthropicStreamAccumulator``; Bedrock and Codex follow the same contract. +""" + +from __future__ import annotations + +from types import SimpleNamespace +from typing import Any + +from agent.chat_completion_helpers import _ToolCallAccumulator +from agent.message_content import flatten_message_text +from agent.reasoning_summaries import separate_glued_reasoning_blocks + + +def _tool_call_delta_view(tc_delta: Any) -> Any: + """Attribute view of a JSON tool-call delta for ``_ToolCallAccumulator.feed`` (written + against SDK objects). Only ``function`` is wrapped: ``feed`` passes ``extra_content`` + (a dict) straight through ``_dump_if_model``, so a recursive view would corrupt it.""" + if not isinstance(tc_delta, dict): + return tc_delta + function = tc_delta.get("function") + return SimpleNamespace(**{**tc_delta, + "function": SimpleNamespace(**function) if isinstance(function, dict) else function}) + + +class RelayChatAccumulator: + """Rebuild a chat.completion from Relay's post-intercept chunk dicts.""" + + def __init__(self) -> None: + self._content: list[str] = [] + self._reasoning: list[str] = [] + self._tool_calls = _ToolCallAccumulator() + self._model = self._usage = self._finish_reason = None + self._role = "assistant" + + def observe(self, chunk: Any) -> None: + if not isinstance(chunk, dict): + return + self._model = chunk.get("model") or self._model + if chunk.get("usage"): + self._usage = chunk["usage"] + choices = chunk.get("choices") or [] + choice = choices[0] if choices else None # Hermes never requests n>1 + if not isinstance(choice, dict): + return + self._finish_reason = choice.get("finish_reason") or self._finish_reason + delta = choice.get("delta") + if not isinstance(delta, dict): + return + if delta.get("role"): + self._role = delta["role"] + text = flatten_message_text(delta.get("content"), sep="") + if text: + self._content.append(text) + reasoning = delta.get("reasoning_content") or delta.get("reasoning") + if reasoning: + self._reasoning.append(separate_glued_reasoning_blocks( + self._reasoning[-1] if self._reasoning else "", reasoning)) + for tc_delta in delta.get("tool_calls") or []: + self._tool_calls.feed(_tool_call_delta_view(tc_delta)) + + def finalize(self) -> dict[str, Any]: + acc = self._tool_calls.materialize() + message = {"role": self._role, "content": "".join(self._content) or None, + "reasoning_content": "".join(self._reasoning) or None, + "tool_calls": [acc[i] for i in sorted(acc)] or None} + # "stop" also covers Nous Portal ``lastOne`` usage frames, which carry no finish_reason. + return {"model": self._model, "usage": self._usage, + "choices": [{"message": message, "finish_reason": self._finish_reason or "stop"}]} diff --git a/agent/codex_responses_adapter.py b/agent/codex_responses_adapter.py index 5f4395b268..bc079c2d32 100644 --- a/agent/codex_responses_adapter.py +++ b/agent/codex_responses_adapter.py @@ -541,16 +541,10 @@ def classify_responses_route(agent: Any) -> ResponsesRouteFlags: ) -def estimate_native_responses_preflight_tokens( - agent: Any, messages: List[Dict[str, Any]], *, system_prompt: str = "", tools: Optional[List[Dict[str, Any]]] = None, -) -> Optional[int]: - """Estimate tokens for the checkpoint-pruned Responses payload (the full transcript overstates a natively compacted - session and fires local compression needlessly). None when native compaction is not proven eligible or conversion fails. - - Automatic preflight previously counted the full durable transcript. On a natively compacted Codex - session that overstates the wire by several times and fires local compression against history the main - request will never send (#96155). - """ +def _native_responses_replay_items( + agent: Any, messages: List[Dict[str, Any]] +) -> Optional[List[Dict[str, Any]]]: + """Build the native-compaction-eligible wire items, or ``None`` when ineligible.""" if getattr(agent, "api_mode", None) != "codex_responses" or not isinstance(messages, list): return None route = classify_responses_route(agent)._asdict() @@ -565,7 +559,37 @@ def estimate_native_responses_preflight_tokens( native_compaction_eligible=True, ) except Exception: - logger.debug("native Responses preflight conversion failed; falling back to generic estimate", exc_info=True) + logger.debug( + "native Responses replay conversion failed; using the generic fallback", + exc_info=True, + ) + return None + return items + + +def has_replayable_native_compaction_checkpoint( + agent: Any, messages: List[Dict[str, Any]] +) -> bool: + """Whether the current route would replay a persisted native checkpoint.""" + items = _native_responses_replay_items(agent, messages) + if items is None: + return False + from agent.native_compaction import has_compaction_checkpoint + return has_compaction_checkpoint(items) + + +def estimate_native_responses_preflight_tokens( + agent: Any, messages: List[Dict[str, Any]], *, system_prompt: str = "", tools: Optional[List[Dict[str, Any]]] = None, +) -> Optional[int]: + """Estimate tokens for the checkpoint-pruned Responses payload (the full transcript overstates a natively compacted + session and fires local compression needlessly). None when native compaction is not proven eligible or conversion fails. + + Automatic preflight previously counted the full durable transcript. On a natively compacted Codex + session that overstates the wire by several times and fires local compression against history the main + request will never send (#96155). + """ + items = _native_responses_replay_items(agent, messages) + if items is None: return None from agent.model_metadata import estimate_request_tokens_rough return estimate_request_tokens_rough(items, system_prompt=system_prompt or "", tools=tools) diff --git a/agent/context_compressor.py b/agent/context_compressor.py index 33838ca390..e9242fa371 100644 --- a/agent/context_compressor.py +++ b/agent/context_compressor.py @@ -2416,6 +2416,20 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine): except (TypeError, ValueError): self._pending_request_rough_tokens = 0 + def note_native_compaction_checkpoint(self) -> None: + """Wait for real usage before trusting a newly checkpointed request. + + Native Responses compaction replaces durable history with an opaque + encrypted checkpoint. Its serialized size is unrelated to the token + count billed by the provider, so the first rough estimate after capture + can jump by more than the whole context window. Reuse the one-response + compaction latch and discard any stale local-compression baseline; the + next provider response then pairs its real usage with the rough estimate + for the checkpointed request. + """ + self.awaiting_real_usage_after_compression = True + self.last_compression_rough_tokens = 0 + def should_defer_preflight_to_real_usage(self, rough_tokens: int) -> bool: """Return True when a high rough preflight estimate is known-noisy. Projects real usage as ``last_real + (rough_now - rough_at_last_real)`` and fires only when the @@ -2426,7 +2440,8 @@ class ContextCompressor(MicroCompactionMixin, ContextEngine): re-runs with the aligned basis.""" if rough_tokens < self.threshold_tokens: return False - # After compaction last_real_prompt_tokens is STALE (above threshold); defer one turn until real usage arrives. + # After local or native compaction, last_real_prompt_tokens is STALE + # (above threshold); defer one turn until real usage arrives. if self.awaiting_real_usage_after_compression: return True if self.last_real_prompt_tokens <= 0 or self.last_real_prompt_tokens >= self.threshold_tokens: diff --git a/agent/context_references.py b/agent/context_references.py index 6fed91b751..ed27d43b6c 100644 --- a/agent/context_references.py +++ b/agent/context_references.py @@ -195,7 +195,13 @@ async def preprocess_context_references_async( # Expand concurrently (each ref is independent; several @url: refs would otherwise # serialize web_extract round-trips). gather preserves order, so warnings/blocks # are assembled in ref order; the token-budget check runs once afterwards. - tasks = (_expand_reference(ref, cwd_path, url_fetcher=url_fetcher, allowed_root=allowed_root_path) for ref in refs) + hard_limit = max(1, int(context_length * 0.50)) + soft_limit = max(1, int(context_length * 0.25)) + tasks = ( + _expand_reference(ref, cwd_path, url_fetcher=url_fetcher, allowed_root=allowed_root_path, + max_inline_tokens=hard_limit) + for ref in refs + ) expanded = await asyncio.gather(*tasks) warnings = [warning for warning, _ in expanded if warning] blocks = [block for _, block in expanded if block] @@ -204,8 +210,6 @@ async def preprocess_context_references_async( message=message, original_message=message, references=refs, warnings=warnings, injected_tokens=injected_tokens ) - hard_limit = max(1, int(context_length * 0.50)) - soft_limit = max(1, int(context_length * 0.25)) if injected_tokens > hard_limit: warnings.append(f"@ context injection refused: {injected_tokens} tokens exceeds the 50% hard limit ({hard_limit}).") result.blocked = True @@ -235,11 +239,12 @@ _GIT_REFERENCE_ARGS: dict[str, Callable[[ContextReference], list[str]]] = { async def _expand_reference( - ref: ContextReference, cwd: Path, *, url_fetcher: UrlFetcher = None, allowed_root: Path | None = None + ref: ContextReference, cwd: Path, *, url_fetcher: UrlFetcher = None, allowed_root: Path | None = None, + max_inline_tokens: int | None = None, ) -> Expansion: try: if ref.kind in ("file", "folder"): - return _expand_path_reference(ref, cwd, allowed_root=allowed_root) + return _expand_path_reference(ref, cwd, allowed_root=allowed_root, max_inline_tokens=max_inline_tokens) if ref.kind in _GIT_REFERENCE_ARGS: git_args = _GIT_REFERENCE_ARGS[ref.kind](ref) return _expand_git_reference(ref, cwd, git_args, "git " + " ".join(git_args)) @@ -261,7 +266,8 @@ async def _expand_reference( return f"{ref.raw}: unsupported reference type", None -def _expand_path_reference(ref: ContextReference, cwd: Path, *, allowed_root: Path | None = None) -> Expansion: +def _expand_path_reference(ref: ContextReference, cwd: Path, *, allowed_root: Path | None = None, + max_inline_tokens: int | None = None) -> Expansion: """``@file:`` / ``@folder:``: resolve, allow-check, then inline text / binary stub / listing.""" is_folder = ref.kind == "folder" path = _resolve_path(cwd, ref.target, allowed_root=allowed_root) @@ -281,7 +287,16 @@ def _expand_path_reference(ref: ContextReference, cwd: Path, *, allowed_root: Pa if ref.line_start is not None: text = "\n".join(text.splitlines()[max(ref.line_start - 1, 0):ref.line_end or ref.line_start]) lang = _FENCE_LANGUAGES.get(path.suffix.lower(), "") - return None, f"📄 {ref.raw} ({estimate_tokens_rough(text)} tokens)\n```{lang}\n{text}\n```" + text_tokens = estimate_tokens_rough(text) + # Check BEFORE building the fenced block: an oversized file is not going to be + # inlined, so don't build a second MB-scale string just to discard it. + if max_inline_tokens is not None and text_tokens > max_inline_tokens: + # One oversized file used to poison the aggregate check and refuse the whole + # turn (#61987); the file stays readable via the agent's tools instead. The + # block alone carries the message (same shape as the binary path) — a warning + # would duplicate it in "--- Context Warnings ---". + return None, _oversized_text_reference_block(ref, path, text_tokens) + return None, f"📄 {ref.raw} ({text_tokens} tokens)\n```{lang}\n{text}\n```" def _run_quiet(cmd: list[str], cwd: Path, timeout: int, env: dict | None = None) -> subprocess.CompletedProcess: @@ -425,12 +440,7 @@ def _iter_visible_entries(path: Path, cwd: Path, limit: int) -> list[Path]: return output -def _binary_reference_block(ref: ContextReference, path: Path) -> str: - mime = mimetypes.guess_type(path.name)[0] or "application/octet-stream" - try: - size = format_bytes(path.stat().st_size) - except OSError: - size = "unknown size" +def _agent_visible_path(path: Path) -> str: # Under a container backend the host path dangles inside the sandbox: translate staged # files to their auto-mounted cache path; fall back to the host path (local backend / # translation failure). Run the idempotent TERMINAL_ENV bridge first so in-process @@ -439,14 +449,42 @@ def _binary_reference_block(ref: ContextReference, path: Path) -> str: from tools.terminal_tool import _ensure_terminal_env_bridged _ensure_terminal_env_bridged() from tools.credential_files import to_agent_visible_cache_path - visible = to_agent_visible_cache_path(str(path)) + return to_agent_visible_cache_path(str(path)) except Exception: - visible = str(path) + return str(path) + + +def _on_disk_reference_block(ref: ContextReference, path: Path, descriptor: str, reason: str, guidance: str) -> str: + """Shared 📎 shape: the file was not inlined, but it IS on disk where the agent's + tools run — hand the model the path and a nudge instead of a dead-end warning.""" + try: + size = format_bytes(path.stat().st_size) + except OSError: + size = "unknown size" return ( - f"📎 {ref.raw} ({mime}, {size}) — binary file, not inlined as text. " - f"It is available on disk at `{visible}`. Use your tools to work with it " - f"(read or convert it, extract its text, or view/render it as needed); " - f"do not tell the user the file type is unsupported." + f"📎 {ref.raw} ({descriptor}, {size}) — {reason} " + f"It is available on disk at `{_agent_visible_path(path)}`. {guidance}" + ) + + +def _binary_reference_block(ref: ContextReference, path: Path) -> str: + mime = mimetypes.guess_type(path.name)[0] or "application/octet-stream" + return _on_disk_reference_block( + ref, path, + descriptor=mime, + reason="binary file, not inlined as text.", + guidance="Use your tools to work with it (read or convert it, extract its text, " + "or view/render it as needed); do not tell the user the file type is unsupported.", + ) + + +def _oversized_text_reference_block(ref: ContextReference, path: Path, text_tokens: int) -> str: + return _on_disk_reference_block( + ref, path, + descriptor=f"text file, approximately {text_tokens} tokens", + reason="too large to inline safely.", + guidance="Use read_file with a narrow line range, search_files, or terminal/code tools " + "to inspect only the relevant parts; do not load the entire file into context.", ) diff --git a/agent/credential_pool.py b/agent/credential_pool.py index 4286903c0e..1f7d5a07cb 100644 --- a/agent/credential_pool.py +++ b/agent/credential_pool.py @@ -791,7 +791,10 @@ def _borrowed_single_use_pool_root() -> Optional[Path]: return None -def _update_root_pool_rows(provider: str, payloads: List[Dict[str, Any]], global_path: Path) -> None: +def _update_root_pool_rows( + provider: str, payloads: List[Dict[str, Any]], global_path: Path, + *, status_cleared_ids: Optional[Iterable[str]] = None, +) -> None: """UPDATE-ONLY merge of *payloads* into the root store's rows for *provider*. A borrower may refresh the root's rows (rotation, cooldown state) but @@ -808,6 +811,7 @@ def _update_root_pool_rows(provider: str, payloads: List[Dict[str, Any]], global existing = pool.get(provider) existing_list = existing if isinstance(existing, list) else [] incoming_by_id = {p.get("id"): p for p in payloads if isinstance(p, dict) and p.get("id")} + cleared = {cid for cid in (status_cleared_ids or ()) if cid} merged: List[Dict[str, Any]] = [] changed = False for disk_entry in existing_list: @@ -816,7 +820,10 @@ def _update_root_pool_rows(provider: str, payloads: List[Dict[str, Any]], global if incoming is None: merged.append(disk_entry) continue - updated = auth_mod._merge_disk_cooldown_state(incoming, disk_entry, provider) + # A deliberately cleared entry has no disk cooldown worth keeping. + updated = auth_mod._merge_disk_cooldown_state( + incoming, None if did in cleared else disk_entry, provider, + ) if updated != disk_entry: changed = True merged.append(updated) @@ -830,6 +837,7 @@ def persist_pool_entries( payloads: List[Dict[str, Any]], *, removed_ids: Optional[Iterable[str]] = None, + status_cleared_ids: Optional[Iterable[str]] = None, ) -> None: """Persist a provider's pool rows to the store that OWNS them. @@ -845,7 +853,10 @@ def persist_pool_entries( global_path = _borrowed_single_use_pool_root() if global_path is not None: try: - _update_root_pool_rows(provider, payloads, global_path) + _update_root_pool_rows( + provider, payloads, global_path, + status_cleared_ids=status_cleared_ids, + ) except Exception as exc: # Fail closed on the FORK, not on the save: never fall back to # writing a local copy (that IS the bug). The in-memory pool @@ -856,7 +867,9 @@ def persist_pool_entries( provider, exc, ) return - write_credential_pool(provider, payloads, removed_ids=removed_ids) + write_credential_pool( + provider, payloads, removed_ids=removed_ids, status_cleared_ids=status_cleared_ids, + ) # --- Per-provider singleton refresh plumbing ------------------------------- @@ -1018,13 +1031,19 @@ class CredentialPool: self._entries[idx] = new return - def _persist(self, *, removed_ids: Optional[List[str]] = None) -> None: + def _persist( + self, + *, + removed_ids: Optional[List[str]] = None, + status_cleared_ids: Optional[List[str]] = None, + ) -> None: # Self-locking: snapshotting self._entries must not race a rotation. with self._lock: persist_pool_entries( self.provider, [entry.to_dict() for entry in self._entries], removed_ids=removed_ids, + status_cleared_ids=status_cleared_ids, ) def _adopt(self, entry: PooledCredential, *, persist: bool = True, **updates: Any) -> PooledCredential: @@ -2139,14 +2158,29 @@ class CredentialPool: return refreshed def reset_statuses(self) -> int: + """Clear exhaustion state on every entry. Returns how many were cleared. + + ``failure_reason`` lives in ``extra``, not a dataclass field, so it is + stripped explicitly. The persist declares the cleared ids because the + disk-recency merge reads a cleared ``last_status_at`` (None -> epoch 0) + as a stale snapshot and would copy a still-binding cooldown back. + """ with self._lock: - stale = [e for e in self._entries if e.last_status or e.last_status_at or e.last_error_code] + stale = [ + e for e in self._entries + if e.last_status or e.last_status_at or e.last_error_code or e.failure_reason + ] if stale: stale_ids = {e.id for e in stale} self._entries = [ - replace(e, **_CLEAR_STATUS) if e.id in stale_ids else e for e in self._entries + replace( + e, **_CLEAR_STATUS, + extra={k: v for k, v in e.extra.items() if k != "failure_reason"}, + ) + if e.id in stale_ids else e + for e in self._entries ] - self._persist() + self._persist(status_cleared_ids=list(stale_ids)) return len(stale) def remove_index(self, index: int) -> Optional[PooledCredential]: diff --git a/agent/error_classifier.py b/agent/error_classifier.py index 7bf3b8716f..ad4e9ae546 100644 --- a/agent/error_classifier.py +++ b/agent/error_classifier.py @@ -113,7 +113,8 @@ _BILLING_ERROR_CODES = frozenset({ # contains an overflow phrase; rate limit is matched first so throttle wins. _RATE_LIMIT_PATTERNS = ( "rate limit", "rate_limit", "too many requests", "throttled", "requests per minute", - "tokens per minute", "requests per day", "try again in", "please retry after", "resource_exhausted", + "tokens per minute", "requests per day", "try again in", "please retry after", + "resource exhausted", "resource_exhausted", "resource-exhausted", "resourceexhausted", "rate increased too quickly", "throttlingexception", "too many concurrent requests", "servicequotaexceededexception", "throttling", ) @@ -486,6 +487,13 @@ def _provider_special_cases(c: _Ctx) -> Optional[Verdict]: # to format_error and a status-less block isn't left retryable (#18028). if any(p in msg for p in _CONTENT_POLICY_BLOCKED_PATTERNS): return _V_CONTENT_BLOCKED + # ChatGPT Codex masks a rejected encrypted-reasoning replay behind the same bare + # ``invalid_prompt: Request blocked.`` it uses for real blocks (#92353). Exact envelope + # + provider only. The verdict keeps format_error's abort-and-fallback hints; the one + # extra thing it buys is turn_recovery's replay strip, which still requires cached + # ``codex_reasoning_items`` — a genuine block with nothing to strip behaves as before. + if _is_codex_masked_replay_rejection(c): + return _v(_R.invalid_encrypted_content, **_ABORT_FALLBACK) # Anthropic thinking-block 400s (signature mismatch after transcript # mutation). Not gated on provider — OpenRouter proxies Anthropic errors. if status == 400 and "thinking" in msg and any(p in msg for p in _THINKING_MUTATION_WORDS): @@ -777,6 +785,22 @@ def _is_server_injected_param_rejection(error_msg: str, provider: str) -> bool: return False +_CODEX_MASKED_REPLAY_MESSAGE = "request blocked." + + +def _is_codex_masked_replay_rejection(c: "_Ctx") -> bool: + """HTTP 400 / status-less ``{code: invalid_prompt, message: "Request blocked."}`` from + ``openai-codex`` — as an SDK error body, a Responses ``error`` SSE frame, or the + ``response.failed`` text ``"invalid_prompt: Request blocked."``.""" + if c.provider_slug != "openai-codex" or c.status_code not in (None, 400): + return False + # The OpenAI SDK unwraps ``body["error"]`` on status errors; stream frames keep the envelope. + body_msg = next((str(m).strip().lower() for m in _body_message_candidates(c.body or {}) if m), "") + return (c.code == "invalid_prompt" and body_msg == _CODEX_MASKED_REPLAY_MESSAGE) or ( + c.msg.strip() == f"invalid_prompt: {_CODEX_MASKED_REPLAY_MESSAGE}" + ) + + def _error_obj(body: Any) -> dict: """``body["error"]`` when it is a dict, else ``{}``.""" err = body.get("error") if isinstance(body, dict) else None diff --git a/agent/interrupt_control.py b/agent/interrupt_control.py index e06f3ac13a..f99aca9dc1 100644 --- a/agent/interrupt_control.py +++ b/agent/interrupt_control.py @@ -9,6 +9,7 @@ import threading from typing import Optional from agent.interrupt_compat import request_hard_interrupt +from tools.interrupt import request_yield as _request_yield from tools.interrupt import set_interrupt as _set_interrupt # Same logger name as the origin module so log records / caplog filters are unchanged. @@ -245,8 +246,20 @@ class InterruptControlMixin: return False # Never kill a tool to deliver guidance; the steer drain puts it on the final tool result. + # A foreground terminal command would park that delivery until it exits (a 5-minute + # `sleep` poller, a build), so ask the tool workers to YIELD: terminal hands the live + # process to the background registry and returns; tools that don't yield are unaffected. if getattr(self, "_executing_tools", False): - return self.steer(cleaned) + accepted = self.steer(cleaned) + if accepted: + tracker = getattr(self, "_tool_worker_threads", None) + tracker_lock = getattr(self, "_tool_worker_threads_lock", None) + if tracker is not None and tracker_lock is not None: + with tracker_lock: + worker_tids = list(tracker) + for tid in worker_tids: + _request_yield(tid) + return accepted _model_active = getattr(self, "_model_request_active", None) with _ic_lock(self, "_pending_redirect_lock"): diff --git a/agent/model_metadata.py b/agent/model_metadata.py index 0900637036..6e0780c9b4 100644 --- a/agent/model_metadata.py +++ b/agent/model_metadata.py @@ -336,6 +336,7 @@ DEFAULT_CONTEXT_LENGTHS = { # OpenAI — direct-API windows (Codex OAuth caps gpt-5.4+/5.5/5.6 at 272K, resolved by # its own branch). 5.4-nano/-mini are 400k, not 1.05M; gpt-5.3-codex-spark is # Codex-OAuth-only and listed so "gpt-5" (400k) doesn't win. + "gpt-6-astra": 1050000, # also matches -pro (verified live on OpenRouter) "gpt-5.6-luna": 1050000, "gpt-5.6-terra": 1050000, "gpt-5.6-sol": 1050000, "gpt-5.5": 1050000, "gpt-5.4-nano": 400000, "gpt-5.4-mini": 400000, "gpt-5.4": 1050000, "gpt-5.3-codex-spark": 128000, "gpt-5.1-chat": 128000, "gpt-5": 400000, diff --git a/agent/native_compaction.py b/agent/native_compaction.py index 0202770c4d..8707cc91af 100644 --- a/agent/native_compaction.py +++ b/agent/native_compaction.py @@ -328,7 +328,12 @@ def is_native_compaction_rejection(error: Any, status_code: Any = None) -> bool: def has_compaction_checkpoint(items: Any) -> bool: """Does this ``codex_reasoning_items`` sidecar carry a compaction checkpoint? A checkpoint is cumulative context living in exactly one place: rewrite/discard the sidecar only after asking.""" - return isinstance(items, list) and any(_is_compaction_item(item) for item in items) + return isinstance(items, list) and any( + _is_compaction_item(item) + and isinstance(item.get("encrypted_content"), str) + and bool(item["encrypted_content"].strip()) + for item in items + ) def merge_interim_reasoning_items(prior_items: Any, new_items: Any) -> List[Dict[str, Any]]: diff --git a/agent/nous_wire.py b/agent/nous_wire.py new file mode 100644 index 0000000000..baf602abde --- /dev/null +++ b/agent/nous_wire.py @@ -0,0 +1,112 @@ +"""Nous Portal ``anthropic/*`` wire selection when ``nous.anthropic_wire`` is ``auto``. + +Portal serves Claude two ways and Hermes cannot tell which from the request: an OpenRouter +passthrough (today, for every ``anthropic/*`` id) or GMI/Vertex (planned once GMI is back). The +native Messages wire is the better transport, but on the OpenRouter path it re-writes the previous +turn's prompt cache on 14-20% of consecutive calls in concurrent tool loops (measured 2026-09-06; +NousResearch/api#227), so the session must ride chat/completions there. On GMI that is untested, +and until it is measured ``auto`` never promotes to native. + +The upstream IS visible in the first RESPONSE: OpenRouter stamps ``provider`` (chat wire) and +mints ``gen--`` ids; GMI/Vertex responses carry neither. So ``auto`` starts every +session on chat (safe on both upstreams), reads the first response, and switches the session to +native only when the upstream is GMI and native has been cleared for GMI. One decision per +session, at call 1, before there is a cache to lose; later calls never flip. + +``classify_upstream`` is pure and unit-tested; ``maybe_switch_wire_after_first_response`` is +the single hook, called from the usage recorder. +""" +from __future__ import annotations + +import logging +import re +from typing import Any, Optional + +logger = logging.getLogger(__name__) + +# Flip to True only after the 20x6 concurrency probe (evals/postmortem/live_ab) is clean on a +# GMI-served anthropic/* id on the native wire. Until then ``auto`` is chat everywhere. +GMI_NATIVE_WIRE_CLEARED = False + +_OPENROUTER_ID = re.compile(r"^gen-\d{9,}-[A-Za-z0-9_-]{8,}$") + + +def classify_upstream(response: Any) -> Optional[str]: + """``"openrouter"`` / ``"gmi"`` / ``None`` (unknown) from a Portal response object. + + Works on both wires: the OpenAI SDK object exposes ``.provider`` (OpenRouter's upstream name, + e.g. ``"Anthropic"``, ``"Amazon Bedrock"``) and an OpenRouter-minted ``.id``; the Anthropic SDK + object has ``.id`` only. GMI/Vertex responses have Anthropic-native ``msg_…`` ids and no + ``provider``. Anything else is unknown, and unknown never triggers a switch. + """ + if response is None: + return None + if isinstance(getattr(response, "provider", None), str) and getattr(response, "provider"): + return "openrouter" + rid = getattr(response, "id", None) + if isinstance(rid, str): + if _OPENROUTER_ID.match(rid): + return "openrouter" + if rid.startswith("msg_"): + return "gmi" + return None + + +def wire_for_upstream(upstream: Optional[str]) -> str: + """The api_mode ``auto`` wants once the upstream is known. Chat unless GMI and cleared.""" + if upstream == "gmi" and GMI_NATIVE_WIRE_CLEARED: + return "anthropic_messages" + return "chat_completions" + + +def maybe_switch_wire_after_first_response(agent: Any, response: Any, api_call_count: int) -> bool: + """Decide the session's wire from its first response; the switch itself is applied at the + start of the next iteration (``apply_pending_wire_switch``), never while a response is being + consumed. Returns True when a switch was scheduled. + + Only for provider=nous, anthropic/* models, ``nous.anthropic_wire: auto``, and only on the + session's first API call. + """ + if api_call_count != 1 or getattr(agent, "_nous_wire_decided", False): + return False + if (getattr(agent, "provider", "") or "").lower() != "nous": + return False + model = str(getattr(agent, "model", "") or "") + if not model.lower().startswith("anthropic/"): + return False + try: + from hermes_cli.providers import _nous_anthropic_wire + if _nous_anthropic_wire() != "auto": + return False + except Exception: + return False + agent._nous_wire_decided = True # one decision per session, whatever it is + upstream = classify_upstream(response) + want = wire_for_upstream(upstream) + if want == getattr(agent, "api_mode", None): + logger.debug("nous wire auto: upstream=%s, staying on %s", upstream, want) + return False + agent._nous_wire_pending = (want, upstream) + return True + + +def apply_pending_wire_switch(agent: Any) -> bool: + """At iteration start, with no response in flight: perform the switch scheduled by + ``maybe_switch_wire_after_first_response``. Reuses ``switch_model`` (same model/provider, new + api_mode) so client rebuild, cache policy and ``_primary_runtime`` stay consistent. A failure + is logged and the session stays on its current wire.""" + pending = getattr(agent, "_nous_wire_pending", None) + if not pending: + return False + agent._nous_wire_pending = None + want, upstream = pending + try: + from agent.agent_runtime_helpers import switch_model + switch_model(agent, agent.model, "nous", api_key=getattr(agent, "api_key", "") or "", + base_url=getattr(agent, "base_url", "") or "", api_mode=want) + except Exception as exc: # never let wire selection break a turn + logger.warning("nous wire auto: switch to %s failed (%s); staying on %s", want, exc, agent.api_mode) + return False + logger.info("nous wire auto: upstream=%s -> %s for the rest of session %s", upstream, want, + getattr(agent, "session_id", "?")) + return True diff --git a/agent/prompt_builder.py b/agent/prompt_builder.py index 6cf5d70dd0..1907e93167 100644 --- a/agent/prompt_builder.py +++ b/agent/prompt_builder.py @@ -98,7 +98,15 @@ def _scan_context_content(content: str, filename: str) -> str: def _find_git_root(start: Path) -> Optional[Path]: """Nearest ancestor (or *start* itself) containing ``.git``, else None.""" current = start.resolve() - return next((p for p in (current, *current.parents) if (p / ".git").exists()), None) + # A parent the process may not stat (locked-down /home on shared hosts) is "no .git here", not a crash. + return next((p for p in (current, *current.parents) if _exists_or_denied(p / ".git")), None) + + +def _exists_or_denied(path: Path) -> bool: + try: + return path.exists() + except OSError: + return False def _find_hermes_md(cwd: Path) -> Optional[Path]: @@ -659,7 +667,12 @@ PLATFORM_HINTS = { "height live, width from the content's first measured span — lay content flush left with no centering wrappers " "or it measures full-bleed. Widgets talk back: data-hermes-send=\"prompt\" on any clickable element (or " "window.hermes.send(\"prompt\")) sends that prompt as a hidden user turn — answer it by updating the widget's " - "file, not with prose." + "file, not with prose. Property/rental listings render as browsable cards: emit a ```listing fence " + "holding JSON — one object, or an array to compare several — with address (required), price, beds, " + "baths, size, note (why it is worth a look), facts[] (short specs), catches[] (risks to verify), " + "images[] (direct https photo URLs, in listing order — the first is the hero), and links[] " + "({label, url} detail pages, never a search-results URL). Use it for every property you present, " + "including follow-ups and re-rankings, so listings stay comparable." ), "sms": ( "You are communicating via SMS. Keep responses concise and use plain text only — no markdown, no " @@ -842,13 +855,6 @@ def _tenv_read(name: str, default: str = "") -> str: _BACKEND_IMAGE_KEYS = {b: f"{b}_image" for b in ("docker", "singularity", "modal", "daytona")} # (config key, default) pairs forwarded to _create_environment's container_config. -_CONTAINER_CONFIG_DEFAULTS = ( - ("container_cpu", 1), ("container_memory", 5120), ("container_disk", 51200), ("container_persistent", True), - ("modal_mode", "auto"), ("docker_volumes", []), ("docker_mount_cwd_to_workspace", False), - ("docker_forward_env", []), ("docker_env", {}), ("docker_run_as_host_user", False), ("docker_extra_args", []), - ("docker_shm_size", "1g"), ("docker_persist_across_processes", True), ("docker_shared_container_key", ""), - ("docker_orphan_reaper", True), -) # Single-line POSIX probe; `2>/dev/null` keeps a missing binary from polluting output. _BACKEND_PROBE_CMD = ( "printf 'os=%s\\nkernel=%s\\nhome=%s\\ncwd=%s\\nuser=%s\\n' \"$(uname -s 2>/dev/null || echo unknown)\" " @@ -859,31 +865,33 @@ _BACKEND_PROBE_CMD = ( def _run_backend_probe(env_type: str, terminal_tool) -> str: """Execute the probe command inside a freshly built backend; "" when it yields nothing.""" - from tools.terminal_tool_backends import _create_environment, _ssh_config_from_config + from tools.terminal_tool_backends import _container_config_from_config, _create_environment, _ssh_config_from_config from tools.terminal_tool_lifecycle import _cleanup_env config = terminal_tool._get_env_config() - # Mirrors tools/terminal_tool.py's live-command assembly (`_create_environment` is the factory). + # Same container_config shaper as the live terminal path: a private copy of the key table here + # drifted (no docker_network) and gave the probe a bridge-networked container under lockdown. env = _create_environment( env_type=env_type, image=config.get(_BACKEND_IMAGE_KEYS[env_type], "") if env_type in _BACKEND_IMAGE_KEYS else "", cwd=config.get("cwd", ""), timeout=config.get("timeout", 180), ssh_config=_ssh_config_from_config(config) if env_type == "ssh" else None, - container_config=({k: config.get(k, d) for k, d in _CONTAINER_CONFIG_DEFAULTS} + container_config=(_container_config_from_config(config) if terminal_tool._is_container_backend(env_type) else None), task_id="prompt-backend-probe", host_cwd=config.get("host_cwd"), + # Only ssh honors this: an isolated ControlMaster socket and no remote dir setup / file sync / + # snapshot. A normal SSHEnvironment would upload the whole ~/.hermes tree just to run `uname`, + # and its later __del__ would sync_back() and close the master shared with the agent's own env. + probe_only=True, ) try: result = env.execute(_BACKEND_PROBE_CMD, timeout=4) finally: # One-shot `uname`; without teardown the backend leaves a second idle sandbox # (task_id="prompt-backend-probe") running for the whole process next to the agent's own. - # ssh is left alone: no task-scoped sandbox, and its cleanup() closes a ControlMaster socket - # (keyed by user@host:port) shared with the agent's real environment; ControlPersist expires it. - if env_type != "ssh": - try: - _cleanup_env(env, force_remove=True) - except Exception: - logger.debug("Backend probe cleanup failed", exc_info=True) + try: + _cleanup_env(env, force_remove=True) + except Exception: + logger.debug("Backend probe cleanup failed", exc_info=True) if result.get("returncode") != 0: logger.debug("Backend probe returned non-zero: %r", result) return "" diff --git a/agent/session_activity.py b/agent/session_activity.py index ddd04c3ab7..1ab19ef978 100644 --- a/agent/session_activity.py +++ b/agent/session_activity.py @@ -4,6 +4,7 @@ only (notification, timeout, kill and retry policy live elsewhere). Provenance i from __future__ import annotations +import sys import time from contextlib import suppress from enum import Enum @@ -45,6 +46,21 @@ def normalize_activity_provenance(provenance: Optional[ActivityProvenance | str] return ActivityProvenance.UNKNOWN +def format_iteration_progress(api_call_count: Any, max_iterations: Any) -> str: + """``iteration N/M`` for user-facing status lines, or ``iteration N`` when the cap is unbounded. + + ``AIAgent.max_iterations`` defaults to ``sys.maxsize`` (unlimited), so printing the pair verbatim + shows ``iteration 3/9223372036854775807`` in busy acks, heartbeats and timeout diagnostics (#102806). + """ + try: + cap = int(max_iterations) + except (TypeError, ValueError): + cap = sys.maxsize + if cap >= sys.maxsize: + return f"iteration {api_call_count}" + return f"iteration {api_call_count}/{cap}" + + def reset_session_activity_persist_window(agent: Any) -> None: """Clear the persist rate-limit so the next stamp writes through (terminal compression labels must not stick on mid-compress text).""" with suppress(Exception): diff --git a/agent/transports/chat_completions.py b/agent/transports/chat_completions.py index 723174f87d..8171aec797 100644 --- a/agent/transports/chat_completions.py +++ b/agent/transports/chat_completions.py @@ -146,8 +146,10 @@ def _build_gemini_thinking_config(model: str, reasoning_config: dict | None) -> return thinking_config if effort not in {"minimal", "low", "medium", "high", "xhigh", "max", "ultra"}: effort = "medium" - # Gemini 3 Flash documents low/medium/high; Gemini 3 Pro only low/high. - if normalized_model.startswith(("gemini-3", "gemini-3.1")): + # Gemini 3 Flash documents low/medium/high thinking levels; Gemini 3 Pro + # is stricter (low/high). Clamp Hermes' wider effort set to what each + # family accepts so we never forward an undocumented level verbatim. + if normalized_model.startswith("gemini-3"): if "flash" in normalized_model: thinking_config["thinkingLevel"] = ( "low" if effort in {"minimal", "low"} else "high" if effort in _HIGH_EFFORTS else "medium" diff --git a/agent/turn_context.py b/agent/turn_context.py index 82bba061a8..82bf522cb4 100644 --- a/agent/turn_context.py +++ b/agent/turn_context.py @@ -525,13 +525,37 @@ def _stage_turn_user_message( def _hydrate_from_history(agent: Any, conversation_history: Optional[List[Any]]) -> None: - """Hydrate the todo store and per-session nudge counters from persisted history.""" + """Hydrate process-local state from persisted history on the first resumed turn.""" if not conversation_history: return if not agent._todo_store.has_items(): agent._hydrate_todo_store(conversation_history) - # Hydrate per-session nudge counters from persisted history. + # A live native checkpoint arms this latch while its response is captured. A + # restarted agent must recover the same one-response deferral before turn-start + # compression can rewrite the restored opaque checkpoint. Reuse the adapter's + # exact route/issuer/replay filtering and tolerate plugin compressors without the + # optional hook. if agent._user_turn_count == 0: + note_checkpoint = getattr( + getattr(agent, "context_compressor", None), + "note_native_compaction_checkpoint", + None, + ) + if callable(note_checkpoint): + try: + from agent.codex_responses_adapter import ( + has_replayable_native_compaction_checkpoint, + ) + + if has_replayable_native_compaction_checkpoint( + agent, conversation_history + ): + note_checkpoint() + except Exception: + logger.debug( + "restored native checkpoint hydration skipped", exc_info=True + ) + # Hydrate per-session nudge counters from persisted history. prior_user_turns = sum(1 for m in conversation_history if m.get("role") == "user") if prior_user_turns > 0: agent._user_turn_count = prior_user_turns diff --git a/agent/turn_context_compaction.py b/agent/turn_context_compaction.py index a9f8a9b703..171f602048 100644 --- a/agent/turn_context_compaction.py +++ b/agent/turn_context_compaction.py @@ -156,6 +156,11 @@ def _idle_compaction( if _idle_gap < _idle_after: return _compressor = agent.context_compressor + # A live or restored native checkpoint must reach its issuer once so real usage, + # rather than an opaque ciphertext estimate, decides whether local compression is + # still needed. Threshold and post-tool preflight honor the same latch. + if bool(getattr(_compressor, "awaiting_real_usage_after_compression", False)): + return # Route-aware pressure: on compacted native-Codex sessions the durable figure # overstates the wire, so reuse the preflight estimator. _idle_tokens = _tc._preflight_request_tokens( diff --git a/agent/turn_iteration_prep.py b/agent/turn_iteration_prep.py index 1b4917758c..96ac96f17c 100644 --- a/agent/turn_iteration_prep.py +++ b/agent/turn_iteration_prep.py @@ -39,6 +39,12 @@ def prepare_iteration(agent: Any,*, messages: Any, api_call_count: Any) -> Itera _INTERRUPT_SCAFFOLD_MARKER, _maybe_inject_run_budget_wrapup ) + # nous.anthropic_wire=auto: a wire switch decided from the previous response lands here, + # before this iteration's request is built and with nothing in flight. + if getattr(agent, "_nous_wire_pending", None): + from agent.nous_wire import apply_pending_wire_switch + apply_pending_wire_switch(agent) + # Fire step_callback for gateway hooks (agent:step event). if agent.step_callback is not None: try: diff --git a/agent/turn_preflight.py b/agent/turn_preflight.py index cc2a5ff011..15d661e53a 100644 --- a/agent/turn_preflight.py +++ b/agent/turn_preflight.py @@ -297,6 +297,9 @@ def compress_after_tool_results( if ( agent.compression_enabled and compression_attempts < max_compression_attempts + and not bool( + getattr(_compressor, "awaiting_real_usage_after_compression", False) + ) and _compressor.should_compress(_real_tokens) ): compression_attempts += 1 diff --git a/agent/turn_recovery.py b/agent/turn_recovery.py index 6fd3fbb25b..fbc9b0731e 100644 --- a/agent/turn_recovery.py +++ b/agent/turn_recovery.py @@ -892,7 +892,8 @@ def log_api_error_attempt( if agent._is_openrouter_url() and "support tool use" in error_msg: _blines(agent, f" 💡 No OpenRouter providers for {_model} support tool calling with your current settings.") - if agent.providers_allowed: + from agent.chat_completion_helpers import _provider_preferences_for_agent + if _provider_preferences_for_agent(agent).get("only"): _blines( agent, " Your provider_routing.only restriction is filtering out tool-capable providers.", @@ -973,30 +974,44 @@ def compute_error_backoff( agent: Any, api_error: Exception, *, retry_count: int, max_retries: int, is_rate_limited: bool, is_zai_coding_overload: bool, base_url: Any, model: Any, ) -> float: - """Pick the wait before the next API retry and announce it. Retry-After wins for rate - limits (capped at 600s: Anthropic Tier 1 buckets reset in ~171s, so a 120s cap re-tripped - the limit); otherwise jittered backoff, replaced by the adaptive policy for 429s / Z.AI - overloads. Normal retries are buffered; long Z.AI Coding waits surface immediately.""" + """Pick the wait before the next API retry and announce it. Retry-After wins for + rate limits and any other retryable error (capped at 600s: Anthropic Tier 1 buckets + reset in ~171s, so a 120s cap re-tripped the limit); otherwise jittered backoff, + replaced by the adaptive policy for 429s / Z.AI overloads. Normal retries are + buffered; long Z.AI Coding waits surface immediately.""" # Imported lazily so tests that patch ``agent.retry_utils.jittered_backoff`` / # ``adaptive_rate_limit_backoff`` (incl. the run_agent conftest fast-backoff fixture) intercept. - from agent.retry_utils import adaptive_rate_limit_backoff, jittered_backoff + from agent.retry_utils import adaptive_rate_limit_backoff, jittered_backoff, parse_retry_after_seconds - _retry_after = None - _resp_headers = getattr(getattr(api_error, "response", None), "headers", None) if is_rate_limited else None - if _resp_headers and hasattr(_resp_headers, "get"): - _ra_raw = _resp_headers.get("retry-after") or _resp_headers.get("Retry-After") - if _ra_raw: - try: - # Cap at 10 minutes. Anthropic Tier 1 input-token buckets reset in ~171s, so a 120s cap - # caused us to retry before the actual reset window and re-trip the limit. 600s covers all - # realistic provider reset windows while still rejecting pathological values. (#26293) - _retry_after = min(float(_ra_raw), 600) - except (TypeError, ValueError): - pass - wait_time = _retry_after if _retry_after else jittered_backoff(retry_count, base_delay=2.0, max_delay=60.0) + # Respect Retry-After on every retryable provider error, not just 429s. Retryable + # 5xx responses (e.g. Cloudflare 520/524) also carry the header or a structured + # ``retry_after`` problem-detail body field; ignoring either turns an origin + # outage into a retry storm. + _retry_after = parse_retry_after_seconds( + getattr(getattr(api_error, "response", None), "headers", None) + ) + if _retry_after is None: + _error_body = getattr(api_error, "body", None) + if isinstance(_error_body, dict): + # Some providers nest it as error.retry_after (the same unwrap + # extract_api_error_context uses), others put it at the top level. + _nested = _error_body.get("error") + _payload = _nested if isinstance(_nested, dict) else _error_body + _retry_after = parse_retry_after_seconds(_payload.get("retry_after")) + if _retry_after is not None: + # Cap at 10 minutes. Anthropic Tier 1 input-token buckets reset in ~171s, so a 120s cap + # caused us to retry before the actual reset window and re-trip the limit. 600s covers all + # realistic provider reset windows while still rejecting pathological values. (#26293) + _retry_after = min(_retry_after, 600) + if _retry_after <= 0: + # A zero/expired cooldown (retry-after: 0, or an HTTP-date in the + # past, which the parser clamps to 0.0) carries no usable wait — + # treat it as absent so we never hot-loop the provider. + _retry_after = None + wait_time = _retry_after if _retry_after is not None else jittered_backoff(retry_count, base_delay=2.0, max_delay=60.0) _backoff_policy = None _adaptive = is_rate_limited or is_zai_coding_overload - if _adaptive and not _retry_after: + if _adaptive and _retry_after is None: wait_time, _backoff_policy = adaptive_rate_limit_backoff( retry_count, base_url=str(base_url), model=model, error=api_error, default_wait=wait_time, ) @@ -1009,7 +1024,16 @@ def compute_error_backoff( else: agent._buffer_status(_rate_limit_status) else: - agent._buffer_status(f"⏳ Retrying in {wait_time:.1f}s (attempt {retry_count}/{max_retries})...") + _retry_status = ( + f"⏳ Retrying in {wait_time:.1f}s (attempt {retry_count}/{max_retries})..." + ) + if _retry_after is not None and _retry_after > 60: + # A 5xx Retry-After can now reach the 600s cap; buffering that wait + # would leave the user silent for minutes, so surface long provider + # cooldowns immediately (mirrors the zai_coding_overload_long path). + agent._emit_status(_retry_status) + else: + agent._buffer_status(_retry_status) logger.warning( "Retrying API call in %ss (attempt %s/%s) %s policy=%s error=%s", wait_time, retry_count, max_retries, agent._client_log_context(), diff --git a/agent/turn_usage.py b/agent/turn_usage.py index 131e2cd348..a535a6f937 100644 --- a/agent/turn_usage.py +++ b/agent/turn_usage.py @@ -181,6 +181,11 @@ def record_response_usage( prompt_tokens, completion_tokens, total_tokens, api_duration, _cache_pct, ) + # nous.anthropic_wire=auto: the session's wire is decided once, from this first response. + if agent.session_api_calls == 1 and (agent.provider or "") == "nous": + with suppress(Exception): + from agent.nous_wire import maybe_switch_wire_after_first_response + maybe_switch_wire_after_first_response(agent, response, agent.session_api_calls) # MoA: agent.model/provider are the virtual preset/"moa" with no pricing entry, silently # dropping aggregator spend. Price at the REAL model/provider from the aggregator slot. diff --git a/agent/usage_pricing.py b/agent/usage_pricing.py index 0f73ce36f9..9830e2e28b 100644 --- a/agent/usage_pricing.py +++ b/agent/usage_pricing.py @@ -191,6 +191,9 @@ _SNAPSHOTS: tuple[tuple[str, Optional[str], str, dict], ...] = ( ("deepseek-chat", "deepseek-reasoner", "deepseek-v4-flash"): ("0.14", "0.28", "0.0028"), "deepseek-v4-pro": ("0.435", "0.87", "0.003625"), }), + ("google", "https://ai.google.dev/gemini-api/docs/pricing", "google-pricing-2026-09-02", { + ("gemini-3.8-flash", "gemini-3.7-flash"): ("0.75", "3.75", "0.075"), + }), ("google", "https://ai.google.dev/gemini-api/docs/pricing", "google-pricing-2026-07-28", { "gemini-3.6-flash": ("1.50", "7.50", "0.15"), "gemini-3.5-flash-lite": ("0.30", "2.50", "0.03"), }), diff --git a/agent/verify/runner.py b/agent/verify/runner.py index eb6f746cfe..4016c1638b 100644 --- a/agent/verify/runner.py +++ b/agent/verify/runner.py @@ -182,6 +182,30 @@ def _run_start_phase( return ReadinessResult(url, ready, status, time.monotonic() - started, error, _tail(output)) +def _compose_live_state_reason(root: Path) -> str | None: + """Why ``docker compose build``/``up`` must not run at *root*, or ``None`` to proceed. + + Read-only ``docker compose ps`` probe. Only a missing docker binary proceeds -- the + build phase would fail the same way, so there is nothing to protect. A hung daemon + or a non-zero probe refuses: containers may be live and unobservable, which is + exactly the #103567 loss window. + """ + try: + result = subprocess.run( + ["docker", "compose", "ps", "--status", "running", "--format", "{{.Name}}"], + cwd=root, capture_output=True, text=True, timeout=15, stdin=subprocess.DEVNULL, + ) + except FileNotFoundError: + return None + except subprocess.TimeoutExpired: + return "docker compose ps timed out after 15s; live containers cannot be ruled out" + if result.returncode != 0: + detail = (result.stderr or result.stdout or "").strip().splitlines() + return f"docker compose ps failed (exit {result.returncode}): {detail[-1] if detail else 'no output'}" + names = [line for line in result.stdout.splitlines() if line.strip()] + return f"this compose project already has running container(s): {', '.join(names)}" if names else None + + def run_verify( root: Path, recipe: Recipe, phases: tuple[str, ...] | list[str] | None = None, phase_timeout: float = DEFAULT_PHASE_TIMEOUT, ready_timeout: float = DEFAULT_READY_TIMEOUT, @@ -189,11 +213,34 @@ def run_verify( on_output: Callable[[str], None] | None = None, ) -> VerifyResult: """Run the selected command phases sequentially, then (unless ``skip_start`` or a - phase failed) boot ``recipe.start``, poll readiness, and tear the process group down.""" + phase failed) boot ``recipe.start``, poll readiness, and tear the process group down. + + A ``compose`` recipe refuses outright when the project already has running + containers: ``docker compose build`` + ``up`` replaces them on an image-hash + change, destroying any container-local state they carry -- this has caused a + real state-loss incident (#103567). The check is best-effort and read-only + (``docker compose ps``); when it cannot run at all, verify proceeds rather than + blocking on an unrelated environment gap, matching every other recipe kind's + behavior when its own tooling is unavailable.""" root = Path(root) selected = tuple(phases) if phases else PHASE_ORDER + ("start",) result = VerifyResult(recipe_name=recipe.name) + mutating = ("build" in selected) or ("start" in selected and not skip_start) + if recipe.kind == "compose" and mutating: + reason = _compose_live_state_reason(root) + if reason: + result.phases.append(PhaseResult( + phase="build", command=recipe.build[0] if recipe.build else "docker compose build", + exit_code=1, duration=0.0, output_tail=( + f"Refusing to run: {reason}. `docker compose build` + `up` would replace " + "live containers on an image-hash change, destroying any container-local " + "state they carry. If you intend to rebuild this live deployment, run " + "`docker compose build`/`up` yourself." + ), + )) + return result + for phase in PHASE_ORDER: if phase not in selected: continue diff --git a/apps/desktop/'/var/folders/5h/qzgt02rn619fttp2d7zdxj600000gn/T/hermes-update-mutex-LMF9y5/home/.hermes-update-in-progress.mutex' b/apps/desktop/'/var/folders/5h/qzgt02rn619fttp2d7zdxj600000gn/T/hermes-update-mutex-LMF9y5/home/.hermes-update-in-progress.mutex' deleted file mode 100644 index e69de29bb2..0000000000 diff --git a/apps/desktop/'/var/folders/5h/qzgt02rn619fttp2d7zdxj600000gn/T/hermes-update-mutex-us8HZu/home/.hermes-update-in-progress.mutex' b/apps/desktop/'/var/folders/5h/qzgt02rn619fttp2d7zdxj600000gn/T/hermes-update-mutex-us8HZu/home/.hermes-update-in-progress.mutex' deleted file mode 100644 index e69de29bb2..0000000000 diff --git a/apps/desktop/electron/backend-dial-claim.test.ts b/apps/desktop/electron/backend-dial-claim.test.ts index 6aa7ce1353..30c7e8195c 100644 --- a/apps/desktop/electron/backend-dial-claim.test.ts +++ b/apps/desktop/electron/backend-dial-claim.test.ts @@ -126,10 +126,10 @@ describe('main.ts wiring for #90812', () => { it('routes the profile-scoped dial IPC through the single-owner claim', () => { const handlerStart = mainSource.indexOf("ipcMain.handle('hermes:connection', ") expect(handlerStart).toBeGreaterThan(-1) - const body = mainSource.slice(handlerStart, handlerStart + 900) + const body = mainSource.slice(handlerStart, handlerStart + 1200) expect(body).toContain('backendDialClaims.run(') - expect(body).toContain('ensureBackend(profile)') + expect(body).toContain('ensureBackend(profile, { spawnPriority })') }) it('routes the registry-scoped dial IPC through the claim keyed by backendScopeKey(connectionId, profile)', () => { @@ -137,8 +137,9 @@ describe('main.ts wiring for #90812', () => { expect(handlerStart).toBeGreaterThan(-1) const body = mainSource.slice(handlerStart, handlerStart + 1_200) - expect(body).toContain('backendDialClaims.run(backendScopeKey(id, profile)') - expect(body).toContain('ensureRegistryBackend(id, profile)') + expect(body).toContain('const scopeKey = backendScopeKey(id, profile)') + expect(body).toContain('backendDialClaims.run(scopeKey, ') + expect(body).toContain("ensureRegistryBackend(id, profile, '', { spawnPriority })") }) // The four IPC/probe surfaces below call ensureRegistryBackend()/ensureBackend() diff --git a/apps/desktop/electron/desktop-remote-route.test.ts b/apps/desktop/electron/desktop-remote-route.test.ts index 2a731b77bb..b2f35ae885 100644 --- a/apps/desktop/electron/desktop-remote-route.test.ts +++ b/apps/desktop/electron/desktop-remote-route.test.ts @@ -3,7 +3,8 @@ import assert from 'node:assert/strict' import { test } from 'vitest' import { normalizeRegistry, REGISTRY_VERSION } from './connection-registry' -import { resolveDesktopRemoteRoute } from './desktop-remote-route' +import { backendScopeKey } from './connection-registry' +import { resolveDesktopRemoteRoute, v1SshTerminalPoolKey } from './desktop-remote-route' const tokenA = { encoding: 'plain', value: 'token-a' } const tokenB = { encoding: 'plain', value: 'token-b' } @@ -153,6 +154,46 @@ test('global SSH treats an omitted port as 22 and checks the primary route', () assert.equal(route?.connectionId, 'ssh-primary') }) +test('v1 settings SSH pool key ignores registry identity tags', () => { + const route = resolveDesktopRemoteRoute({ + config: { mode: 'ssh', remote: { mode: 'ssh', host: 'box.test', user: 'hermes' } }, + profile: 'worker', + registry: registry('ssh-primary', [ + { id: 'ssh-primary', kind: 'ssh', label: 'SSH primary', host: 'box.test', user: 'hermes', port: 22 } + ]) + }) + + assert.ok(route) + assert.equal(route.kind, 'ssh') + assert.equal(route.connectionId, 'ssh-primary') + assert.equal(v1SshTerminalPoolKey(route, 'worker'), '') + assert.notEqual(v1SshTerminalPoolKey(route, 'worker'), backendScopeKey(route.connectionId, 'worker')) +}) + +test('v1 profile SSH pool key is the profile, not conn:id::profile', () => { + const ssh = { + mode: 'ssh', + host: 'box.test', + user: 'hermes', + port: 2222, + keyPath: '/keys/a', + remoteHermesPath: '/srv/hermes', + remoteProfile: 'worker' + } + + const route = resolveDesktopRemoteRoute({ + config: { mode: 'local', profiles: { worker: ssh } }, + profile: 'worker', + registry: registry('local', [{ id: 'worker-ssh', kind: 'ssh', label: 'Worker SSH', ...ssh }]) + }) + + assert.ok(route) + assert.equal(route.kind, 'ssh') + assert.equal(route.connectionId, 'worker-ssh') + assert.equal(v1SshTerminalPoolKey(route, 'worker'), 'worker') + assert.notEqual(v1SshTerminalPoolKey(route, 'worker'), backendScopeKey(route.connectionId, 'worker')) +}) + test('profile route omits identity when two registry entries match exactly', () => { const block = { mode: 'remote', url: 'https://worker.test', authMode: 'token', token: tokenA } diff --git a/apps/desktop/electron/desktop-remote-route.ts b/apps/desktop/electron/desktop-remote-route.ts index e41df9b199..56ca6b74a8 100644 --- a/apps/desktop/electron/desktop-remote-route.ts +++ b/apps/desktop/electron/desktop-remote-route.ts @@ -51,6 +51,22 @@ function withConnectionId(route: T, connectionId?: string): T return connectionId ? { ...route, connectionId } : route } +/** + * Pool key used by v1 settings/profile SSH (`sshConnections`). + * + * Bootstrap stores the tunnel under sshScopeKey(profile or null). + * resolveDesktopRemoteRoute may also stamp a registry connectionId when + * the host matches a v2 row. That id is identity, not the pool key. + * Looking up backendScopeKey(connectionId, profile) misses the live tunnel. + */ +export function v1SshTerminalPoolKey(route: { source: string }, profile?: null | string): string { + if (route.source === 'profile') { + return connectionScopeKey(profile) || '' + } + + return '' +} + /** * Select one remote route with the existing precedence and freeze any exact * registry identity before I/O. A null result means the profile resolves diff --git a/apps/desktop/electron/link-title-window.test.ts b/apps/desktop/electron/link-title-window.test.ts index f81190f47d..a460de7c38 100644 --- a/apps/desktop/electron/link-title-window.test.ts +++ b/apps/desktop/electron/link-title-window.test.ts @@ -10,13 +10,16 @@ import { } from './link-title-window' function makeFakeBrowserWindow() { - const calls = { audioMuted: [] } + const calls = { audioMuted: [], windowOpenHandlers: [] } const FakeBrowserWindow = function (options) { this.options = options this.webContents = { setAudioMuted(value) { calls.audioMuted.push(value) + }, + setWindowOpenHandler(handler) { + calls.windowOpenHandlers.push(handler) } } } @@ -45,6 +48,9 @@ test('createLinkTitleWindow mutes audio so historical links never autoplay sound assert.ok(window instanceof FakeBrowserWindow) assert.deepEqual(calls.audioMuted, [true]) + // GHSA-9f4c-93c8-jc8g: a page loaded for its title must not be able to pop a window. + assert.equal(calls.windowOpenHandlers.length, 1) + assert.deepEqual(calls.windowOpenHandlers[0]({ url: 'https://attacker.test/popup' }), { action: 'deny' }) }) test('createLinkTitleWindow still returns the window if muting throws', () => { diff --git a/apps/desktop/electron/link-title-window.ts b/apps/desktop/electron/link-title-window.ts index 226ac3ce60..1dd9ad142c 100644 --- a/apps/desktop/electron/link-title-window.ts +++ b/apps/desktop/electron/link-title-window.ts @@ -3,6 +3,8 @@ // in an offscreen window and read its title. That window loads arbitrary // user-linked pages, so it must never emit sound or trigger real downloads. +import { createWindowOpenHandler } from './window-open-policy' + export function linkTitleWindowOptions(partitionSession) { return { show: false, @@ -34,6 +36,9 @@ export function createLinkTitleWindow(BrowserWindow, partitionSession) { try { window.webContents.setAudioMuted(true) + // Loads arbitrary user-linked pages on render; it only needs the title, so + // a popup from that page never has a reason to exist (GHSA-9f4c-93c8-jc8g). + window.webContents.setWindowOpenHandler(createWindowOpenHandler()) } catch { // webContents may be unavailable in degraded/headless environments; muting // is best-effort and the window is destroyed within a few seconds anyway. diff --git a/apps/desktop/electron/main.ts b/apps/desktop/electron/main.ts index ac44220756..79f2d3b610 100644 --- a/apps/desktop/electron/main.ts +++ b/apps/desktop/electron/main.ts @@ -160,7 +160,7 @@ import { describeCrashReason, installCrashForensics } from './crash-forensics' import { adoptServedDashboardToken } from './dashboard-token' import { loadOrCreateInstallationId, sshOwnershipId } from './desktop-installation' import { formatDesktopLogLine } from './desktop-log-line' -import { resolveDesktopRemoteRoute } from './desktop-remote-route' +import { resolveDesktopRemoteRoute, v1SshTerminalPoolKey } from './desktop-remote-route' import { buildPosixCleanupScript, buildWindowsCleanupScript, @@ -293,7 +293,9 @@ import { import { selectPoolEvictions } from './pool-eviction' import { clampPoolLimits, parsePoolLimits, POOL_LIMITS_DEFAULTS } from './pool-limits' import { + isBackgroundSlotWaitTimeout, LocalBackendSpawnCoordinator, + type LocalBackendSpawnPriority, type LocalBackendSpawnRequest, releaseLocalBackendSlotAfterExit } from './pool-spawn-coordinator' @@ -416,7 +418,12 @@ import { isHermesOwnedVenvDaemon } from './venv-holder-select' import { fetchMarketplaceThemes, searchMarketplaceThemes } from './vscode-marketplace' import { createWakeIndicatorWindowController } from './wake-indicator-window' import { enumerateWindowsFrontToBack, enumerationFailed, readWindowBelow } from './window-below' -import { registrySshScopeForWindowRoute, WindowConnectionRouteRegistry } from './window-connection-route' +import { + registrySshPoolScopeByConnectionId, + registrySshScopeForWindowRoute, + WindowConnectionRouteRegistry +} from './window-connection-route' +import { createWindowOpenHandler } from './window-open-policy' import { installWindowRendererLifecycle } from './window-renderer-lifecycle' import { createWindowRevealController } from './window-reveal' import { @@ -1501,6 +1508,70 @@ const localBackendSpawnCoordinator = new LocalBackendSpawnCoordinator(poolLimits // the queued ticket fails before the renderer does and the user sees why. const POOL_SLOT_WAIT_MS = 30_000 +function spawnPriorityFrom(value: unknown): LocalBackendSpawnPriority { + return value === 'foreground' ? 'foreground' : 'background' +} + +// Foreground intent for a dial whose pool entry does not exist yet: a user +// click that joins an in-flight backendDialClaims claim never re-enters +// ensureBackend(), and the claim owner may still be awaiting poolStopper / +// registry resolution before backendPool.set(). The local spawn takes the mark +// right before its slot request; the IPC handler that set it clears it once +// the claim settles, so a dial that never reaches a slot request (primary +// route, remote scope, a guard rejection) cannot leave it for a later +// hydration spawn of the same key to pick up. +const pendingForegroundSpawns = new Set() + +function takeForegroundSpawn(...poolKeys: string[]): boolean { + let marked = false + + for (const poolKey of poolKeys) { + marked = pendingForegroundSpawns.delete(poolKey) || marked + } + + return marked +} + +// Upgrade a pooled entry (running, spawning, or queued for a slot) to +// foreground so a queued slot wait can take the reserved foreground slot. +function promotePoolEntry(entry: any): void { + entry.spawnPriority = 'foreground' + entry.localBackendSpawnRequest?.promote?.('foreground') +} + +// Land a spawn failure in desktop.log. A background slot-wait timeout is +// routine under a saturated pool (the next hydration pass retries), so it is +// logged as such instead of as a backend-start failure. +function logPoolSpawnFailure(label: string, error: unknown): void { + if (isBackgroundSlotWaitTimeout(error)) { + rememberLog(`Profile backend ${label} slot wait timed out (background); will retry on the next hydration`) + } else { + rememberLog( + `Hermes backend for profile ${label} failed to start: ${error instanceof Error ? error.message : String(error)}` + ) + } +} + +// Apply foreground intent to the dial claim for `scopeKey`: an entry already +// in the pool is promoted directly, otherwise the intent is marked for the +// spawn the claim owner is about to start. Returns the cleanup that clears a +// mark the dial never consumed. +function applySpawnPriority(scopeKey: string, spawnPriority: LocalBackendSpawnPriority): () => void { + if (spawnPriority !== 'foreground') { + return () => undefined + } + + const existing = backendPool.get(scopeKey) + + if (existing) { + promotePoolEntry(existing) + } else { + pendingForegroundSpawns.add(scopeKey) + } + + return () => void pendingForegroundSpawns.delete(scopeKey) +} + function poolMaxBackends() { return poolLimits.maxBackends } @@ -10054,7 +10125,18 @@ function activeSshTerminalTarget(webContentsId?: number) { const state = sshConnections.get(scope) - return state && state.ssh ? { ssh: state.ssh, scope } : 'pending' + if (state && state.ssh) { + return { ssh: state.ssh, scope } + } + + // The pool's single writer publishes under the per-profile bootstrap key + // while stamping the entry with its registry connection id (#97345), so a + // composite-key miss must still resolve the live tunnel by that identity + // instead of reporting 'pending' forever. + const pooledScope = registrySshPoolScopeByConnectionId(sshConnections, windowRoute.connectionId) + const pooledState = pooledScope === null ? null : sshConnections.get(pooledScope) + + return pooledState && pooledState.ssh ? { ssh: pooledState.ssh, scope: pooledScope } : 'pending' } const profile = windowRoute?.profile ?? primaryProfileKey() @@ -10074,9 +10156,7 @@ function activeSshTerminalTarget(webContentsId?: number) { return null } - const scope = route.connectionId - ? backendScopeKey(route.connectionId, profile) - : sshScopeKey(route.source === 'profile' ? profile : null) + const scope = v1SshTerminalPoolKey(route, profile) const state = sshConnections.get(scope) @@ -11108,8 +11188,9 @@ function profileRouteOptions(profile, request?) { // Resolve a backend connection for the given profile, per the routing table in // resolveProfileBackendRoute(). An empty / unknown profile resolves to the // primary, so legacy callers are unchanged. -async function ensureBackend(profile) { +async function ensureBackend(profile, opts: { spawnPriority?: LocalBackendSpawnPriority } = {}) { const key = profile && String(profile).trim() ? String(profile).trim() : primaryProfileKey() + const spawnPriority = spawnPriorityFrom(opts.spawnPriority) profileDeletionGate.assertCanStart(key) @@ -11141,6 +11222,11 @@ async function ensureBackend(profile) { if (existing) { existing.lastActiveAt = Date.now() + + if (spawnPriority === 'foreground') { + promotePoolEntry(existing) + } + const connection = await existing.connectionPromise setWslBridgeProfileState(key, connection.mode !== 'remote') @@ -11158,16 +11244,15 @@ async function ensureBackend(profile) { remoteBaseUrl: null, releaseLocalBackendSlot: null, localBackendSlotKey: null, - localBackendSpawnRequest: null + localBackendSpawnRequest: null, + spawnPriority } entry.connectionPromise = spawnPoolBackend(key, entry).catch(async error => { // Land the failure in desktop.log: without this a spawn that dies before // its child exists (guard rejection, runtime resolution) leaves no trace // beyond renderer-side rejections users never see in a bundle. - rememberLog( - `Hermes backend for profile "${key}" failed to start: ${error instanceof Error ? error.message : String(error)}` - ) + logPoolSpawnFailure(`"${key}"`, error) await teardownFailedLocalBackend(key, entry) throw error @@ -11188,7 +11273,13 @@ async function ensureBackend(profile) { // a genuinely-local child when the v1 mode says remote; non-local connections // pool under the composite key from backendScopeKey() and reuse the same pool // entry lifecycle (LRU, idle reaper, touch) as per-profile local backends. -async function ensureRegistryBackend(connectionId, profile, managedUpdateCorrelation = '') { +async function ensureRegistryBackend( + connectionId, + profile, + managedUpdateCorrelation = '', + opts: { spawnPriority?: LocalBackendSpawnPriority } = {} +) { + const spawnPriority = spawnPriorityFrom(opts.spawnPriority) const registry = readDesktopConnectionsRegistry() const id = String(connectionId || '').trim() || registry.primary const source = registry.connections.find(c => c.id === id) @@ -11245,7 +11336,7 @@ async function ensureRegistryBackend(connectionId, profile, managedUpdateCorrela const primary = await reuseMatchingPrimarySshBackend({ connectionId: id, effectiveFingerprint: resolveRegistryEffectiveFingerprint, - ensurePrimary: () => ensureBackend(profile), + ensurePrimary: () => ensureBackend(profile, { spawnPriority }), profile, registry, source @@ -11296,7 +11387,7 @@ async function ensureRegistryBackend(connectionId, profile, managedUpdateCorrela }) if (localRoute.delegate) { - return ensureBackend(profile) + return ensureBackend(profile, { spawnPriority }) } const stoppingLocal = poolStopper.inFlight(localRoute.poolKey) @@ -11310,6 +11401,10 @@ async function ensureRegistryBackend(connectionId, profile, managedUpdateCorrela if (existingLocal) { existingLocal.lastActiveAt = Date.now() + if (spawnPriority === 'foreground') { + promotePoolEntry(existingLocal) + } + return existingLocal.connectionPromise } @@ -11324,7 +11419,8 @@ async function ensureRegistryBackend(connectionId, profile, managedUpdateCorrela remoteBaseUrl: null, releaseLocalBackendSlot: null, localBackendSlotKey: null, - localBackendSpawnRequest: null + localBackendSpawnRequest: null, + spawnPriority } localEntry.connectionPromise = spawnPoolBackend(profileKey, localEntry, { @@ -11333,9 +11429,7 @@ async function ensureRegistryBackend(connectionId, profile, managedUpdateCorrela }).catch(async error => { // Same trace rule as the v1 pool path: a forced-local child whose spawn // rejects before the child exists must still land in desktop.log. - rememberLog( - `Hermes backend for profile "${profileKey}" (forced-local) failed to start: ${error instanceof Error ? error.message : String(error)}` - ) + logPoolSpawnFailure(`"${profileKey}" (forced-local)`, error) await teardownFailedLocalBackend(localRoute.poolKey, localEntry) throw error @@ -12173,11 +12267,23 @@ async function spawnPoolBackend(profile, entry, opts: { forceLocal?: boolean; po // pool-idle window (10 min) would hold the pool key hostage and every // later click on the profile would join that stale wait. Failing here // surfaces the "all N slots busy" reason instead of a generic boot timeout. - const spawnRequest = localBackendSpawnCoordinator.request(poolKey, { timeoutMs: POOL_SLOT_WAIT_MS }) + // The caller stamped entry.spawnPriority from its own request; a foreground + // dial that joined the claim before this entry existed left a mark instead. + if (takeForegroundSpawn(poolKey, profile)) { + entry.spawnPriority = 'foreground' + } + + const spawnPriority: LocalBackendSpawnPriority = spawnPriorityFrom(entry.spawnPriority) + + const spawnRequest = localBackendSpawnCoordinator.request(poolKey, { + timeoutMs: POOL_SLOT_WAIT_MS, + priority: spawnPriority + }) + entry.localBackendSlotKey = poolKey entry.localBackendSpawnRequest = spawnRequest - if (localBackendSpawnCoordinator.activeCount >= poolMaxBackends()) { + if (spawnRequest.queued) { rememberLog( `Profile backend "${profile}" waiting for a free local slot (${localBackendSpawnCoordinator.activeCount}/${poolMaxBackends()} busy, ${localBackendSpawnCoordinator.queuedCount} queued)` ) @@ -13019,11 +13125,11 @@ function wireCommonWindowHandlers(win, { zoom = true }: { zoom?: boolean } = {}) } installContextMenuBridge(win) - win.webContents.setWindowOpenHandler(details => { - void openExternalUrl(details.url) - - return { action: 'deny' } - }) + // Always deny, never open as a side effect: GHSA-9f4c-93c8-jc8g. Trusted + // links arrive via `hermes:openExternal`, not here. See window-open-policy.ts. + win.webContents.setWindowOpenHandler( + createWindowOpenHandler(origin => rememberLog(`[window-open] denied: ${origin}`)) + ) win.webContents.on('will-navigate', (event, url) => { if ((DEV_SERVER && url.startsWith(DEV_SERVER)) || (!DEV_SERVER && url.startsWith('file:'))) { return @@ -14474,13 +14580,27 @@ function createWindow() { }) } -ipcMain.handle('hermes:connection', async (_event, profile) => { +ipcMain.handle('hermes:connection', async (_event, profile, extra) => { // Coalesce concurrent renderer dials for one profile scope (#90812): the // renderer-side reconnect lock is per-window, so two windows waking at once // both land here. The claim key mirrors ensureBackend()'s own profile // normalization so every spelling of the primary coalesces onto one dial. const profileKey = profile && String(profile).trim() ? String(profile).trim() : primaryProfileKey() - const connection = await backendDialClaims.run(backendScopeKey(null, profileKey), () => ensureBackend(profile)) + // A user click may join an in-flight hydration claim; the foreground intent + // is applied to that claim so its slot wait can take the reserved slot. + const spawnPriority = spawnPriorityFrom(extra?.priority) + + const scopeKey = backendScopeKey(null, profileKey) + const clearSpawnPriority = applySpawnPriority(scopeKey, spawnPriority) + + let connection + + try { + connection = await backendDialClaims.run(scopeKey, () => ensureBackend(profile, { spawnPriority })) + } finally { + clearSpawnPriority() + } + const connectionId = resolvedConnectionId(readDesktopConnectionsRegistry(), connection) return connectionId ? { ...connection, connectionId } : connection @@ -14491,13 +14611,24 @@ ipcMain.handle('hermes:connection', async (_event, profile) => { // forces a genuinely-local child when the v1 global mode is remote (the // registry 'local' entry always means this machine). ipcMain.handle('hermes:connection:for', async (_event, payload) => { - const { connectionId, profile } = payload && typeof payload === 'object' ? (payload as any) : ({} as any) + const { connectionId, profile, priority } = payload && typeof payload === 'object' ? (payload as any) : ({} as any) const registry = readDesktopConnectionsRegistry() const id = String(connectionId || '').trim() || registry.primary + const spawnPriority = spawnPriorityFrom(priority) + // Same single-owner claim as 'hermes:connection', keyed by the composite // (connectionId, profile) scope (#90812): concurrent registry dials for one // scope share the first spawn instead of bootstrapping duplicate remotes. - const connection = await backendDialClaims.run(backendScopeKey(id, profile), () => ensureRegistryBackend(id, profile)) + const scopeKey = backendScopeKey(id, profile) + const clearSpawnPriority = applySpawnPriority(scopeKey, spawnPriority) + + let connection + + try { + connection = await backendDialClaims.run(scopeKey, () => ensureRegistryBackend(id, profile, '', { spawnPriority })) + } finally { + clearSpawnPriority() + } return { ...connection, connectionId: id, registryScoped: true } }) diff --git a/apps/desktop/electron/plugin-compat-notice.test.ts b/apps/desktop/electron/plugin-compat-notice.test.ts index e4032690ed..e474879f15 100644 --- a/apps/desktop/electron/plugin-compat-notice.test.ts +++ b/apps/desktop/electron/plugin-compat-notice.test.ts @@ -73,6 +73,7 @@ test('dismissal is remembered for the same report and forgotten for a different ...REPORT, plugins: { ...REPORT.plugins, gamma: [{ file: 'g.py', line: 1, old: 'x.y', new: 'z.y' }] } } + fs.writeFileSync(path.join(home, REPORT_FILE), JSON.stringify(grown)) const second = pendingNotice(home, userData) assert.ok(second) diff --git a/apps/desktop/electron/plugin-compat-notice.ts b/apps/desktop/electron/plugin-compat-notice.ts index 2e01f3e8b6..5cd56f3f60 100644 --- a/apps/desktop/electron/plugin-compat-notice.ts +++ b/apps/desktop/electron/plugin-compat-notice.ts @@ -87,6 +87,7 @@ export function recordDismissed(userData: string, key: string): void { if (!keys.includes(key)) { keys.push(key) } + fs.mkdirSync(userData, { recursive: true }) fs.writeFileSync(file, JSON.stringify({ keys: keys.slice(-20) }, null, 2)) } @@ -105,11 +106,13 @@ export function pendingNotice(hermesHome: string, userData: string): PendingNoti if (!report) { return null } + const key = reportKey(report) if (wasDismissed(userData, key)) { return null } + const names = Object.keys(report.plugins).sort() const list = names diff --git a/apps/desktop/electron/pool-spawn-coordinator.test.ts b/apps/desktop/electron/pool-spawn-coordinator.test.ts index a9207ccfc9..83f7884766 100644 --- a/apps/desktop/electron/pool-spawn-coordinator.test.ts +++ b/apps/desktop/electron/pool-spawn-coordinator.test.ts @@ -6,7 +6,11 @@ import { fileURLToPath } from 'node:url' import { test } from 'vitest' -import { LocalBackendSpawnCoordinator, releaseLocalBackendSlotAfterExit } from './pool-spawn-coordinator' +import { + LocalBackendSlotWaitTimeoutError, + LocalBackendSpawnCoordinator, + releaseLocalBackendSlotAfterExit +} from './pool-spawn-coordinator' const deferred = () => { let resolve!: () => void @@ -306,6 +310,197 @@ test('setLimit rejects a non-positive or fractional cap', () => { assert.equal(coordinator.limit, 2) }) +test('cap 3: two background leases leave a reserved slot for foreground', async () => { + const coordinator = new LocalBackendSpawnCoordinator(3) + const bg1 = await coordinator.request('bg-1', { priority: 'background' }).acquired + const bg2 = await coordinator.request('bg-2', { priority: 'background' }).acquired + assert.equal(coordinator.activeCount, 2) + assert.equal(coordinator.queuedCount, 0) + + let fgGranted = false + + const fgPromise = coordinator.request('fg', { priority: 'foreground' }).acquired.then(release => { + fgGranted = true + + return release + }) + + await flush() + assert.equal(fgGranted, true) + assert.equal(coordinator.activeCount, 3) + + const releaseFg = await fgPromise + bg1() + bg2() + releaseFg() + assert.equal(coordinator.activeCount, 0) +}) + +test('untagged acquire still fills the cap (foreground default)', async () => { + const coordinator = new LocalBackendSpawnCoordinator(3) + const releases = await Promise.all(['a', 'b', 'c'].map(key => coordinator.acquire(key))) + assert.equal(coordinator.activeCount, 3) + assert.equal(coordinator.queuedCount, 0) + + for (const release of releases) { + release() + } + + assert.equal(coordinator.activeCount, 0) +}) + +test('foreground is granted the reserved slot ahead of a background hydration queue', async () => { + const coordinator = new LocalBackendSpawnCoordinator(3) + + const bgRunning = await Promise.all( + ['bg-run-1', 'bg-run-2'].map(key => coordinator.request(key, { priority: 'background' }).acquired) + ) + + const queued = Array.from({ length: 20 }, (_, index) => + coordinator.request(`bg-wait-${index}`, { priority: 'background', timeoutMs: 5_000 }) + ) + + await flush() + assert.equal(coordinator.activeCount, 2) + assert.equal(coordinator.queuedCount, 20) + + const started = Date.now() + const releaseFg = await coordinator.request('user-click', { priority: 'foreground', timeoutMs: 100 }).acquired + assert.ok(Date.now() - started < 80, 'foreground must not wait behind the background queue') + assert.equal(coordinator.activeCount, 3) + + for (const request of queued) { + request.cancel() + } + + releaseFg() + + for (const release of bgRunning) { + release() + } + + await Promise.all( + queued.map(request => + request.acquired.then( + () => undefined, + () => undefined + ) + ) + ) + assert.equal(coordinator.activeCount, 0) + assert.equal(coordinator.queuedCount, 0) +}) + +test('drain prefers a foreground waiter over an earlier background waiter', async () => { + const coordinator = new LocalBackendSpawnCoordinator(1) + const releaseHolder = await coordinator.acquire('holder') + const background = coordinator.request('background', { priority: 'background' }) + const foreground = coordinator.request('foreground', { priority: 'foreground' }) + await flush() + assert.equal(coordinator.queuedCount, 2) + + let backgroundEntered = false + let foregroundEntered = false + + const backgroundGrant = background.acquired.then(release => { + backgroundEntered = true + + return release + }) + + const foregroundGrant = foreground.acquired.then(release => { + foregroundEntered = true + + return release + }) + + releaseHolder() + await flush() + assert.equal(foregroundEntered, true) + assert.equal(backgroundEntered, false) + assert.equal(coordinator.activeCount, 1) + + const releaseForeground = await foregroundGrant + releaseForeground() + const releaseBackground = await backgroundGrant + assert.equal(backgroundEntered, true) + releaseBackground() + assert.equal(coordinator.activeCount, 0) +}) + +test('background slot-wait timeout is distinguishable; foreground keeps a user-facing message', async () => { + const coordinator = new LocalBackendSpawnCoordinator(1) + const releaseFirst = await coordinator.acquire('first') + + const background = coordinator.request('bg', { priority: 'background', timeoutMs: 10 }) + await assert.rejects(background.acquired, error => { + assert.ok(error instanceof LocalBackendSlotWaitTimeoutError) + assert.equal(error.name, 'LocalBackendSlotWaitTimeoutError') + assert.equal(error.priority, 'background') + assert.equal(error.silent, true) + assert.match(error.message, /timed out while waiting for a free slot/) + assert.match(error.message, /\(background\)/) + + return true + }) + + const foreground = coordinator.request('fg', { priority: 'foreground', timeoutMs: 10 }) + await assert.rejects(foreground.acquired, error => { + assert.ok(error instanceof Error) + assert.match(error.message, /timed out while waiting for a free slot/) + assert.doesNotMatch(error.message, /\(background\)/) + assert.notEqual(error.name, 'LocalBackendSlotWaitTimeoutError') + + return true + }) + + releaseFirst() + assert.equal(coordinator.activeCount, 0) +}) + +test('request() reports whether the caller actually waited behind the queue', async () => { + const coordinator = new LocalBackendSpawnCoordinator(3) + const bg1 = coordinator.request('bg-1', { priority: 'background' }) + const bg2 = coordinator.request('bg-2', { priority: 'background' }) + const bgWait = coordinator.request('bg-3', { priority: 'background' }) + assert.equal(bg1.queued, false) + assert.equal(bg2.queued, false) + assert.equal(bgWait.queued, true) + + // The reserved slot is free: a foreground request is granted immediately + // even though a background waiter is queued. + const fg = coordinator.request('fg', { priority: 'foreground' }) + assert.equal(fg.queued, false) + assert.equal(coordinator.activeCount, 3) + + bgWait.cancel() + await bgWait.acquired.catch(() => undefined) + ;(await fg.acquired)() + ;(await bg1.acquired)() + ;(await bg2.acquired)() + assert.equal(coordinator.activeCount, 0) +}) + +test('promoting a queued background waiter lets it take the reserved foreground slot', async () => { + const coordinator = new LocalBackendSpawnCoordinator(3) + const bg1 = await coordinator.request('bg-1', { priority: 'background' }).acquired + const bg2 = await coordinator.request('bg-2', { priority: 'background' }).acquired + const queued = coordinator.request('same-bot', { priority: 'background' }) + await flush() + assert.equal(coordinator.activeCount, 2) + assert.equal(coordinator.queuedCount, 1) + + assert.equal(queued.promote('foreground'), true) + const releasePromoted = await queued.acquired + assert.equal(coordinator.activeCount, 3) + assert.equal(coordinator.queuedCount, 0) + + releasePromoted() + bg1() + bg2() + assert.equal(coordinator.activeCount, 0) +}) + // ── main.ts wiring ────────────────────────────────────────────────────────── // The coordinator is only as good as the timeout main.ts hands it. A queued // ticket that outlives the renderer's backend-boot budget holds the pool key @@ -330,7 +525,10 @@ test('setLimit rejects a non-positive or fractional cap', () => { assert.ok(Number.isFinite(slotWait) && slotWait > 0, 'POOL_SLOT_WAIT_MS must be a literal in main.ts') assert.ok(Number.isFinite(bootBudget), 'BACKEND_BOOT_WAIT_TIMEOUT_MS must be a literal') assert.ok(slotWait < bootBudget, `slot wait ${slotWait}ms must be below the boot budget ${bootBudget}ms`) - assert.match(mainSource, /localBackendSpawnCoordinator\.request\(poolKey, \{ timeoutMs: POOL_SLOT_WAIT_MS \}\)/) + assert.match( + mainSource, + /localBackendSpawnCoordinator\.request\(poolKey, \{\s*timeoutMs: POOL_SLOT_WAIT_MS,\s*priority: spawnPriority\s*\}\)/ + ) assert.doesNotMatch(mainSource, /request\(poolKey, \{ timeoutMs: POOL_IDLE_MS \}\)/) }) diff --git a/apps/desktop/electron/pool-spawn-coordinator.ts b/apps/desktop/electron/pool-spawn-coordinator.ts index 8e565ed133..5f80af332d 100644 --- a/apps/desktop/electron/pool-spawn-coordinator.ts +++ b/apps/desktop/electron/pool-spawn-coordinator.ts @@ -1,17 +1,47 @@ export type ReleaseLocalBackendSlot = () => void +export type LocalBackendSpawnPriority = 'foreground' | 'background' + export type LocalBackendSpawnRequest = { acquired: Promise cancel: () => boolean + promote: (priority: LocalBackendSpawnPriority) => boolean + /** False when the slot was granted without waiting behind the queue. */ + queued: boolean } type Waiter = { key: string + priority: LocalBackendSpawnPriority resolve: (release: ReleaseLocalBackendSlot) => void reject: (error: Error) => void timer: ReturnType | null } +const SLOT_WAIT_TIMEOUT_MESSAGE = (key: string) => + `Local backend start for "${key}" timed out while waiting for a free slot.` + +/** + * Slot-wait timeout. Background hydrations set `silent` so call sites can fail + * quiet instead of toasting a user-visible backend-start failure. + */ +export class LocalBackendSlotWaitTimeoutError extends Error { + readonly priority: LocalBackendSpawnPriority + readonly silent: boolean + + constructor(key: string, priority: LocalBackendSpawnPriority) { + const suffix = priority === 'background' ? ' (background)' : '' + super(`${SLOT_WAIT_TIMEOUT_MESSAGE(key)}${suffix}`) + this.name = 'LocalBackendSlotWaitTimeoutError' + this.priority = priority + this.silent = priority === 'background' + } +} + +export function isBackgroundSlotWaitTimeout(error: unknown): boolean { + return error instanceof LocalBackendSlotWaitTimeoutError && error.silent +} + export async function releaseLocalBackendSlotAfterExit( release: ReleaseLocalBackendSlot, waitForExit: () => Promise @@ -25,10 +55,15 @@ export async function releaseLocalBackendSlotAfterExit( * * A lease is acquired immediately before local start work and is held until * the child exits or the start fails. Remote descriptors never call request(). + * + * When the cap is at least 2, one slot is reserved for foreground (user-open) + * requests so background roster hydration cannot occupy the whole pool. + * Untagged acquire() is foreground, so existing cap tests still fill `limit`. */ export class LocalBackendSpawnCoordinator { #limit: number - #active = 0 + #activeForeground = 0 + #activeBackground = 0 #queue: Waiter[] = [] constructor(limit: number) { @@ -40,7 +75,7 @@ export class LocalBackendSpawnCoordinator { } get activeCount(): number { - return this.#active + return this.#activeForeground + this.#activeBackground } get limit(): number { @@ -66,39 +101,47 @@ export class LocalBackendSpawnCoordinator { return this.#queue.length } - request(key: string, options: { timeoutMs?: number } = {}): LocalBackendSpawnRequest { + request( + key: string, + options: { timeoutMs?: number; priority?: LocalBackendSpawnPriority } = {} + ): LocalBackendSpawnRequest { if (options.timeoutMs !== undefined && (!Number.isFinite(options.timeoutMs) || options.timeoutMs < 1)) { throw new RangeError('Local backend spawn timeout must be a positive number.') } - if (this.#active < this.#limit) { + const priority: LocalBackendSpawnPriority = options.priority === 'background' ? 'background' : 'foreground' + + if (this.#queue.length === 0 && this.#canGrant(priority)) { return { - acquired: Promise.resolve(this.#grant()), - cancel: () => false + acquired: Promise.resolve(this.#grant(priority)), + cancel: () => false, + promote: () => false, + queued: false } } let waiter!: Waiter const acquired = new Promise((resolve, reject) => { - waiter = { key, resolve, reject, timer: null } + waiter = { key, priority, resolve, reject, timer: null } this.#queue.push(waiter) if (options.timeoutMs !== undefined) { waiter.timer = setTimeout(() => { - this.#rejectWaiter( - waiter, - new Error(`Local backend start for "${key}" timed out while waiting for a free slot.`) - ) + this.#rejectWaiter(waiter, this.#timeoutError(waiter)) }, options.timeoutMs) waiter.timer.unref?.() } }) + this.#drain() + return { acquired, cancel: () => - this.#rejectWaiter(waiter, new Error(`Local backend start for "${key}" was cancelled while queued.`)) + this.#rejectWaiter(waiter, new Error(`Local backend start for "${key}" was cancelled while queued.`)), + promote: (nextPriority: LocalBackendSpawnPriority) => this.#promoteWaiter(waiter, nextPriority), + queued: this.#queue.includes(waiter) } } @@ -106,6 +149,45 @@ export class LocalBackendSpawnCoordinator { return this.request(key).acquired } + #timeoutError(waiter: Waiter): Error { + if (waiter.priority === 'background') { + return new LocalBackendSlotWaitTimeoutError(waiter.key, 'background') + } + + return new Error(SLOT_WAIT_TIMEOUT_MESSAGE(waiter.key)) + } + + #backgroundLimit(): number { + return this.#limit >= 2 ? this.#limit - 1 : this.#limit + } + + #canGrant(priority: LocalBackendSpawnPriority): boolean { + if (this.activeCount >= this.#limit) { + return false + } + + if (priority === 'background' && this.#activeBackground >= this.#backgroundLimit()) { + return false + } + + return true + } + + #promoteWaiter(waiter: Waiter, priority: LocalBackendSpawnPriority): boolean { + if (!this.#queue.includes(waiter)) { + return false + } + + if (waiter.priority === priority) { + return false + } + + waiter.priority = priority + this.#drain() + + return true + } + #rejectWaiter(waiter: Waiter, error: Error): boolean { const index = this.#queue.indexOf(waiter) @@ -127,8 +209,13 @@ export class LocalBackendSpawnCoordinator { } } - #grant(): ReleaseLocalBackendSlot { - this.#active += 1 + #grant(priority: LocalBackendSpawnPriority): ReleaseLocalBackendSlot { + if (priority === 'background') { + this.#activeBackground += 1 + } else { + this.#activeForeground += 1 + } + let released = false return () => { @@ -137,17 +224,49 @@ export class LocalBackendSpawnCoordinator { } released = true - this.#active -= 1 + + if (priority === 'background') { + this.#activeBackground -= 1 + } else { + this.#activeForeground -= 1 + } + this.#drain() } } - /** Hand free slots to queued waiters while under the (possibly lowered) cap. */ + #takeWaiter(priority: LocalBackendSpawnPriority): Waiter | undefined { + const index = this.#queue.findIndex(waiter => waiter.priority === priority) + + if (index === -1) { + return undefined + } + + return this.#queue.splice(index, 1)[0] + } + + /** Hand free slots to queued waiters. Foreground waiters always go first. */ #drain(): void { - while (this.#active < this.#limit && this.#queue.length > 0) { - const next = this.#queue.shift()! + while (this.#canGrant('foreground')) { + const next = this.#takeWaiter('foreground') + + if (!next) { + break + } + this.#clearTimer(next) - next.resolve(this.#grant()) + next.resolve(this.#grant('foreground')) + } + + while (this.#canGrant('background')) { + const next = this.#takeWaiter('background') + + if (!next) { + break + } + + this.#clearTimer(next) + next.resolve(this.#grant('background')) } } } diff --git a/apps/desktop/electron/preload.ts b/apps/desktop/electron/preload.ts index 28cc94cd5c..32d17e48de 100644 --- a/apps/desktop/electron/preload.ts +++ b/apps/desktop/electron/preload.ts @@ -16,7 +16,7 @@ contextBridge.exposeInMainWorld('hermesDesktop', { glassSupported: translucencySupport?.glass === true, translucencySupported: translucencySupport?.translucency === true, localModelsEnabled: featureFlags?.localModels === true, - getConnection: profile => ipcRenderer.invoke('hermes:connection', profile), + getConnection: (profile, opts) => ipcRenderer.invoke('hermes:connection', profile, opts), // Registry-scoped backend resolution: { connectionId, profile } → descriptor. getConnectionFor: payload => ipcRenderer.invoke('hermes:connection:for', payload), getProfileRoutes: profiles => ipcRenderer.invoke('hermes:plugin-profile-routes', profiles), diff --git a/apps/desktop/electron/remote-lifecycle.test.ts b/apps/desktop/electron/remote-lifecycle.test.ts index dffc716193..408cd9e27a 100644 --- a/apps/desktop/electron/remote-lifecycle.test.ts +++ b/apps/desktop/electron/remote-lifecycle.test.ts @@ -1849,3 +1849,66 @@ test('cleanupStale keeps the lockfile when even SIGKILL cannot confirm the pid d // The record must survive so the next connect's reap pass retries. assert.ok(!ssh.calls.some(c => /rm -f .*backend\.lock\.json/.test(c))) }) +test.skipIf(process.platform === 'win32')( + 'buildSpawnCommand quotes expandRemotePath fragments exactly once (real sh parse)', + async () => { + // expandRemotePath() output is pre-quoted; a second shq() ships literal quote + // characters to the remote python. Parse the composed command with a real sh, + // as the remote login shell does, and require every path to come out clean. + const cmd = buildSpawnCommand('/x/hermes', 'work', { + hermesHome: '~/.hermes', + logPath: spawnLogPath(OWNERSHIP_ID, SPAWN_NONCE), + ownershipId: OWNERSHIP_ID, + reservationNonce: SPAWN_NONCE, + spawnNonce: SPAWN_NONCE, + tokenFilePath: spawnTokenPath(OWNERSHIP_ID, SPAWN_NONCE), + lockMetadata: { ownershipId: OWNERSHIP_ID, spawnNonce: SPAWN_NONCE } + }) + + // Capture the argv a remote shell would hand to python3, via a shim on PATH. + const root = await mkdtemp(path.join(os.tmpdir(), 'hermes-argv-shim-')) + + try { + const shimDir = path.join(root, 'shim') + const fakeHome = path.join(root, 'home') + await mkdir(shimDir) + await mkdir(fakeHome) + const argvFile = path.join(shimDir, 'argv') + await writeFile(path.join(shimDir, 'python3'), `#!/bin/sh\nprintf '%s\\0' "$@" > '${argvFile}'\n`, { + mode: 0o755 + }) + await exec(cmd, { + env: { ...process.env, PATH: `${shimDir}:${process.env.PATH}`, HOME: fakeHome } + }) + const argv = (await readFile(argvFile, 'utf8')).split('\0') + + // argv: ['-c', , , ] + assert.equal( + argv[2], + `${fakeHome}/.hermes/.hermes-update-in-progress.mutex`, + 'mutex path must reach python fully expanded, with no quote characters' + ) + + // The payload assigns reservation/lock/owner_file before its mkdir loop. + // Evaluate only that prefix the way the remote sh does; never the loop itself. + const payload = argv[3] + const loopStart = payload.indexOf('i=0;') + assert.ok(loopStart > 0, 'payload prefix sentinel missing') + + const { stdout } = await exec( + `${payload.slice(0, loopStart)} printf '%s\\n' "$reservation" "$lock" "$owner_file"`, + { + env: { ...process.env, HOME: fakeHome } + } + ) + + const [reservation, lock, ownerFile] = stdout.split('\n') + const base = `${fakeHome}/.hermes/desktop-ssh/${OWNERSHIP_ID}` + assert.equal(reservation, `${base}/.connect.lock`) + assert.equal(lock, `${base}/backend.lock.json`) + assert.equal(ownerFile, `${base}/.connect.lock/owner`) + } finally { + await rm(root, { recursive: true, force: true }) + } + } +) diff --git a/apps/desktop/electron/remote-lifecycle.ts b/apps/desktop/electron/remote-lifecycle.ts index d096705c95..40108383cd 100644 --- a/apps/desktop/electron/remote-lifecycle.ts +++ b/apps/desktop/electron/remote-lifecycle.ts @@ -895,7 +895,9 @@ finally: // the marker check, spawns the backend, and publishes its initial lockfile. // Python keeps the descriptor close-on-exec by default and passes it explicitly // only to the intended outer shell; each detached child closes it before -// execing Hermes. +// execing Hermes. mutexPath is expandRemotePath() output — a complete shell +// word ("$HOME"'/…' or '/abs/…') embedded raw so $HOME expands remotely; a +// second shq() would hand python the quote characters as part of the path. function withRemoteUpdateMutex(command, mutexPath) { const script = ` import fcntl,os,subprocess,sys @@ -913,7 +915,7 @@ finally: sys.exit(result.returncode if result is not None else 1) `.trim() - return `python3 -c ${shq(script)} ${shq(mutexPath)} ${shq(command)}` + return `python3 -c ${shq(script)} ${mutexPath} ${shq(command)}` } /** diff --git a/apps/desktop/electron/window-connection-route.test.ts b/apps/desktop/electron/window-connection-route.test.ts index 17a4a488f4..400f151db7 100644 --- a/apps/desktop/electron/window-connection-route.test.ts +++ b/apps/desktop/electron/window-connection-route.test.ts @@ -4,6 +4,7 @@ import { test } from 'vitest' import { normalizeWindowConnectionRoute, + registrySshPoolScopeByConnectionId, registrySshScopeForWindowRoute, WindowConnectionRouteRegistry } from './window-connection-route' @@ -120,3 +121,20 @@ test('uses the canonical default profile scope when a registry SSH route has no 'conn:source-b::default' ) }) + +test('recovers the bootstrap pool key a registry SSH tunnel was published under', () => { + const pool = new Map([['', { registryConnectionId: 'source-b', ssh: { alive: true } }]]) + + assert.equal(registrySshPoolScopeByConnectionId(pool, 'source-b'), '') +}) + +test('does not match another connection, an unlabelled entry, or a torn-down tunnel', () => { + const pool = new Map([ + ['research', { registryConnectionId: 'source-a', ssh: { alive: true } }], + ['', { registryConnectionId: '', ssh: { alive: true } }], + ['worker', { registryConnectionId: 'source-b' }] + ]) + + assert.equal(registrySshPoolScopeByConnectionId(pool, 'source-b'), null) + assert.equal(registrySshPoolScopeByConnectionId(pool, 'source-c'), null) +}) diff --git a/apps/desktop/electron/window-connection-route.ts b/apps/desktop/electron/window-connection-route.ts index 7c0a83cee5..806f690508 100644 --- a/apps/desktop/electron/window-connection-route.ts +++ b/apps/desktop/electron/window-connection-route.ts @@ -40,6 +40,30 @@ export function registrySshScopeForWindowRoute( return backendScopeKey(route.connectionId, route.profile) } +export interface RegistrySshPoolEntry { + registryConnectionId?: null | string + ssh?: unknown +} + +// The sshConnections pool has a single writer that publishes every tunnel under +// its per-profile bootstrap key while stamping the entry with the registry +// connection that owns it. A registry-scoped lookup under the composite +// backendScopeKey can therefore miss a live tunnel; this recovers the writer's +// actual key from the stamped identity, the same match managedSshScopeRole +// applies to pool entries. +export function registrySshPoolScopeByConnectionId( + entries: Iterable, + connectionId: string +): null | string { + for (const [scope, entry] of entries) { + if (entry?.ssh && entry.registryConnectionId === connectionId) { + return scope + } + } + + return null +} + export class WindowConnectionRouteRegistry { private readonly routes = new Map() diff --git a/apps/desktop/electron/window-open-policy.ts b/apps/desktop/electron/window-open-policy.ts new file mode 100644 index 0000000000..48821b9188 --- /dev/null +++ b/apps/desktop/electron/window-open-policy.ts @@ -0,0 +1,58 @@ +/** + * Window-open policy for every BrowserWindow's webContents. + * + * Every external URL the desktop opens on purpose goes through the audited + * `hermes:openExternal` IPC channel (`openExternalUrl` in main.ts: http/https/ + * mailto allowlist, guarded file:). The `window.open` / `target=_blank` path + * that reaches `setWindowOpenHandler` is therefore only ever driven by content + * we did NOT initiate — most dangerously untrusted HTML in sandboxed + * `allow-scripts` iframes (artifact previews, inline preview directives). + * + * GHSA-9f4c-93c8-jc8g (CVE-2026-70608): a sandboxed iframe without + * `allow-popups` and without a user gesture can still reach this handler via + * the OpenURL navigation path. If the handler opens `details.url` as a side + * effect, a malicious artifact forces the user's OS browser to an attacker URL. + * There is no fixed Electron 40.x, so the defence lives here regardless of the + * pin: deny every request and never open a URL from this handler. + */ + +export interface WindowOpenRequestLike { + url: string +} + +export interface WindowOpenDecision { + action: 'deny' +} + +/** + * `origin` only — a denied URL can carry query credentials, signed-URL tokens + * or attacker-controlled text, none of which belongs in a persisted log. + */ +export function describeDeniedUrl(url: string): string { + try { + const parsed = new URL(url) + + return parsed.origin === 'null' ? parsed.protocol : parsed.origin + } catch { + return '' + } +} + +/** + * Build a `setWindowOpenHandler` callback that denies unconditionally. + * `onDenied` is logging-only and receives the sanitized origin; a throwing + * observer must not be able to change the decision. + */ +export function createWindowOpenHandler( + onDenied?: (origin: string) => void +): (details: WindowOpenRequestLike) => WindowOpenDecision { + return details => { + try { + onDenied?.(describeDeniedUrl(details.url)) + } catch { + // observer failure is not a reason to reconsider the decision + } + + return { action: 'deny' } + } +} diff --git a/apps/desktop/scripts/desktop-update-ui.test.mjs b/apps/desktop/scripts/desktop-update-ui.test.mjs new file mode 100644 index 0000000000..bc69320a4a --- /dev/null +++ b/apps/desktop/scripts/desktop-update-ui.test.mjs @@ -0,0 +1,100 @@ +import assert from 'node:assert/strict' +import fs from 'node:fs' +import { JSDOM } from 'jsdom' +import { afterEach, test, vi } from 'vitest' + +// Execute the shipped page, including its inline script, rather than matching +// source strings or testing a second implementation of the progress client. +const html = fs.readFileSync(new URL('../../../scripts/desktop-update/ui.html', import.meta.url), 'utf8') +const windows = [] + +function openPage(fetch) { + vi.useFakeTimers() + const dom = new JSDOM(html, { + url: 'http://127.0.0.1:12345/', + runScripts: 'dangerously', + beforeParse(window) { + window.fetch = fetch + window.AbortController = AbortController + window.setTimeout = setTimeout + window.clearTimeout = clearTimeout + window.requestAnimationFrame = () => 1 + window.cancelAnimationFrame = () => {} + } + }) + windows.push(dom.window) + return dom.window.document +} + +afterEach(() => { + windows.splice(0).forEach(window => window.close()) + vi.useRealTimers() +}) + +test.each(['done', 'manual', 'error'])('renders %s before acknowledging terminal delivery', async status => { + let document + const requests = [] + const receipt = '550e8400-e29b-41d4-a716-446655440000' + const fetch = vi.fn(async (url, options) => { + requests.push(url) + if (url.startsWith('/ack/')) { + assert.equal(options.method, 'POST') + assert.equal(document.body.className, status === 'error' ? 'error' : 'done') + assert.notEqual(document.getElementById('title').textContent, 'Updating Hermes') + return { ok: true } + } + return { ok: true, json: async () => ({ status, receipt, message: 'The updater result' }) } + }) + document = openPage(fetch) + await vi.advanceTimersByTimeAsync(1000) + assert.equal(document.body.className, status === 'error' ? 'error' : 'done') + assert.notEqual(document.getElementById('title').textContent, 'Updating Hermes') + assert.deepEqual(requests, ['/progress', `/ack/${receipt}`]) +}) + +test.each(['disconnect', 'hung', 'hung-body', 'http', 'invalid'])('bounds %s progress failures without inventing an update outcome', async failure => { + let attempts = 0 + const fetch = vi.fn((_url, options) => { + attempts++ + if (attempts === 1) { + return Promise.resolve({ ok: true, json: async () => ({ status: 'running', message: 'Installing dependencies' }) }) + } + if (failure === 'hung' || failure === 'hung-body') { + const pending = () => new Promise((_resolve, reject) => { + options.signal?.addEventListener('abort', () => reject(new Error('timeout')), { once: true }) + }) + return failure === 'hung' ? pending() : Promise.resolve({ ok: true, json: pending }) + } + if (failure === 'http') return Promise.resolve({ ok: false }) + if (failure === 'invalid') return Promise.resolve({ ok: true, json: async () => ({}) }) + return Promise.reject(new Error('connection refused')) + }) + const document = openPage(fetch) + await vi.advanceTimersByTimeAsync(20_000) + assert.equal(document.body.className, 'disconnected') + assert.equal(document.getElementById('title').textContent, 'Update status unavailable') + assert.match(document.getElementById('line').textContent, /Check Hermes/) + assert.ok(attempts <= 4, `unbounded retry loop: ${attempts}`) +}) + +test.each(['legacy', 'transient', 'ack-failure'])('preserves terminal truth with %s servers', async mode => { + let attempts = 0 + const fetch = vi.fn(async url => { + if (url.startsWith('/ack/')) throw new Error('server already stopped') + if (++attempts === 1 && mode === 'transient') throw new Error('temporary disconnect') + return { ok: true, json: async () => ({ status: 'done', ...(mode === 'ack-failure' ? { receipt: 'test-receipt' } : {}) }) } + }) + const document = openPage(fetch) + await vi.advanceTimersByTimeAsync(20_000) + assert.equal(document.body.className, 'done') + assert.equal(document.getElementById('title').textContent, 'Update complete') +}) + +test('continues displaying a healthy long update while progress remains reachable', async () => { + const fetch = vi.fn(async () => ({ ok: true, json: async () => ({ status: 'running', message: 'Building Desktop' }) })) + const document = openPage(fetch) + await vi.advanceTimersByTimeAsync(60_000) + assert.equal(document.body.className, '') + assert.equal(document.getElementById('title').textContent, 'Updating Hermes') + assert.equal(document.getElementById('line').textContent, 'Building Desktop') +}) diff --git a/apps/desktop/src/app/chat/composer/hooks/use-composer-voice.ts b/apps/desktop/src/app/chat/composer/hooks/use-composer-voice.ts index 9a7db494ee..eb8307b234 100644 --- a/apps/desktop/src/app/chat/composer/hooks/use-composer-voice.ts +++ b/apps/desktop/src/app/chat/composer/hooks/use-composer-voice.ts @@ -4,7 +4,7 @@ import { useCallback, useEffect, useRef, useState } from 'react' import { useI18n } from '@/i18n' import { chatMessageText, collectUnspokenTurnSpeech } from '@/lib/chat-messages' import { triggerHaptic } from '@/lib/haptics' -import { markAssistantIdSpoken, resolveSpokenReply } from '@/lib/spoken-reply' +import { adoptSpokenReplySession, markAssistantIdSpoken, resolveSpokenReply } from '@/lib/spoken-reply' import { CONVERSATION_LEASE, READ_ALOUD_LEASE, syncTtsLease } from '@/lib/tts-lease' import { clearWakeIndicator, syncWakeIndicatorWithVoice } from '@/lib/wake-indicator' import { $voiceConversationStartRequest, takeVoiceConversationStart } from '@/store/composer' @@ -65,8 +65,15 @@ export function useComposerVoice({ const { $messages } = useComposerScope() const [voiceConversationActive, setVoiceConversationActive] = useState(false) const ownsWakeIndicatorRef = useRef(false) + const previousSessionIdRef = useRef(sessionId) const voiceStartRequest = useStore($voiceConversationStartRequest) + // eslint-disable-next-line no-restricted-syntax -- session-id adopt token, not an atom mirror + useEffect(() => { + adoptSpokenReplySession(previousSessionIdRef.current, sessionId) + previousSessionIdRef.current = sessionId + }, [sessionId]) + const { dictate, voiceActivityState, voiceStatus } = useVoiceRecorder({ focusInput, maxRecordingSeconds, diff --git a/apps/desktop/src/app/chat/composer/hooks/use-status-presence.test.ts b/apps/desktop/src/app/chat/composer/hooks/use-status-presence.test.ts new file mode 100644 index 0000000000..1c5fb5f0b7 --- /dev/null +++ b/apps/desktop/src/app/chat/composer/hooks/use-status-presence.test.ts @@ -0,0 +1,200 @@ +import { act, cleanup, renderHook } from '@testing-library/react' +import { afterEach, beforeEach, describe, expect, it } from 'vitest' + +import { $composerActionsBySession } from '@/store/composer-actions' +import { $previewStatusBySession } from '@/store/preview-status' +import { $sessionControlBySession, type SessionControlEntry } from '@/store/session-control' +import { $todosBySession } from '@/store/todos' + +import { useSessionStatusPresence } from './use-status-presence' + +const SID = 'presence-session-1' + +const mockEntry = (overrides?: Partial): SessionControlEntry => ({ + capability: 'supported', + error: null, + loading: false, + pendingAction: null, + snapshot: { + goal: null, + heartbeat: null, + loop: null, + revision: 'rev-1', + updated_at: 1000 + }, + ...overrides +}) + +describe('useSessionStatusPresence', () => { + beforeEach(() => { + $todosBySession.set({}) + $composerActionsBySession.set({}) + $previewStatusBySession.set({}) + $sessionControlBySession.set({}) + }) + + afterEach(() => { + cleanup() + $todosBySession.set({}) + $composerActionsBySession.set({}) + $previewStatusBySession.set({}) + $sessionControlBySession.set({}) + }) + + it('returns false when session is null or empty', () => { + const { result } = renderHook(() => useSessionStatusPresence(null)) + expect(result.current).toBe(false) + + const { result: emptyResult } = renderHook(() => useSessionStatusPresence(SID)) + expect(emptyResult.current).toBe(false) + }) + + it('returns true when legacy status items exist', () => { + const { result } = renderHook(() => useSessionStatusPresence(SID)) + expect(result.current).toBe(false) + + act(() => { + $todosBySession.set({ + [SID]: [{ content: 'task 1', id: '1', status: 'in_progress' }] + }) + }) + + expect(result.current).toBe(true) + }) + + it('returns true when structured goal exists in session control', () => { + const { result } = renderHook(() => useSessionStatusPresence(SID)) + expect(result.current).toBe(false) + + act(() => { + $sessionControlBySession.set({ + [SID]: mockEntry({ + snapshot: { + goal: { + contract: { + boundaries: '', + constraints: '', + outcome: 'test outcome', + stop_when: '', + verification: '' + }, + gates: [], + max_turns: 10, + status: 'active', + subgoals: [], + title: 'Structured Goal', + turns_used: 1 + }, + heartbeat: null, + loop: null, + revision: 'rev-2', + updated_at: 2000 + } + }) + }) + }) + + expect(result.current).toBe(true) + }) + + it('returns true when structured loop exists in session control', () => { + const { result } = renderHook(() => useSessionStatusPresence(SID)) + expect(result.current).toBe(false) + + act(() => { + $sessionControlBySession.set({ + [SID]: mockEntry({ + snapshot: { + goal: null, + heartbeat: null, + loop: { + awaiting_response: false, + created_at: 1000, + current_delay: 60, + deferred_by_goal: false, + interval_seconds: 60, + last_fired_at: 1000, + max_ticks: 10, + mode: 'interval', + next_due_at: 2000, + prompt: 'Run loop', + status: 'active', + ticks_fired: 0, + times: 5, + until: '' + }, + revision: 'rev-3', + updated_at: 2000 + } + }) + }) + }) + + expect(result.current).toBe(true) + }) + + it('returns true when structured heartbeat exists in session control', () => { + const { result } = renderHook(() => useSessionStatusPresence(SID)) + expect(result.current).toBe(false) + + act(() => { + $sessionControlBySession.set({ + [SID]: mockEntry({ + snapshot: { + goal: null, + heartbeat: { + created_at: 1000, + fire_count: 1, + interval_seconds: 300, + last_fired_at: 1000, + prompt: 'Heartbeat check', + status: 'active' + }, + loop: null, + revision: 'rev-4', + updated_at: 2000 + } + }) + }) + }) + + expect(result.current).toBe(true) + }) + + it('returns false when session control entry has empty snapshot (null goal, loop, heartbeat)', () => { + const { result } = renderHook(() => useSessionStatusPresence(SID)) + expect(result.current).toBe(false) + + act(() => { + $sessionControlBySession.set({ + [SID]: mockEntry({ + snapshot: { + goal: null, + heartbeat: null, + loop: null, + revision: 'rev-empty', + updated_at: 2000 + } + }) + }) + }) + + expect(result.current).toBe(false) + }) + + it('returns true when session control entry has only an error', () => { + const { result } = renderHook(() => useSessionStatusPresence(SID)) + expect(result.current).toBe(false) + + act(() => { + $sessionControlBySession.set({ + [SID]: mockEntry({ + error: 'Gateway connection failed', + snapshot: null + }) + }) + }) + + expect(result.current).toBe(true) + }) +}) diff --git a/apps/desktop/src/app/chat/composer/hooks/use-status-presence.ts b/apps/desktop/src/app/chat/composer/hooks/use-status-presence.ts index b4655ffd50..b2c059ba35 100644 --- a/apps/desktop/src/app/chat/composer/hooks/use-status-presence.ts +++ b/apps/desktop/src/app/chat/composer/hooks/use-status-presence.ts @@ -3,8 +3,9 @@ import { useSyncExternalStore } from 'react' import { $composerActionsBySession } from '@/store/composer-actions' import { $statusItemsBySession } from '@/store/composer-status' import { $previewStatusBySession } from '@/store/preview-status' +import { $sessionControlBySession } from '@/store/session-control' -/** Structural view of the three per-session feeds — they hold different item +/** Structural view of the per-session feeds — they hold different item * types, and all this hook needs from each is "does this key have rows". */ interface PresenceFeed { get(): Record @@ -14,7 +15,7 @@ interface PresenceFeed { const FEEDS: PresenceFeed[] = [$statusItemsBySession, $composerActionsBySession, $previewStatusBySession] const subscribe = (onChange: () => void) => { - const offs = FEEDS.map(feed => feed.listen(onChange)) + const offs = [...FEEDS.map(feed => feed.listen(onChange)), $sessionControlBySession.listen(onChange)] return () => { for (const off of offs) { @@ -24,8 +25,9 @@ const subscribe = (onChange: () => void) => { } /** - * Whether a session has any status items, micro actions, or previews, as a - * coarse *edge*: the boolean only flips when the stack appears/disappears. + * Whether a session has any status items, micro actions, previews, or + * structured session controls (goal, loop, heartbeat), as a coarse *edge*: + * the boolean only flips when the stack appears/disappears. * ChatBar uses it to toggle a styling data-attr — subscribing to the whole * `$statusItemsBySession` (a `computed` that rebuilds the entire map) / * `$previewStatusBySession` maps re-rendered the ~1.4k ChatBar on every @@ -39,6 +41,16 @@ export function useSessionStatusPresence(sessionId: string | null): boolean { return false } - return FEEDS.some(feed => (feed.get()[sessionId]?.length ?? 0) > 0) + if (FEEDS.some(feed => (feed.get()[sessionId]?.length ?? 0) > 0)) { + return true + } + + const control = $sessionControlBySession.get()[sessionId] + + return Boolean( + control?.error || + (control?.snapshot && + (control.snapshot.goal !== null || control.snapshot.loop !== null || control.snapshot.heartbeat !== null)) + ) }) } diff --git a/apps/desktop/src/app/chat/composer/index.tsx b/apps/desktop/src/app/chat/composer/index.tsx index b01381e41d..371bf32b07 100644 --- a/apps/desktop/src/app/chat/composer/index.tsx +++ b/apps/desktop/src/app/chat/composer/index.tsx @@ -1195,6 +1195,7 @@ export function ChatBar({ grows upward over the thread and the dock's own measurement covers it. Collapses to nothing when every status is empty. */} 0 ? ( group.type === 'todo' && group.items.some(item => item.todoStatus === 'in_progress' && item.state === 'running') interface ComposerStatusStackProps { + onSubmit?: (value: string, options?: SubmitTextOptions) => Promise | boolean /** The queue, built by the composer (it owns the queue's callbacks). Rendered * as the last group so it stays fused to the composer like before. */ queue: ReactNode @@ -85,7 +89,7 @@ interface ComposerStatusStackProps { * every session-scoped status — subagents, background tasks, queue — grouped by * type and separated by light dividers. Collapses to nothing when empty. */ -export function ComposerStatusStack({ queue, sessionId }: ComposerStatusStackProps) { +export function ComposerStatusStack({ onSubmit, queue, sessionId }: ComposerStatusStackProps) { const { t } = useI18n() const navigate = useNavigate() // Subscribe to THIS session's slice only. Both maps churn on other @@ -96,10 +100,22 @@ export function ComposerStatusStack({ queue, sessionId }: ComposerStatusStackPro // items actually changed. const items = useSessionSlice($statusItemsBySession, sessionId) const previews = useSessionSlice($previewStatusBySession, sessionId) + const controlEntry = useSessionValue($sessionControlBySession, sessionId) + const scrolledUp = useStore($threadScrolledUp) const billing = useStore($billingBlock) - const groups = useMemo(() => groupStatusItems(items), [items]) + const isStructuredSupported = controlEntry?.capability === 'supported' + + const groups = useMemo(() => { + const raw = groupStatusItems(items) + + if (isStructuredSupported) { + return raw.filter(g => g.type !== 'goal') + } + + return raw + }, [items, isStructuredSupported]) // Seed from the registry on session open; event-driven refreshes (terminal / // process tool completions) live in use-message-stream. This must NOT reset @@ -112,7 +128,7 @@ export function ComposerStatusStack({ queue, sessionId }: ComposerStatusStackPro useEffect(() => { if (sessionId) { void refreshBackgroundProcesses(sessionId) - void refreshSessionGoal(sessionId) + void refreshSessionControl(sessionId) } }, [sessionId]) @@ -160,6 +176,21 @@ export function ComposerStatusStack({ queue, sessionId }: ComposerStatusStackPro sections.push({ key: 'billing', node: }) } + const hasControlContent = Boolean( + controlEntry && + (controlEntry.error || + controlEntry.snapshot?.goal || + controlEntry.snapshot?.loop || + controlEntry.snapshot?.heartbeat) + ) + + if (sessionId && controlEntry && hasControlContent) { + sections.push({ + key: 'session-control', + node: + }) + } + for (const group of groups) { sections.push({ key: group.type, diff --git a/apps/desktop/src/app/chat/composer/status-stack/session-control-goal.tsx b/apps/desktop/src/app/chat/composer/status-stack/session-control-goal.tsx new file mode 100644 index 0000000000..59d55dbd4a --- /dev/null +++ b/apps/desktop/src/app/chat/composer/status-stack/session-control-goal.tsx @@ -0,0 +1,639 @@ +import { memo, useCallback, useState } from 'react' + +import { queueKickoffIfSessionBusy } from '@/app/session/hooks/use-prompt-actions/queue-if-busy' +import type { SubmitTextOptions } from '@/app/session/hooks/use-prompt-actions/utils' +import { StatusSection } from '@/components/chat/status-section' +import { Button } from '@/components/ui/button' +import { Codicon } from '@/components/ui/codicon' +import { ConfirmDialog } from '@/components/ui/confirm-dialog' +import { + ContextMenu, + ContextMenuContent, + ContextMenuItem, + ContextMenuSeparator, + ContextMenuTrigger +} from '@/components/ui/context-menu' +import { + Dialog, + DialogContent, + DialogDescription, + DialogFooter, + DialogHeader, + DialogTitle +} from '@/components/ui/dialog' +import { + DropdownMenu, + DropdownMenuContent, + DropdownMenuItem, + DropdownMenuSeparator, + DropdownMenuTrigger +} from '@/components/ui/dropdown-menu' +import { Tip } from '@/components/ui/tooltip' +import { useI18n } from '@/i18n' +import { + runSessionControlAction, + type SessionControlAction, + type SessionControlActionArgs, + type SessionControlGoal +} from '@/store/session-control' + +import type { ConfirmState } from './session-control-utils' + +interface GoalSectionProps { + goal: SessionControlGoal + sessionId: string + pendingAction: SessionControlAction | null + onSubmit?: (value: string, options?: SubmitTextOptions) => Promise | boolean + onFeedback: (error: string | null, success: string | null) => void +} + +export const SessionControlGoalSection = memo(function SessionControlGoalSection({ + goal, + sessionId, + pendingAction, + onSubmit, + onFeedback +}: GoalSectionProps) { + const { t } = useI18n() + const s = t.statusStack + const ctrl = s.control + + const [detailsOpen, setDetailsOpen] = useState(false) + const [addCriterionOpen, setAddCriterionOpen] = useState(false) + const [addCriterionError, setAddCriterionError] = useState(null) + const [newCriterionText, setNewCriterionText] = useState('') + const [confirmState, setConfirmState] = useState(null) + const [menuOpen, setMenuOpen] = useState(false) + + const isBusy = Boolean(pendingAction) + + const handleAction = useCallback( + async ( + action: SessionControlAction, + args?: SessionControlActionArgs, + onFailure?: (message: string) => void + ): Promise => { + onFeedback(null, null) + + try { + const dispatch = await runSessionControlAction(sessionId, action, args) + + if (dispatch.type === 'send') { + if (!dispatch.message || !onSubmit) { + onFeedback(ctrl.continuationFailed, null) + onFailure?.(ctrl.continuationFailed) + + return false + } + + // The backend has already resumed the goal; if a turn is running the + // kickoff must queue (same as a typed `/goal resume`, slash.ts), not + // read as a failure. + const queued = queueKickoffIfSessionBusy({ + displayText: dispatch.display ?? undefined, + sessionId, + text: dispatch.message + }) + + if (queued !== 'idle') { + const copy = queued === 'queued' ? ctrl.continuationQueued : ctrl.continuationBusy + onFeedback(queued === 'queued' ? null : copy, queued === 'queued' ? copy : null) + + return queued === 'queued' + } + + const submitted = await onSubmit(dispatch.message, { + displayKind: 'hidden', + sessionId + }) + + if (!submitted) { + onFeedback(ctrl.continuationFailed, null) + onFailure?.(ctrl.continuationFailed) + + return false + } + } + + onFeedback(null, ctrl.actionSucceeded) + + return true + } catch (err) { + const msg = err instanceof Error ? err.message : String(err) + const failure = ctrl.actionFailed(msg) + onFeedback(failure, null) + onFailure?.(failure) + + return false + } + }, + [sessionId, onSubmit, onFeedback, ctrl] + ) + + const copyCriterionText = useCallback( + async (text: string) => { + try { + await navigator.clipboard.writeText(text) + onFeedback(null, ctrl.copySuccess) + } catch { + onFeedback(ctrl.copyFailure, null) + } + }, + [onFeedback, ctrl] + ) + + const visibleState: 'waiting' | 'active' | 'paused' | 'done' = goal.wait_barrier + ? 'waiting' + : goal.status === 'paused' + ? 'paused' + : goal.status === 'done' + ? 'done' + : 'active' + + const iconClass = + goal.last_verdict === 'blocked' + ? 'text-red-500' + : visibleState === 'done' + ? 'text-muted-foreground/70' + : visibleState === 'active' + ? 'text-emerald-500' + : 'text-amber-500' + + const stateLabel = + goal.last_verdict === 'blocked' + ? s.goalBlocked + : visibleState === 'waiting' + ? s.goalWaiting + : visibleState === 'paused' + ? s.goalPaused + : visibleState === 'done' + ? s.goalDone + : s.goalActive + + const headerLabel = + visibleState === 'done' + ? `${stateLabel} · ${ctrl.goalDoneTurns(goal.turns_used)}` + : goal.max_turns > 0 + ? `${stateLabel} · ${ctrl.goalActiveTurns(goal.turns_used, goal.max_turns)}` + : goal.turns_used > 0 + ? `${stateLabel} · ${ctrl.goalTurn(goal.turns_used)}` + : stateLabel + + const confirmClearGoal = () => { + setConfirmState({ + title: ctrl.clearGoalConfirmTitle, + description: ctrl.clearGoalConfirmBody, + destructive: true, + confirmLabel: ctrl.clearGoal, + onConfirm: async () => { + await handleAction('goal.clear') + } + }) + } + + const confirmRemoveCriterion = (index: number) => { + setConfirmState({ + title: ctrl.removeCriterionConfirmTitle(index), + description: ctrl.removeCriterionConfirmBody(index), + destructive: true, + confirmLabel: ctrl.removeCriterion(index), + onConfirm: async () => { + await handleAction('subgoal.remove', { index }) + } + }) + } + + const confirmClearCriteria = () => { + setConfirmState({ + title: ctrl.clearCriteriaConfirmTitle, + description: ctrl.clearCriteriaConfirmBody, + destructive: true, + confirmLabel: ctrl.clearCriteria, + onConfirm: async () => { + await handleAction('subgoal.clear') + } + }) + } + + const openAddCriterion = () => { + setAddCriterionError(null) + setAddCriterionOpen(true) + } + + const renderMenuItems = (isContext = false) => { + const Item = isContext ? ContextMenuItem : DropdownMenuItem + const Sep = isContext ? ContextMenuSeparator : DropdownMenuSeparator + + return ( + <> + {hasDetails && ( + setDetailsOpen(true)}> + + {ctrl.viewDetails} + + )} + {visibleState !== 'done' && ( + + + {ctrl.addCriterion} + + )} + {(hasDetails || visibleState !== 'done') && } + {visibleState === 'active' && ( + void handleAction('goal.pause')}> + + {ctrl.pauseGoal} + + )} + {visibleState === 'paused' && ( + void handleAction('goal.resume')}> + + {ctrl.resumeGoal} + + )} + {visibleState === 'waiting' && ( + <> + void handleAction('goal.unwait')}> + + {ctrl.resumeNow} + + void handleAction('goal.pause')}> + + {ctrl.pauseGoal} + + + )} + + + {ctrl.clearGoal} + + + ) + } + + const hasDetails = Boolean( + goal.contract.outcome || + goal.contract.verification || + goal.contract.constraints || + goal.contract.boundaries || + goal.contract.stop_when || + goal.wait_barrier || + goal.gates.length > 0 + ) + + return ( + <> + + +
+ + + + + + + + + + {renderMenuItems(false)} + + + } + defaultCollapsed={false} + icon={} + label={headerLabel} + > +
+ {/* Full goal title */} +
{goal.title}
+ + {/* Optional reasons */} + {!detailsOpen && goal.wait_barrier && ( +
+ {goal.wait_barrier.reason + ? `${ctrl.waitBarrierTitle}: ${goal.wait_barrier.reason}` + : ctrl.waitBarrierTitle} +
+ )} + {!goal.wait_barrier && goal.paused_reason && ( +
{goal.paused_reason}
+ )} + {!goal.wait_barrier && !goal.paused_reason && goal.last_reason && ( +
{goal.last_reason}
+ )} + + {/* View details button */} + {hasDetails && ( +
+ +
+ )} + + {/* Criteria subsection */} +
+
+
+ + {ctrl.criteriaHeader(goal.subgoals.length)} +
+
+ + {goal.subgoals.length > 0 && ( + + )} +
+
+ + {goal.subgoals.length > 0 ? ( +
+ {goal.subgoals.map((subgoal, idx) => { + const index = idx + 1 + + return ( +
+
+ {index}. + {subgoal} +
+
+ + + + + + +
+
+ ) + })} +
+ ) : null} +
+
+
+
+
+ {renderMenuItems(true)} +
+ + {/* Add Criterion Dialog */} + { + if (!open) { + setAddCriterionOpen(false) + setAddCriterionError(null) + } + }} + open={addCriterionOpen} + > + +
{ + e.preventDefault() + const text = newCriterionText.trim() + + if (!text || isBusy) { + return + } + + setAddCriterionError(null) + const ok = await handleAction('subgoal.add', { text }, setAddCriterionError) + + if (ok) { + setNewCriterionText('') + setAddCriterionOpen(false) + } + }} + > + + {ctrl.addCriterionDialogTitle} + {ctrl.addCriterionPlaceholder} + +
+