From 80d23f06e67dfe82b59d2cc9efa79e48654e7f1f Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Wed, 16 Sep 2026 01:14:44 +0800 Subject: [PATCH 1/4] perf(chat): bound completed event retention Signed-off-by: duanjialing.777 --- loopx/chat_store.py | 70 +++++++------ tests/test_chat_event_retention.py | 155 +++++++++++++++++++++++++++++ 2 files changed, 192 insertions(+), 33 deletions(-) create mode 100644 tests/test_chat_event_retention.py diff --git a/loopx/chat_store.py b/loopx/chat_store.py index fb17b6df0a..7ead472c8b 100644 --- a/loopx/chat_store.py +++ b/loopx/chat_store.py @@ -11,6 +11,7 @@ import threading from typing import Any import uuid +from weakref import WeakValueDictionary from .chat import require_matching_replay, resolve_attached_completion_replay from .file_lock import exclusive_file_lock @@ -26,6 +27,8 @@ CHAT_SESSION_MODE_ATTACHED = "attached_host" RESUMABLE_SESSION_STATES = {"ready", "busy", "stale", "resuming"} TERMINAL_TURN_STATES = {"completed", "interrupted", "timed_out", "failed"} +TERMINAL_EVENT_KINDS = {"turn.completed", "turn.failed", "turn.interrupted"} +REPLAY_ONLY_EVENT_KINDS = {"answer.delta", "assistant.delta", "agent.phase", "turn.activity"} SESSION_QUEUE_MAX_PENDING = 20 SESSION_QUEUE_TTL_SECONDS = 3600 @@ -138,12 +141,12 @@ def __init__(self, runtime_root: Path) -> None: self.root = runtime_root.expanduser().resolve() / "chat" self.sessions_root = self.root / "sessions" self._session_lock_guard = threading.Lock() - self._session_locks: dict[str, threading.Lock] = {} + self._session_locks: WeakValueDictionary[str, threading.Lock] = WeakValueDictionary() self._event_lock = threading.RLock() self._event_cache: dict[tuple[str, str], list[dict[str, Any]]] = {} self._event_cache_revision: dict[tuple[str, str], tuple[int, int, int] | None] = {} self._event_pending: dict[tuple[str, str], list[dict[str, Any]]] = {} - self._event_flush_locks: dict[tuple[str, str], threading.Lock] = {} + self._event_flush_locks: WeakValueDictionary[tuple[str, str], threading.Lock] = WeakValueDictionary() self.sessions_root.mkdir(parents=True, exist_ok=True, mode=0o700) os.chmod(self.root, 0o700) os.chmod(self.sessions_root, 0o700) @@ -1312,6 +1315,12 @@ def _event_flush_lock(self, key: tuple[str, str]) -> threading.Lock: with self._event_lock: return self._event_flush_locks.setdefault(key, threading.Lock()) + def _drop_event_cache(self, key: tuple[str, str], expected: list[dict[str, Any]] | None = None) -> None: + with self._event_lock: + if expected is None or self._event_cache.get(key) is expected: + self._event_cache.pop(key, None) + self._event_cache_revision.pop(key, None) + def append_event( self, session_id: str, @@ -1358,9 +1367,12 @@ def flush_events(self, session_id: str, turn_id: str) -> int: event["event_id"] = str(sequence) event["sequence"] = sequence _append_jsonl_rows(path, pending) - with self._event_lock: - self._event_cache[key] = [*rows, *pending] - self._event_cache_revision[key] = self._event_revision(path) + if any(row["kind"] in TERMINAL_EVENT_KINDS for row in pending): + self._drop_event_cache(key) + else: + with self._event_lock: + self._event_cache[key] = [*rows, *pending] + self._event_cache_revision[key] = self._event_revision(path) except Exception: with self._event_lock: later = self._event_pending.get(key, []) @@ -1379,22 +1391,17 @@ def events_after(self, session_id: str, turn_id: str, event_id: str | None) -> l revision = self._event_revision(path) with self._event_lock: cached = self._event_cache.get(key) - rows = ( - cached - if cached is not None and self._event_cache_revision.get(key) == revision - else None - ) + rows = cached if cached is not None and self._event_cache_revision.get(key) == revision else None if rows is None: with exclusive_file_lock(path, agent_id="loopx-chat", operation="read_chat_events"): rows = self._event_rows_locked(session_id, turn_id) # Relies on the sequence ordering maintained by flush_events and compaction; # gaps are valid, but inserting or rewriting rows must preserve that order. - start = bisect_right( - rows, - after, - key=lambda row: int(row.get("sequence") or 0), - ) - return rows[start:] + start = bisect_right(rows, after, key=lambda row: int(row.get("sequence") or 0)) + result = rows[start:] + if rows and rows[-1].get("kind") in TERMINAL_EVENT_KINDS: + self._drop_event_cache(key, rows) + return result def compact_completed_events(self, *, older_than_hours: float = 24.0) -> int: """Drop replay-only deltas after the durable final message is old enough.""" @@ -1415,25 +1422,22 @@ def compact_completed_events(self, *, older_than_hours: float = 24.0) -> int: turn_id = turn_path.stem self.flush_events(session_id, turn_id) event_path = turn_path.with_name(f"{turn_path.stem}.events.jsonl") + revision = self._event_revision(event_path) + if turn.get("event_compaction_revision") == list(revision or ()): + continue with exclusive_file_lock(event_path, agent_id="loopx-chat", operation="compact_chat_events"): rows = self._event_rows_locked(session_id, turn_id) - retained = [ - row for row in rows - if row.get("kind") not in { - "answer.delta", - "assistant.delta", - "agent.phase", - "turn.activity", - } - ] - if len(retained) == len(rows): - continue - _replace_jsonl(event_path, retained) - with self._event_lock: - key = (session_id, turn_id) - self._event_cache[key] = retained - self._event_cache_revision[key] = self._event_revision(event_path) - compacted += 1 + retained = [row for row in rows if row.get("kind") not in REPLAY_ONLY_EVENT_KINDS] + if len(retained) != len(rows): + _replace_jsonl(event_path, retained) + compacted += 1 + revision = self._event_revision(event_path) + self._drop_event_cache((session_id, turn_id)) + with exclusive_file_lock(turn_path, agent_id="loopx-chat", operation="mark_chat_events_compacted"): + current = _read_json(turn_path) + if current.get("status") in TERMINAL_TURN_STATES: + current["event_compaction_revision"] = list(revision or ()) + _atomic_write_json(turn_path, current, preserve_mode=True) return compacted def public_session(self, payload: dict[str, Any]) -> dict[str, Any]: diff --git a/tests/test_chat_event_retention.py b/tests/test_chat_event_retention.py new file mode 100644 index 0000000000..811415ac0a --- /dev/null +++ b/tests/test_chat_event_retention.py @@ -0,0 +1,155 @@ +from __future__ import annotations + +from datetime import datetime, timedelta, timezone +import gc +import json +from pathlib import Path + +import loopx.chat_store as chat_store +from loopx.chat_store import CHAT_TURN_SCHEMA_VERSION, ChatSessionStore + + +def _write_completed_turn(root: Path, *, with_events: bool = True) -> tuple[Path, Path]: + turn_path = root / "chat" / "sessions" / "session" / "turns" / "turn.json" + turn_path.parent.mkdir(parents=True) + turn_path.write_text( + json.dumps( + { + "schema_version": CHAT_TURN_SCHEMA_VERSION, + "session_id": "session", + "turn_id": "turn", + "status": "completed", + "completed_at": ( + datetime.now(timezone.utc) - timedelta(days=2) + ).isoformat(), + } + ), + encoding="utf-8", + ) + event_path = turn_path.with_name("turn.events.jsonl") + if with_events: + rows = [ + {"kind": "answer.delta", "sequence": 1, "event_id": "1"}, + {"kind": "turn.completed", "sequence": 2, "event_id": "2"}, + ] + event_path.write_text( + "\n".join(json.dumps(row) for row in rows) + "\n", + encoding="utf-8", + ) + return turn_path, event_path + + +def test_terminal_event_history_does_not_remain_in_the_hot_cache( + tmp_path: Path, +) -> None: + store = ChatSessionStore(tmp_path) + key = ("session", "turn") + store.append_event( + *key, + kind="assistant.delta", + payload={"text": "visible delta"}, + buffered=True, + ) + store.append_event(*key, kind="turn.completed", payload={}) + + assert key not in store._event_cache + assert key not in store._event_cache_revision + assert [row["kind"] for row in store.events_after(*key, None)] == [ + "assistant.delta", + "turn.completed", + ] + assert key not in store._event_cache + assert key not in store._event_cache_revision + + +def test_keyed_locks_are_released_after_callers_drop_them(tmp_path: Path) -> None: + store = ChatSessionStore(tmp_path) + session_lock = store._session_lock("session") + event_lock = store._event_flush_lock(("session", "turn")) + + assert store._session_lock("session") is session_lock + assert store._event_flush_lock(("session", "turn")) is event_lock + + del session_lock, event_lock + gc.collect() + + assert "session" not in store._session_locks + assert ("session", "turn") not in store._event_flush_locks + + +def test_completed_event_compaction_is_skipped_until_the_file_changes( + tmp_path: Path, + monkeypatch, +) -> None: + turn_path, event_path = _write_completed_turn(tmp_path) + first = ChatSessionStore(tmp_path) + turn = json.loads(turn_path.read_text(encoding="utf-8")) + + assert turn["event_compaction_revision"] == list( + first._event_revision(event_path) or () + ) + assert ("session", "turn") not in first._event_cache + + original_read_jsonl = chat_store._read_jsonl + + def reject_redundant_event_read(path: Path): + if path == event_path: + raise AssertionError("an unchanged compacted event stream was read again") + return original_read_jsonl(path) + + monkeypatch.setattr(chat_store, "_read_jsonl", reject_redundant_event_read) + restarted = ChatSessionStore(tmp_path) + + assert not restarted._event_cache + assert not restarted._event_flush_locks + + +def test_completed_event_compaction_marker_is_invalidated_by_a_new_event( + tmp_path: Path, +) -> None: + turn_path, event_path = _write_completed_turn(tmp_path) + ChatSessionStore(tmp_path) + with event_path.open("a", encoding="utf-8") as handle: + handle.write( + json.dumps({"kind": "answer.delta", "sequence": 3, "event_id": "3"}) + ) + handle.write("\n") + + restarted = ChatSessionStore(tmp_path) + rows = chat_store._read_jsonl(event_path) + turn = json.loads(turn_path.read_text(encoding="utf-8")) + + assert [row["kind"] for row in rows] == ["turn.completed"] + assert turn["event_compaction_revision"] == list( + restarted._event_revision(event_path) or () + ) + assert ("session", "turn") not in restarted._event_cache + + +def test_missing_event_stream_is_marked_without_retaining_empty_state( + tmp_path: Path, +) -> None: + turn_path, _event_path = _write_completed_turn(tmp_path, with_events=False) + + store = ChatSessionStore(tmp_path) + turn = json.loads(turn_path.read_text(encoding="utf-8")) + + assert turn["event_compaction_revision"] == [] + assert not store._event_cache + assert not store._event_flush_locks + + +def test_compaction_marker_does_not_make_legacy_terminal_turns_unreadable( + tmp_path: Path, +) -> None: + turn_path, event_path = _write_completed_turn(tmp_path) + turn = json.loads(turn_path.read_text(encoding="utf-8")) + turn.pop("schema_version") + turn_path.write_text(json.dumps(turn), encoding="utf-8") + + store = ChatSessionStore(tmp_path) + persisted = json.loads(turn_path.read_text(encoding="utf-8")) + + assert persisted["event_compaction_revision"] == list( + store._event_revision(event_path) or () + ) From df698a03cba37e0935aa70c3055b7d72794062dd Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Wed, 16 Sep 2026 15:43:12 +0800 Subject: [PATCH 2/4] fix(steward): preserve the static type-check boundary Signed-off-by: duanjialing.777 --- .../steward_executor/machine_defaults.py | 15 +++++++-------- 1 file changed, 7 insertions(+), 8 deletions(-) diff --git a/loopx/capabilities/steward_executor/machine_defaults.py b/loopx/capabilities/steward_executor/machine_defaults.py index 3efb352a4a..35482c8a67 100644 --- a/loopx/capabilities/steward_executor/machine_defaults.py +++ b/loopx/capabilities/steward_executor/machine_defaults.py @@ -23,6 +23,7 @@ from __future__ import annotations from collections.abc import Mapping +import importlib from pathlib import Path from typing import Any @@ -50,17 +51,15 @@ def steward_executor_endpoints() -> frozenset[str]: namespace, and a module-level import would close that loop. """ - from ...chat_manager import MANAGER_ENDPOINT_KINDS - - return frozenset(MANAGER_ENDPOINT_KINDS) + manager = importlib.import_module("loopx.chat_manager") + return frozenset(manager.MANAGER_ENDPOINT_KINDS) def steward_reasoning_efforts() -> tuple[str, ...]: """Return the reasoning-effort vocabulary a selected executor accepts.""" - from ...chat_manager import MANAGER_REASONING_EFFORTS - - return tuple(MANAGER_REASONING_EFFORTS) + manager = importlib.import_module("loopx.chat_manager") + return tuple(manager.MANAGER_REASONING_EFFORTS) def _optional_text(value: Any, *, field: str) -> str | None: @@ -238,10 +237,10 @@ def load_effective_steward_executor_defaults(runtime_root: Path) -> dict[str, An a typed reason instead of failing the channel a person is talking to. """ - from ..machine_configuration.store import read_stored_machine_configuration + store = importlib.import_module("loopx.capabilities.machine_configuration.store") try: - configuration = read_stored_machine_configuration(runtime_root) + configuration = store.read_stored_machine_configuration(runtime_root) except (OSError, TypeError, ValueError): return _projection( status="unavailable", From 0db32a76821177b91063f716ff213962f605c0ca Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Wed, 16 Sep 2026 16:47:04 +0800 Subject: [PATCH 3/4] fix(lark): keep manager routing within module budget Signed-off-by: duanjialing.777 --- loopx/extensions/lark/goal_topic_runtime.py | 35 ++++++--------------- loopx/extensions/lark/manager_routing.py | 31 +++++++++++++++++- tests/test_manager_channel_binding.py | 2 +- 3 files changed, 40 insertions(+), 28 deletions(-) diff --git a/loopx/extensions/lark/goal_topic_runtime.py b/loopx/extensions/lark/goal_topic_runtime.py index 45ba4065af..e165f783bf 100644 --- a/loopx/extensions/lark/goal_topic_runtime.py +++ b/loopx/extensions/lark/goal_topic_runtime.py @@ -14,14 +14,12 @@ from pathlib import Path from typing import Any -from ...chat_manager import ( - MANAGER_AGENT_OBJECTIVE, - manager_executor_endpoint_default, - steward_machine_defaults, -) +from ...chat_manager import MANAGER_AGENT_OBJECTIVE from .manager_routing import ( has_manager_binding, invalid_manager_authority_result, + manager_session_requires_executor_rebind, + manager_turn_executor, parse_manager_authority_mode, unavailable_manager_context_result, ManagerAuthorityMode, @@ -791,20 +789,12 @@ def answer_lark_goal_topic( runtime_controller: Any, ) -> str: """Deliver one Topic message using its exact Agent ingress contract.""" - goal_id = str(route.get("goal_id") or "") ingress_mode = str(route.get("ingress_mode") or "direct_session") session_id = str(route.get("session_id") or "") manager = route.get("conversation_kind") == "manager" - # The machine owns its manager channel's executor, so a manager Turn runs on - # the machine's current selection even when the connection record still - # carries the endpoint that was the default on the day it was created. - agent_id = ( - manager_executor_endpoint_default( - machine_defaults=steward_machine_defaults(runtime_controller) - ) - if manager - else str(route.get("agent_id") or "codex") + agent_id = manager_turn_executor(runtime_controller) if manager else str( + route.get("agent_id") or "codex" ) expected_channel = ( str(route.get("manager_channel_id") or "") if manager else f"goal.{goal_id}" @@ -833,18 +823,11 @@ def answer_lark_goal_topic( or session.get("channel_id") != expected_channel or session.get("status") == "closed" ): - if ( - manager - and session is not None - and session.get("channel_id") == expected_channel - and session.get("agent_id") != agent_id - and session.get("status") != "closed" + if manager and manager_session_requires_executor_rebind( + session, + expected_channel=expected_channel, + agent_id=agent_id, ): - # A machine that changes its steward executor leaves the older - # bound Session behind on the same audience. Name that state - # instead of failing the Turn under the opaque "manager failed" - # label, because one connection re-apply repairs it. Every other - # mismatch keeps the outcome it had. raise LarkGoalTopicTurnFailed( "manager_channel_executor_rebind_required", _session_turn_effect(route), diff --git a/loopx/extensions/lark/manager_routing.py b/loopx/extensions/lark/manager_routing.py index 80ef512dfe..5e3743ce81 100644 --- a/loopx/extensions/lark/manager_routing.py +++ b/loopx/extensions/lark/manager_routing.py @@ -7,7 +7,12 @@ from typing import Any from pathlib import Path -from ...chat_manager import manager_channel, manager_connection_executor_endpoint +from ...chat_manager import ( + manager_channel, + manager_connection_executor_endpoint, + manager_executor_endpoint_default, + steward_machine_defaults, +) from ..external_connector_runtime import project_external_connector_status from .goal_channel_contracts import bindings_for_goal from .goal_channel_targets import goal_channel_target_for_name @@ -22,6 +27,30 @@ class ManagerAuthorityMode(str, Enum): TURN_AUTHORIZED = "turn_authorized" +def manager_turn_executor(runtime_controller: Any) -> str: + """Resolve the manager executor from this machine's current selection.""" + + return manager_executor_endpoint_default( + machine_defaults=steward_machine_defaults(runtime_controller) + ) + + +def manager_session_requires_executor_rebind( + session: Mapping[str, Any] | None, + *, + expected_channel: str, + agent_id: str, +) -> bool: + """Identify a live manager Session left on the machine's old executor.""" + + return bool( + session is not None + and session.get("channel_id") == expected_channel + and session.get("agent_id") != agent_id + and session.get("status") != "closed" + ) + + def parse_manager_authority_mode(value: object) -> ManagerAuthorityMode | None: """Parse a persisted route mode without coercing unknown values.""" diff --git a/tests/test_manager_channel_binding.py b/tests/test_manager_channel_binding.py index 1ec61900b6..88149c7a62 100644 --- a/tests/test_manager_channel_binding.py +++ b/tests/test_manager_channel_binding.py @@ -838,7 +838,7 @@ def test_every_production_steward_caller_passes_the_machine_defaults() -> None: "loopx/chat_manager.py", "loopx/chat_manager_context.py", "loopx/chat_runtime.py", - "loopx/extensions/lark/goal_topic_runtime.py", + "loopx/extensions/lark/manager_routing.py", } undocumented = [ f"{path}:{line} {name}" From 2858bbffd451aa2909484bd31929a640ebaa67e2 Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Wed, 16 Sep 2026 16:52:50 +0800 Subject: [PATCH 4/4] fix(quota): stay within decision complexity budget Signed-off-by: duanjialing.777 --- loopx/control_plane/quota/should_run_prepare.py | 17 ++++++++--------- 1 file changed, 8 insertions(+), 9 deletions(-) diff --git a/loopx/control_plane/quota/should_run_prepare.py b/loopx/control_plane/quota/should_run_prepare.py index 85f21e7bc1..ceb8430ba2 100644 --- a/loopx/control_plane/quota/should_run_prepare.py +++ b/loopx/control_plane/quota/should_run_prepare.py @@ -789,18 +789,17 @@ def _prepare_quota_should_run_item( # which left the caller no legal way to bind its own quota guard to it. # The builder still applies every eligibility predicate, so this widens # reachability of the lookup, not what may be selected. - explicit_action_selection_items = [ - *quota_runnable_action_candidates( - agent_id=agent_frontier_id or "", - agent_todo_summary=agent_todo_summary, - capability_gate=capability_gate, - ), - *agent_todo_planning_source_items, - ] requested_action_candidate = ( build_explicit_advancement_next_action( agent_identity=agent_identity, - agent_todo_items=explicit_action_selection_items, + agent_todo_items=[ + *quota_runnable_action_candidates( + agent_id=agent_frontier_id or "", + agent_todo_summary=agent_todo_summary, + capability_gate=capability_gate, + ), + *agent_todo_planning_source_items, + ], available_capabilities=effective_available_capabilities, todo_id=requested_action_todo_id, selection_binding="pending_action_selection",