Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
15 changes: 7 additions & 8 deletions loopx/capabilities/steward_executor/machine_defaults.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
from __future__ import annotations

from collections.abc import Mapping
import importlib
from pathlib import Path
from typing import Any

Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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",
Expand Down
70 changes: 37 additions & 33 deletions loopx/chat_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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

Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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, [])
Expand All @@ -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."""
Expand All @@ -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]:
Expand Down
17 changes: 8 additions & 9 deletions loopx/control_plane/quota/should_run_prepare.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
35 changes: 9 additions & 26 deletions loopx/extensions/lark/goal_topic_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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}"
Expand Down Expand Up @@ -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),
Expand Down
31 changes: 30 additions & 1 deletion loopx/extensions/lark/manager_routing.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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."""

Expand Down
Loading