From 58da93b3f404e210c031dacf8ae277403da9f5cf Mon Sep 17 00:00:00 2001 From: "ang-fleet-workers[bot]" <333956806+ang-fleet-workers[bot]@users.noreply.github.com> Date: Sun, 27 Sep 2026 19:16:39 -0700 Subject: [PATCH 1/2] fix: C3 re-bucket follow-up, 8 FleetReview rows (#942 x2, #970, #976, #1032, #1043, #1095, #1198, #1254:121) (t_7da6cadf) - #942: desktop hydration accepts a tool-call confab notice on an empty system row (mirrors notice_from_display_row); per-kind label. - #970: switch_model(probe_catalog=False) passes allow_network through get_label / determine_api_mode; cold models.dev cache opens no socket. - #976: _build_hyg_agent / hygiene get_session use the _hyg_old_sid snapshot. - #1043: post-turn clear_resume_pending keeps a mark written during the turn. - #1095: an early-imported bundled provider module is re-registered in the bundled discovery step (filesystem over pip precedence). - #1198: cooling_providers ignores runtime-stage pin refusals. - #1254:121: process_env_files overlay pre-values ride to child agent processes (HERMES_PROCESS_ENV_OVERLAY) so strip_overlay works there. - #1032: store-level model resolution that changes the provider re-runs the base_url exfil guard at write time. Verified: new tests fail on fork/main (8/8), pass on head; neighbouring suites pass (see PR). --- .../src/lib/chat-messages/hydration.ts | 8 +- .../src/lib/confab-notice-hydration.test.ts | 27 +++ apps/shared/src/confab-notice.ts | 40 +++- cron/jobs.py | 27 +++ gateway/run.py | 29 ++- gateway/session.py | 13 +- hermes_cli/kanban_db.py | 6 + hermes_cli/model_switch.py | 12 +- hermes_cli/process_env_files.py | 34 ++- hermes_cli/providers.py | 15 +- providers/__init__.py | 25 ++- tests/gateway/test_c3_followup_t_7da6cadf.py | 199 ++++++++++++++++++ .../test_no_network_cold_models_dev.py | 55 +++++ tests/hermes_cli/test_process_env_files.py | 1 + 14 files changed, 460 insertions(+), 31 deletions(-) create mode 100644 tests/gateway/test_c3_followup_t_7da6cadf.py create mode 100644 tests/gateway/test_no_network_cold_models_dev.py diff --git a/apps/desktop/src/lib/chat-messages/hydration.ts b/apps/desktop/src/lib/chat-messages/hydration.ts index b10869eb8f5c..c70536353871 100644 --- a/apps/desktop/src/lib/chat-messages/hydration.ts +++ b/apps/desktop/src/lib/chat-messages/hydration.ts @@ -1,5 +1,5 @@ import { skillInvocationText } from '@hermes/shared' -import { CONFAB_NOTICE_EVENT_TEXT, confabNoticeFromRow } from '@hermes/shared/confab-notice' +import { confabNoticeEventText, confabNoticeFromRow } from '@hermes/shared/confab-notice' import { extractImageRefs } from '@/lib/embedded-images' import { dedupeGeneratedImageEchoesInParts } from '@/lib/generated-images' @@ -207,12 +207,14 @@ export function toChatMessages(messages: SessionMessage[]): ChatMessage[] { // Gated on role + re-validated metadata (confabNoticeFromRow): the // `display_kind` string alone is open-ended, so a malformed or imported // non-assistant row carrying it must not be shown as a confirmed catch. - if (confabNoticeFromRow(message)) { + const confabNotice = confabNoticeFromRow(message) + + if (confabNotice) { flushPendingTools(index) result.push({ id: `confab-notice-${message.timestamp || Date.now()}-${index}`, role: 'system', - parts: [textPart(CONFAB_NOTICE_EVENT_TEXT, message.timestamp)], + parts: [textPart(confabNoticeEventText(confabNotice), message.timestamp)], timestamp: message.timestamp }) activeAssistantIndex = null diff --git a/apps/desktop/src/lib/confab-notice-hydration.test.ts b/apps/desktop/src/lib/confab-notice-hydration.test.ts index fcb31731d2f1..1942113d4efb 100644 --- a/apps/desktop/src/lib/confab-notice-hydration.test.ts +++ b/apps/desktop/src/lib/confab-notice-hydration.test.ts @@ -131,3 +131,30 @@ describe('desktop hydration gating', () => { expect(markers([noticeRow({ display_metadata: '{not json' })])).toHaveLength(0) }) }) + +describe('tool-call notice on a metadata-only system row (FleetReview #942)', () => { + // conversation_loop persists a tool-call notice as role=system, content ''. + // The Python reader (notice_from_display_row) accepts that row; an + // assistant-only gate here dropped it on every reload. + const TOOL_NOTICE = { grammar: null, kind: 'tool_call_as_text', request_id: 'r1', scope: 'visible', version: 1 } + const toolRow = (over: Record = {}) => + noticeRow({ content: '', display_metadata: { confab_notice: TOOL_NOTICE }, role: 'system', ...over }) + + it('surfaces the tool-call label from an empty system row', () => { + const out = texts([toolRow()]).filter(m => m.text.includes('Tool call not executed')) + + expect(out).toHaveLength(1) + expect(out[0]?.role).toBe('system') + expect(out[0]?.text).toContain('text was sent instead of a native tool call') + }) + + it('still refuses a scaffold catch on a system row', () => { + expect(markers([toolRow({ display_metadata: { confab_notice: NOTICE } })])).toHaveLength(0) + }) + + it('refuses a tool-call notice with a non-visible scope', () => { + const row = toolRow({ display_metadata: { confab_notice: { ...TOOL_NOTICE, scope: 'both' } } }) + + expect(texts([row]).filter(m => m.text.includes('Tool call not executed'))).toHaveLength(0) + }) +}) diff --git a/apps/shared/src/confab-notice.ts b/apps/shared/src/confab-notice.ts index 7254f67fe9a3..13fe39a7e5c7 100644 --- a/apps/shared/src/confab-notice.ts +++ b/apps/shared/src/confab-notice.ts @@ -13,7 +13,8 @@ * row. Presenting the confirmed-confabulation claim off that string would put * a false accusation on a user or system turn. So a row must clear two gates: * - * 1. `role === 'assistant'` — only a model reply can carry a catch; + * 1. `role === 'assistant'` — only a model reply can carry a scaffold catch; + * a metadata-only `system` event row may carry a TOOL-CALL notice only; * 2. `display_metadata.confab_notice` re-validates against the same * fail-closed v1 schema the wire payload had to pass. * @@ -32,6 +33,17 @@ export const CONFAB_NOTICE_VERSION = 1 /** The only catch kind defined by v1 of the contract. */ export const CONFAB_NOTICE_KIND = 'scaffold_confab_removed' +/** + * Tool-call notice kinds (Python `TOOL_CALL_NOTICE_TEXT`) and their fixed + * labels (`confab_notice_status`). These persist as an EMPTY `system` row, so + * a reader that only accepted assistant rows dropped them on reload + * (FleetReview #942). + */ +export const TOOL_CALL_NOTICE_EVENT_TEXT: Readonly> = { + tool_call_as_text: 'Tool call not executed: text was sent instead of a native tool call.', + tool_call_unparseable: 'Tool call not executed: tool-call JSON could not be parsed.' +} + /** Allowed `scope` values. */ export const CONFAB_NOTICE_SCOPES = ['visible', 'intermediate', 'both'] as const @@ -76,7 +88,9 @@ export function validateConfabNotice(raw: unknown): ConfabNotice | null { return null } - if (candidate.kind !== CONFAB_NOTICE_KIND) { + const isToolCallKind = typeof candidate.kind === 'string' && Object.hasOwn(TOOL_CALL_NOTICE_EVENT_TEXT, candidate.kind) + + if (candidate.kind !== CONFAB_NOTICE_KIND && !isToolCallKind) { return null } @@ -88,6 +102,10 @@ export function validateConfabNotice(raw: unknown): ConfabNotice | null { return null } + if (isToolCallKind && candidate.scope !== 'visible') { + return null + } + // `grammar` is the detector's bounded label, or null when several catches // cannot be represented by one label. Absent is treated as null. const grammar = candidate.grammar @@ -98,7 +116,7 @@ export function validateConfabNotice(raw: unknown): ConfabNotice | null { return { grammar: grammar === undefined || grammar === null ? null : grammar, - kind: CONFAB_NOTICE_KIND, + kind: candidate.kind as string, request_id: candidate.request_id, scope: candidate.scope, version: CONFAB_NOTICE_VERSION @@ -112,7 +130,7 @@ export function validateConfabNotice(raw: unknown): ConfabNotice | null { * JSON text, so parse a string form before reading into it. */ export function confabNoticeFromRow(row: ConfabNoticeRow | null | undefined): ConfabNotice | null { - if (!row || row.role !== 'assistant' || row.display_kind !== CONFAB_NOTICE_DISPLAY_KIND) { + if (!row || (row.role !== 'assistant' && row.role !== 'system') || row.display_kind !== CONFAB_NOTICE_DISPLAY_KIND) { return null } @@ -130,5 +148,17 @@ export function confabNoticeFromRow(row: ConfabNoticeRow | null | undefined): Co return null } - return validateConfabNotice((metadata as Record)[CONFAB_NOTICE_KEY]) + const notice = validateConfabNotice((metadata as Record)[CONFAB_NOTICE_KEY]) + + // Mirrors Python: a system row may only carry a tool-call notice. + if (row.role === 'system' && (!notice || !Object.hasOwn(TOOL_CALL_NOTICE_EVENT_TEXT, notice.kind))) { + return null + } + + return notice +} + +/** The fixed event label for a validated notice. */ +export function confabNoticeEventText(notice: ConfabNotice): string { + return TOOL_CALL_NOTICE_EVENT_TEXT[notice.kind] ?? CONFAB_NOTICE_EVENT_TEXT } diff --git a/cron/jobs.py b/cron/jobs.py index 556ebd77ea3e..893bc11421d9 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -2552,6 +2552,26 @@ def _resolve_auto_model_sentinel( return resolved_model, resolved_provider, True +def _revalidate_resolved_provider(before: Any, provider: Any, base_url: Any) -> None: + """Re-run the tool's base_url/provider exfil guard when store-level model + resolution (auto sentinel, alias, legacy heal) CHANGED the provider. + + ``cronjob_tools`` validates the pair it was handed; the resolution here can + rewrite the provider afterwards, and the stored pair was then never checked + (FleetReview #1032). The scheduler's fire-time backstop still refuses such a + pair, so the job would only fail on every run -- refuse it at write time. + """ + before = _normalize_job_optional_text(before) + after = _normalize_job_optional_text(provider) + if before == after or not _normalize_job_optional_text(base_url): + return + from tools.cronjob_tools import _validate_cron_base_url + + err = _validate_cron_base_url(after, base_url) + if err: + raise ValueError(err) + + def _normalize_job_optional_text(value: Any, *, strip_trailing_slash: bool = False) -> Optional[str]: if not isinstance(value, str): return None @@ -2809,6 +2829,7 @@ def create_job( normalized_model, allow_flagship_reason=allow_flagship_reason ) normalized_base_url = _normalize_job_optional_text(base_url, strip_trailing_slash=True) + _revalidate_resolved_provider(provider, normalized_provider, normalized_base_url) normalized_script = str(script).strip() if isinstance(script, str) else None normalized_script = normalized_script or None normalized_toolsets = [str(t).strip() for t in enabled_toolsets if str(t).strip()] if enabled_toolsets else None @@ -3142,6 +3163,9 @@ def update_job(job_id: str, updates: Dict[str, Any]) -> Optional[Dict[str, Any]] # `--model ""` clears the pin and is deliberately untouched here. # The "auto" sentinel is resolved (LLM job) or dropped (script job) # against the EFFECTIVE no_agent of this update — never stored. + _pre_resolution_provider = ( + updates["provider"] if "provider" in updates else job.get("provider") + ) if "model" in updates and _is_auto_model_sentinel(updates["model"]): _eff_no_agent = updates["no_agent"] if "no_agent" in updates else job.get("no_agent") _am, _ap, _pinned = _resolve_auto_model_sentinel( @@ -3191,6 +3215,9 @@ def update_job(job_id: str, updates: Dict[str, Any]) -> Optional[Dict[str, Any]] if _reason: updated["allow_flagship_reason"] = _reason _legacy_auto_healed = True + _revalidate_resolved_provider( + _pre_resolution_provider, updated.get("provider"), updated.get("base_url") + ) if "allow_flagship_reason" in updates and "model" not in updates: from hermes_cli.model_policy import validate_worker_model diff --git a/gateway/run.py b/gateway/run.py index 3e49fbcb9848..35dbdee85624 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -5191,6 +5191,10 @@ def _is_gateway_hidden_reasoning_incomplete_turn(agent_result: dict) -> bool: return not final_response or final_response == error_text +#: Sentinel: clear whatever resume mark is present (legacy callers). +_ANY_RESUME_MARK = object() + + def _should_clear_resume_pending_after_turn(agent_result: dict) -> bool: """Return True only when a gateway turn really completed successfully. @@ -14466,7 +14470,7 @@ def _clear_restart_replay_marks(self, session_key: str) -> None: entry["armed"] = False counts[session_key] = entry - def _apply_post_turn_resume_gate(self, session_key: str) -> None: + def _apply_post_turn_resume_gate(self, session_key: str, *, marked_at=_ANY_RESUME_MARK) -> None: """Post-(clean-turn) replay-loop gate for the F2 circuit-breaker. Called when a turn completed cleanly enough to clear ``resume_pending``. @@ -14513,9 +14517,10 @@ def _apply_post_turn_resume_gate(self, session_key: str) -> None: ) except Exception: deferred = False + _clear_kw = {} if marked_at is _ANY_RESUME_MARK else {"marked_at": marked_at} if deferred: getattr(self, "_session_initiated_restart", {}).pop(session_key, None) - self.session_store.clear_resume_pending(session_key) + self.session_store.clear_resume_pending(session_key, **_clear_kw) return flag = bool( @@ -14533,7 +14538,7 @@ def _apply_post_turn_resume_gate(self, session_key: str) -> None: # the two must not leave a set marker beside a zeroed counter, which # would hand the session a whole extra budget of unattended replays. try: - self.session_store.clear_resume_pending(session_key) + self.session_store.clear_resume_pending(session_key, **_clear_kw) except Exception as exc: logger.debug("clear_resume_pending failed for %s: %s", session_key, exc) if initiated_restart: @@ -26657,7 +26662,7 @@ async def _notify_stale_lease_holder( _hyg_old_sid = session_entry.session_id try: _hyg_session_row = await self._session_db.get_session( - session_entry.session_id + _hyg_old_sid ) except Exception as exc: _hyg_session_row = None @@ -26704,7 +26709,12 @@ def _build_hyg_agent(): quiet_mode=True, skip_memory=not _hyg_checkpoint_required, enabled_toolsets=["memory"], - session_id=session_entry.session_id, + # The snapshot, never the live entry: this + # runs in a worker thread after awaits, and + # a /new or rotation can move + # session_entry.session_id meanwhile + # (FleetReview #976). + session_id=_hyg_old_sid, session_db=_hyg_session_db, ) _seed_hygiene_system_prompt(_a, _hyg_session_row) @@ -27619,6 +27629,10 @@ def _build_hyg_agent(): # below; a /new or another lifecycle transition may move # session_entry.session_id while the old run is still unwinding. _run_start_session_id = session_entry.session_id + # The resume mark this turn recovers from. A mark written while the + # turn runs (a concurrent shutdown drain) is newer and must survive + # the post-turn clear (FleetReview #1043). + _run_start_resume_marked_at = getattr(session_entry, "last_resume_marked_at", None) # Same rule as the context-prompt pin: an internal event reuses # the last human turn's channel_prompt / parent_chat_id, which # its rebuilt source lacks, so combined_ephemeral cannot toggle. @@ -27736,7 +27750,10 @@ def _build_hyg_agent(): # succeeded and subsequent messages should no longer receive # the restart-interruption system note. if session_key and _should_clear_resume_pending_after_turn(agent_result): - await asyncio.to_thread(self._apply_post_turn_resume_gate, session_key) + await asyncio.to_thread( + self._apply_post_turn_resume_gate, session_key, + marked_at=_run_start_resume_marked_at, + ) # Normalize empty responses: surface errors, partial failures, and # the case where agent did work but returned no text. Fix for #18765. diff --git a/gateway/session.py b/gateway/session.py index 1fd220a2d480..a0d14b2ef609 100644 --- a/gateway/session.py +++ b/gateway/session.py @@ -4706,13 +4706,18 @@ def flush(self) -> None: # durable write, so wait for the queue to drain before returning. self.drain_sessions_json_writes() - def clear_resume_pending(self, session_key: str) -> bool: + def clear_resume_pending(self, session_key: str, **kw) -> bool: """Clear the resume-pending flag after a successful resumed turn. Called from the gateway after ``run_conversation()`` returns a final response for a session that had ``resume_pending=True``, signalling that recovery succeeded. + ``marked_at`` (optional): the ``last_resume_marked_at`` the caller's + turn started from. When given, a mark written AFTER that snapshot (a + concurrent shutdown drain) is left in place — it is not the mark this + turn recovered from (FleetReview #1043). + Returns True if a flag was cleared. """ with self._lock: @@ -4720,6 +4725,12 @@ def clear_resume_pending(self, session_key: str) -> bool: entry = self._entries.get(session_key) if entry is None or not entry.resume_pending: return False + if "marked_at" in kw and entry.last_resume_marked_at != kw["marked_at"]: + logger.info( + "Keeping resume mark for %s: re-marked during the turn (%s -> %s)", + session_key, kw["marked_at"], entry.last_resume_marked_at, + ) + return False self._clear_resume_pending_entry(entry) self._save() return True diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index 3d84ef4ee0b0..f7628b8b3e16 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -16182,6 +16182,12 @@ def cooling_providers( continue if not isinstance(payload, dict) or payload.get("rate_limited") is not True: continue + # A runtime-stage refusal is one mid-turn 429 on a pinned (often + # pooled) provider; the worker keeps retrying it. Only an auth-stage + # refusal (the provider was unusable at startup) cools the provider + # (FleetReview #1198). Legacy events without a stage were auth. + if payload.get("stage", "auth") != "auth": + continue provider = payload.get("provider") if not isinstance(provider, str) or not provider.strip(): continue diff --git a/hermes_cli/model_switch.py b/hermes_cli/model_switch.py index c0de75f73e0b..7c53b17283e5 100644 --- a/hermes_cli/model_switch.py +++ b/hermes_cli/model_switch.py @@ -1941,7 +1941,7 @@ def switch_model( is_global=is_global, error_message=( f"Provider '{_explicit_norm}' is an alias that routes " - f"through {get_label(target_provider)}, which " + f"through {get_label(target_provider, allow_network=probe_catalog)}, which " f"has no credentials configured.{_hint}" ), ) @@ -2191,7 +2191,7 @@ def switch_model( # ================================================================= provider_changed = target_provider != current_provider - provider_label = get_label(target_provider) + provider_label = get_label(target_provider, allow_network=probe_catalog) if target_provider == "custom" and current_base_url: provider_label = "Custom endpoint" if target_provider.startswith("custom:"): @@ -2259,7 +2259,7 @@ def switch_model( elif target_provider == "custom" and current_base_url: api_key = current_api_key base_url = current_base_url - api_mode = determine_api_mode(target_provider, base_url) + api_mode = determine_api_mode(target_provider, base_url, allow_network=probe_catalog) else: try: runtime = resolve_runtime_provider( @@ -2318,7 +2318,7 @@ def switch_model( # provider, causing validation to probe the wrong model-list URL. api_key = current_api_key or "no-key-required" base_url = current_base_url - api_mode = determine_api_mode(current_provider, base_url) + api_mode = determine_api_mode(current_provider, base_url, allow_network=probe_catalog) validation_headers = ollama_headers else: try: @@ -2394,7 +2394,7 @@ def switch_model( if _mandated_mode is not None: api_mode = _mandated_mode elif not api_mode: - api_mode = determine_api_mode(target_provider, base_url) + api_mode = determine_api_mode(target_provider, base_url, allow_network=probe_catalog) # --- Normalize model name for target provider --- new_model = _resolve_named_custom_model_id( @@ -2564,7 +2564,7 @@ def switch_model( # --- Determine api_mode if not already set --- if not api_mode: api_mode = determine_api_mode( - target_provider, base_url, model=new_model + target_provider, base_url, model=new_model, allow_network=probe_catalog ) # OpenCode base URLs end with /v1 for OpenAI-compatible models, but the diff --git a/hermes_cli/process_env_files.py b/hermes_cli/process_env_files.py index 717fd14da5c1..8cc07e3186a4 100644 --- a/hermes_cli/process_env_files.py +++ b/hermes_cli/process_env_files.py @@ -16,7 +16,8 @@ untouched and logs a warning. Process start never fails because of it. * The applied diff is remembered so a spawn site that must NOT carry it (cron script children are plain scripts, not agent processes) can undo it - with :func:`strip_overlay`. + with :func:`strip_overlay`. The pre-apply values ride to child agent + processes in ``HERMES_PROCESS_ENV_OVERLAY`` so a child can undo them too. """ from __future__ import annotations @@ -37,6 +38,23 @@ # key -> (value before apply or None if unset, value after apply or None if unset) _OVERLAY: Dict[str, Tuple[Optional[str], Optional[str]]] = {} +# Carries the overlay's pre-apply values ({key: old}) to child agent processes. +# A child inherits the ALREADY-overlaid env, so its own diff is empty for every +# key the parent set and strip_overlay could never undo them (FleetReview #1254 +# :121). Only pre-overlay values ride here; the overlaid ones are in the env. +_INHERITED_ENV = "HERMES_PROCESS_ENV_OVERLAY" + + +def _inherited_overlay(env: MutableMapping[str, str]) -> Dict[str, Tuple[Optional[str], Optional[str]]]: + try: + raw = json.loads(env.get(_INHERITED_ENV) or "{}") + except (TypeError, ValueError): + return {} + if not isinstance(raw, dict): + return {} + return {str(k): (None if v is None else str(v), env.get(str(k))) + for k, v in raw.items() if v is None or isinstance(v, str)} + def configured_files(cfg: Optional[dict]) -> List[str]: """Expanded, existing paths from ``agent.process_env_files`` (order kept).""" @@ -114,16 +132,23 @@ def apply_process_env_files( (default ``os.environ``). Returns the applied diff; ``{}`` when nothing is configured or sourcing failed.""" target = os.environ if env is None else env + inherited = _inherited_overlay(target) diff = compute_overlay(configured_files(cfg), target) for key, (_old, new) in diff.items(): if new is None: target.pop(key, None) else: target[key] = new - if env is None and diff: + if env is None and (diff or inherited): + # An inherited key keeps the ORIGINAL pre-overlay value, even if this + # process's files changed it again. + merged = {**diff, **{k: (old, diff[k][1] if k in diff else cur) + for k, (old, cur) in inherited.items()}} _OVERLAY.clear() - _OVERLAY.update(diff) - logger.info("agent.process_env_files applied: %s", ",".join(sorted(diff))) + _OVERLAY.update(merged) + target[_INHERITED_ENV] = json.dumps({k: old for k, (old, _new) in merged.items()}) + if diff: + logger.info("agent.process_env_files applied: %s", ",".join(sorted(diff))) return diff @@ -132,6 +157,7 @@ def strip_overlay(env: MutableMapping[str, str]) -> MutableMapping[str, str]: A key is restored only while it still holds the value the overlay set, so a later deliberate change to that key is never reverted.""" + env.pop(_INHERITED_ENV, None) for key, (old, new) in _OVERLAY.items(): if key == "PATH" and new is not None and env.get("PATH"): # PATH is routinely re-edited after start (venv/tool dirs), so match diff --git a/hermes_cli/providers.py b/hermes_cli/providers.py index 71524f6905d5..379527054578 100644 --- a/hermes_cli/providers.py +++ b/hermes_cli/providers.py @@ -600,8 +600,11 @@ def get_provider(name: str, *, allow_network: bool = True) -> Optional[ProviderD return None -def get_label(provider_id: str) -> str: - """Get a human-readable display name for a provider.""" +def get_label(provider_id: str, *, allow_network: bool = True) -> str: + """Get a human-readable display name for a provider. + + ``allow_network=False`` never fetches models.dev (a cold cache falls back + to the canonical id); see ``switch_model(probe_catalog=False)``.""" canonical = normalize_provider(provider_id) # Check label overrides first @@ -609,7 +612,7 @@ def get_label(provider_id: str) -> str: return _LABEL_OVERRIDES[canonical] # Try models.dev - pdef = get_provider(canonical) + pdef = get_provider(canonical) if allow_network else get_provider(canonical, allow_network=False) if pdef: return pdef.name @@ -777,7 +780,9 @@ def nous_api_mode(model: str = "") -> str: return "chat_completions" -def determine_api_mode(provider: str, base_url: str = "", model: str = "") -> str: +def determine_api_mode( + provider: str, base_url: str = "", model: str = "", *, allow_network: bool = True +) -> str: """Determine the API mode (wire protocol) for a provider/endpoint. Resolution order: @@ -802,7 +807,7 @@ def determine_api_mode(provider: str, base_url: str = "", model: str = "") -> st if provider_norm in {"nous", "nous-portal", "nousresearch"}: return nous_api_mode(model) - pdef = get_provider(provider) + pdef = get_provider(provider) if allow_network else get_provider(provider, allow_network=False) if pdef is not None: return TRANSPORT_TO_API_MODE.get(pdef.transport, "chat_completions") diff --git a/providers/__init__.py b/providers/__init__.py index 449837549491..eee832d706ea 100644 --- a/providers/__init__.py +++ b/providers/__init__.py @@ -100,11 +100,27 @@ def register_provider(profile: ProviderProfile) -> None: plugins under ``$HERMES_HOME/plugins/model-providers/`` can override bundled profiles without editing repo code. """ + _register(profile) + # Remember which module registered it, so discovery can replay the + # registration for a module that was imported before discovery ran. + try: + mod = sys._getframe(1).f_globals.get("__name__") + except ValueError: # pragma: no cover - no caller frame + mod = None + if mod: + _PROFILES_BY_MODULE.setdefault(mod, []).append(profile) + + +def _register(profile: ProviderProfile) -> None: _REGISTRY[profile.name] = profile for alias in profile.aliases: _ALIASES[alias] = profile.name +#: module name -> profiles it registered at import (see _import_plugin_dir). +_PROFILES_BY_MODULE: dict[str, list[ProviderProfile]] = {} + + def get_provider_profile(name: str) -> ProviderProfile | None: """Look up a provider profile by name or alias. @@ -170,7 +186,14 @@ def _import_plugin_dir(plugin_dir: Path, source: str) -> None: module_name = f"_hermes_user_provider_{safe_name}" if module_name in sys.modules: - return # already imported + # Already imported (e.g. a direct import of a bundled profile class + # before discovery). Its import-time registration ran BEFORE the + # entry-point scan, so a pip plugin of the same name may have replaced + # it since. Replay it here to keep filesystem-over-pip precedence + # (FleetReview #1095). + for profile in _PROFILES_BY_MODULE.get(module_name, ()): + _register(profile) + return try: spec = importlib.util.spec_from_file_location( diff --git a/tests/gateway/test_c3_followup_t_7da6cadf.py b/tests/gateway/test_c3_followup_t_7da6cadf.py new file mode 100644 index 000000000000..29eb240ec231 --- /dev/null +++ b/tests/gateway/test_c3_followup_t_7da6cadf.py @@ -0,0 +1,199 @@ +"""C3 re-bucket follow-up (card t_7da6cadf): one pin per fixed FleetReview row. + +#970 has its own file (test_no_network_cold_models_dev.py) and #942 is a +desktop vitest (apps/desktop/src/lib/confab-notice-hydration.test.ts). +""" +from __future__ import annotations + +import ast +import json +import os +import sqlite3 +import sys +import textwrap +import types +from datetime import datetime, timedelta +from pathlib import Path + +import pytest + +REPO = Path(__file__).resolve().parents[2] + + +# --------------------------------------------------------------------------- #976 +def test_976_hygiene_agent_build_uses_the_session_id_snapshot(): + """_build_hyg_agent runs in a worker thread after awaits: it must read the + _hyg_old_sid snapshot, never the live session_entry (a /new or rotation can + move session_entry.session_id meanwhile).""" + tree = ast.parse((REPO / "gateway" / "run.py").read_text(encoding="utf-8")) + builds = [n for n in ast.walk(tree) if isinstance(n, ast.FunctionDef) and n.name == "_build_hyg_agent"] + assert len(builds) == 1 + live = [n.lineno for n in ast.walk(builds[0]) + if isinstance(n, ast.Name) and n.id == "session_entry"] + assert live == [], f"_build_hyg_agent reads the live session_entry at gateway/run.py:{live}" + + +# --------------------------------------------------------------------------- #1043 +def _store(tmp_path): + from gateway.config import GatewayConfig + from gateway.session import SessionSource, SessionStore + from gateway.config import Platform + + store = SessionStore(sessions_dir=tmp_path, config=GatewayConfig()) + entry = store.get_or_create_session(SessionSource(platform=Platform.TELEGRAM, chat_id="1", user_id="u")) + return store, entry.session_key + + +def test_1043_post_turn_clear_keeps_a_mark_written_during_the_turn(tmp_path): + store, key = _store(tmp_path) + assert store.mark_resume_pending(key, "restart_timeout") + at_turn_start = store._entries[key].last_resume_marked_at + # A concurrent shutdown drain re-marks the session while the turn runs. + store._entries[key].last_resume_marked_at = at_turn_start + timedelta(seconds=5) + assert store.clear_resume_pending(key, marked_at=at_turn_start) is False + assert store._entries[key].resume_pending is True + # The mark the turn started from is cleared as before. + assert store.clear_resume_pending(key, marked_at=store._entries[key].last_resume_marked_at) is True + assert store._entries[key].resume_pending is False + + +def test_1043_gate_passes_the_turn_start_mark(tmp_path): + """The runner hands the snapshot through; legacy callers (no marked_at) + keep the old clear-whatever-is-there behaviour.""" + store, key = _store(tmp_path) + store.mark_resume_pending(key, "restart_timeout") + fresh = store._entries[key].last_resume_marked_at + assert store.clear_resume_pending(key) is True # legacy path unchanged + store.mark_resume_pending(key, "restart_timeout") + assert store.clear_resume_pending(key, marked_at=fresh - timedelta(seconds=1)) is False + src = (REPO / "gateway" / "run.py").read_text(encoding="utf-8") + assert "marked_at=_run_start_resume_marked_at" in src + + +# --------------------------------------------------------------------------- #1095 +def test_1095_early_imported_bundled_profile_is_re_registered(tmp_path, monkeypatch): + import providers + from providers.base import ProviderProfile + + name = "c3f-probe" + plugin = tmp_path / "c3f_probe" + plugin.mkdir() + (plugin / "__init__.py").write_text(textwrap.dedent(f""" + from providers import register_provider + from providers.base import ProviderProfile + bundled = ProviderProfile(name="{name}", base_url="https://bundled.example/v1") + register_provider(bundled) + """)) + mod_name = "plugins.model_providers.c3f_probe" + monkeypatch.delitem(sys.modules, mod_name, raising=False) + # 1. a direct early import (before discovery) registers the bundled profile + import importlib.util + spec = importlib.util.spec_from_file_location(mod_name, plugin / "__init__.py", + submodule_search_locations=[str(plugin)]) + mod = importlib.util.module_from_spec(spec) + monkeypatch.setitem(sys.modules, mod_name, mod) + spec.loader.exec_module(mod) + try: + # 2. a pip entry-point plugin of the same name registers during discovery + providers.register_provider(ProviderProfile(name=name, base_url="https://pip.example/v1")) + assert providers._REGISTRY[name].base_url == "https://pip.example/v1" + # 3. the bundled step reaches the already-imported module: bundled wins again + providers._import_plugin_dir(plugin, "bundled") + assert providers._REGISTRY[name].base_url == "https://bundled.example/v1" + finally: + # Same restore as tests/providers/test_entry_point_discovery.py: the + # registry is additive, so reset it and re-run real discovery. + from hermes_cli import provider_seam + + monkeypatch.delitem(sys.modules, mod_name, raising=False) + providers._PROFILES_BY_MODULE.pop(mod_name, None) + provider_seam._reset("_REGISTRY", "_ALIASES") + providers._discovered = False + for m in [m for m in sys.modules if m.startswith("plugins.model_providers.")]: + del sys.modules[m] + providers._discover_providers() + + +# --------------------------------------------------------------------------- #1198 +def test_1198_runtime_stage_refusals_do_not_cool_the_provider(): + from hermes_cli import kanban_db as kb + + conn = sqlite3.connect(":memory:") + conn.row_factory = sqlite3.Row + conn.execute("CREATE TABLE task_runs (id INTEGER PRIMARY KEY, ended_at INTEGER)") + conn.execute("CREATE TABLE task_events (run_id INTEGER, kind TEXT, payload TEXT, created_at INTEGER)") + now = 10_000 + conn.execute("INSERT INTO task_runs VALUES (1, NULL), (2, NULL), (3, NULL)") + ev = "worker_route_pin_refused" + rows = [ + (1, ev, json.dumps({"stage": "runtime", "provider": "claude-bpr", "rate_limited": True}), now - 10), + (2, ev, json.dumps({"stage": "auth", "provider": "openai-codex", "rate_limited": True}), now - 10), + (3, ev, json.dumps({"provider": "legacy-prov", "rate_limited": True}), now - 10), # pre-stage events + ] + conn.executemany("INSERT INTO task_events VALUES (?, ?, ?, ?)", rows) + cooling = kb.cooling_providers(conn, now=now, window=600) + assert "claude-bpr" not in cooling + assert cooling["openai-codex"] == now - 10 + 600 + assert "legacy-prov" in cooling + + +# --------------------------------------------------------------------------- #1254 :121 +@pytest.fixture +def _pef(): + from hermes_cli import process_env_files as pef + + saved = dict(pef._OVERLAY) + pef._OVERLAY.clear() + yield pef + pef._OVERLAY.clear() + pef._OVERLAY.update(saved) + + +def test_1254_child_process_can_strip_an_inherited_overlay(tmp_path, monkeypatch, _pef): + f = tmp_path / "lane.sh" + f.write_text("export PEF_C3F_TOKEN=lane\n") + cfg = {"agent": {"process_env_files": [str(f)]}} + monkeypatch.delenv("PEF_C3F_TOKEN", raising=False) + monkeypatch.delenv(_pef._INHERITED_ENV, raising=False) + # Parent agent process applies the overlay. + _pef.apply_process_env_files(cfg) + try: + assert os.environ["PEF_C3F_TOKEN"] == "lane" + inherited = dict(os.environ) # what a child agent process starts with + # Child agent process: same config, already-overlaid env -> its own diff is empty. + _pef._OVERLAY.clear() + monkeypatch.setattr(os, "environ", inherited) + assert _pef.apply_process_env_files(cfg) == {} + # A plain script spawned by the CHILD must still lose the overlay. + script_env = dict(inherited) + _pef.strip_overlay(script_env) + assert "PEF_C3F_TOKEN" not in script_env + assert _pef._INHERITED_ENV not in script_env + finally: + monkeypatch.undo() + os.environ.pop("PEF_C3F_TOKEN", None) + os.environ.pop(_pef._INHERITED_ENV, None) + + +# --------------------------------------------------------------------------- #1032 +def test_1032_resolution_that_changes_the_provider_revalidates_base_url(monkeypatch): + from cron import jobs + import tools.cronjob_tools as ct + + seen = [] + + def _validate(provider, base_url): + seen.append((provider, base_url)) + return "blocked" if provider == "anthropic" else None + + monkeypatch.setattr(ct, "_validate_cron_base_url", _validate) + # provider rewritten by resolution (custom -> anthropic) with an off-host base_url + with pytest.raises(ValueError, match="blocked"): + jobs._revalidate_resolved_provider("my-custom", "anthropic", "https://evil.example/v1") + # unchanged provider: the tool already validated this pair + jobs._revalidate_resolved_provider("anthropic", "anthropic", "https://evil.example/v1") + # no base_url: nothing to exfiltrate to + jobs._revalidate_resolved_provider("my-custom", "anthropic", None) + assert seen == [("anthropic", "https://evil.example/v1")] + src = (REPO / "cron" / "jobs.py").read_text(encoding="utf-8") + assert src.count("_revalidate_resolved_provider(") == 3 # def + create_job + update_job diff --git a/tests/gateway/test_no_network_cold_models_dev.py b/tests/gateway/test_no_network_cold_models_dev.py new file mode 100644 index 000000000000..fb6e6c5fb324 --- /dev/null +++ b/tests/gateway/test_no_network_cold_models_dev.py @@ -0,0 +1,55 @@ +"""``switch_model(probe_catalog=False)`` opens no socket on a COLD models.dev cache. + +FleetReview #970 (C3 re-bucket, t_7da6cadf): the sibling arm in +test_no_network_reachable_from_loop.py stubs ``fetch_models_dev`` out, so the +cold-cache path (no memory or disk cache: stage 4, a singleflight FOREGROUND +fetch) was never exercised. That path is reachable from the event loop on a +fresh install or after the cache file is removed. Here the real +``fetch_models_dev`` runs against an empty cache and every connect raises. +""" +from __future__ import annotations + +import socket + + +def test_switch_model_probe_off_cold_models_dev_cache_opens_no_socket(monkeypatch, tmp_path): + import agent.models_dev as mdev + from hermes_cli import model_switch + + attempts: list = [] + + def _no_network(*args, **kwargs): + attempts.append(args[:1]) + raise OSError("network disabled in test") + + # Cold: no in-memory cache, no disk cache, no failure backoff. + monkeypatch.setattr(mdev, "_models_dev_cache", {}, raising=False) + monkeypatch.setattr(mdev, "_models_dev_cache_time", 0, raising=False) + monkeypatch.setattr(mdev, "_models_dev_retry_after", 0, raising=False) + monkeypatch.setattr(mdev, "_get_cache_path", lambda: tmp_path / "models_dev_cache.json", raising=False) + fetched: list = [] + real_fetch = mdev.fetch_models_dev + + def _spy_fetch(*a, **k): + fetched.append(k.get("allow_network", True)) + return real_fetch(*a, **k) + + monkeypatch.setattr(mdev, "fetch_models_dev", _spy_fetch) + monkeypatch.setattr(socket, "create_connection", _no_network) + monkeypatch.setattr(socket.socket, "connect", lambda self, addr: _no_network(addr)) + + user_providers = { + "local-test": { + "name": "local-test", + "base_url": "http://127.0.0.1:9/v1", + "api_key": "sk-test", + "models": {"m-1": {}}, + } + } + model_switch.switch_model( + raw_input="m-1", current_provider="local-test", current_model="m-1", + explicit_provider="local-test", user_providers=user_providers, + custom_providers=[], probe_catalog=False, + ) + assert True not in fetched, f"probe_catalog=False fetched models.dev with network allowed: {fetched}" + assert attempts == [], f"probe_catalog=False opened a socket on a cold cache: {attempts}" diff --git a/tests/hermes_cli/test_process_env_files.py b/tests/hermes_cli/test_process_env_files.py index e24a1c331374..d1bfd3a33a6e 100644 --- a/tests/hermes_cli/test_process_env_files.py +++ b/tests/hermes_cli/test_process_env_files.py @@ -16,6 +16,7 @@ def _clean_overlay(): yield pef._OVERLAY.clear() pef._OVERLAY.update(saved) + os.environ.pop(pef._INHERITED_ENV, None) # written by an env=None apply def _write(path, text): From 99cc8fcdd958e085f324d501685d971429a5310c Mon Sep 17 00:00:00 2001 From: "ang-fleet-workers[bot]" <333956806+ang-fleet-workers[bot]@users.noreply.github.com> Date: Sun, 27 Sep 2026 19:28:03 -0700 Subject: [PATCH 2/2] test: count the marked_at form of the single resume-gate call site (#1043, t_7da6cadf) --- tests/gateway/test_restart_cascade.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/tests/gateway/test_restart_cascade.py b/tests/gateway/test_restart_cascade.py index 33b43de2b15d..543f91d135e3 100644 --- a/tests/gateway/test_restart_cascade.py +++ b/tests/gateway/test_restart_cascade.py @@ -792,9 +792,10 @@ def test_single_gate_call_site(): src = inspect.getsource(gr) # Direct call or offloaded via asyncio.to_thread(self._apply_..., key). + # (t_7da6cadf: the offloaded call now also passes marked_at=..., FleetReview #1043.) n = src.count("self._apply_post_turn_resume_gate(session_key)") + src.count( "self._apply_post_turn_resume_gate, session_key)" - ) + ) + src.count("self._apply_post_turn_resume_gate, session_key,") assert n == 1, f"expected 1 gate call site, found {n}"