From e77e24a6a650abe1882449dc94aaf8ffe1594768 Mon Sep 17 00:00:00 2001 From: teknium1 <127238744+teknium1@users.noreply.github.com> Date: Sat, 19 Sep 2026 12:24:17 -0700 Subject: [PATCH] fix: persist the codex thread id per session and thread/resume it across an API-server restart MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit After a codex app-server turn's projected rows are durable in the session DB, store the codex thread id as ``codex_thread_id`` in the session row's model_config (atomic merge via patch_session_model_config; never for a retired thread). The FIRST CodexAppServerSession an AIAgent builds for that session passes the stored id as resume_thread_id, so a rebuilt agent — the next /api/sessions/{id}/chat request, or the first turn after the API server or gateway restarts — resumes the model-side thread before turn/start instead of starting an empty one while Hermes' own transcript continues. Fail closed when codex cannot hand the thread back (rollout gone, CODEX_HOME changed, previous app-server killed mid-write): drop the stored id, start a fresh thread on the same client, and say so once — "Codex thread could not be resumed; starting a new one." — through _emit_diagnostic_status, the lifecycle status rail every surface renders (CLI vprint, TUI/Desktop and gateway status_callback). No other lifecycle change: a retired or prompt-recreated session in the same process keeps today's fresh-thread behaviour and overwrites the binding once its turn commits. Why: CodexAppServerSession kept the thread id in memory only, so every API-server restart (and every per-request agent) silently reset the model's memory of the conversation (#100531). Supersedes the persistence half of #100528 (_persist_projected_messages now reports durability) and the refuse-and-raise policy of #103352 with the maintainer-approved bounded slice. --- agent/codex_runtime.py | 61 ++++++++++-- .../test_codex_app_server_thread_resume.py | 95 +++++++++++++++++++ .../test_codex_app_server_session.py | 28 ++++++ .../features/codex-app-server-runtime.md | 3 +- 4 files changed, 180 insertions(+), 7 deletions(-) create mode 100644 tests/agent/test_codex_app_server_thread_resume.py diff --git a/agent/codex_runtime.py b/agent/codex_runtime.py index 09009f5996..4292367754 100644 --- a/agent/codex_runtime.py +++ b/agent/codex_runtime.py @@ -459,10 +459,51 @@ def _codex_developer_instructions(agent) -> str: return developer_instructions +# Durable codex thread binding: ``sessions.model_config.codex_thread_id`` (hermes_state), written after the +# turn's projected rows were committed, read by the next AIAgent built for the same Hermes session so an +# API-server restart (or the per-request agents of /api/sessions/{id}/chat) resumes the model-side thread +# instead of starting an empty one while Hermes' own transcript continues (#100531). +_CODEX_THREAD_ID_KEY = "codex_thread_id" +_CODEX_THREAD_RESUME_NOTICE = "Codex thread could not be resumed; starting a new one." + + +def _stored_codex_thread_id(agent) -> str | None: + db, session_id = getattr(agent, "_session_db", None), getattr(agent, "session_id", None) + if db is None or not session_id: + return None + thread_id = db.get_session_model_config_value(session_id, _CODEX_THREAD_ID_KEY) + return thread_id if isinstance(thread_id, str) and thread_id else None + + +def _store_codex_thread_id(agent, thread_id: str | None) -> None: + """Merge (``None`` clears) the binding into the session row; a failed write only logs — the turn is done.""" + db, session_id = getattr(agent, "_session_db", None), getattr(agent, "session_id", None) + if db is None or not session_id: + return + _call_guarded(db.patch_session_model_config, "codex thread id could not be stored on the session row", + args=(session_id, {_CODEX_THREAD_ID_KEY: thread_id})) + + +def _start_codex_thread(agent) -> str: + """``ensure_started`` with the fail-closed resume policy: a stored thread that codex cannot hand back + (unknown id, rollout locked by a killed app-server, different thread) is dropped from the session row, + the user is told once on the status rail, and a fresh thread starts on the same client.""" + from agent.transports.codex_app_server_session import CodexThreadResumeError + try: + return agent._codex_session.ensure_started() + except CodexThreadResumeError as exc: + logger.warning("%s; starting a new codex thread (session=%s)", exc.message, getattr(agent, "session_id", None)) + _store_codex_thread_id(agent, None) + agent._emit_diagnostic_status(_CODEX_THREAD_RESUME_NOTICE) + return agent._codex_session.ensure_started() + + def _ensure_codex_session(agent) -> None: """Lazily spawn one CodexAppServerSession per AIAgent (reused across turns, closed by the _cleanup hook). A live session whose thread was started with a different prompt composition (TUI/Desktop ``/personality`` - or a prompt mirror mutate the agent in place) is retired first so the new thread carries the current one.""" + or a prompt mirror mutate the agent in place) is retired first so the new thread carries the current one. + Only the FIRST session of an AIAgent resumes the stored codex thread: a retired/recreated one keeps + today's fresh-thread behaviour and overwrites the binding once its turn is committed.""" developer_instructions = _codex_developer_instructions(agent) if getattr(agent, "_codex_session", None) is not None: # Only a session whose recorded composition differs is stale; one attached without a record is kept. @@ -470,6 +511,7 @@ def _ensure_codex_session(agent) -> None: if recorded is None or recorded == developer_instructions: return _close_codex_session(agent) + resume_thread_id = None if getattr(agent, "_codex_session_prompt", None) is not None else _stored_codex_thread_id(agent) from agent.runtime_cwd import resolve_agent_cwd from agent.transports.codex_app_server_session import CodexAppServerSession, _ServerRequestRouting from hermes_cli.codex_runtime_switch import get_configured_codex_binary @@ -511,17 +553,19 @@ def _ensure_codex_session(agent) -> None: on_event=make_codex_app_server_event_bridge(agent), developer_instructions=developer_instructions or None, model=getattr(agent, "model", None) if model_provider else None, model_provider=model_provider, + resume_thread_id=resume_thread_id, ) -def _persist_projected_messages(agent, turn, messages: List[Dict[str, Any]]) -> None: - """Splice the projected messages into ``messages`` and flush them to the session DB. +def _persist_projected_messages(agent, turn, messages: List[Dict[str, Any]]) -> bool: + """Splice the projected messages into ``messages`` and flush them to the session DB; True when the + rows are durable in the session DB (the codex thread binding may then be published). Bypasses conversation_loop's per-step _persist_session(); the flush dedups via _DB_PERSISTED_MARKER so only the new codex rows are written. The agent stays the sole persister (agent_persisted=True): a gateway re-write would re-INSERT the user turn.""" if not turn.projected_messages: - return + return False from agent.message_metadata import append_message projected_messages = turn.projected_messages # Turn-start persistence owns the accepted input. Codex's leading user item @@ -534,7 +578,7 @@ def _persist_projected_messages(agent, turn, messages: List[Dict[str, Any]]) -> for projected_message in projected_messages: append_message(messages, projected_message) if getattr(agent, "_session_db", None) is None: - return + return False flush_ok = False try: flush_ok = agent._flush_messages_to_session_db(messages) @@ -544,6 +588,7 @@ def _persist_projected_messages(agent, turn, messages: List[Dict[str, Any]]) -> # Output already streamed and agent_persisted cannot flip to False: surface the gap loudly. logger.warning("codex app-server turn was delivered but could NOT be persisted to the session DB " "(session=%s) — this turn will be missing after restart/resume", getattr(agent, "session_id", None)) + return flush_ok is True def _finish_codex_turn(agent, turn, messages: List[Dict[str, Any]], *, original_user_message: Any, @@ -583,6 +628,7 @@ def run_codex_app_server_turn(agent, *, user_message: str, original_user_message "without a truthful pre-compaction transcript boundary") _ensure_codex_session(agent) try: + _start_codex_thread(agent) turn = agent._codex_session.run_turn(user_input=user_message) except Exception as exc: logger.exception("codex app-server turn failed") @@ -597,7 +643,10 @@ def run_codex_app_server_turn(agent, *, user_message: str, original_user_message if getattr(turn, "should_retire", False): logger.warning("codex app-server session retired (turn error: %s)", turn.error) _close_codex_session(agent) - _persist_projected_messages(agent, turn, messages) + # The binding is published only once the transcript it belongs to is durable, and never for a + # retired thread (the next agent would only resume into the same wedge). + if _persist_projected_messages(agent, turn, messages) and not getattr(turn, "should_retire", False): + _store_codex_thread_id(agent, turn.thread_id) usage_result = _finish_codex_turn( agent, turn, messages, original_user_message=original_user_message, should_review_memory=should_review_memory, ) diff --git a/tests/agent/test_codex_app_server_thread_resume.py b/tests/agent/test_codex_app_server_thread_resume.py new file mode 100644 index 0000000000..17dd60015a --- /dev/null +++ b/tests/agent/test_codex_app_server_thread_resume.py @@ -0,0 +1,95 @@ +"""#100531 — the codex app-server thread survives an AIAgent rebuild (API-server restart, per-request agents). + +``CodexAppServerSession`` keeps the codex thread id in memory only, so every new ``AIAgent`` for the same +Hermes session used to ``thread/start`` an empty thread while Hermes' own transcript continued. The runtime +now publishes ``codex_thread_id`` into the session row's ``model_config`` once the turn's projected rows are +durable, the next agent for that session issues ``thread/resume`` for it, and a stored id codex cannot hand +back fails closed: fresh thread, binding dropped, one status-rail notice. +""" + +from pathlib import Path + +from agent.transports import codex_app_server_session as session_mod +from agent.transports.codex_app_server import CodexAppServerError +from agent.transports.codex_app_server_session import CodexAppServerSession, TurnResult +from hermes_state import SessionDB + +SID = "sess-codex-restart" + + +class _WireClient: + """Minimal app-server stand-in: answers thread/start with a fresh id, thread/resume with the requested + id (or refuses ids in ``dead``), and records every JSON-RPC method it saw.""" + + dead: set[str] = set() + instances: list["_WireClient"] = [] + counter = 0 + + def __init__(self, **kwargs): + self.requests: list[tuple[str, dict]] = [] + _WireClient.instances.append(self) + + def initialize(self, **kwargs): + return {} + + def request(self, method, params=None, timeout=30.0): + params = params or {} + self.requests.append((method, params)) + if method == "thread/resume": + if params["threadId"] in _WireClient.dead: + raise CodexAppServerError(code=-32600, message=f"no rollout found for thread id {params['threadId']}") + return {"thread": {"id": params["threadId"]}} + _WireClient.counter += 1 + return {"thread": {"id": f"thread-{_WireClient.counter}"}} + + def close(self): + pass + + +def _agent(db, **kwargs): + from run_agent import AIAgent + agent = AIAgent(api_key="stub", base_url="https://stub.invalid", provider="openai", api_mode="codex_app_server", + quiet_mode=True, skip_context_files=True, skip_memory=True, session_db=db, session_id=SID, **kwargs) + agent._spawn_background_review = lambda **kw: None + return agent + + +def _run_turn(self, user_input, **kwargs): + return TurnResult(final_text=f"echo {user_input}", thread_id=self._thread_id, turn_id="turn-1", + projected_messages=[{"role": "assistant", "content": f"echo {user_input}"}]) + + +def test_rebuilt_agent_resumes_the_stored_codex_thread_and_an_unresumable_one_fails_closed(monkeypatch, tmp_path): + monkeypatch.setattr(session_mod, "CodexAppServerClient", _WireClient) + monkeypatch.setattr(CodexAppServerSession, "run_turn", _run_turn) + _WireClient.instances, _WireClient.dead, _WireClient.counter = [], set(), 0 + db = SessionDB(Path(tmp_path) / "state.db") + try: + # Process 1: first turn publishes the binding only after the transcript is durable. + first = _agent(db) + assert first.run_conversation("Remember the word amber.")["completed"] + assert db.get_session_model_config_value(SID, "codex_thread_id") == "thread-1" + assert [r["content"] for r in db.get_messages(SID)][-1] == "echo Remember the word amber." + + # Process 2 (restart): a new AIAgent for the same session resumes that thread before turn/start. + notices: list[str] = [] + def status_callback(kind, message): # the lifecycle rail every surface renders; other kinds are noise + if kind == "lifecycle": + notices.append(str(message)) + second = _agent(db) + second.status_callback = status_callback + assert second.run_conversation("Which word?", conversation_history=[])["completed"] + assert [m for m, _ in _WireClient.instances[1].requests] == ["thread/resume"] + assert _WireClient.instances[1].requests[0][1]["threadId"] == "thread-1" + assert notices == [] + + # Control: the stored thread is gone on the codex side -> fresh thread, binding rotated, one notice. + _WireClient.dead = {"thread-1"} + third = _agent(db) + third.status_callback = status_callback + assert third.run_conversation("And now?", conversation_history=[])["completed"] + assert [m for m, _ in _WireClient.instances[2].requests] == ["thread/resume", "thread/start"] + assert db.get_session_model_config_value(SID, "codex_thread_id") == "thread-2" + assert notices == ["Codex thread could not be resumed; starting a new one."] + finally: + db.close() diff --git a/tests/agent/transports/test_codex_app_server_session.py b/tests/agent/transports/test_codex_app_server_session.py index de483613f8..748ffad000 100644 --- a/tests/agent/transports/test_codex_app_server_session.py +++ b/tests/agent/transports/test_codex_app_server_session.py @@ -227,6 +227,34 @@ class TestLifecycle: assert thread_start_params(provider="openai-codex", requested_provider="openai-codex", model="gpt-5.4") == base assert thread_start_params(provider="custom", requested_provider="custom", model="gpt-5.4") == base + def test_stored_thread_is_resumed_and_an_unresumable_one_falls_back_to_a_fresh_start(self): + """#100531: a stored id goes out as ``thread/resume`` (same params as thread/start, never a + second ``thread/start``); when codex cannot hand it back the failure is typed and the NEXT + ensure_started() starts a fresh thread on the same handshaken client.""" + from agent.transports.codex_app_server import CodexAppServerError + from agent.transports.codex_app_server_session import CodexThreadResumeError + + client = FakeClient() + client._request_handler = lambda method, params: ( + {"thread": {"id": params["threadId"]}} if method == "thread/resume" else {"thread": {"id": "fresh-1"}}) + s = make_session(client, resume_thread_id="stored-1", developer_instructions="SOUL") + assert s.ensure_started() == s.ensure_started() == "stored-1" + assert [m for m, _ in client.requests] == ["thread/resume"] + assert client.requests[0][1] == {"threadId": "stored-1", "cwd": "/tmp", "personality": "none", "developerInstructions": "SOUL"} + + def refuse(method, params): + if method == "thread/resume": + raise CodexAppServerError(code=-32600, message=f"no rollout found for thread id {params['threadId']}") + return {"thread": {"id": "fresh-2"}} + client = FakeClient() + client._request_handler = refuse + s = make_session(client, resume_thread_id="gone-1") + with pytest.raises(CodexThreadResumeError) as exc_info: + s.ensure_started() + assert exc_info.value.thread_id == "gone-1" + assert s.ensure_started() == "fresh-2" + assert [m for m, _ in client.requests] == ["thread/resume", "thread/start"] + def test_close_idempotent(self): client = FakeClient() s = make_session(client) diff --git a/website/docs/user-guide/features/codex-app-server-runtime.md b/website/docs/user-guide/features/codex-app-server-runtime.md index be054ebaeb..dbcbd3366d 100644 --- a/website/docs/user-guide/features/codex-app-server-runtime.md +++ b/website/docs/user-guide/features/codex-app-server-runtime.md @@ -470,6 +470,7 @@ Known limitations: - **No inline patch preview in approval prompts when codex doesn't track the changeset.** Codex's `fileChange` approval params don't always carry the changeset. Hermes caches the data from the corresponding `item/started` notification when possible, but if approval arrives before the item has streamed, the prompt falls back to whatever `reason` codex provides. - **`fallback_providers` fail over only on quota and rate-limit failures.** When a codex app-server turn fails with a billing / usage-limit / rate-limit error, Hermes switches to the configured [fallback provider](./fallback-providers.md) and retries the same turn on it; auth failures (`codex login` expired), turn timeouts and unknown-model errors do not fail over on this runtime and surface as the turn's error instead. - **Conversation history is not projected into the codex thread.** The codex thread receives Hermes' system prompt when it starts plus each new user message; prior Hermes history (e.g. from a resumed session) is not replayed into it. When the composed prompt changes mid-session (for example `/personality` in the TUI or Desktop), the next turn retires the running thread and starts a new one carrying the updated prompt; that new thread does not inherit the retired thread's history. +- **The codex thread itself does survive a restart.** After each committed turn Hermes stores the codex thread id on the session row (`codex_thread_id` in the session's `model_config`, `hermes sessions` / `state.db`). The next agent built for that same Hermes session — a later `/api/sessions/{id}/chat` request, or the first turn after the API server or gateway restarts — issues `thread/resume` for the stored id before `turn/start`, so the model keeps its own memory of the earlier turns even though Hermes never replays its transcript. When codex cannot hand the thread back (its rollout was deleted, `CODEX_HOME` changed, the previous app-server was killed while still writing it), Hermes fails closed: it drops the stored id, starts a fresh thread and shows one line — `Codex thread could not be resumed; starting a new one.` — on the status rail of the surface you are on (CLI, TUI/Desktop, messaging gateway). A `/new` session never resumes an older thread. - **Sub-second cancellation isn't guaranteed.** Mid-stream interrupts (Ctrl+C while codex is responding) are sent via `turn/interrupt`, but if codex has already flushed the final message, you get the response anyway. If you find a bug, [open an issue](https://github.com/NousResearch/hermes-agent/issues) with the output of `hermes logs --since 5m`. Mention `codex-runtime` in the title so it's easy to triage. @@ -493,7 +494,7 @@ If you find a bug, [open an issue](https://github.com/NousResearch/hermes-agent/ ▼ │ ┌──────────────────────────────────┐ │ │ codex app-server (subprocess) │──────────────┘ - │ thread/start, turn/start │ + │ thread/start|resume, turn/start │ │ item/* notifications │ │ shell + apply_patch + update_plan│ │ view_image + sandbox │