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
13 changes: 13 additions & 0 deletions docs/architecture/rfcs/harness-selection-dsh-pi-v0.md
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,19 @@ where the endpoint came from (`executor_endpoint_source`, now including
(`executor_endpoint_default_reason`), so an operator reads a decided default
instead of inferring it from the resolved host name.

A manager connection does not keep a second copy of that decision. The
connection record stores the resolved endpoint as an **observation** with its
source, and every read path -- the Lark route, the authorized-connection
resolution, and the Turn that answers on the channel -- re-resolves from the
machine. A record written while a different default was in force therefore
cannot keep answering on an endpoint the operator has since replaced, which is
what previously let a machine whose readback said `dsh` keep running its
steward on `codex`. When the machine does change the selection, the Session
bound to the channel still runs on the earlier endpoint; that Turn is refused
with the typed `manager_channel_executor_rebind_required` receipt, and the reply
names the one action that repairs it -- re-applying the connection, which opens
the channel Session on the endpoint the machine now selects.

Both managed surfaces resolve their **execution profile** from one owner
(`loopx/control_plane/turn_driver/execution_profile.py`): provider
`deepseek-official`, model `deepseek-v4-flash` (DeepSeek V4.1 Flash) and reasoning
Expand Down
8 changes: 8 additions & 0 deletions docs/architecture/rfcs/harness-selection-dsh-pi-v0.zh-CN.md
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,14 @@ LoopX **选择**托管有界 Turn 的默认宿主,而从不由启动时的意
(`executor_endpoint_default_reason`),因此运维方读到的是一个已决定的默认值,而不是
从解析出的宿主名去反推。

管家连接不会保存这份决定的第二份副本。连接记录把解析出的端点连同来源作为**观测值**
存下来;所有读取路径——Lark 路由、授权连接解析、以及在该通道上应答的 Turn——都重新
从本机解析。因此在另一个默认值仍生效时写下的记录,无法继续在被运维方替换过的端点上
应答——而这正是"回读说 `dsh`、管家却仍在 `codex` 上跑"的来路。当本机确实改了选择时,
通道上已绑定的 Session 仍跑在旧端点上;该 Turn 会以类型化回执
`manager_channel_executor_rebind_required` 被拒绝,回复直接给出唯一能修复它的动作——
重新应用一次该连接,把通道 Session 开在本机当前选择的端点上。

两个托管面从同一个所有者解析**执行档位**
(`loopx/control_plane/turn_driver/execution_profile.py`):provider
`deepseek-official`、模型 `deepseek-v4-flash`(DeepSeek V4.1 Flash)、推理档位
Expand Down
18 changes: 13 additions & 5 deletions loopx/chat_lark_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
)
from .chat_agent import CodexChatAgentError
from .chat_manager import (
controller_runtime_root,
manager_channel,
manager_executor_endpoint_default,
open_manager_session,
Expand Down Expand Up @@ -483,15 +484,19 @@ def _lark_connect(self) -> None:
or stored_routing.get("conversation_kind")
or "goal"
)
executor_endpoint_id = _compact_text(
body.get("executor_endpoint_id"), limit=100
) or stored_routing.get("executor_endpoint_id")
if conversation_kind == "manager" and not executor_endpoint_id:
executor_endpoint_id = manager_executor_endpoint_default(
# The machine owns its manager channel's executor, so the machine
# setting -- not a stored connection field or a request field --
# decides which endpoint this connection runs on and which Session
# it binds. The connection write below records the resolution.
executor_endpoint_id = (
manager_executor_endpoint_default(
machine_defaults=steward_machine_defaults(
self.server.runtime_controller
)
)
if conversation_kind == "manager"
else None
)
session_id: str | None = None
session_ids_by_agent: dict[str, str] = {}
if conversation_kind == "manager":
Expand Down Expand Up @@ -579,6 +584,9 @@ def _lark_connect(self) -> None:
"ingress_mode": ingress_mode or "async_inbox",
"reply_mode": reply_mode,
"registry_path": binding_path.parent / "registry.json",
"runtime_root": controller_runtime_root(
getattr(self.server, "runtime_controller", None)
),
"execute": body.get("execute") is True,
"runner": self._lark_runner(),
"cli_bin": cli_bin,
Expand Down
27 changes: 27 additions & 0 deletions loopx/chat_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
MANAGED_TURN_HOST,
managed_executor_binding,
)
from .capabilities.steward_executor import load_effective_steward_executor_defaults
from .chat_store import (
CHAT_SESSION_MODE_ATTACHED,
CHAT_SESSION_MODE_MANAGED,
Expand Down Expand Up @@ -344,6 +345,32 @@ def manager_executor_endpoint_default(
)[0]


def manager_connection_executor_endpoint(
runtime_root: Path | str | None,
*,
environ: dict[str, str] | None = None,
) -> tuple[str, str]:
"""Return the endpoint a manager connection runs on, and the source of it.

A manager conversation is one machine-level channel, so the machine owns
which executor answers there. A connection record therefore stores this
resolution as an observation instead of a decision that would outlive the
machine setting that made it: reading the connection's endpoint back as
authority is what let a machine that had selected a managed executor keep
answering on the interactive CLI endpoint that was the default when the
connection was created.
"""

machine_defaults = (
load_effective_steward_executor_defaults(Path(runtime_root))
if runtime_root is not None
else None
)
return selected_manager_executor_endpoint(
environ, machine_defaults=machine_defaults
)


# The channel's readback quotes the mode and the status of the Session it is an
# entry point to. The execution-mode RFC makes the binding, not the endpoint,
# the transport or the audience, the unit of mode ownership, so this projection
Expand Down
22 changes: 17 additions & 5 deletions loopx/control_plane/quota/should_run_prepare.py
Original file line number Diff line number Diff line change
Expand Up @@ -781,14 +781,26 @@ def _prepare_quota_should_run_item(
recovery_allowed = False
reason = str(projection_gap_repair.get("reason") or reason)
boundary_projection_repair = None
# An explicit `--todo-id` names one row the caller already chose, so it is
# resolved against the non-terminal rows this Goal records as the shared
# planning inventory. The bounded suggestion lanes above are a presentation
# budget: seeding the by-id lookup from them alone made an owned, open,
# typed advancement Todo unselectable whenever the display lanes were full,
# 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=quota_runnable_action_candidates(
agent_id=agent_frontier_id or "",
agent_todo_summary=agent_todo_summary,
capability_gate=capability_gate,
),
agent_todo_items=explicit_action_selection_items,
available_capabilities=effective_available_capabilities,
todo_id=requested_action_todo_id,
selection_binding="pending_action_selection",
Expand Down
2 changes: 2 additions & 0 deletions loopx/extensions/lark/goal_topic_batch.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ def connect_lark_goal_topics(
ingress_mode: str = "async_inbox",
reply_mode: str = "topic_reply",
registry_path: Path | None = None,
runtime_root: str | Path | None = None,
execute: bool = True,
runner: CommandRunner = default_subprocess_runner,
cli_bin: str = DEFAULT_CLI_BIN,
Expand Down Expand Up @@ -66,6 +67,7 @@ def connect_lark_goal_topics(
"ingress_mode": ingress_mode,
"reply_mode": reply_mode,
"registry_path": registry_path,
"runtime_root": runtime_root,
"runner": runner,
"cli_bin": cli_bin,
}
Expand Down
18 changes: 16 additions & 2 deletions loopx/extensions/lark/goal_topic_connections.py
Original file line number Diff line number Diff line change
Expand Up @@ -298,6 +298,7 @@ def connect_lark_goal_topic(
ingress_mode: str | None = None,
conversation_kind: str | None = None,
executor_endpoint_id: str | None = None,
runtime_root: str | Path | None = None,
reply_mode: str = "topic_reply",
registry_path: Path | None = None,
execute: bool = True,
Expand All @@ -321,11 +322,17 @@ def connect_lark_goal_topic(
editing = edit.binding
app_ref, chat_id, chat_name = edit.app_ref, edit.chat_id, edit.chat_name
agent_id, capture_scope = edit.agent_id, edit.capture_scope
conversation_kind, executor_endpoint_id, ingress_mode = resolve_conversation_policy(
(
conversation_kind,
executor_endpoint_id,
executor_endpoint_source,
ingress_mode,
) = resolve_conversation_policy(
editing=editing,
conversation_kind=conversation_kind,
executor_endpoint_id=executor_endpoint_id,
ingress_mode=ingress_mode,
runtime_root=runtime_root,
)
if conversation_kind == "manager":
agent_id = MANAGER_AGENT_GOAL_ID
Expand Down Expand Up @@ -763,6 +770,7 @@ def connect_lark_goal_topic(
{
"conversation_kind": "manager",
"executor_endpoint_id": executor_endpoint_id,
"executor_endpoint_source": executor_endpoint_source,
}
if conversation_kind == "manager"
else {}
Expand Down Expand Up @@ -1066,11 +1074,15 @@ def decide_lark_topic_event(
target_payload: Mapping[str, Any],
binding_payloads: Mapping[str, Mapping[str, Any]],
event: Mapping[str, Any],
runtime_root: str | Path | None = None,
) -> dict[str, Any]:
"""Return a content-free routing decision for one provider event."""

manager_decision = decide_manager_event(
target_payload=target_payload, binding_payloads=binding_payloads, event=event
target_payload=target_payload,
binding_payloads=binding_payloads,
event=event,
runtime_root=runtime_root,
)
if manager_decision is not None:
return manager_decision
Expand Down Expand Up @@ -1277,11 +1289,13 @@ def route_lark_topic_event(
target_payload: Mapping[str, Any],
binding_payloads: Mapping[str, Mapping[str, Any]],
event: Mapping[str, Any],
runtime_root: str | Path | None = None,
) -> dict[str, str] | None:
decision = decide_lark_topic_event(
target_payload=target_payload,
binding_payloads=binding_payloads,
event=event,
runtime_root=runtime_root,
)
route = decision.get("route")
return dict(route) if isinstance(route, Mapping) else None
Expand Down
27 changes: 23 additions & 4 deletions loopx/extensions/lark/goal_topic_edit.py
Original file line number Diff line number Diff line change
Expand Up @@ -121,23 +121,42 @@ def resolve_conversation_policy(
conversation_kind: str | None,
executor_endpoint_id: str | None,
ingress_mode: str | None,
) -> tuple[str, str | None, str | None]:
runtime_root: str | Path | None = None,
) -> tuple[str, str | None, str | None, str | None]:
"""Resolve one connection's purpose, executor observation and ingress mode.

The machine owns its manager channel's executor. A manager connection
therefore records the machine's current resolution and may restate it, but
it may not override it: a caller asking for a different endpoint is refused
with the owner named, instead of being written into a connection record
that the route would then ignore.
"""

from ...chat_manager import manager_connection_executor_endpoint
from .goal_channel_transport import SAFE_PROFILE_PATTERN

prior = (editing or {}).get("routing") or {}
kind = conversation_kind or prior.get("conversation_kind") or "goal"
if kind not in {"goal", "manager"}:
raise ValueError("conversation_kind must be goal or manager")
if kind == "manager":
endpoint = executor_endpoint_id or prior.get("executor_endpoint_id") or "codex"
endpoint, endpoint_source = manager_connection_executor_endpoint(
runtime_root
)
requested = str(executor_endpoint_id or "").strip()
if requested and requested != endpoint:
raise ValueError(
"the machine steward executor setting owns this connection's "
f"endpoint; change it there instead of requesting {requested}"
)
if not SAFE_PROFILE_PATTERN.fullmatch(endpoint):
raise ValueError("executor_endpoint_id must be a safe endpoint reference")
if ingress_mode and ingress_mode != "session_queue":
raise ValueError(
"the machine manager uses synchronous session_queue delivery"
)
return kind, endpoint, "session_queue"
return kind, None, ingress_mode
return kind, endpoint, endpoint_source, "session_queue"
return kind, None, None, ingress_mode


def _unregister_async_inbox(
Expand Down
27 changes: 22 additions & 5 deletions loopx/extensions/lark/goal_topic_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -796,12 +796,12 @@ def answer_lark_goal_topic(
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 = (
str(
route.get("executor_endpoint_id")
or manager_executor_endpoint_default(
machine_defaults=steward_machine_defaults(runtime_controller)
)
manager_executor_endpoint_default(
machine_defaults=steward_machine_defaults(runtime_controller)
)
if manager
else str(route.get("agent_id") or "codex")
Expand Down Expand Up @@ -833,6 +833,22 @@ 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"
):
# 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),
)
raise RuntimeError(
"bound Agent session is unavailable or no longer matches"
)
Expand Down Expand Up @@ -989,6 +1005,7 @@ def process_lark_goal_topic_event(
target_payload=target_payload,
binding_payloads=binding_payloads,
event=event,
runtime_root=runtime_root,
)
route = decision.get("route")
if route is None:
Expand Down
3 changes: 3 additions & 0 deletions loopx/extensions/lark/manager_context.py
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,9 @@ def manager_failure_reply(error: Exception) -> tuple[str, str]:
"hard_timeout": "处理超过时间限制",
"interrupted": "处理已中断",
"manager_authorization_unavailable": "当前连接的授权范围不可用",
"manager_channel_executor_rebind_required": (
"管家的执行器已由本机设置更改,需要重新应用一次管家连接"
),
}
code = str(getattr(error, "error_code", ""))
code = code if code in labels else "processing_failed"
Expand Down
20 changes: 16 additions & 4 deletions loopx/extensions/lark/manager_routing.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
from typing import Any
from pathlib import Path

from ...chat_manager import manager_channel
from ...chat_manager import manager_channel, manager_connection_executor_endpoint
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 Down Expand Up @@ -76,8 +76,15 @@ def decide_manager_event(
target_payload: Mapping[str, Any],
binding_payloads: Mapping[str, Any],
event: Mapping[str, Any],
runtime_root: str | Path | None = None,
) -> dict[str, Any] | None:
"""A manager receives addressed messages; exact worker Topics keep priority."""
"""A manager receives addressed messages; exact worker Topics keep priority.

The route carries the executor this machine selected for its manager
channel rather than the one the connection recorded when it was created, so
a machine that changes its steward executor does not keep answering on the
endpoint that happened to be the default on the day of the connection.
"""
chat_id, message_id = (
str(event.get("chat_id") or ""),
str(event.get("message_id") or ""),
Expand Down Expand Up @@ -131,6 +138,9 @@ def ignored(reason: str) -> dict[str, Any]:
return ignored("invalid_routing_state")
profile = str(identity.get("sender_profile") or "default")
turn_authorized = is_event_addressed_to_bot(event, identity)
executor_endpoint_id, executor_endpoint_source = (
manager_connection_executor_endpoint(runtime_root)
)
return {
"matched": True,
"reason": "matched" if turn_authorized else "context_only",
Expand All @@ -140,7 +150,8 @@ def ignored(reason: str) -> dict[str, Any]:
"agent_id": binding["agent_id"],
"session_id": binding["session_id"],
"conversation_kind": "manager",
"executor_endpoint_id": routing.get("executor_endpoint_id") or "codex",
"executor_endpoint_id": executor_endpoint_id,
"executor_endpoint_source": executor_endpoint_source,
"manager_channel_id": manager_channel(
provider="lark", audience=f"{profile}\0{chat_id}"
),
Expand Down Expand Up @@ -219,7 +230,8 @@ def authorized_manager_goal_ids(
goal_id, binding, routing = candidates[0]
if (
binding.get("session_id") != session.get("session_id")
or (routing.get("executor_endpoint_id") or "codex") != session.get("agent_id")
or manager_connection_executor_endpoint(runtime_root)[0]
!= session.get("agent_id")
or not _valid_manager_binding(goal_id, binding, routing)
):
return []
Expand Down
Loading
Loading