Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
54 commits
Select commit Hold shift + click to select a range
eadc85d
feat(gateway): allow operators to suppress engine warning notifications
victor-kyriazakos Sep 15, 2026
ef42c37
fix(gateway): close warning policy coverage and malformed config gaps
victor-kyriazakos Sep 15, 2026
6eb988e
feat(notifications): classify diagnostics and carry suppression throu…
victor-kyriazakos Sep 16, 2026
d90c76f
fix(cron): preserve admitted Bot Chat outcomes across policy changes
victor-kyriazakos Sep 16, 2026
4ab3358
fix(api): separate diagnostic presentation from execution outcomes
victor-kyriazakos Sep 16, 2026
4a96c44
fix(api): bind notification category to replay identity
victor-kyriazakos Sep 16, 2026
3e8f8c0
test(notifications): cover Signal pacing and Slack upload producers
victor-kyriazakos Sep 16, 2026
4419894
fix(cron): preserve suppression disposition and replay identity
victor-kyriazakos Sep 16, 2026
21723e0
fix(gateway): preserve owner scope and truthful failure outcomes
victor-kyriazakos Sep 16, 2026
9ec6a84
fix(notifications): preserve diagnostic observers at presentation bou…
victor-kyriazakos Sep 16, 2026
6ac5ee6
fix(discord): scope operator alerts through policy and delivery
victor-kyriazakos Sep 16, 2026
fe0d4f3
test(gateway): cover watchdog policy and update turn fixtures
victor-kyriazakos Sep 16, 2026
11b6e0c
fix(adapters): preserve captions when suppressing media diagnostics
victor-kyriazakos Sep 16, 2026
d738176
docs: describe opt-in diagnostic notification suppression
victor-kyriazakos Sep 16, 2026
7b11661
test(gateway): supply real source to watchdog review fixture
victor-kyriazakos Sep 16, 2026
57505cf
test(cli): align notification scope fixtures with real turns
victor-kyriazakos Sep 16, 2026
703b651
fix(telegram): bind inbound media diagnostics to routed owner
victor-kyriazakos Sep 16, 2026
2e6c67f
fix(agent): classify remaining core and retry diagnostics
victor-kyriazakos Sep 16, 2026
79ac511
fix(gateway): retain diagnostic observers on muted turns
victor-kyriazakos Sep 16, 2026
3b1fd71
fix(notifications): snapshot CLI and TUI foreground policy
victor-kyriazakos Sep 16, 2026
fa4e3d2
test(notifications): verify TUI snapshot cleanup after errors
victor-kyriazakos Sep 16, 2026
0d49ec3
test(notifications): exercise native media owner policy matrix
victor-kyriazakos Sep 16, 2026
1b4f376
test(notifications): cover credit capture through frontend sinks
victor-kyriazakos Sep 16, 2026
d3c6f4f
fix(cron): reconcile deferred delivery projections by execution
victor-kyriazakos Sep 16, 2026
75958b7
fix(weixin): suppress optional voice downgrade notice
victor-kyriazakos Sep 16, 2026
bae30df
fix(cron): fence late alerts to incident occurrence
victor-kyriazakos Sep 16, 2026
5268815
test(notifications): cover hygiene abort and fallback lifecycle
victor-kyriazakos Sep 16, 2026
f657c9e
test(notifications): compose warnings with native Slack stream
victor-kyriazakos Sep 16, 2026
ccf3992
fix(cron): preserve pending siblings and unprojected receipts
victor-kyriazakos Sep 16, 2026
063da93
fix(cron): bind incident occurrence atomically and fence old writers
victor-kyriazakos Sep 16, 2026
4c40395
refactor(notifications): centralize channel warning emission
victor-kyriazakos Sep 16, 2026
e64063e
refactor(notifications): share UI rendering policy boundary
victor-kyriazakos Sep 16, 2026
9730a3f
refactor(notifications): route gateway renderers through shared boundary
victor-kyriazakos Sep 16, 2026
71bb1b3
refactor(notifications): project mixed adapter diagnostics through wa…
victor-kyriazakos Sep 16, 2026
2aa8465
refactor(notifications): present lifecycle diagnostics through shared…
victor-kyriazakos Sep 16, 2026
a3885be
refactor(notifications): route remaining foreground sinks through ren…
victor-kyriazakos Sep 16, 2026
1d219e0
refactor(notifications): single admission rule for diagnostic-muted t…
victor-kyriazakos Sep 16, 2026
8024bde
fix(cron): park delivery manifest on receipts when the ledger write f…
victor-kyriazakos Sep 16, 2026
fb759b6
fix(cron,notifications): close Salt R3 findings B1-B4, S1-S4
victor-kyriazakos Sep 17, 2026
37afcab
fix(cron): pre-send manifest intent is the authoritative projection f…
victor-kyriazakos Sep 17, 2026
f6fb4c0
fix(cron): every exit after manifest intent settles; legacy journals …
victor-kyriazakos Sep 17, 2026
cb4ab15
fix(cron): migration-completion fence is the legacy-intent adoption a…
victor-kyriazakos Sep 17, 2026
cab8a64
fix(cron): account for pre-flag writers that publish after adoption
victor-kyriazakos Sep 17, 2026
b632076
fix(cron): late-intent adoption rechecks the row at write time
victor-kyriazakos Sep 17, 2026
3224cb1
fix(cron): finish derives its terminal decision from the authoritativ…
victor-kyriazakos Sep 17, 2026
f5b807e
fix(delegate): sync child failure echo honors warning-notification po…
victor-kyriazakos Sep 17, 2026
24ab0a9
fix(gateway): route remaining automatic gateway diagnostics through t…
victor-kyriazakos Sep 17, 2026
66d1cd6
fix(agent,cli,tui): classify remaining CLI/TUI/stream automatic diagn…
victor-kyriazakos Sep 17, 2026
83e6cba
fix(tui): preview-restart status rows honor warning-notification policy
victor-kyriazakos Sep 17, 2026
7301691
fix(hindsight): root-guard startup warning honors warning-notificatio…
victor-kyriazakos Sep 17, 2026
29043ab
fix(cli): import render_notification for the vision fallback notices
victor-kyriazakos Sep 17, 2026
43fae0a
Merge origin/main into feat/user-channel-warning-suppression
victor-kyriazakos Sep 17, 2026
3ebf42f
test(agent): welcome-identity stub agent grows the diagnostic status …
victor-kyriazakos Sep 17, 2026
f45c640
test(cron): R7 initializer probe stops asserting WAL journal mode
victor-kyriazakos Sep 17, 2026
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
24 changes: 19 additions & 5 deletions agent/agent_init.py
Original file line number Diff line number Diff line change
Expand Up @@ -952,7 +952,9 @@ def _init_openai_client(agent, api_key, base_url, fallback_model, _provider_time
print(f"🤖 AI Agent initialized with model: {agent.model}")
if base_url:
print(f"🔗 Using custom base URL: {base_url}")
_print_key_banner(client_kwargs.get("api_key", "none"), "API key", warn_missing=True)
from gateway.warning_notifications import warning_notifications_enabled
_print_key_banner(client_kwargs.get("api_key", "none"), "API key",
warn_missing=warning_notifications_enabled(agent.platform))
except Exception as e:
raise RuntimeError(f"Failed to initialize OpenAI client: {e}")

Expand Down Expand Up @@ -1097,7 +1099,7 @@ def _load_tools(agent, enabled_toolsets, disabled_toolsets):
requirements = model_tools.check_toolset_requirements()
missing_reqs = [name for name, available in requirements.items() if not available]
if missing_reqs:
print(f"⚠️ Some tools may not work due to missing requirements: {missing_reqs}")
agent._safe_print(f"⚠️ Some tools may not work due to missing requirements: {missing_reqs}", diagnostic=True)
else:
print("🛠️ No tools loaded (all tools filtered out or unavailable)")
if agent.save_trajectories:
Expand Down Expand Up @@ -1522,12 +1524,23 @@ def _parse_compression_config(agent, _agent_cfg) -> CompressionSettings:

def _warn_invalid_config_int(
what: str, value: Any, requirement: str, fallback: str, print_fallback: str = "",
agent: Any = None,
) -> None:
"""Log + stderr-print an invalid integer config value (``print_fallback``: user-facing
wording where it differs from the log line)."""
wording where it differs from the log line). The print is an automatic diagnostic and
honors the warning-notification policy; the log line never does."""
_ra().logger.warning(
"Invalid %s: %r — %s. Falling back to %s.", what, value, requirement, fallback,
)
from gateway.warning_notifications import warning_notifications_enabled
try:
if not warning_notifications_enabled(
getattr(agent, "_notification_platform", getattr(agent, "platform", "cli")),
getattr(agent, "_notification_config", None),
):
return
except Exception:
pass
print(
f"\n⚠ Invalid {what}: {value!r}\n"
f" {requirement[0].upper() + requirement[1:]}.\n"
Expand Down Expand Up @@ -1681,6 +1694,7 @@ def _warn_invalid_custom_provider_context_length(agent, _custom_providers) -> No
_warn_invalid_config_int(
f"context_length for model {agent.model!r} in custom_providers",
_cp_ctx, _CTX_LEN_REQUIREMENT, "auto-detection", "auto-detected context window",
agent=agent,
)
return

Expand Down Expand Up @@ -1708,7 +1722,7 @@ def _resolve_context_length(agent, _agent_cfg, base_url):
_warn_invalid_config_int(
"model.context_length in config.yaml", _config_context_length,
"must be a plain integer (e.g. 256000, not '256K')",
"auto-detection", "auto-detected context window",
"auto-detection", "auto-detected context window", agent=agent,
)
_config_context_length = None

Expand Down Expand Up @@ -2083,7 +2097,7 @@ def _emit_compression_summary(agent, cs):
print(f"📊 Context limit: {_cc.context_length:,} tokens (auto-compression disabled)")
# Gateway users get the same text via _compression_warning on turn 1.
if _autoraise_notice:
print(_autoraise_notice)
agent._safe_print(_autoraise_notice, diagnostic=True)

# status_callback isn't wired yet: stash for replay on the first turn; mark shown so
# repeated inits stay silent.
Expand Down
4 changes: 2 additions & 2 deletions agent/agent_runtime_helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -970,7 +970,7 @@ def try_recover_primary_transport(
wait_time = min(3 + retry_count, 8)
agent._vprint(
f"{agent.log_prefix}🔁 Transient {error_type} on {agent.provider} — "
f"rebuilt client, waiting {wait_time}s before one last primary attempt.", force=True,
f"rebuilt client, waiting {wait_time}s before one last primary attempt.", force=True, diagnostic=True,
)
time.sleep(wait_time)
return True
Expand Down Expand Up @@ -1195,7 +1195,7 @@ def _load_primary_pool():
if provider_fallback_active:
# Notification surfaces are best-effort and must never undo a successful restore.
with contextlib.suppress(Exception):
agent._emit_status(
agent._emit_diagnostic_status(
f"✅ Primary model restored: {agent.model} via {agent.provider}; "
f"fallback {previous_model} via {previous_provider} is no longer active."
)
Expand Down
27 changes: 15 additions & 12 deletions agent/chat_completion_helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -589,7 +589,7 @@ def _report_stale_nonstream_kill(agent, api_kwargs: dict, elapsed: float, stale_
"model=%s context=~%s tokens. Killing connection.", "Inline n" if inline else "N", elapsed,
stale_timeout, model, f"{estimate_request_context_tokens(api_kwargs):,}")
try:
agent._buffer_status(
agent._buffer_diagnostic_status(
f"⚠️ No response from provider for {int(elapsed)}s (non-streaming, model: {model}). {hint or 'Aborting call.'}")
except Exception:
logger.debug("stale status buffering failed", exc_info=True)
Expand Down Expand Up @@ -1894,7 +1894,7 @@ def _rescope_fallback_extra_body(agent, old_model: str, old_provider: str, old_b
def _buffer_fallback_notice(agent, notice: str) -> None:
"""Buffer the switch notice for terminal failure AND retain it as a durable one-shot for
_emit_pending_fallback_notice (a successful fallback clears retry chatter)."""
agent._buffer_status(notice)
agent._buffer_diagnostic_status(notice)
pending = getattr(agent, "_pending_fallback_notice", None)
if isinstance(pending, list):
pending.append(notice)
Expand Down Expand Up @@ -2164,7 +2164,7 @@ def handle_max_iterations(agent, messages: list, api_call_count: int) -> str:
# stdout so wrappers receive only the final assistant content (#93220 class).
logger.warning(warning)
else:
agent._safe_print(warning)
agent._safe_print(warning, diagnostic=True)

summary_api_request_id = f"iteration-summary:{uuid.uuid4()}"
summary_call_outcome = "failed"
Expand Down Expand Up @@ -2396,7 +2396,7 @@ def _fall_back_to_converse(self, client, final_kwargs: dict, exc: Exception):
self.agent._disable_streaming = True
self.agent._safe_print("\n⚠ AWS IAM denied bedrock:InvokeModelWithResponseStream — "
"falling back to non-streaming InvokeModel.\n"
" Grant that action to restore streaming output.\n")
" Grant that action to restore streaming output.\n", diagnostic=True)
logger.info("bedrock: converse_stream denied by IAM (%s) — "
"using non-streaming converse() for this session.", type(exc).__name__)
return normalize_converse_response(client.converse(**final_kwargs))
Expand Down Expand Up @@ -2464,7 +2464,7 @@ def _on_stale(self, stale_elapsed: float) -> None:
logger.warning("Bedrock stream stale for %.0fs (threshold %.0fs) — no events "
"received. region=%s model=%s. Aborting call.", stale_elapsed, self.stale_timeout, self.region,
self._model())
agent._buffer_status(f"⚠️ No events from Bedrock for {int(stale_elapsed)}s (model: {self._model()}). Aborting...")
agent._buffer_diagnostic_status(f"⚠️ No events from Bedrock for {int(stale_elapsed)}s (model: {self._model()}). Aborting...")
_bump_stale_streak(agent)
# Evict the region's cached client so the NEXT call gets a fresh pool.
# This does NOT abort the in-flight botocore EventStream (no external
Expand Down Expand Up @@ -3234,7 +3234,8 @@ def _maybe_disable_streaming(self, e) -> None:
" Grant that action to restore streaming output.\n"
if _is_bedrock_stream_denied else
"\n⚠ Streaming is not supported for this model/provider. Switching to non-streaming.\n"
" To avoid this delay, set display.streaming: false in config.yaml\n"
" To avoid this delay, set display.streaming: false in config.yaml\n",
diagnostic=True,
)

def _handle_stream_error(self, e: Exception, attempt: int, max_retries: int) -> bool:
Expand Down Expand Up @@ -3293,7 +3294,8 @@ def _handle_stream_error(self, e: Exception, attempt: int, max_retries: int) ->
return False
# Marker explains the re-streamed preamble (``_emit_stream_drop`` logs the WARNING);
# reset the streamed-text buffer so it isn't double-recorded; fresh accumulators.
self._quiet(self.agent._fire_stream_delta, "\n\n⚠ Connection dropped mid tool-call; reconnecting…\n\n")
if self.agent._warning_presentation_enabled():
self._quiet(self.agent._fire_stream_delta, "\n\n⚠ Connection dropped mid tool-call; reconnecting…\n\n")
self._quiet(self.agent._reset_stream_delivery_tracking)
self.deltas_were_sent["yes"] = False
self.first_delta_fired["done"] = False
Expand All @@ -3312,7 +3314,7 @@ def _handle_stream_error(self, e: Exception, attempt: int, max_retries: int) ->
_what = ("Provider returned malformed streaming data after" if _is_stream_parse_err
else "Provider returned an empty response stream after" if _is_empty_stream
else "Connection to provider failed after")
self.agent._buffer_status(
self.agent._buffer_diagnostic_status(
f"❌ {_what} {max_retries + 1} attempts. The provider may be experiencing issues — try again in a moment.")
else:
self._maybe_disable_streaming(e)
Expand Down Expand Up @@ -3421,7 +3423,7 @@ def _kill_stale_stream(self, elapsed: float) -> None:
"Stream stale for %.0fs (threshold %.0fs) — no chunks received. model=%s context=~%s tokens. Killing connection.",
elapsed, self._stream_stale_timeout, self.api_kwargs.get("model", "unknown"), f"{_est_ctx:,}",
)
self.agent._buffer_status(
self.agent._buffer_diagnostic_status(
f"⚠️ No response from provider for {int(elapsed)}s (model: {self.api_kwargs.get('model', 'unknown')}, "
f"context: ~{_est_ctx:,} tokens). Reconnecting...")
# Captured BEFORE the cancel/abort: the pool sweep can miss a checked-out
Expand All @@ -3435,7 +3437,7 @@ def _kill_stale_stream(self, elapsed: float) -> None:
_bump_stale_streak(self.agent) # circuit breaker, see ``_stale_streak()``
# Reset the timer so we don't kill repeatedly while the worker unwinds.
self.last_chunk_time["t"] = time.time()
self.agent._emit_wait_notice(f"⚠ no output from provider for {int(elapsed)}s — reconnecting...")
self.agent._emit_diagnostic_wait(f"⚠ no output from provider for {int(elapsed)}s — reconnecting...")
self.agent._touch_activity(f"stale stream detected after {int(elapsed)}s, reconnecting")

def _abort_for_interrupt(self, stale_elapsed: float) -> None:
Expand Down Expand Up @@ -3499,8 +3501,9 @@ def _partial_stream_stub(self):
_name_str += f", +{len(_partial_names) - 3} more"
_warn = (f"\n\n⚠ Stream stalled mid tool-call ({_name_str}); the action was not executed. "
f"Ask me to retry if you want to continue.")
_partial_text = (_partial_text or "") + _warn
self._quiet(self.agent._fire_stream_delta, _warn) # visible immediately
_partial_text = (_partial_text or "") + _warn # model/result bookkeeping, never gated
if self.agent._warning_presentation_enabled():
self._quiet(self.agent._fire_stream_delta, _warn) # visible immediately
logger.warning(
"Partial stream dropped tool call(s) %s after %s chars of text; surfaced warning to user: %s",
_partial_names, len(_partial_text or ""), error)
Expand Down
6 changes: 3 additions & 3 deletions agent/chat_completion_nonstream.py
Original file line number Diff line number Diff line change
Expand Up @@ -173,11 +173,11 @@ def _ttfb_kill(self, elapsed: float) -> None:
"(%.0fs > %.0fs, model=%s). Backend accepted the connection "
"but sent no stream events. Killing connection so the retry loop can reconnect.", elapsed,
wd.ttfb_timeout, self._model())
agent._buffer_status(
agent._buffer_diagnostic_status(
f"⚠️ No first stream event from provider in {int(elapsed)}s (codex stream, model: {self._model()}). "
f"Reconnecting." + (f" {silent_hint}" if silent_hint else ""))
self._abort_request("codex_ttfb_kill")
agent._emit_wait_notice(f"⚠ no response from provider in {int(elapsed)}s — reconnecting...")
agent._emit_diagnostic_wait(f"⚠ no response from provider in {int(elapsed)}s — reconnecting...")
agent._touch_activity(f"codex stream killed after {int(elapsed)}s with no first stream event")
self._await_worker_after_kill(
f"Codex stream produced no parsed stream event within {int(elapsed)}s "
Expand All @@ -197,7 +197,7 @@ def _idle_kill(self, event_stale_elapsed: float) -> None:
"(threshold %.0fs, model=%s, context=~%s tokens). Killing "
"connection so the retry loop can reconnect.", event_stale_elapsed, arm_point, wd.idle_timeout,
self._model(), f"{wd.est_tokens:,}")
agent._buffer_status(
agent._buffer_diagnostic_status(
f"⚠️ Codex stream sent no events for {int(event_stale_elapsed)}s after {arm_point} "
f"(model: {self._model()}). Reconnecting.")
self._abort_request("codex_stream_idle_kill")
Expand Down
10 changes: 6 additions & 4 deletions agent/conversation_compression.py
Original file line number Diff line number Diff line change
Expand Up @@ -1791,7 +1791,7 @@ def _lower_threshold_to_aux_context(
f"{recomputed_threshold:,} tokens, still above the compression model's {aux_context:,}.)"
)
agent._compression_warning = msg
agent._emit_status(msg)
agent._emit_diagnostic_status(msg)
logger.warning(
"Auxiliary compression model %s has %d token context, below the main model's compression threshold of %d "
"tokens — auto-lowered session threshold to %d to keep compression working.", aux_model, aux_context,
Expand Down Expand Up @@ -1848,7 +1848,7 @@ def check_compression_model_feasibility(agent: Any) -> None:
"long chats, so older messages will be cut without a summary. Run `hermes setup` to add one."
)
agent._compression_warning = msg
agent._emit_status(msg)
agent._emit_diagnostic_status(msg)
logger.warning("No auxiliary LLM provider for compression — summaries will be unavailable.")
return
aux_base_url = str(getattr(client, "base_url", ""))
Expand Down Expand Up @@ -1902,8 +1902,10 @@ def replay_compression_warning(agent: Any) -> None:
``__init__``) is finally wired."""
msg = getattr(agent, "_compression_warning", None)
if msg and agent.status_callback:
# Replayed as a classified diagnostic so every sink applies its own policy snapshot.
from gateway.warning_notifications import DiagnosticText
with contextlib.suppress(Exception):
agent.status_callback("lifecycle", msg)
agent.status_callback("lifecycle", DiagnosticText(msg))


