forked from NousResearch/hermes-agent
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathrun_agent.py
More file actions
1589 lines (1390 loc) · 86.6 KB
/
Copy pathrun_agent.py
File metadata and controls
1589 lines (1390 loc) · 86.6 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
#!/usr/bin/env python3
"""AIAgent: the tool-calling agent runner (conversation loop, tool execution, session lifecycle).
from run_agent import AIAgent
agent = AIAgent(base_url="http://localhost:30000/v1", model="claude-opus-4-20250514")
response = agent.run_conversation("Tell me about the latest Python updates")
"""
# hermes_bootstrap must be the very first import (UTF-8 stdio on Windows; no-op on POSIX).
try:
import hermes_bootstrap # noqa: F401
except ModuleNotFoundError:
pass # partial `hermes update` — only skips the Windows UTF-8 stdio setup
import json
import logging
logger = logging.getLogger(__name__)
import os
import re
import sys
import time
import threading
import uuid
import warnings
from typing import List, Dict, Any, Optional, Callable
from datetime import datetime
from pathlib import Path
from hermes_constants import get_hermes_home
def _launch_cwd_for_session(source: str) -> Optional[str]:
"""cwd to stamp on a new session row (``hermes -c`` / ``--resume``), or None.
Only local CLI sessions record one: gateway/cron/remote backends (non-"local" ``TERMINAL_ENV``) have no
stable host cwd for the agent's tools.
"""
if source not in CLI_FAMILY_SOURCES or (os.environ.get("TERMINAL_ENV") or "local").strip().lower() not in ("", "local"):
return None
try:
return os.getcwd()
except OSError: # cwd was unlinked out from under us
return None
# Sources that label the human conversation an interactive UI transport hosts. A finite ``hermes chat -q`` /
# one-shot child spawned from such a session inherits HERMES_SESSION_SOURCE (the terminal tool bridges the
# session env into child processes) but is NOT that conversation: labelling it ``tui``/``desktop`` lists it
# in the TUI/WebUI pickers as a resumable chat and lets ``hermes -c`` in the TUI continue it (#112550).
# Automation sources (kanban, tool, cron, a2a, ...) are inherited on purpose.
_UI_TRANSPORT_SOURCES = frozenset({"tui", "desktop"})
# Finite non-interactive CLI runs (``hermes chat -q``/``--oneshot``, ``hermes -z``) get their own source so human
# pickers hide them without title/cwd heuristics; ``hermes -c`` still treats them as CLI history.
ONESHOT_SOURCE = "oneshot"
CLI_FAMILY_SOURCES = frozenset({"cli", ONESHOT_SOURCE})
def _session_source_for_agent(platform: Optional[str]) -> str:
try:
from gateway.session_context import get_session_env
except Exception:
get_session_env = os.environ.get
source = str(get_session_env("HERMES_SESSION_SOURCE", "") or "").strip()
single_query = get_session_env("HERMES_SINGLE_QUERY_SESSION", "") == "1"
explicit = get_session_env("HERMES_SESSION_SOURCE_EXPLICIT", "") == "1"
if single_query and not explicit and source in _UI_TRANSPORT_SOURCES:
source = ""
if single_query and not source and (platform or "cli") == "cli":
return ONESHOT_SOURCE
return source or platform or "cli"
def _gateway_origin_json(agent: "AIAgent") -> Optional[str]:
"""Gateway routing ``origin_json`` for a session row; None when the agent carries no gateway identity.
Mirrors ``SessionSource.to_dict()`` so state.db consumers see the same fields ``record_gateway_session_peer`` writes.
"""
chat_id = getattr(agent, "_chat_id", None)
session_key = getattr(agent, "_gateway_session_key", None)
user_id = getattr(agent, "_user_id", None)
if not (chat_id or session_key or user_id):
return None
origin: Dict[str, Any] = {
"platform": getattr(agent, "platform", None) or "", "chat_id": chat_id,
"chat_name": getattr(agent, "_chat_name", None), "chat_type": getattr(agent, "_chat_type", None) or "dm",
"user_id": user_id, "user_name": getattr(agent, "_user_name", None), "thread_id": getattr(agent, "_thread_id", None),
}
if getattr(agent, "_user_id_alt", None):
origin["user_id_alt"] = agent._user_id_alt
profile = getattr(agent, "_profile_name", None)
if not profile:
try:
from hermes_cli.profiles import get_active_profile_name
profile = get_active_profile_name()
except Exception:
profile = None
if profile == "default":
profile = None
if profile:
origin["profile"] = profile
try:
return json.dumps(origin)
except Exception:
return None
from agent.iteration_budget import IterationBudget
from hermes_cli.env_loader import load_hermes_dotenv
from hermes_cli.timeouts import get_provider_request_timeout, get_provider_stale_timeout
_hermes_home = get_hermes_home() # read by agent_init via _ra()._hermes_home
_loaded_env_paths = load_hermes_dotenv(hermes_home=_hermes_home, project_env=Path(__file__).parent / '.env')
for _env_path in _loaded_env_paths:
logger.info("Loaded environment variables from %s", _env_path)
if not _loaded_env_paths:
logger.info("No .env file found. Using system environment variables.")
from model_tools import get_toolset_for_tool
from tools.terminal_tool_lifecycle import cleanup_vm, get_active_env
from tools.interrupt import set_interrupt as _set_interrupt
from tools.browser_tool_lifecycle import cleanup_browser
from agent.memory_provider import is_trivial_prompt
from agent.client_lifecycle import ClientLifecycleMixin
from agent.stream_delivery import StreamDeliveryMixin
from agent.status_output import StatusOutputMixin
from agent.api_request_hooks import ApiRequestHooksMixin
from agent.api_error_summary import PROVIDER_STREAM_PARSE_MARKERS, ApiErrorSummaryMixin
from agent.interrupt_control import InterruptControlMixin
from agent.turn_explainers import TurnExplainersMixin
from agent.activity_tracking import ActivityTrackingMixin
from agent.rate_limit_credits import RateLimitCreditsMixin
from agent.session_persistence import SessionPersistenceMixin
from agent.compression_facade import CompressionFacadeMixin
from agent.turn_facade import TurnFacadeMixin
from agent.vision_message_prep import VisionMessagePrepMixin
from agent.reasoning_params import ReasoningParamsMixin
from agent.lazy_forward import forward as _forward, forward_static as _forward_static
from agent.session_activity import ActivityProvenance
from agent.model_metadata import is_local_endpoint
from agent.message_sanitization import (
coalesce_tool_call_id as _sanitize_coalesce_tool_call_id,
deterministic_call_id as _codex_deterministic_call_id,
uniquify_tool_call_ids as _sanitize_uniquify_tool_call_ids,
)
from agent.codex_responses_adapter import (
_derive_responses_function_call_id as _codex_derive_responses_function_call_id,
_split_responses_tool_id as _codex_split_responses_tool_id,
_summarize_user_message_for_log,
)
from agent.tool_guardrails import ToolGuardrailDecision, append_toolguard_guidance, toolguard_synthetic_result
from utils import base_url_host_matches, base_url_hostname, env_float, model_forces_max_completion_tokens
_MAX_TOOL_WORKERS = 8
# Spawn the OpenRouter pre-warm thread once per process, not per AIAgent (gateway thread leak).
_openrouter_prewarm_done = threading.Event()
def _quietly(fn: Callable, *args, **kwargs) -> None:
"""Run one teardown step, swallowing any exception so sibling steps still run."""
try:
fn(*args, **kwargs)
except Exception:
pass
def _call_engine_hook(engine: Any, hook: str, *args, **kwargs) -> None:
"""Invoke an optional context-engine lifecycle hook; failures are logged, never raised."""
if not hasattr(engine, hook):
return
try:
getattr(engine, hook)(*args, **kwargs)
except Exception as exc:
logger.debug("context engine %s during transition: %s", hook, exc)
def _positive_int(value: Any) -> Optional[int]:
"""``value`` when it is a real positive int (bools excluded), else None."""
return value if isinstance(value, int) and not isinstance(value, bool) and value > 0 else None
def _review_should_defer(agent: Any, task_cfg: Optional[Dict[str, Any]]) -> bool:
"""True when an automatic background review targets the managed local runtime under ``defer: auto``."""
from agent.review_idle_queue import defer_mode, review_targets_managed_local
return defer_mode(task_cfg) == "auto" and review_targets_managed_local(agent, task_cfg)
def _review_queue_key(agent: Any) -> str:
return str(getattr(agent, "session_id", None) or id(agent))
def _notify_context_engine_session_end(agent: Any, messages: Optional[list]) -> None:
"""Tell the context engine the session ended (flush DAG, close DBs) at the same lifecycle moment as the
memory manager, so per-session engine state never leaks into the next session."""
engine = getattr(agent, "context_compressor", None)
if engine:
_quietly(lambda: engine.on_session_end(agent.session_id or "", messages or []))
def _pool_may_recover_from_rate_limit(pool) -> bool:
"""Wait for credential-pool rotation (True) or fall back to ``fallback_model`` (False) after a 429.
Rotation only helps when the pool has somewhere to go; a single-credential pool would retry the same quota.
See issues #11314 and #13636.
"""
return pool is not None and pool.has_available() and len(pool.entries()) > 1
class _StreamErrorEvent(Exception):
"""Provider error synthesized from a standalone Responses ``type=error`` SSE frame (Codex-style backends).
Gives ``_summarize_api_error`` / the entitlement detector the familiar ``.body`` / ``.status_code`` shape.
"""
def __init__(self, message: str, *, code: Optional[str] = None, param: Optional[str] = None,
status_code: Optional[int] = None) -> None:
super().__init__(message)
self.message, self.code, self.param, self.status_code = message, code, param, status_code
# OpenAI SDK-shaped body so _extract_api_error_context / _summarize_api_error / classify_api_error pick it up.
self.body: Dict[str, Any] = {"error": {"message": message, "code": code, "param": param, "type": "error"}}
class AIAgent(
ClientLifecycleMixin, StreamDeliveryMixin, StatusOutputMixin, ApiRequestHooksMixin, ApiErrorSummaryMixin,
InterruptControlMixin, TurnExplainersMixin, ActivityTrackingMixin, RateLimitCreditsMixin,
SessionPersistenceMixin, CompressionFacadeMixin, TurnFacadeMixin, VisionMessagePrepMixin, ReasoningParamsMixin,
):
"""AI Agent with tool calling capabilities."""
_TOOL_CALL_ARGUMENTS_CORRUPTION_MARKER = (
"[hermes-agent: tool call arguments were corrupted in this session and "
"have been dropped to keep the conversation alive. See issue #15236.]"
)
@property
def base_url(self) -> str:
return self._base_url
@base_url.setter
def base_url(self, value: str) -> None:
self._base_url = value
self._base_url_lower = value.lower() if value else ""
self._base_url_hostname = base_url_hostname(value)
def __init__(
self,
base_url: str = None, api_key: str = None, provider: str = None, api_mode: str = None,
acp_command: str = None, acp_args: list[str] | None = None, command: str = None, args: list[str] | None = None,
model: str = "",
max_iterations: int = sys.maxsize, # unlimited tool-calling iterations by default (shared with subagents)
tool_delay: float = None, # deprecated: accepted for compatibility, ignored
enabled_toolsets: List[str] = None, disabled_toolsets: List[str] = None,
save_trajectories: bool = False, verbose_logging: bool = False, quiet_mode: bool = False,
tool_progress_mode: str = "all", ephemeral_system_prompt: str = None,
log_prefix_chars: int = 100, log_prefix: str = "",
providers_allowed: List[str] = None, providers_ignored: List[str] = None, providers_order: List[str] = None,
provider_sort: str = None, provider_require_parameters: bool = False, provider_data_collection: str = None,
openrouter_min_coding_score: Optional[float] = None,
session_id: str = None,
tool_progress_callback: callable = None, tool_start_callback: callable = None,
tool_complete_callback: callable = None, thinking_callback: callable = None,
reasoning_callback: callable = None, clarify_callback: callable = None,
read_terminal_callback: callable = None, read_preview_callback: callable = None,
drive_preview_callback: callable = None, read_window_below_callback: callable = None,
connection_callback: callable = None, tour_callback: callable = None, step_callback: callable = None,
stream_delta_callback: callable = None, interim_assistant_callback: callable = None,
tool_gen_callback: callable = None, status_callback: callable = None,
notice_callback: callable = None, notice_clear_callback: callable = None,
event_callback: Optional[Callable[[str, dict], None]] = None,
reaction_callback: Optional[Callable[[str], None]] = None,
max_tokens: int = None, reasoning_config: Dict[str, Any] = None, service_tier: str = None,
request_overrides: Dict[str, Any] = None, prefill_messages: List[Dict[str, Any]] = None,
platform: str = None, user_id: str = None, user_id_alt: str = None, user_name: str = None,
chat_id: str = None, chat_name: str = None, chat_type: str = None, thread_id: str = None,
gateway_session_key: str = None,
skip_context_files: bool = False, load_soul_identity: bool = False,
skip_memory: bool = False, skip_background_review: bool = False,
session_db=None, parent_session_id: str = None,
iteration_budget: "IterationBudget" = None, run_budget_seconds: Optional[float] = None,
fallback_model: Dict[str, Any] = None, credential_pool=None,
checkpoints_enabled: bool = False, checkpoint_max_snapshots: int = 20,
checkpoint_max_total_size_mb: int = 500, checkpoint_max_file_size_mb: int = 10,
pass_session_id: bool = False, requested_provider: str = None,
capabilities: Dict[str, bool] | None = None, cwd: str | None = None,
):
"""Forwarder — see ``agent.agent_init.init_agent`` (same keyword parameters, minus ``tool_delay``)."""
init_kwargs = {k: v for k, v in locals().items() if k not in ("self", "tool_delay")}
if tool_delay is not None:
warnings.warn("tool_delay is deprecated and ignored; sequential tool calls "
"no longer sleep between executions.", DeprecationWarning, stacklevel=2)
from agent.agent_init import init_agent
init_agent(self, **init_kwargs)
def _get_session_db_for_recall(self):
"""SessionDB for recall, opening the default state DB when no ``session_db`` was passed so the
advertised ``session_search`` tool stays usable."""
# Persistence-isolated forks (background review) must not lazily open the canonical state DB —
# that would re-arm the flush to write the fork's harness turn into the user's real session.
if getattr(self, "_persist_disabled", False):
return None
if self._session_db is not None:
return self._session_db
try:
from hermes_state_registry import acquire
self._session_db = acquire()
self._owns_session_db = True # we opened it, so close() must release it
return self._session_db
except Exception:
logger.debug("SessionDB unavailable for recall", exc_info=True)
return None
def _session_row_model_config(self) -> Any:
"""``model_config`` for the session row: the init config plus the live YOLO bypass.
The row is created lazily on the first turn, so this is the only chance to record a pre-first-turn
/yolo toggle for ``hermes --resume``.
"""
model_config = self._session_init_model_config
try:
from tools.approval import is_session_yolo_enabled
if is_session_yolo_enabled(self.session_id):
model_config = dict(model_config or {})
model_config["yolo_mode"] = True
except Exception:
pass
return model_config
def _ensure_db_session(self) -> None:
"""Create the session DB row on first use; a transient failure leaves it to retry next turn."""
if getattr(self, "_persist_disabled", False) or self._session_db_created or not self._session_db:
return
source = _session_source_for_agent(self.platform)
try:
# Persist the profile name explicitly, including "default": profile-keyed consumers treat NULL
# as unowned.
try:
from hermes_cli.profiles import get_active_profile_name
profile_for_session = get_active_profile_name()
except Exception:
# Persist the profile name EXPLICITLY, including "default". NULL used to stand in for the
# default profile, but the #94724 legacy-owner backfill already stamps literal "default"
# onto old rows, and profile-keyed consumers (sidebar scope matching,
# @session:<profile>/<id> deep links) treat NULL as unowned — rows minted NULL after the
# one-shot backfill vanished from the sidebar (#99222).
profile_for_session = None
# Carry the gateway routing identity: when the gateway SessionStore degraded to JSONL (corrupt
# state.db) this lazy create is the ONLY durable write, and an identity-less row is unrecoverable.
self._session_db.create_session(
session_id=self.session_id, source=source, model=self.model,
model_config=self._session_row_model_config(), system_prompt=self._cached_system_prompt,
user_id=getattr(self, "_user_id", None), session_key=getattr(self, "_gateway_session_key", None),
chat_id=getattr(self, "_chat_id", None), chat_type=getattr(self, "_chat_type", None),
thread_id=getattr(self, "_thread_id", None),
display_name=getattr(self, "_chat_name", None) or getattr(self, "_user_name", None),
origin_json=_gateway_origin_json(self), parent_session_id=self._parent_session_id,
cwd=_launch_cwd_for_session(source), profile_name=profile_for_session,
)
self._session_db_created = True
except Exception as e:
# Transient failure (e.g. SQLite lock): _session_db_created stays False so the next turn retries.
logger.warning("Session DB creation failed (will retry next turn): %s", e)
def _transition_context_engine_session(
self, *, old_session_id: Optional[str] = None, new_session_id: Optional[str] = None,
previous_messages: Optional[list] = None, carry_over_context: bool = False, reset_engine: bool = True,
**extra_context,
) -> None:
"""Drive the context engine's session transition: on_session_end → on_session_reset → on_session_start
→ carry_over_new_session_context. Each hook is optional (the built-in compressor only resets)."""
engine = getattr(self, "context_compressor", None)
if not engine:
return
if old_session_id and previous_messages is not None:
_call_engine_hook(engine, "on_session_end", old_session_id, previous_messages)
if reset_engine:
_call_engine_hook(engine, "on_session_reset")
should_start = bool(old_session_id or previous_messages is not None or carry_over_context or extra_context)
target_session_id = new_session_id or getattr(self, "session_id", "") or ""
if should_start and target_session_id and hasattr(engine, "on_session_start"):
start_context = {
"old_session_id": old_session_id, "carry_over_context": carry_over_context,
"platform": _session_source_for_agent(getattr(self, "platform", None)),
"model": getattr(self, "model", ""), "context_length": getattr(engine, "context_length", None),
"conversation_id": getattr(self, "_gateway_session_key", None), **extra_context,
}
start_context = {k: v for k, v in start_context.items() if v not in (None, "")}
_call_engine_hook(engine, "on_session_start", target_session_id, **start_context)
if carry_over_context and old_session_id and target_session_id:
_call_engine_hook(engine, "carry_over_new_session_context", old_session_id, target_session_id)
def reset_session_state(self, previous_messages: Optional[list] = None, old_session_id: Optional[str] = None,
carry_over_context: bool = False):
"""Reset session-scoped token/cost counters and compressor state for a fresh session.
With ``previous_messages`` / ``old_session_id`` / ``carry_over_context`` the context engine gets the
full transition lifecycle instead of a bare reset.
"""
for counter in (
"session_total_tokens", "session_input_tokens", "session_output_tokens", "session_prompt_tokens",
"session_completion_tokens", "session_cache_read_tokens", "session_cache_write_tokens",
"session_reasoning_tokens", "session_api_calls",
):
setattr(self, counter, 0)
self.session_estimated_cost_usd = 0.0
self.session_cost_status = "unknown"
self.session_cost_source = "none"
# Session boundary: the usage anchor describes the OLD transcript; fall back to full estimation.
self._usage_anchor = None
self._turn_base_usage_anchor = None
# The workspace snapshot is pinned per session (agent/system_prompt.py::_coding_parts); a
# /new, /resume or /branch on the same agent must re-snapshot at its own session start.
self._frozen_workspace_snapshot = None
# Turn counter (added after reset_session_state was first written — #2635)
self._user_turn_count = 0
# The drifted-prompt compaction INFO is once per session, so a /new or /resume re-arms it.
self._compaction_prompt_drift_logged = False
# Who wrote the current turn. build_turn_context() sets it at the start of every turn.
self._turn_author = None
# Copilot x-initiator: True for the first API call of a user turn, False for tool-loop follow-ups.
self._is_user_initiated_turn = False
self._transition_context_engine_session(
old_session_id=old_session_id, new_session_id=getattr(self, "session_id", None),
previous_messages=previous_messages, carry_over_context=carry_over_context, reset_engine=True,
)
# Reset-only switches (/new, /resume, /branch) change session_id before this call; rebind the
# built-in compressor's session-keyed cooldown state when no full start hook ran.
engine = getattr(self, "context_compressor", None)
target_session_id = getattr(self, "session_id", "") or ""
if (engine is not None and hasattr(engine, "bind_session_state") and target_session_id
and target_session_id != getattr(engine, "_session_id", "")):
try:
engine.bind_session_state(getattr(self, "_session_db", None), target_session_id)
except Exception as exc:
logger.debug("context engine bind_session_state during reset: %s", exc)
@staticmethod
def _effective_lmstudio_context_length(config_context_length: Optional[int], runtime_context_length: Any) -> Optional[int]:
"""Return a safe context budget from explicit intent and verified runtime."""
explicit = _positive_int(config_context_length)
runtime = _positive_int(getattr(runtime_context_length, "context_length", runtime_context_length))
if bool(getattr(runtime_context_length, "rejected", False)) or (
bool(getattr(runtime_context_length, "load_attempted", False)) and runtime is None
):
return None
if runtime is not None and explicit is not None:
return min(runtime, explicit)
return runtime if runtime is not None else explicit
@staticmethod
def _lmstudio_load_was_unverified(load_result: Any) -> bool:
"""Return true when a management load was rejected or unverifiable."""
return bool(getattr(load_result, "rejected", False)) or (
bool(getattr(load_result, "load_attempted", False)) and getattr(load_result, "context_length", None) is None
)
def _ensure_lmstudio_runtime_loaded(self, config_context_length: Optional[int] = None) -> Any:
"""Preload LM Studio unless configured to rely on JIT loading."""
if (self.provider or "").strip().lower() != "lmstudio":
return None
if (getattr(self, "lmstudio_load_mode", "explicit") or "explicit").strip().lower() == "jit":
logger.debug("LM Studio explicit preload skipped: lmstudio_load_mode=jit")
return None
from hermes_cli.models_local import ensure_lmstudio_model_loaded
if config_context_length is None:
config_context_length = getattr(self, "_config_context_length", None)
return ensure_lmstudio_model_loaded(
self.model, self.base_url, getattr(self, "api_key", ""), config_context_length, return_load_result=True,
)
switch_model = _forward("agent.agent_runtime_helpers", "switch_model")
def _disable_codex_reasoning_replay(self, messages: Optional[List[Dict[str, Any]]] = None) -> Dict[str, int]:
"""On HTTP 400 ``invalid_encrypted_content``: disable Responses reasoning replay and pop
``codex_reasoning_items`` from every assistant message. Returns ``{"messages", "items"}`` counts."""
stripped_messages = stripped_items = 0
for msg in (messages if isinstance(messages, list) else []):
if not isinstance(msg, dict) or msg.get("role") != "assistant":
continue
items = msg.pop("codex_reasoning_items", None)
if isinstance(items, list) and items:
stripped_messages += 1
stripped_items += len(items)
self._codex_reasoning_replay_enabled = False
return {"messages": stripped_messages, "items": stripped_items}
_stream_diag_init = _forward_static("agent.stream_diag", "stream_diag_init")
_stream_diag_capture_response = _forward("agent.stream_diag", "stream_diag_capture_response")
_flatten_exception_chain = _forward_static("agent.stream_diag", "flatten_exception_chain")
def _is_provider_stream_parse_error(self, error: BaseException) -> bool:
"""True for a malformed Anthropic event-stream frame (surfaced by the SDK as a plain ``ValueError``);
that is wire trouble, not local validation, so it follows the truncated-JSON retry path."""
return (getattr(self, "api_mode", None) == "anthropic_messages" and isinstance(error, ValueError)
and not isinstance(error, (UnicodeEncodeError, json.JSONDecodeError))
and any(marker in str(error).strip().lower() for marker in PROVIDER_STREAM_PARSE_MARKERS))
_log_stream_retry = _forward("agent.stream_diag", "log_stream_retry")
_emit_stream_drop = _forward("agent.stream_diag", "emit_stream_drop")
def _emit_auxiliary_failure(self, task: str, exc: BaseException) -> None:
"""Surface a compact warning for failed auxiliary work."""
try:
detail = self._summarize_api_error(exc)
except Exception:
detail = str(exc)
detail = (detail or exc.__class__.__name__).strip()
if len(detail) > 220:
detail = detail[:217].rstrip() + "..."
self._emit_warning(f"⚠ Auxiliary {task} failed: {detail}")
def _current_main_runtime(self) -> Dict[str, str]:
"""Return the live main runtime for session-scoped auxiliary routing."""
return {
key: getattr(self, key, "") or ""
for key in ("model", "provider", "base_url", "api_key", "api_mode", "auth_mode", "session_id")
}
_check_compression_model_feasibility = _forward("agent.conversation_compression", "check_compression_model_feasibility")
_replay_compression_warning = _forward("agent.conversation_compression", "replay_compression_warning")
def _hostname_for(self, base_url: Optional[str]) -> str:
"""Hostname of ``base_url``, or of the agent's own base URL when None."""
if base_url is not None:
return base_url_hostname(base_url)
return getattr(self, "_base_url_hostname", "") or base_url_hostname(getattr(self, "_base_url_lower", ""))
def _is_direct_openai_url(self, base_url: str = None) -> bool:
"""Return True when a base URL targets OpenAI's native API."""
return self._hostname_for(base_url) == "api.openai.com"
def _is_azure_openai_url(self, base_url: str = None) -> bool:
"""True when a base URL targets Azure OpenAI (standard client, but NO Responses API support)."""
url = str(base_url).lower() if base_url is not None else (getattr(self, "_base_url_lower", "") or "")
return base_url_host_matches(url, "openai.azure.com")
def _is_github_copilot_url(self, base_url: str = None) -> bool:
"""Return True when a base URL targets GitHub Copilot's OpenAI-compatible API."""
hostname = self._hostname_for(base_url)
return bool(hostname) and (hostname == "api.githubcopilot.com" or hostname.endswith(".githubcopilot.com"))
def _resolved_api_call_timeout(self) -> float:
"""Per-call request timeout: per-model ``timeout_seconds`` > provider ``request_timeout_seconds`` >
``HERMES_API_TIMEOUT`` > 1800s."""
cfg = get_provider_request_timeout(self.provider, self.model)
return cfg if cfg is not None else env_float("HERMES_API_TIMEOUT", 1800.0)
def _resolved_api_call_stale_timeout_base(self) -> tuple[float, bool]:
"""Base non-stream stale timeout: per-model ``stale_timeout_seconds`` > provider-wide >
``HERMES_API_CALL_STALE_TIMEOUT`` > reasoning floor > 90s.
Returns ``(seconds, uses_implicit_default)``; the implicit flag lets callers auto-disable the detector
for local endpoints only when the user configured nothing.
"""
cfg = get_provider_stale_timeout(self.provider, self.model)
if cfg is not None:
return cfg, False
env_timeout = os.getenv("HERMES_API_CALL_STALE_TIMEOUT")
if env_timeout is not None:
return float(env_timeout), False
# Reasoning-model floor (cloud gateways idle-kill mid-think); not "implicit" so the local-endpoint
# short-circuit does not disable stale detection here.
from agent.reasoning_timeouts import get_reasoning_stale_timeout_floor
reasoning_floor = get_reasoning_stale_timeout_floor(self.model)
if reasoning_floor is not None:
return reasoning_floor, False
return 90.0, True
def _compute_non_stream_stale_timeout(self, api_payload: Any) -> float:
"""Effective non-stream stale timeout for ``api_payload`` (an ``api_kwargs`` dict or legacy ``messages``
list), scaled by estimated context size and capped by the run budget."""
stale_base, uses_implicit_default = self._resolved_api_call_stale_timeout_base()
base_url = getattr(self, "_base_url", None) or self.base_url or ""
if uses_implicit_default and base_url and is_local_endpoint(base_url):
return float("inf")
from agent.chat_completion_helpers import _high_effort_silence_floor, estimate_request_context_tokens
est_tokens = estimate_request_context_tokens(api_payload)
timeout = max(stale_base, 240.0) if est_tokens > 100_000 else max(stale_base, 150.0) if est_tokens > 50_000 else stale_base
explicit = self._stale_timeout_is_explicit()
# High-effort Codex reasoning (#112909) floors the IMPLICIT stale timeout before the run-budget
# cap below, so the floor can never outlive the run budget.
if self.api_mode == "codex_responses" and not explicit:
timeout = max(timeout, _high_effort_silence_floor(self))
# Run-budget cap: an implicit stale timeout is capped at half the remaining budget (>= 60s) so one
# hung call cannot eat the run. Never raises the timeout; explicit user config still wins.
run_budget = getattr(self, "run_budget_seconds", None)
started = getattr(self, "_run_budget_started_at", None)
if run_budget and started and not explicit:
remaining = float(run_budget) - (time.time() - started)
timeout = min(timeout, max(60.0, remaining * 0.5))
return timeout
def _stale_timeout_is_explicit(self) -> bool:
"""True when the user explicitly configured the stale timeout (config or env var); implicit values
(reasoning floors, the 90s default) yield to the run-budget cap, explicit ones never do."""
return (get_provider_stale_timeout(self.provider, self.model) is not None
or os.getenv("HERMES_API_CALL_STALE_TIMEOUT") is not None)
def _codex_silent_hang_hint(self, model: Optional[str] = None) -> Optional[str]:
"""Actionable hint when the request matches a known Codex silent-reject shape (currently the ``gpt-5.5``
family: connection accepted, no events, no error), else None. Makes the stale timeout actionable."""
if self.api_mode != "codex_responses":
return None
from agent.codex_responses_adapter import classify_responses_route
if not classify_responses_route(self).is_codex_backend:
return None
eff_model = (model if model is not None else self.model) or ""
# Match the gpt-5.5 family at word boundaries (bare, -codex, vendor-prefixed) but not gpt-5.50.
if not re.search(r"(?:^|[/\-_])gpt-5\.5(?:$|[\-_])", eff_model.lower()):
return None
return (
f"Codex backend appears to be silently rejecting {eff_model!r} "
"on chatgpt.com/backend-api/codex (no stream events, no error). "
"This is a known backend-side pattern that has affected ChatGPT "
"Plus accounts intermittently. "
"Workaround: try `gpt-5.4` on the same OAuth profile, or `gpt-5.3-codex`, "
"or switch to a different model/provider in your fallback chain. "
"Some ChatGPT Codex accounts do not support `gpt-5.4-codex`. "
"See hermes-agent#21444 for symptom history."
)
def _is_openrouter_url(self) -> bool:
"""Return True when the base URL targets OpenRouter."""
return base_url_host_matches(self._base_url_lower, "openrouter.ai")
def _is_copilot_url(self) -> bool:
"""Return True when the base URL targets GitHub Copilot or GitHub Models."""
return any(base_url_host_matches(self._base_url_lower, h) for h in ("api.githubcopilot.com", "models.github.ai"))
def _is_copilot_provider(self) -> bool:
"""True when the active provider is GitHub Copilot under any alias (``copilot`` / ``github-copilot`` /
``github``) or by base URL; a bare equality check would silently skip credential recovery."""
return (self.provider or "").strip().lower() in {"copilot", "github-copilot", "github"} or self._is_copilot_url()
def _is_codex_backend(self) -> bool:
"""Return True for the ChatGPT OAuth Codex Responses backend."""
return (getattr(self, "api_mode", None) == "codex_responses"
and getattr(self, "_base_url_hostname", "") == "chatgpt.com"
and "/backend-api/codex" in (getattr(self, "_base_url_lower", "") or ""))
_anthropic_prompt_cache_policy = _forward("agent.agent_runtime_helpers", "anthropic_prompt_cache_policy")
_direct_native_anthropic_tool_cache_capability = _forward("agent.agent_runtime_helpers", "_direct_native_anthropic_tool_cache_capability")
@staticmethod
def _model_requires_responses_api(model: str) -> bool:
"""True for GPT-5.x, which OpenAI and OpenRouter reject on /v1/chat/completions
(``unsupported_api_for_model``)."""
return model.lower().rsplit("/", 1)[-1].startswith("gpt-5") # strip vendor prefix ("openai/gpt-5.4")
@staticmethod
def _provider_model_requires_responses_api(model: str, *, provider: Optional[str] = None) -> bool:
"""Return True when this provider/model pair should use Responses API."""
from hermes_cli.providers import is_actual_route
normalized_provider = (provider or "").strip().lower()
# Nous serves GPT-5.x via chat completions (its /v1/responses returns 404); generic custom endpoints
# may relay GPT-5 without full Responses semantics — only direct OpenAI/xAI URLs auto-upgrade.
if normalized_provider in ("nous", "custom") or is_actual_route(provider):
return False
if normalized_provider == "copilot":
try:
from hermes_cli.models import _should_use_copilot_responses_api
return _should_use_copilot_responses_api(model)
except Exception:
pass # fall back to the generic GPT-5 rule
return AIAgent._model_requires_responses_api(model)
def _max_tokens_param(self, value: int) -> dict:
"""``max_completion_tokens`` for newer OpenAI families (and Azure / Copilot serving them), else
``max_tokens``. URL-first, then model-name fallback for third-party endpoints fronting those models."""
if (self._is_direct_openai_url() or self._is_azure_openai_url() or self._is_github_copilot_url()
or model_forces_max_completion_tokens(self.model)):
return {"max_completion_tokens": value}
return {"max_tokens": value}
@staticmethod
def _requested_output_cap_from_api_kwargs(api_kwargs: Any) -> Optional[int]:
"""Extract the outgoing response token cap from a prepared request."""
if not isinstance(api_kwargs, dict):
return None
for key in ("max_output_tokens", "max_completion_tokens", "max_tokens"):
try:
value = int(api_kwargs.get(key))
except (TypeError, ValueError):
continue
if value > 0:
return value
return None
def _has_content_after_think_block(self, content: str) -> bool:
"""True when text remains after stripping reasoning blocks (reasoning-only output is retried)."""
return bool(content) and bool(self._strip_think_blocks(content).strip())
_strip_think_blocks = _forward("agent.agent_runtime_helpers", "strip_think_blocks")
@staticmethod
def _has_natural_response_ending(content: str) -> bool:
"""Heuristic: does visible assistant text look intentionally finished?"""
stripped = (content or "").rstrip()
if not stripped:
return False
last = stripped[-1]
# Closing punctuation/brackets, a fenced-code close, or an emoji (Misc Symbols, Dingbats, Emoticons, ...).
return stripped.endswith("```") or last in '.!?:)"\']}。!?:)】」』》^' or ord(last) >= 0x1F300
def _is_ollama_glm_backend(self) -> bool:
"""Ollama-hosted GLM models misreport finish_reason='stop'. Matches only explicit Ollama signatures
(port 11434, "ollama" in URL, provider ollama), never arbitrary local proxies; excludes Ollama Cloud
(``ollama.com`` / ``:cloud``), which reports faithfully — rewriting it would manufacture truncations.
Crucially it does NOT match arbitrary local/private endpoints (LiteLLM/sglang/vLLM/LM Studio
proxies, Tailscale boxes), which report finish_reason correctly and were the source of #13971's
false-positive truncation continuations.
Two signatures identify it: the ``ollama.com`` host (provider ``ollama-cloud``) and the ``:cloud``
model suffix (cloud generation proxied through a local 11434 endpoint, #98406). Applying the
stop→length rewrite to them manufactures false truncations and causes the continuation nudge to
consume the model's output budget on the next retry, making further false-positives more likely.
"""
model_lower = (self.model or "").lower()
provider_lower = (self.provider or "").lower()
if "glm" not in model_lower and provider_lower != "zai":
return False
base = self._base_url_lower
# Ollama Cloud (hosted service or :cloud proxy) forwards finish_reason faithfully — do not rewrite.
if "ollama.com" in base or ":cloud" in model_lower:
return False
if "ollama" in base or ":11434" in base:
return True
return provider_lower == "ollama"
def _should_treat_stop_as_truncated(self, finish_reason: str, assistant_message, messages: Optional[list] = None) -> bool:
"""Detect conservative stop->length misreports for Ollama-hosted GLM models."""
if finish_reason != "stop" or self.api_mode != "chat_completions" or not self._is_ollama_glm_backend():
return False
if not any(isinstance(msg, dict) and msg.get("role") == "tool" for msg in (messages or [])):
return False
if assistant_message is None or getattr(assistant_message, "tool_calls", None):
return False
content = getattr(assistant_message, "content", None)
if not isinstance(content, str):
return False
visible_text = self._strip_think_blocks(content).strip()
if len(visible_text) < 20 or not re.search(r"\s", visible_text):
return False
return not self._has_natural_response_ending(visible_text)
_looks_like_codex_intermediate_ack = _forward("agent.agent_runtime_helpers", "looks_like_codex_intermediate_ack")
_extract_reasoning = _forward("agent.agent_runtime_helpers", "extract_reasoning")
_cleanup_task_resources = _forward("agent.chat_completion_helpers", "cleanup_task_resources")
# Background memory/skill review — prompts live in agent.background_review.
from agent.background_review import _MEMORY_REVIEW_PROMPT, _SKILL_REVIEW_PROMPT, _COMBINED_REVIEW_PROMPT
_summarize_background_review_actions = _forward_static("agent.background_review", "summarize_background_review_actions")
def _spawn_background_review(self, messages_snapshot: List[Dict], review_memory: bool = False,
review_skills: bool = False, focus: Optional[str] = None, explicit: bool = False) -> None:
"""Post-turn review entry point: decide WHEN, then spawn.
A review whose runtime is the MANAGED LOCAL llama-server is queued for machine idle (``defer: auto``)
instead of hitting the user's GPU mid-session; everything else spawns immediately. ``explicit``
(/refine) is never deferred but does not touch the ``focus``-keyed delegate/enabled gates.
"""
# Gates run at enqueue/spawn time; the idle dispatcher re-checks `enabled` at dispatch time.
if focus is None and getattr(self, "_delegate_depth", 0) > 0:
return
task_cfg = None
if focus is None:
from agent.background_review import load_background_review_settings
enabled, task_cfg = load_background_review_settings()
if not enabled:
return
# Structural clone at the single chokepoint: the fork sanitizes in place, and a shallow copy would
# alias the live history's nested tool_calls/content.
# Structural clone at the single chokepoint every review path (automatic, /refine, idle-queue
# deferral) goes through. See #100795.
from agent.turn_finalizer import _clone_background_review_messages
kwargs = dict(messages_snapshot=_clone_background_review_messages(messages_snapshot),
review_memory=review_memory, review_skills=review_skills, focus=focus, task_cfg=task_cfg,
explicit=explicit)
if focus is None and not explicit and _review_should_defer(self, task_cfg):
from agent.review_idle_queue import QUEUE
QUEUE.enqueue(self, _review_queue_key(self), kwargs)
return
self._spawn_background_review_now(**kwargs)
def _spawn_background_review_now(self, messages_snapshot: List[Dict], review_memory: bool = False,
review_skills: bool = False, focus: Optional[str] = None,
task_cfg: Optional[Dict[str, Any]] = None, _requeue_attempts: int = 0,
explicit: bool = False) -> None:
"""Spawn the background memory/skill review thread.
``threading.Thread`` is constructed here so tests patching ``run_agent.threading.Thread`` keep working.
``focus`` is /refine steering text; ``task_cfg`` is the pre-loaded config block (None on direct calls).
``explicit`` (/refine) forks under the ``refine_review`` write origin, keeping the full
memory operation set. A deferred review preempted by a live turn is requeued (bounded)
rather than lost.
"""
from agent.background_review import (
finish_background_review_run, prepare_background_review_run, spawn_background_review_thread,
)
from tools.thread_context import propagate_context_to_thread
review_run = prepare_background_review_run(self)
if review_run is None:
return
try:
target, _prompt = spawn_background_review_thread(
self, messages_snapshot, review_memory=review_memory, review_skills=review_skills,
focus=focus, task_cfg=task_cfg, review_run=review_run, explicit=explicit,
)
def _target_with_requeue() -> None:
target()
self._maybe_requeue_preempted_review(review_run, dict(
messages_snapshot=messages_snapshot, review_memory=review_memory, review_skills=review_skills,
focus=focus, task_cfg=task_cfg, _requeue_attempts=_requeue_attempts + 1,
explicit=explicit))
# Carry the active profile into the review thread so MEMORY.md / skill review writes land in the
# right profile.
threading.Thread(target=propagate_context_to_thread(_target_with_requeue), daemon=True, name="bg-review").start()
except Exception:
finish_background_review_run(self, review_run)
raise
_REVIEW_REQUEUE_MAX_ATTEMPTS = 3
def _maybe_requeue_preempted_review(self, review_run, kwargs) -> None:
"""Requeue a deferred-mode review that a live turn cancelled.
Only for automatic reviews on the managed local runtime; bounded attempts stop a busy box cycling
forever.
"""
try:
# Not cancelled == ran to completion (or was never admitted).
if not review_run.cancel_requested.is_set() or kwargs.get("focus") is not None:
return
if kwargs.get("_requeue_attempts", 0) > self._REVIEW_REQUEUE_MAX_ATTEMPTS:
logger.info("Preempted background review dropped after %d requeues", self._REVIEW_REQUEUE_MAX_ATTEMPTS)
return
if not _review_should_defer(self, kwargs.get("task_cfg")):
return
from agent.review_idle_queue import QUEUE
# kwargs carries the incremented _requeue_attempts through the queue so the cap survives.
QUEUE.enqueue(self, _review_queue_key(self), dict(kwargs))
except Exception: # noqa: BLE001 — requeue is best-effort
logger.debug("Preempted-review requeue failed", exc_info=True)
_build_memory_write_metadata = _forward("agent.background_review", "build_memory_write_metadata")
_apply_pending_steer_to_tool_results = _forward("agent.agent_runtime_helpers", "apply_pending_steer_to_tool_results")
def get_activity_summary(self) -> dict:
"""Diagnostic snapshot: ``last_activity_*`` plus the short aliases gateway and delegate readers use."""
from agent.session_activity import build_activity_snapshot
provenance = getattr(self, "_last_activity_provenance", None)
return build_activity_snapshot(
last_activity_at=getattr(self, "_last_activity_ts", None),
last_activity_description=getattr(self, "_last_activity_desc", None) or "",
last_activity_provenance=provenance if provenance is not None else ActivityProvenance.UNKNOWN,
extra={
"current_tool": self._current_tool, "api_call_count": self._api_call_count,
"max_iterations": self.max_iterations, "budget_used": self.iteration_budget.used,
"budget_max": self.iteration_budget.max_total,
},
)
def shutdown_memory_provider(self, messages: list = None) -> None:
"""Shut down the memory provider and context engine at session end (idempotent: gateway cleanup and
``close()`` may both call it)."""
if getattr(self, "_memory_provider_shutdown", False):
return
self._memory_provider_shutdown = True
if self._memory_manager:
try:
self._memory_manager.on_session_end(messages or [])
except Exception as e:
logger.warning("Memory provider on_session_end failed during shutdown: %s", e, exc_info=True)
_quietly(lambda: self._memory_manager.shutdown_all())
_notify_context_engine_session_end(self, messages)
def commit_memory_session(self, messages: list = None) -> None:
"""Flush end-of-session extraction on session_id rotation (/new, compression) without tearing providers
down."""
if self._memory_manager:
_quietly(lambda: self._memory_manager.on_session_end(messages or []))
_notify_context_engine_session_end(self, messages)
def _sync_external_memory_for_turn(self, *, original_user_message: Any, final_response: Any, interrupted: bool,
messages: list | None = None) -> None:
"""Mirror a completed turn into external memory providers (``sync_all`` + ``queue_prefetch_all``).
Uses ``original_user_message`` (``user_message`` may carry injected skill content). Interrupted turns
are skipped: partial output is not durable truth. Best-effort — an offline backend never blocks.
A partial assistant output, an aborted tool chain, or a mid-stream reset is not durable
conversational truth — mirroring it into an external memory backend pollutes future recall with
state the user never saw completed. The prefetch is gated on the same flag: the user's next message
is almost certainly a retry of the same intent, and a prefetch keyed on the interrupted turn would
fire against stale context. See #15218.
"""
if interrupted or not (self._memory_manager and final_response and original_user_message):
return
# Flatten multimodal parts to text (newline-joined for memory).
user_text = _summarize_user_message_for_log(original_user_message, sep="\n")
response_text = _summarize_user_message_for_log(final_response, sep="\n")
if not (user_text and response_text):
return
try:
sync_kwargs = {"session_id": self.session_id or "", **({"messages": messages} if messages is not None else {})}
# Stashed by build_turn_context() for this turn, None on a human turn.
turn_author = getattr(self, "_turn_author", None)
if turn_author is not None:
sync_kwargs["turn_author"] = turn_author
self._memory_manager.sync_all(user_text, response_text, **sync_kwargs)
# Sibling of the build_turn_context() prefetch gate: don't key recall on zero-signal prompts.
if not is_trivial_prompt(user_text):
self._memory_manager.queue_prefetch_all(user_text, session_id=self.session_id or "")
except Exception:
pass
def release_clients(self) -> None:
"""Release LLM clients and child agents WITHOUT tearing down session tool state (gateway cache
eviction: the session may resume on the same task_id, so processes, sandbox, browser, computer-use and
memory provider are kept). Idempotent; distinct from ``close()``."""
self._close_active_children(soft=True)
# Retire (don't hard-close) the shared client: eviction runs on the gateway memory-manager thread,
# and a cross-thread close can release TLS FDs under a still-unwinding worker.
_quietly(self._drop_shared_client, lambda c: self._retire_shared_openai_client(c, reason="cache_evict"))
self._close_request_clients("cache_evict")
def close(self) -> None:
"""Release every resource this agent holds (idempotent); each phase is guarded so one failure never
blocks the rest."""
# close() is the hard owner boundary; shutdown_memory_provider() is idempotent so gateway pre-calls
# never double-extract.
session_messages = getattr(self, "_session_messages", None)
_quietly(self.shutdown_memory_provider, session_messages if isinstance(session_messages, list) else None)
self._close_task_resources(getattr(self, "session_id", None) or "")
self._close_active_children(soft=False)
_quietly(self._drop_shared_client, lambda c: self._close_openai_client(c, reason="agent_close", shared=True))
self._close_request_clients("agent_close")
_quietly(self._close_codex_session)
# Free conversation history proactively: callers may still hold the closed agent. The DB-flush
# settled-prefix snapshot and the streamed-text accumulator are shadow copies of the same transcript;
# on a closed delegate child they were the only remaining owners, pinning its history in the parent heap.
self._session_messages = []
self._db_flush_scan_prefix = None
self._streamed_assistant_text_parts = []
_quietly(self._trim_process_memory)
_quietly(self._finalize_owned_session_row)
# -- close()/release_clients() phases -------------------------------------------------------------
def _close_active_children(self, *, soft: bool) -> None:
"""Detach and close per-turn child agents; ``soft`` releases their clients first, falling back to close()."""
try:
with self._active_children_lock:
children = list(self._active_children)
self._active_children.clear()
except Exception:
return
for child in children:
if soft:
try:
child.release_clients()
continue
except Exception:
pass
_quietly(lambda: child.close())
def _drop_shared_client(self, close_fn: Callable[[Any], None]) -> None:
"""Hand the shared OpenAI/httpx client to ``close_fn`` and clear the attribute."""
# Retire the OpenAI/httpx client to release sockets immediately. #70773: eviction runs on the
# gateway's memory-manager thread — a cross-thread hard close of the shared client can release TLS
# FDs under a still-unwinding worker (FD-recycle → SQLite corruption). Retirement shuts the pooled
# sockets down (the memory/socket win we want here) and lets GC release the FDs once no thread holds
# them.
client = getattr(self, "client", None)
if client is not None:
close_fn(client)
self.client = None
def _close_request_clients(self, reason: str) -> None:
"""Drop the cached per-request wire clients (reused across sequential LLM calls)."""
_quietly(self._close_cached_request_openai_client, reason=reason)