From a9ef4a7625303370f80454de8138fa2ea7b2d234 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Mon, 7 Sep 2026 01:58:51 -0700 Subject: [PATCH] fix(codex): keep transport echoes out of durable user history Port the exact submitted-wire-text ownership boundary from #93546 onto current topical runtime code. Do not add the candidate's mocked-result fallback or storage-level content deduplication. Preserve later distinct and identical user events, separate identical accepted turns, and keyless inputs. Add two regression invariants and offline subprocess-wire A/B. Local wire A/B: 4/8 control matrix passing on base, 8/8 after. Broader tests queued behind campaign lock; not ready for merge. Refs #104653 Original diagnosis: @gitszabolcs (#38254) Original implementation: #43127, submitted by @vashkartik Focused salvage and wire-text correction: @fancyboi999 (#93546) Current-main carry-forward considered: #104698 Co-authored-by: Xinmin Zeng <135568692+fancyboi999@users.noreply.github.com> Co-authored-by: VECTOR --- agent/codex_runtime.py | 10 ++- agent/transports/codex_app_server_session.py | 5 +- evals/codex_echo/fixture_server.py | 38 ++++++++++++ evals/codex_echo/probe.py | 61 +++++++++++++++++++ tests/agent/test_codex_echo_ownership.py | 50 +++++++++++++++ .../test_codex_app_server_session.py | 15 +++++ .../docs/developer-guide/session-storage.md | 11 ++++ 7 files changed, 188 insertions(+), 2 deletions(-) create mode 100755 evals/codex_echo/fixture_server.py create mode 100644 evals/codex_echo/probe.py create mode 100644 tests/agent/test_codex_echo_ownership.py diff --git a/agent/codex_runtime.py b/agent/codex_runtime.py index f15732f708..e58638cca8 100644 --- a/agent/codex_runtime.py +++ b/agent/codex_runtime.py @@ -419,7 +419,15 @@ def _persist_projected_messages(agent, turn, messages: List[Dict[str, Any]]) -> if not turn.projected_messages: return from agent.message_metadata import append_message - for projected_message in turn.projected_messages: + projected_messages = turn.projected_messages + # Turn-start persistence owns the accepted input. Codex's leading user item + # only echoes its coerced wire text; later/nonmatching events remain real. + submitted_user_text = getattr(turn, "submitted_user_text", None) + first = projected_messages[0] + if (submitted_user_text is not None and first.get("role") == "user" + and first.get("content") == submitted_user_text): + projected_messages = projected_messages[1:] + for projected_message in projected_messages: append_message(messages, projected_message) if getattr(agent, "_session_db", None) is None: return diff --git a/agent/transports/codex_app_server_session.py b/agent/transports/codex_app_server_session.py index 5ff11c73e4..effc28e69d 100644 --- a/agent/transports/codex_app_server_session.py +++ b/agent/transports/codex_app_server_session.py @@ -46,6 +46,8 @@ class TurnResult: error: Optional[str] = None # non-recoverable turn error turn_id: Optional[str] = None thread_id: Optional[str] = None + # Exact turn/start text distinguishes the input echo from a new user event. + submitted_user_text: Optional[str] = None token_usage_last: Optional[dict[str, Any]] = None model_context_window: Optional[int] = None compacted: bool = False @@ -341,9 +343,10 @@ class CodexAppServerSession: if self._interrupt_event.is_set(): result.interrupted = True else: + result.submitted_user_text = _coerce_turn_input_text(user_input) ts = self._request_for( result, "turn/start", - {"threadId": self._thread_id, "input": [{"type": "text", "text": _coerce_turn_input_text(user_input)}]}, + {"threadId": self._thread_id, "input": [{"type": "text", "text": result.submitted_user_text}]}, "turn/start", ) if ts is not None: diff --git a/evals/codex_echo/fixture_server.py b/evals/codex_echo/fixture_server.py new file mode 100755 index 0000000000..fd01e1a664 --- /dev/null +++ b/evals/codex_echo/fixture_server.py @@ -0,0 +1,38 @@ +#!/usr/bin/env python3 +"""Offline Codex JSON-RPC peer; no provider or authentication calls.""" +import json +import os +import sys + + +def send(value): + print(json.dumps(value), flush=True) + + +for line in sys.stdin: + request = json.loads(line) + if "id" not in request: + continue + method = request["method"] + result = {} + if method == "thread/start": + result = {"thread": {"id": "fixture-thread"}} + elif method == "turn/start": + result = {"turn": {"id": "fixture-turn"}} + send({"jsonrpc": "2.0", "id": request["id"], "result": result}) + if method != "turn/start": + continue + text = request["params"]["input"][0]["text"] + mode = os.environ.get("ECHO_FIXTURE_MODE", "echo") + items = [] + if mode != "assistant_only": + items.append({"type": "userMessage", "id": "input", "content": [ + {"type": "text", "text": "different" if mode == "different" else text}]}) + items.append({"type": "agentMessage", "id": "reply", "text": "fixture reply"}) + if mode == "later_equal": + items.append({"type": "userMessage", "id": "steer", "content": [{"type": "text", "text": text}]}) + for item in items: + send({"method": "item/completed", "params": { + "threadId": "fixture-thread", "turnId": "fixture-turn", "item": item}}) + send({"method": "turn/completed", "params": { + "threadId": "fixture-thread", "turn": {"id": "fixture-turn", "status": "completed", "error": None}}}) diff --git a/evals/codex_echo/probe.py b/evals/codex_echo/probe.py new file mode 100644 index 0000000000..1b82e1aeee --- /dev/null +++ b/evals/codex_echo/probe.py @@ -0,0 +1,61 @@ +"""Run from the checkout: PYTHONPATH=. python evals/codex_echo/probe.py. + +Real subprocess pipes, protocol parser/projector, persistence and SQLite replay. +The peer is an offline fixture, not a Codex/provider validation. +""" +import json +import os +from pathlib import Path +import tempfile +import threading + + +def main(): + with tempfile.TemporaryDirectory(prefix="hermes-echo-") as tmp: + os.environ.update(HOME=tmp, HERMES_HOME=tmp, CODEX_HOME=tmp, HERMES_DISABLE_PLUGINS="1") + from agent import codex_runtime + from agent.message_metadata import append_message + from agent.session_persistence import SessionPersistenceMixin + from agent.transports.codex_app_server_session import CodexAppServerSession + from hermes_state import SessionDB + + observations = [] + for mode in ["echo", "assistant_only", "different", "later_equal"]: + for rich in [False, True]: + os.environ["ECHO_FIXTURE_MODE"] = mode + sid = f"{mode}-{rich}" + db = SessionDB(Path(tmp) / f"{sid}.db") + db.create_session(session_id=sid, source="telegram", model="fixture") + agent = SessionPersistenceMixin() + agent.session_id = sid + agent._session_db = db + agent._session_db_created = True + agent._last_flushed_db_idx = 0 + agent._session_persist_lock = threading.RLock() + session = CodexAppServerSession(cwd=tmp, codex_home=tmp, + codex_bin=str(Path(__file__).with_name("fixture_server.py").resolve())) + messages = [] + wire_input = [{"type": "text", "text": "caption"}, + {"type": "image_url", "image_url": {"url": "data:image/png;base64,abc"}}] if rich else "caption" + try: + for accepted in range(2): + append_message(messages, {"role": "user", "content": "caption", + "platform_message_id": str(accepted) if not rich else None}) + assert agent._flush_messages_to_session_db(messages) + turn = session.run_turn(wire_input, turn_timeout=10) + assert turn.error is None, turn.error + assert turn.final_text == "fixture reply", turn + codex_runtime._persist_projected_messages(agent, turn, messages) + assert agent._flush_messages_to_session_db(messages) + users = [m["content"] for m in db.get_messages_as_conversation(sid) if m["role"] == "user"] + expected_count = 4 if mode in {"different", "later_equal"} else 2 + observations.append({"mode": mode, "rich_keyless": rich, "users": users, + "expected_user_rows": expected_count, "user_rows": len(users), "pass": len(users) == expected_count}) + finally: + session.close() + db.close() + print(json.dumps({"module": codex_runtime.__file__, "cases": observations}, indent=2)) + + +if __name__ == "__main__": + main() diff --git a/tests/agent/test_codex_echo_ownership.py b/tests/agent/test_codex_echo_ownership.py new file mode 100644 index 0000000000..7b9beebae4 --- /dev/null +++ b/tests/agent/test_codex_echo_ownership.py @@ -0,0 +1,50 @@ +"""The accepted input owns its row; transport echoes do not own new rows.""" +import threading +from types import SimpleNamespace + +import pytest + +from agent.codex_runtime import _persist_projected_messages +from agent.message_metadata import append_message +from agent.session_persistence import SessionPersistenceMixin +from agent.transports.codex_event_projector import CodexEventProjector +from hermes_state import SessionDB + + +@pytest.mark.parametrize("platform_id", ["2146", None]) +@pytest.mark.parametrize("projection", ["echo", "assistant_only", "different", "later_equal"]) +def test_only_submitted_leading_echo_is_excluded(tmp_path, platform_id, projection): + db = SessionDB(tmp_path / "state.db") + try: + db.create_session(session_id="echo", source="telegram", model="codex") + agent = SessionPersistenceMixin() + agent.session_id = "echo" + agent._session_db = db + agent._session_db_created = True + agent._last_flushed_db_idx = 0 + agent._session_persist_lock = threading.RLock() + messages = [] + # Two independently accepted identical inputs must both survive. + for turn_index in range(2): + append_message(messages, {"role": "user", "content": "accepted", "platform_message_id": platform_id}) + assert agent._flush_messages_to_session_db(messages) + projector = CodexEventProjector() + items = [] + if projection != "assistant_only": + items.append({"type": "userMessage", "id": f"u{turn_index}", "content": [ + {"type": "text", "text": "different" if projection == "different" else "wire caption"}]}) + items.append({"type": "agentMessage", "id": f"a{turn_index}", "text": "reply"}) + if projection == "later_equal": + items.append({"type": "userMessage", "id": f"s{turn_index}", "content": [ + {"type": "text", "text": "wire caption"}]}) + projected = [] + for item in items: + projected.extend(projector.project({"method": "item/completed", "params": {"item": item}}).messages) + _persist_projected_messages(agent, SimpleNamespace( + projected_messages=projected, submitted_user_text="wire caption"), messages) + assert agent._flush_messages_to_session_db(messages) + users = [row["content"] for row in db.get_messages_as_conversation("echo") if row["role"] == "user"] + extra = {"different": ["different"], "later_equal": ["wire caption"]}.get(projection, []) + assert users == (["accepted"] + extra) * 2 + 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 5a82567ba6..7548e59e29 100644 --- a/tests/agent/transports/test_codex_app_server_session.py +++ b/tests/agent/transports/test_codex_app_server_session.py @@ -211,6 +211,21 @@ class TestRunTurn: + def test_result_records_the_exact_submitted_input(self): + client = FakeClient() + client.queue_notification( + "turn/completed", threadId="t", + turn={"id": "turn-fake-001", "status": "completed", "error": None}, + ) + rich_input = [ + {"type": "text", "text": "caption"}, + {"type": "image_url", "image_url": {"url": "data:image/png;base64,abc"}}, + ] + result = make_session(client).run_turn(rich_input, turn_timeout=2.0) + _, params = next(request for request in client.requests if request[0] == "turn/start") + assert result.submitted_user_text == params["input"][0]["text"] + assert result.submitted_user_text != rich_input + def test_foreign_completion_in_server_request_drain_is_ignored(self): """Approval draining must not project a child result into the parent.""" client = FakeClient() diff --git a/website/docs/developer-guide/session-storage.md b/website/docs/developer-guide/session-storage.md index 5329d862b6..09fcaca450 100644 --- a/website/docs/developer-guide/session-storage.md +++ b/website/docs/developer-guide/session-storage.md @@ -31,6 +31,17 @@ history that appears to revert. +## Codex app-server input ownership + +The agent persists an accepted user input before starting its Codex turn. Codex +then projects that input as a leading `userMessage` notification. At the runtime +splice boundary, Hermes excludes only that leading item when it exactly matches +the text serialized into `turn/start`, including rich-input coercion. Later or +nonmatching user events remain intact, as do separately accepted identical turns. +This also applies to synthetic/keyless input; it does not depend on a platform +message ID. Existing historical duplicates are not rewritten. The gateway skips +its transcript write when the agent reports that it owns persistence. + ## Architecture Overview ```