fix: persist the codex thread id per session and thread/resume it across an API-server restart
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.
This commit is contained in:
@@ -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,
|
||||
)
|
||||
|
||||
95
tests/agent/test_codex_app_server_thread_resume.py
Normal file
95
tests/agent/test_codex_app_server_thread_resume.py
Normal file
@@ -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()
|
||||
@@ -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)
|
||||
|
||||
@@ -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 │
|
||||
|
||||
Reference in New Issue
Block a user