def conversation_history_after_compression(
Expand Down Expand Up @@ -3148,7 +3150,7 @@ def _finish_compaction_boundary(
f"{agent.log_prefix}⚠️ Session compressed {_cc} times — accuracy may degrade. Consider /new to start fresh."
)
agent._compression_warning = _cc_msg
agent._emit_status(_cc_msg)
agent._emit_diagnostic_status(_cc_msg)

# session:compress lets hooks ingest the old session before it's lost;
# in_place=True tells them the same id was compacted rather than rotated.
Expand Down
2 changes: 1 addition & 1 deletion agent/conversation_loop.py
Original file line number Diff line number Diff line change
Expand Up @@ -465,7 +465,7 @@ def _print_guidance(agent, message: str) -> bool:
if not message:
return False
for line in message.splitlines():
agent._vprint(f"{agent.log_prefix} 💡 {line}", force=True)
agent._vprint(f"{agent.log_prefix} 💡 {line}", force=True, diagnostic=True)
return True


Expand Down
2 changes: 1 addition & 1 deletion agent/fallback_cooldown.py
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,7 @@ def _mark_entitlement_rejected_model(agent, api_error) -> bool:
"treating it as unavailable for this session",
model, provider,
)
agent._buffer_status(
agent._buffer_diagnostic_status(
f"🚫 This account is not entitled to {model} via {provider}; it will be skipped "
"until restart. Switch to an entitled model via /model or `hermes model`."
)
Expand Down
83 changes: 83 additions & 0 deletions agent/notification_presentation.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
"""Turn-local presentation of diagnostic-only wakes; execution and controls stay live."""
from contextlib import contextmanager
from contextvars import ContextVar


_muted_surface: ContextVar[str | None] = ContextVar("muted_notification_surface", default=None)
# These are text/media UI events, not approval/clarify/connection requests or outcomes.
_FREEFORM_EVENTS = frozenset({
"message.start", "message.delta", "message.interim", "message.complete",
"reasoning.delta", "thinking.delta", "status.update", "notification.show",
"tool.start", "tool.complete", "tool.generating", "error", "reaction",
})
_PRESENTATION_CALLBACKS = (
"stream_delta_callback", "interim_assistant_callback", "reasoning_callback",
"tool_progress_callback",
"tool_start_callback", "tool_complete_callback", "tool_gen_callback", "reaction_callback",
)


def event_presentation_muted(event: str, session_id: str) -> bool:
return _muted_surface.get() == session_id and event in _FREEFORM_EVENTS


def diagnostic_process_event(event: dict) -> bool:
"""Early failure/monitor diagnostics, not the explicitly requested final result."""
return bool(event.get("task_failure_notice")) or event.get("type") in {
"watch_disabled", "watch_overflow_tripped", "watch_overflow_released",
}


def notification_config_snapshot():
"""Read the canonical owning effective config once; presentation fails open."""
from copy import deepcopy
from hermes_cli.config_effective import load_user_config_effective
try:
config = load_user_config_effective()
return deepcopy(config) if isinstance(config, dict) else {}
except Exception:
return {}


@contextmanager
def notification_policy_snapshot(agent, platform, config):
"""Bind one foreground policy for callbacks, including worker threads."""
from copy import deepcopy
missing = object()
saved = {key: getattr(agent, key, missing)
for key in ("_notification_config", "_notification_platform")}
try:
agent._notification_config = deepcopy(config)
agent._notification_platform = platform
yield
finally:
for key, value in saved.items():
if value is missing:
delattr(agent, key)
else:
setattr(agent, key, value)


@contextmanager
def notification_turn(agent, *, muted: bool, session_id: str = ""):
"""Freeze the current turn's presentation without touching prompts or tool schemas."""
if not muted:
yield
return
missing = object()
keys = (*_PRESENTATION_CALLBACKS, "suppress_status_output", "_mute_notification_reply")
saved = {key: getattr(agent, key, missing) for key in keys}
token = _muted_surface.set(session_id)
try:
for key in _PRESENTATION_CALLBACKS:
setattr(agent, key, None)
agent.suppress_status_output = True
agent._mute_notification_reply = True
yield
finally:
for key, value in saved.items():
if value is missing:
delattr(agent, key)
else:
setattr(agent, key, value)
_muted_surface.reset(token)
Loading
Loading