diff --git a/docs/architecture/rfcs/harness-selection-dsh-pi-v0.md b/docs/architecture/rfcs/harness-selection-dsh-pi-v0.md index e0ec1dba0a..8a7c173f52 100644 --- a/docs/architecture/rfcs/harness-selection-dsh-pi-v0.md +++ b/docs/architecture/rfcs/harness-selection-dsh-pi-v0.md @@ -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 diff --git a/docs/architecture/rfcs/harness-selection-dsh-pi-v0.zh-CN.md b/docs/architecture/rfcs/harness-selection-dsh-pi-v0.zh-CN.md index 1eb5193c8e..cbc192c647 100644 --- a/docs/architecture/rfcs/harness-selection-dsh-pi-v0.zh-CN.md +++ b/docs/architecture/rfcs/harness-selection-dsh-pi-v0.zh-CN.md @@ -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)、推理档位 diff --git a/loopx/chat_lark_api.py b/loopx/chat_lark_api.py index 868c171e10..9f330d4496 100644 --- a/loopx/chat_lark_api.py +++ b/loopx/chat_lark_api.py @@ -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, @@ -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": @@ -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, diff --git a/loopx/chat_manager.py b/loopx/chat_manager.py index 91c11f1eb0..176feb4e7f 100644 --- a/loopx/chat_manager.py +++ b/loopx/chat_manager.py @@ -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, @@ -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 diff --git a/loopx/control_plane/quota/should_run_prepare.py b/loopx/control_plane/quota/should_run_prepare.py index 159e938af6..85f21e7bc1 100644 --- a/loopx/control_plane/quota/should_run_prepare.py +++ b/loopx/control_plane/quota/should_run_prepare.py @@ -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", diff --git a/loopx/extensions/lark/goal_topic_batch.py b/loopx/extensions/lark/goal_topic_batch.py index 20421d4a02..27fcd7fd21 100644 --- a/loopx/extensions/lark/goal_topic_batch.py +++ b/loopx/extensions/lark/goal_topic_batch.py @@ -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, @@ -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, } diff --git a/loopx/extensions/lark/goal_topic_connections.py b/loopx/extensions/lark/goal_topic_connections.py index fb20fce567..9152b7967f 100644 --- a/loopx/extensions/lark/goal_topic_connections.py +++ b/loopx/extensions/lark/goal_topic_connections.py @@ -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, @@ -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 @@ -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 {} @@ -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 @@ -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 diff --git a/loopx/extensions/lark/goal_topic_edit.py b/loopx/extensions/lark/goal_topic_edit.py index 77f34173db..75a3c43ef7 100644 --- a/loopx/extensions/lark/goal_topic_edit.py +++ b/loopx/extensions/lark/goal_topic_edit.py @@ -121,7 +121,18 @@ 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 {} @@ -129,15 +140,23 @@ def resolve_conversation_policy( 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( diff --git a/loopx/extensions/lark/goal_topic_runtime.py b/loopx/extensions/lark/goal_topic_runtime.py index b05277a19a..45ba4065af 100644 --- a/loopx/extensions/lark/goal_topic_runtime.py +++ b/loopx/extensions/lark/goal_topic_runtime.py @@ -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") @@ -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" ) @@ -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: diff --git a/loopx/extensions/lark/manager_context.py b/loopx/extensions/lark/manager_context.py index b9304e0147..f564ac8264 100644 --- a/loopx/extensions/lark/manager_context.py +++ b/loopx/extensions/lark/manager_context.py @@ -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" diff --git a/loopx/extensions/lark/manager_routing.py b/loopx/extensions/lark/manager_routing.py index 2e8b0f31f6..80ef512dfe 100644 --- a/loopx/extensions/lark/manager_routing.py +++ b/loopx/extensions/lark/manager_routing.py @@ -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 @@ -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 ""), @@ -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", @@ -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}" ), @@ -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 [] diff --git a/tests/control_plane/test_quota_settlement_cli.py b/tests/control_plane/test_quota_settlement_cli.py index 6d9b4f6453..cfc236cf70 100644 --- a/tests/control_plane/test_quota_settlement_cli.py +++ b/tests/control_plane/test_quota_settlement_cli.py @@ -28,6 +28,7 @@ ALTERNATIVE_TODO_ID = "todo_fixture_alternative" SECOND_ALTERNATIVE_TODO_ID = "todo_fixture_second_alternative" OUTSIDE_BOUNDED_PORTFOLIO_TODO_ID = "todo_fixture_outside_portfolio" +DEEP_ALTERNATIVE_TODO_ID = "todo_fixture_deep_alternative" REENTRY_TODO_ID = "todo_fixture_network_reentry" DUE_MONITOR_TODO_ID = "todo_fixture_due_monitor" TURN_ID = "turn-settlement-cli-1" @@ -2689,6 +2690,86 @@ def test_agent_can_select_eligible_todo_outside_bounded_suggestions( assert _heartbeat_receipt_count(runtime, turn_instance_id) == 2 +def _configure_deep_alternative(project: Path, *, fillers: int = 9) -> None: + """Add one owned advancement Todo beyond every bounded suggestion lane.""" + + state_path = project / f".codex/goals/{GOAL_ID}/ACTIVE_GOAL_STATE.md" + filler_rows = "".join( + f"- [ ] [P1] Advance filler delivery {index}.\n" + f" \n" + for index in range(fillers) + ) + state_path.write_text( + state_path.read_text(encoding="utf-8").rstrip() + + "\n" + + filler_rows + + "- [ ] [P1] Advance the deep alternative delivery.\n" + + " \n", + encoding="utf-8", + ) + + +def test_agent_can_select_an_owned_todo_outside_every_bounded_lane( + tmp_path: Path, +) -> None: + """The display lanes are a presentation budget, not the eligible Todo set.""" + + project, runtime, registry_path = _write_fixture(tmp_path) + _configure_deep_alternative(project) + turn_instance_id = "turn-agent-selection-beyond-lanes" + guard_args = ( + "quota", + "should-run", + "--codex-app", + "--goal-id", + GOAL_ID, + "--agent-id", + AGENT_ID, + "--turn-instance-id", + turn_instance_id, + "--scan-path", + str(project), + ) + first_rc, first = _run_cli(registry_path, runtime, *guard_args) + selected_rc, selected = _run_cli( + registry_path, + runtime, + *guard_args, + "--todo-id", + DEEP_ALTERNATIVE_TODO_ID, + ) + + assert first_rc == 0, first + # The row is a real, open, typed advancement Todo that this Agent owns, and + # no bounded presentation lane lists it. + assert DEEP_ALTERNATIVE_TODO_ID not in { + item["todo_id"] for item in first["action_portfolio"]["suggested_actions"] + } + # Every bounded lane this payload publishes must not list it either; the + # count guard keeps the check from passing vacuously on an empty payload. + bounded_lanes = [ + (key, items) + for key, items in first["agent_todo_summary"].items() + if key.endswith("executable_items") and isinstance(items, list) + ] + assert bounded_lanes, sorted(first["agent_todo_summary"]) + for key, items in bounded_lanes: + assert DEEP_ALTERNATIVE_TODO_ID not in { + item["todo_id"] for item in items + }, key + assert selected_rc == 0, selected + assert selected["ok"] is True + assert selected["selected_todo"]["todo_id"] == DEEP_ALTERNATIVE_TODO_ID + assert selected["selected_todo"]["selection_binding"] == "heartbeat_receipt" + assert selected["heartbeat_receipt"]["status"] == "upgraded" + assert selected["heartbeat_receipt"]["settlement_identity"]["todo_id"] == ( + DEEP_ALTERNATIVE_TODO_ID + ) + + def test_same_turn_can_select_eligible_todo_created_after_unbound_receipt( tmp_path: Path, ) -> None: diff --git a/tests/extensions/test_lark_goal_topic_connections.py b/tests/extensions/test_lark_goal_topic_connections.py index bc4462321a..718981ab4a 100644 --- a/tests/extensions/test_lark_goal_topic_connections.py +++ b/tests/extensions/test_lark_goal_topic_connections.py @@ -2894,3 +2894,155 @@ def test_manager_explicit_audience_scope_still_requires_live_binding(tmp_path): assert authorized_manager_goal_ids(snapshot,session,runtime_root=tmp_path)==['goal-alpha','goal-beta'] session['session_id']='different-session' assert authorized_manager_goal_ids(snapshot,session,runtime_root=tmp_path)==[] + + +def _machine_selects_steward_executor(runtime_root: Path, endpoint: str = "dsh") -> None: + """Store one machine-level steward executor, as a machine surface would.""" + + from loopx.capabilities.machine_configuration.builtins import ( + build_builtin_machine_configuration_registry, + ) + from loopx.capabilities.machine_configuration.store import ( + configure_machine_configuration, + ) + + registry = build_builtin_machine_configuration_registry() + configuration = { + "schema_version": "loopx_machine_configuration_v0", + "namespaces": { + "steward_executor": { + "schema_version": "steward_executor_machine_defaults_v0", + "executor_endpoint": endpoint, + "executor_model": None, + "executor_reasoning_effort": None, + } + }, + } + preview = configure_machine_configuration( + runtime_root=runtime_root, + configuration=configuration, + registry=registry, + execute=False, + ) + configure_machine_configuration( + runtime_root=runtime_root, + configuration=configuration, + registry=registry, + execute=True, + expected_plan_revision=str(preview["plan_revision"]), + ) + + +def test_a_manager_connection_runs_on_the_executor_this_machine_selected( + tmp_path: Path, +) -> None: + """The machine, not the connection record, decides where the steward answers.""" + + kwargs, _state, bindings = _manager_fixture(tmp_path) + runtime_root = tmp_path / "runtime" + _machine_selects_steward_executor(runtime_root) + + def decision() -> dict[str, Any]: + return decide_lark_topic_event( + target_payload=read_goal_channel_targets(kwargs["target_path"]), + binding_payloads={"goal-alpha": bindings}, + event={ + "chat_id": CHAT_ID, + "message_id": "om_manager_machine_endpoint", + "mentions": [{"id": APP_ID}], + }, + runtime_root=runtime_root, + )["route"] + + # The connection was created while the shipped default was still the + # interactive CLI endpoint, and it keeps that value on disk. + connection = binding_for_goal(bindings, "goal-alpha") + assert connection is not None + assert connection["routing"]["executor_endpoint_id"] == "codex" + + route = decision() + + assert route["executor_endpoint_id"] == "dsh" + assert route["executor_endpoint_source"] == "machine_configuration" + # A stale record cannot outrank the machine selection on a later event either. + assert decision()["executor_endpoint_id"] == "dsh" + + +def test_a_manager_connection_write_records_the_machine_resolution( + tmp_path: Path, +) -> None: + kwargs, _state = _upgrade_fixture(tmp_path, agent_id="agent-alpha", peers=True) + runtime_root = tmp_path / "runtime" + _machine_selects_steward_executor(runtime_root) + + connected = connect_lark_goal_topic( + **kwargs, + conversation_kind="manager", + session_id="manager-session", + runtime_root=runtime_root, + ) + + assert connected["ok"] is True + stored = read_goal_channel_binding(kwargs["binding_path"]) + connection = binding_for_goal(stored, "goal-alpha") + assert connection is not None + routing = stored["bindings"]["goal-alpha"]["connections"][ + connection["connection_id"] + ]["routing"] + assert routing["executor_endpoint_id"] == "dsh" + assert routing["executor_endpoint_source"] == "machine_configuration" + + # A caller may restate the machine selection, but it may not overrule it + # from a connection record the route would then have to ignore. + restated = connect_lark_goal_topic( + **kwargs, + conversation_kind="manager", + session_id="manager-session", + runtime_root=runtime_root, + executor_endpoint_id="dsh", + ) + assert restated["ok"] is True + with pytest.raises(ValueError, match="machine steward executor setting owns"): + connect_lark_goal_topic( + **kwargs, + conversation_kind="manager", + session_id="manager-session", + runtime_root=runtime_root, + executor_endpoint_id="codex", + ) + + +def test_an_authorized_manager_session_must_run_on_the_machine_executor( + tmp_path: Path, +) -> None: + from loopx.extensions.lark.manager_routing import authorized_manager_goal_ids + + kwargs, _state, bindings = _manager_fixture(tmp_path) + runtime_root = tmp_path / "runtime" + _machine_selects_steward_executor(runtime_root) + targets = read_goal_channel_targets(kwargs["target_path"]) + route = decide_lark_topic_event( + target_payload=targets, + binding_payloads={"goal-alpha": bindings}, + event={ + "chat_id": CHAT_ID, + "message_id": "om_manager_stale_session", + "mentions": [{"id": APP_ID}], + }, + runtime_root=runtime_root, + )["route"] + snapshot = {"target_payload": targets, "binding_payloads": {"goal-alpha": bindings}} + bound = { + "session_id": "manager-session", + "channel_id": route["manager_channel_id"], + } + + assert ( + authorized_manager_goal_ids( + snapshot, {**bound, "agent_id": "codex"}, runtime_root=runtime_root + ) + == [] + ) + assert authorized_manager_goal_ids( + snapshot, {**bound, "agent_id": "dsh"}, runtime_root=runtime_root + ) == ["goal-alpha"] diff --git a/tests/extensions/test_lark_goal_topic_runtime.py b/tests/extensions/test_lark_goal_topic_runtime.py index 5da9021d02..6f04a462e0 100644 --- a/tests/extensions/test_lark_goal_topic_runtime.py +++ b/tests/extensions/test_lark_goal_topic_runtime.py @@ -26,6 +26,67 @@ def test_goal_topic_runtime_exposes_the_inbox_bridge() -> None: assert callable(getattr(module, "process_lark_goal_topic_event", None)) +def test_a_stale_steward_session_names_the_rebind_instead_of_a_generic_failure( + tmp_path: Path, +) -> None: + """A machine that changed its steward executor must say what repairs it.""" + + from types import SimpleNamespace + + from loopx.extensions.lark.goal_topic_runtime import ( + LarkGoalTopicTurnFailed, + answer_lark_goal_topic, + ) + from loopx.extensions.lark.manager_context import manager_failure_reply + + controller = SimpleNamespace( + steward_executor_defaults=lambda: { + "schema_version": "steward_executor_effective_defaults_v0", + "status": "ready", + "source": "machine_configuration", + "executor_endpoint": "dsh", + "executor_model": "deepseek-v4-flash", + "executor_reasoning_effort": "high", + }, + store=SimpleNamespace( + load_session=lambda _session_id: { + "session_id": "manager-session", + "agent_id": "codex", + "channel_id": "manager.external.public_fixture", + "status": "ready", + } + ), + ) + route = { + "schema_version": "lark_goal_topic_route_v0", + "goal_id": "goal-alpha", + "conversation_kind": "manager", + "ingress_mode": "session_queue", + "session_id": "manager-session", + "manager_channel_id": "manager.external.public_fixture", + "message_id": "om_manager_stale_endpoint", + "topic_root_message_id": "om_manager_topic_root", + "app_ref": "cli_public_fixture", + "target_ref": "public_fixture_target", + } + + with pytest.raises(LarkGoalTopicTurnFailed) as failure: + answer_lark_goal_topic( + route=route, + text="@LoopX 管家 status", + work_dir=str(tmp_path), + objective="worker objective", + runtime_controller=controller, + ) + + assert failure.value.error_code == "manager_channel_executor_rebind_required" + code, text = manager_failure_reply(failure.value) + assert code == "manager_channel_executor_rebind_required" + # The reply names the operator action instead of the opaque manager label. + assert "重新应用" in text + assert "管家处理失败" not in text + + def test_existing_collector_uses_the_real_compact_event_schema() -> None: projection = _jq_projection("oc_public_fixture")