Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 5 additions & 3 deletions apps/desktop/src/lib/chat-messages/hydration.ts
Original file line number Diff line number Diff line change
@@ -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'
Expand Down Expand Up @@ -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
Expand Down
27 changes: 27 additions & 0 deletions apps/desktop/src/lib/confab-notice-hydration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown> = {}) =>
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)
})
})
40 changes: 35 additions & 5 deletions apps/shared/src/confab-notice.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*
Expand All @@ -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<Record<string, string>> = {
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

Expand Down Expand Up @@ -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
}

Expand All @@ -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
Expand All @@ -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
Expand All @@ -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
}

Expand All @@ -130,5 +148,17 @@ export function confabNoticeFromRow(row: ConfabNoticeRow | null | undefined): Co
return null
}

return validateConfabNotice((metadata as Record<string, unknown>)[CONFAB_NOTICE_KEY])
const notice = validateConfabNotice((metadata as Record<string, unknown>)[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
}
27 changes: 27 additions & 0 deletions cron/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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

Expand Down
29 changes: 23 additions & 6 deletions gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down Expand Up @@ -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``.
Expand Down Expand Up @@ -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(
Expand All @@ -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:
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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.
Expand Down
13 changes: 12 additions & 1 deletion gateway/session.py
Original file line number Diff line number Diff line change
Expand Up @@ -4706,20 +4706,31 @@ 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:
self._ensure_loaded_locked()
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
Expand Down
6 changes: 6 additions & 0 deletions hermes_cli/kanban_db.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
12 changes: 6 additions & 6 deletions hermes_cli/model_switch.py
Original file line number Diff line number Diff line change
Expand Up @@ -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}"
),
)
Expand Down Expand Up @@ -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:"):
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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
Expand Down
Loading
Loading