Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
70e1ee7
fix(chat): preserve attached session ownership
Duang777 Sep 9, 2026
e918856
fix(coordination): stabilize fence read errors
Duang777 Sep 10, 2026
ff8c0de
test(chat): cover attached interrupt endpoint error code
Duang777 Sep 11, 2026
44719d0
Merge remote-tracking branch 'origin/main' into codex/fix-attached-se…
Duang777 Sep 11, 2026
4246bc1
fix(chat): preserve interrupt error codes
Duang777 Sep 11, 2026
1c0b4c1
Merge remote-tracking branch 'origin/main' into codex/fix-attached-se…
Duang777 Sep 11, 2026
df950ea
Merge remote-tracking branch 'origin/main' into codex/fix-attached-se…
Duang777 Sep 11, 2026
6d9f8be
Merge remote-tracking branch 'origin/main' into codex/fix-attached-se…
Duang777 Sep 11, 2026
e9516ae
Merge remote-tracking branch 'origin/main' into codex/fix-attached-se…
Duang777 Sep 11, 2026
b8d4903
Merge remote-tracking branch 'origin/main' into codex/fix-attached-se…
Duang777 Sep 11, 2026
5eb8d46
test(dashboard): disambiguate recovery pending state
Duang777 Sep 11, 2026
e79db35
Merge remote-tracking branch 'origin/main' into codex/fix-attached-se…
Duang777 Sep 11, 2026
84a76e1
test(dashboard): replay completed mock turn after reload
Duang777 Sep 11, 2026
d3e42ab
Merge remote-tracking branch 'origin/main' into codex/fix-attached-se…
Duang777 Sep 11, 2026
00d1936
Merge remote-tracking branch 'origin/main' into codex/fix-attached-se…
Duang777 Sep 11, 2026
9908571
test(dashboard): target exact recovery pending state
Duang777 Sep 11, 2026
596716c
fix(todo): type effect decision results
Duang777 Sep 11, 2026
f3302b3
Merge remote-tracking branch 'origin/main' into codex/fix-attached-se…
Duang777 Sep 11, 2026
9af4ac3
Merge main into codex/fix-attached-session-lifecycle
Duang777 Sep 12, 2026
71114ae
Merge main into codex/fix-attached-session-lifecycle
Duang777 Sep 12, 2026
b3f8d58
Merge main into codex/fix-attached-session-lifecycle
Duang777 Sep 12, 2026
c929e44
Merge main into codex/fix-attached-session-lifecycle
Duang777 Sep 12, 2026
60facc5
Merge main into codex/fix-attached-session-lifecycle
Duang777 Sep 13, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 32 additions & 9 deletions loopx/chat_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -434,10 +434,6 @@ def open_session(
agent_goal_id: str | None = None,
) -> tuple[dict[str, Any], bool]:
capability = next((item for item in self.capabilities() if item["agent_id"] == agent_id), None)
if capability is None:
raise ValueError(f"unknown Agent endpoint: {agent_id}")
if not capability["available"]:
raise ValueError(f"Agent endpoint is unavailable: {agent_id}")
if mode not in {"resume_latest", "new"}:
raise ValueError("mode must be resume_latest or new")
selected_channel = channel_id or f"goal.{goal_id}"
Expand All @@ -452,15 +448,22 @@ def open_session(
with self.lock:
route_lock = self.session_open_locks.setdefault(route_key, threading.Lock())
with route_lock:
latest = None
if mode == "resume_latest":
latest = self.store.latest_session(
goal_id=None if is_manager_channel(selected_channel) else goal_id,
agent_id=agent_id,
channel_id=selected_channel,
)
if latest is not None:
self._ensure_adapter(latest, work_dir=work_dir, objective=objective)
return self.store.load_session(latest["session_id"]) or latest, True
if latest is not None and latest.get("session_mode") == CHAT_SESSION_MODE_ATTACHED:
return latest, True
if capability is None:
raise ValueError(f"unknown Agent endpoint: {agent_id}")
if not capability["available"]:
raise ValueError(f"Agent endpoint is unavailable: {agent_id}")
if latest is not None:
self._ensure_adapter(latest, work_dir=work_dir, objective=objective)
return self.store.load_session(latest["session_id"]) or latest, True
adapter = self._start_adapter(
agent_id=agent_id,
work_dir=work_dir,
Expand Down Expand Up @@ -1114,8 +1117,13 @@ def scope_valid() -> bool:
self._fail_turn(session_id, turn_id, exc.error_code, str(exc), status="failed", gate=exc.gate)
if not adapter.healthcheck():
with self.lock:
self.adapters.pop(session_id, None)
self.store.update_session(session_id, status="stale", last_error_code="transport_disconnected")
if self.adapters.get(session_id) is adapter:
self.adapters.pop(session_id, None)
self.store.update_session(
session_id,
status="stale",
last_error_code="transport_disconnected",
)
except Exception as exc: # noqa: BLE001 - preserve compact runtime failure.
event_buffer.close()
if consume_interrupted():
Expand Down Expand Up @@ -1170,6 +1178,21 @@ def interrupt_turn(self, *, session_id: str, turn_id: str) -> dict[str, Any]:
raise KeyError("chat turn was not found")
if turn.get("status") in TERMINAL_TURN_STATES:
return turn
session = self.store.load_session(session_id)
if (
session
and session.get("session_mode") == CHAT_SESSION_MODE_ATTACHED
and session.get("active_turn_id") == turn_id
):
raise CodexChatAgentError(
"The attached host does not expose interrupt control to LoopX Chat.",
error_code="attached_session_interrupt_unavailable",
gate={
"kind": "host_tool_gate",
"summary": "The active Turn is owned by the attached host.",
"next_action": "Stop the Turn in the attached host, then retry.",
},
)
interrupting = self.store.update_turn(
session_id,
turn_id,
Expand Down
2 changes: 1 addition & 1 deletion loopx/chat_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -833,7 +833,7 @@ def _interrupt_turn(self, session_id: str, turn_id: str) -> None:
self._send_error("chat turn was not found", status=404)
return
except CodexChatAgentError as exc:
self._send_error(str(exc), status=424, gate=exc.gate)
self._send_error(str(exc), status=424, error_code=exc.error_code, gate=exc.gate)
return
self._send_json(
{
Expand Down
113 changes: 113 additions & 0 deletions tests/test_attached_session_broker.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
claim_attached_agent_turn,
complete_attached_agent_turn,
)
from loopx.chat_agent import CodexChatAgentError
from loopx.chat_runtime import ChatRuntimeController
from loopx.chat_server import ChatRequestHandler
from loopx.chat_store import CHAT_SESSION_MODE_ATTACHED, ChatSessionStore
Expand Down Expand Up @@ -362,6 +363,118 @@ def test_web_and_lark_share_ordered_attached_session_without_spawning(
]


def test_resume_latest_reuses_attached_session_without_local_adapter(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
store = ChatSessionStore(tmp_path)
session = _bind(store)["session"]
runtime = ChatRuntimeController(store=store, codex_bin="missing-codex")
monkeypatch.setattr(
runtime,
"_start_adapter",
lambda **_kwargs: pytest.fail("attached Session must not start a local adapter"),
)

resumed, reused = runtime.open_session(
goal_id=GOAL_ID,
agent_id=AGENT_ID,
work_dir=tmp_path,
objective="sample objective",
mode="resume_latest",
)

assert reused is True
assert resumed["session_id"] == session["session_id"]
assert runtime.adapters == {}


def test_attached_claimed_turn_interrupt_fails_closed(tmp_path: Path) -> None:
store = ChatSessionStore(tmp_path)
session_id = str(_bind(store)["session"]["session_id"])
runtime = ChatRuntimeController(store=store, codex_bin="missing-codex")
turn, _created = store.create_queued_turn(
session_id,
client_turn_id="claimed-interrupt",
message="keep host ownership",
)
claim_attached_agent_turn(
store=store,
session_id=session_id,
host_surface=HOST_SURFACE,
host_session_id=HOST_SESSION_ID,
claim_id="claimed-interrupt",
)

with pytest.raises(CodexChatAgentError) as raised:
runtime.interrupt_turn(session_id=session_id, turn_id=str(turn["turn_id"]))

assert raised.value.error_code == "attached_session_interrupt_unavailable"
current_turn = store.load_turn(session_id, str(turn["turn_id"]))
current_session = store.load_session(session_id)
assert current_turn is not None and current_turn["status"] == "running"
assert current_session is not None and current_session["status"] == "busy"
assert current_session["active_turn_id"] == turn["turn_id"]
assert (
claim_attached_agent_turn(
store=store,
session_id=session_id,
host_surface=HOST_SURFACE,
host_session_id=HOST_SESSION_ID,
claim_id="next-claim",
)["claimed"]
is False
)


def test_attached_interrupt_endpoint_preserves_typed_failure(tmp_path: Path) -> None:
store = ChatSessionStore(tmp_path)
session_id = str(_bind(store)["session"]["session_id"])
runtime = ChatRuntimeController(store=store, codex_bin="missing-codex")
turn, _created = store.create_queued_turn(
session_id,
client_turn_id="endpoint-interrupt",
message="keep host ownership",
)
turn_id = str(turn["turn_id"])
claim_attached_agent_turn(
store=store,
session_id=session_id,
host_surface=HOST_SURFACE,
host_session_id=HOST_SESSION_ID,
claim_id="endpoint-interrupt",
)
responses: list[dict[str, object]] = []

class Handler:
server = SimpleNamespace(runtime_controller=runtime)

def _send_error(self, message: str, **kwargs: object) -> None:
responses.append({"error": message, **kwargs})

def _send_json(self, *_args: object, **_kwargs: object) -> None:
raise AssertionError("interrupt should fail closed")

ChatRequestHandler._interrupt_turn(Handler(), session_id, turn_id) # type: ignore[arg-type]

assert responses == [
{
"error": "The attached host does not expose interrupt control to LoopX Chat.",
"status": 424,
"error_code": "attached_session_interrupt_unavailable",
"gate": {
"kind": "host_tool_gate",
"summary": "The active Turn is owned by the attached host.",
"next_action": "Stop the Turn in the attached host, then retry.",
},
}
]
assert store.load_turn(session_id, turn_id)["status"] == "running" # type: ignore[index]
session = store.load_session(session_id)
assert session is not None and session["status"] == "busy"
assert session["active_turn_id"] == turn_id


def test_attached_completion_uses_canonical_response_and_terminal_events(
tmp_path: Path,
) -> None:
Expand Down
68 changes: 68 additions & 0 deletions tests/test_chat_session_active_turn.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
import pytest

import loopx.chat_store as chat_store
from loopx.chat_agent import CodexChatAgentError
from loopx.chat_runtime import ChatRuntimeController
from loopx.chat_store import (
SESSION_QUEUE_MAX_PENDING,
Expand Down Expand Up @@ -56,6 +57,19 @@ def close_session(self) -> None:
self.closed = True


class _FailingChatAdapter(_HealthyChatAdapter):
def start_turn(self, message: str, event_sink) -> dict[str, object]:
del message, event_sink
raise CodexChatAgentError(
"transport disconnected",
error_code="transport_disconnected",
gate={"kind": "host_tool_gate"},
)

def healthcheck(self) -> bool:
return False


def _slow_new_turn_writes(monkeypatch) -> set[str]:
original_atomic_write = chat_store._atomic_write_json
turn_write_threads: set[str] = set()
Expand Down Expand Up @@ -1222,6 +1236,60 @@ def test_managed_close_rejects_pending_queue_without_stranding_it(
assert store.load_turn(session_id, str(queued["turn_id"]))["status"] == "queued" # type: ignore[index]


def test_failed_old_turn_does_not_remove_replacement_adapter(
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
store = ChatSessionStore(tmp_path)
session = store.create_session(
goal_id="goal-one",
agent_id="codex",
executor_endpoint_id="codex",
adapter_kind="codex_app_server",
upstream_thread_id="thread-one",
upstream_mode="chat",
)
session_id = str(session["session_id"])
first_turn, _created = store.create_turn(
session_id,
client_turn_id="failed-old-turn",
message="fail on the old adapter",
)
runtime = ChatRuntimeController(store=store, codex_bin="missing-codex")
old_adapter = _FailingChatAdapter()
replacement_adapter = _HealthyChatAdapter()
runtime.adapters[session_id] = old_adapter
original_fail_turn = runtime._fail_turn
replacement_turn: dict[str, object] = {}

def fail_then_replace(*args: object, **kwargs: object) -> None:
original_fail_turn(*args, **kwargs) # type: ignore[arg-type]
with runtime.lock:
runtime.adapters[session_id] = replacement_adapter # type: ignore[assignment]
replacement_turn.update(
store.create_turn(
session_id,
client_turn_id="replacement-turn",
message="continue on the replacement adapter",
)[0]
)

monkeypatch.setattr(runtime, "_fail_turn", fail_then_replace)

runtime._run_turn(
session_id=session_id,
turn_id=str(first_turn["turn_id"]),
message="fail on the old adapter",
attachments=[],
adapter=old_adapter,
)

current = store.load_session(session_id)
assert runtime.adapters[session_id] is replacement_adapter
assert current is not None and current["status"] == "busy"
assert current["active_turn_id"] == replacement_turn["turn_id"]


@pytest.mark.parametrize("terminal_status", sorted(TERMINAL_TURN_STATES))
def test_managed_close_clears_terminal_active_turn_left_by_crash(
tmp_path: Path,
Expand Down