diff --git a/examples/capability-extension-registry-smoke.py b/examples/capability-extension-registry-smoke.py index 22110a3e0a..a71c766744 100644 --- a/examples/capability-extension-registry-smoke.py +++ b/examples/capability-extension-registry-smoke.py @@ -79,6 +79,7 @@ def run_cli(runtime_root: Path, *args: str) -> dict[str, object]: "deep-research", "public-safe-outbound", "connector-registry", + "reliability-diagnostics", ] assert all(item["provider_id"] == "loopx-core" for item in builtin_capabilities) value_summary = next( diff --git a/examples/lark-goal-topic-connection-smoke.py b/examples/lark-goal-topic-connection-smoke.py index 3d0a6e0369..08e2c5c3f5 100644 --- a/examples/lark-goal-topic-connection-smoke.py +++ b/examples/lark-goal-topic-connection-smoke.py @@ -17,9 +17,10 @@ REPO_ROOT = Path(__file__).resolve().parents[1] sys.path.insert(0, str(REPO_ROOT)) -from loopx.chat_server import ChatHTTPServer, ChatRequestHandler -from loopx.extensions.lark.goal_channel_contracts import read_goal_channel_binding -from loopx.extensions.lark.goal_channel_targets import read_goal_channel_targets +# Standalone smoke imports follow the repository-path bootstrap above. +from loopx.chat_server import ChatHTTPServer, ChatRequestHandler # noqa: E402 +from loopx.extensions.lark.goal_channel_contracts import binding_for_goal, read_goal_channel_binding # noqa: E402 +from loopx.extensions.lark.goal_channel_targets import read_goal_channel_targets # noqa: E402 APP_ID = "cli_public_fixture" @@ -136,6 +137,7 @@ def run(args: list[str], _cwd: object, _timeout: object) -> dict[str, Any]: "chats": [ { "message_id": message_id, + "chat_id": CHAT_ID, "body": {"content": state.get("messages", {}).get(message_id, "")}, } ] @@ -232,6 +234,9 @@ def main() -> None: assert preview["status"] == "preview_ready", preview connected = request(base, "/api/chat/lark/connections", method="POST", body={**body, "execute": True}) assert connected["status"] == "connected" and connected["readback_verified"] is True, connected + reconnected = request(base, "/api/chat/lark/connections", method="POST", body={**body, "execute": True}) + assert reconnected["status"] == "connected" and reconnected["readback_verified"] is True, reconnected + assert reconnected["external_write_performed"] is False, reconnected connections = request(base, "/api/chat/lark/connections") assert len(connections["connections"]) == 2, connections @@ -243,12 +248,13 @@ def main() -> None: target_path = runtime / "goal-channel-targets.json" binding_path = registry_path.parent / "goal-channel.json" assert len(read_goal_channel_targets(target_path)["targets"]) == 1 - bindings = read_goal_channel_binding(binding_path)["bindings"] + payload = read_goal_channel_binding(binding_path) + bindings = {goal_id: binding_for_goal(payload, goal_id, agent_id="agent-alpha") for goal_id in ("goal-alpha", "goal-beta")} assert bindings["goal-alpha"]["topic"]["root_message_id"] != bindings["goal-beta"]["topic"]["root_message_id"] disconnected = request( base, - "/api/chat/lark/connections?goal_id=goal-alpha", + "/api/chat/lark/connections?goal_id=goal-alpha&connection_id=" + bindings["goal-alpha"]["connection_id"], method="DELETE", ) assert disconnected["status"] == "disconnected", disconnected diff --git a/loopx/extensions/lark/goal_channel_contracts.py b/loopx/extensions/lark/goal_channel_contracts.py index d0e57d85e8..531b41fc7f 100644 --- a/loopx/extensions/lark/goal_channel_contracts.py +++ b/loopx/extensions/lark/goal_channel_contracts.py @@ -718,3 +718,53 @@ def gate_message( if kanban_url: lines.extend(["", f"Kanban: {kanban_url}"]) return "\n".join(lines), question + + +def reusable_goal_topic_root( + payload: Mapping[str, Any], + goal_id: str, + *, + connection_id: str, + provider_target: Mapping[str, Any] | None, + chat_id: str, +) -> str: + """Return the established topic root when reconnect may adopt it. + + A reconnect may skip sending a fresh Goal Topic only when the stored + connection for this ``connection_id`` is enabled, still resolves through + ``provider_target`` (target_ref/provider/chat all validated by the typed + reader), and its stored root is a well-formed message id for this chat. + Anything else returns "" so the caller sends a new topic. + """ + + from .goal_channel_transport import MESSAGE_ID_PATTERN + + try: + existing = binding_for_goal( + payload, + goal_id, + connection_id=connection_id, + ) + if existing is not None: + existing = _resolve_goal_binding(existing, provider_target=provider_target) + except ValueError: + return "" + if not existing or existing.get("enabled") is not True: + return "" + prior_topic = ( + existing.get("topic") if isinstance(existing.get("topic"), Mapping) else {} + ) + prior_channel = ( + existing.get("channel") if isinstance(existing.get("channel"), Mapping) else {} + ) + candidate_root = str( + prior_topic.get("root_message_id") + or prior_channel.get("pinned_message_id") + or "" + ) + if ( + MESSAGE_ID_PATTERN.fullmatch(candidate_root) + and str(prior_channel.get("chat_id") or "") == chat_id + ): + return candidate_root + return "" diff --git a/loopx/extensions/lark/goal_channel_transport.py b/loopx/extensions/lark/goal_channel_transport.py index d0206b2753..872380333c 100644 --- a/loopx/extensions/lark/goal_channel_transport.py +++ b/loopx/extensions/lark/goal_channel_transport.py @@ -298,6 +298,7 @@ def message_readback_verified( identity: str, message_id: str, expected_text: str, + expected_chat_id: str | None = None, ) -> bool: result = call( runner, @@ -322,6 +323,10 @@ def message_readback_verified( result.get("returncode") == 0 and contains_exact_field(payload, "message_id", message_id) and payload_contains_text(payload, expected_text) + and ( + expected_chat_id is None + or contains_exact_field(payload, "chat_id", expected_chat_id) + ) ) diff --git a/loopx/extensions/lark/goal_topic_connections.py b/loopx/extensions/lark/goal_topic_connections.py index f7c537940c..6265334fd2 100644 --- a/loopx/extensions/lark/goal_topic_connections.py +++ b/loopx/extensions/lark/goal_topic_connections.py @@ -43,6 +43,7 @@ operation_packet, provider_idempotency_key, read_goal_channel_binding, + reusable_goal_topic_root, save_goal_connection, serialize_goal_binding_mutation, semantic_key, @@ -558,6 +559,13 @@ def connect_lark_goal_topic( objective = goal_objective(goal) topic_text = f"LoopX Goal Topic: {objective}\nGoal ID: {goal_id}" + reusable_root = reusable_goal_topic_root( + read_goal_channel_binding(binding_path), + goal_id, + connection_id=connection_id, + provider_target=matched[1] if matched is not None else None, + chat_id=safe_chat_id, + ) key = provider_idempotency_key( semantic_key( goal_id, @@ -569,60 +577,63 @@ def connect_lark_goal_topic( topic_text, ) ) - sent = call( - runner, - lark_args( - cli_bin=effective_cli_bin, - profile=profile, - tail=[ - "im", - "+messages-send", - "--chat-id", - safe_chat_id, - "--text", - topic_text, - "--idempotency-key", - key, - "--as", - "bot", - "--format", - "json", - ], - ), - ) - root_message_id = find_first_string( - json_payload(sent), - {"message_id"}, - MESSAGE_ID_PATTERN, - ) - if sent.get("returncode") != 0 or not root_message_id: - return operation_packet( - ok=False, - goal_id=goal_id, - operation="connect_topic", - execute=True, - status="failed", - blocker="provider_api_failed", - public_summary="the Goal Topic root message could not be sent", + if reusable_root: + root_message_id = reusable_root + else: + sent = call( + runner, + lark_args( + cli_bin=effective_cli_bin, + profile=profile, + tail=[ + "im", + "+messages-send", + "--chat-id", + safe_chat_id, + "--text", + topic_text, + "--idempotency-key", + key, + "--as", + "bot", + "--format", + "json", + ], + ), + ) + root_message_id = find_first_string( + json_payload(sent), + {"message_id"}, + MESSAGE_ID_PATTERN, ) - verified = message_readback_verified( + if sent.get("returncode") != 0 or not root_message_id: + return operation_packet( + ok=False, + goal_id=goal_id, + operation="connect_topic", + execute=True, + status="failed", + blocker="provider_api_failed", + public_summary="the Goal Topic root message could not be sent", + ) + if not message_readback_verified( runner=runner, cli_bin=effective_cli_bin, profile=profile, identity="bot", message_id=root_message_id, - expected_text=topic_text, - ) - if not verified: + expected_text=f"Goal ID: {goal_id}" if reusable_root else topic_text, + expected_chat_id=safe_chat_id if reusable_root else None, + ): return operation_packet( ok=False, goal_id=goal_id, operation="connect_topic", execute=True, - status="sent_unverified", + status="blocked" if reusable_root else "sent_unverified", blocker="readback_mismatch", - public_summary="the Goal Topic was sent but could not be verified", - external_write_performed=True, + public_summary="the Goal Topic could not be verified; binding preserved", + external_write_performed=not bool(reusable_root), ) connector_binding: dict[str, Any] | None = None @@ -683,7 +694,7 @@ def connect_lark_goal_topic( status="sent_verified_registration_failed", blocker="agent_inbox_registration_failed", public_summary="the Agent-scoped inbox could not be registered", - external_write_performed=True, + external_write_performed=not bool(reusable_root), readback_verified=True, ) @@ -732,7 +743,7 @@ def connect_lark_goal_topic( execute=True, status="connected", public_summary="connected one Goal to a dedicated Lark topic", - external_write_performed=True, + external_write_performed=not bool(reusable_root), readback_verified=True, idempotency_key=key, details={ diff --git a/tests/extensions/test_lark_goal_topic_connections.py b/tests/extensions/test_lark_goal_topic_connections.py index 3b07db8661..f5adbc806b 100644 --- a/tests/extensions/test_lark_goal_topic_connections.py +++ b/tests/extensions/test_lark_goal_topic_connections.py @@ -10,10 +10,13 @@ from loopx.control_plane.quota.goal_boundary import goal_boundary from loopx.extensions.lark.goal_channel_contracts import ( + GOAL_CHANNEL_BINDING_SCHEMA_VERSION, + GOAL_CHANNEL_CONNECTION_SET_SCHEMA_VERSION, binding_for_goal, bindings_for_goal, goal_channel_connection_id, read_goal_channel_binding, + write_goal_channel_binding, ) from loopx.extensions.lark.goal_channel_targets import ( add_lark_goal_channel_target, @@ -108,6 +111,7 @@ def run(args: list[str], _cwd: object, _timeout: object) -> dict[str, Any]: "chats": [ { "message_id": message_id, + "chat_id": CHAT_ID, "body": { "content": state.get("sent", {}).get( message_id, "reply ok" @@ -1394,9 +1398,7 @@ def test_disconnect_async_inbox_unregisters_only_selected_agent( assert result["ok"] is True assert result["details"]["agent_inbox_unregistered"] is True - remaining = bindings_for_goal( - read_goal_channel_binding(binding_path), "goal-alpha" - ) + remaining = bindings_for_goal(read_goal_channel_binding(binding_path), "goal-alpha") assert [item["agent_id"] for item in remaining] == ["agent-beta"] registry = json.loads(registry_path.read_text(encoding="utf-8")) goal = registry["goals"][0] @@ -1407,9 +1409,12 @@ def test_disconnect_async_inbox_unregisters_only_selected_agent( assert alpha_boundary is None or "lark_event_inbox" not in alpha_boundary.get( "capabilities", {} ) - assert goal_boundary( - goal, agent_id="agent-beta", registry_path=registry_path - )["capabilities"]["lark_event_inbox"]["enabled"] is True + assert ( + goal_boundary(goal, agent_id="agent-beta", registry_path=registry_path)[ + "capabilities" + ]["lark_event_inbox"]["enabled"] + is True + ) def test_disconnect_reports_agent_inbox_cleanup_failure( @@ -1462,8 +1467,240 @@ def test_disconnect_reports_agent_inbox_cleanup_failure( "agent_inbox_unregistered": False, "agent_id": "agent-alpha", } - assert binding_for_goal( + assert ( + binding_for_goal( + read_goal_channel_binding(binding_path), + "goal-alpha", + connection_id=str(connection["connection_id"]), + ) + is None + ) + + +def _legacy_v0_binding_payload( + root_message_id: str, agent_id: str, target_ref: str = "mew-product" +) -> dict[str, Any]: + """A v0 single-binding file as written by pre-#3969 releases.""" + + return { + "schema_version": GOAL_CHANNEL_BINDING_SCHEMA_VERSION, + "bindings": { + "goal-alpha": { + "enabled": True, + "provider": "lark", + "target_ref": target_ref, + "agent_id": agent_id, + "session_id": "", + "topic": {"name": "Alpha delivery", "root_message_id": root_message_id}, + "channel": {"chat_id": CHAT_ID, "chat_name": "Product group"}, + "routing": { + "incoming_mode": "mentions", + "capture_scope": "addressed_only", + "ingress_mode": "direct_session", + "reply_mode": "topic_reply", + }, + "receipts": {}, + "automation": {"human_gate_auto_notify": True}, + } + }, + } + + +def _prep_goal_channel_target(root: Path) -> Path: + target_path = root / "targets.json" + add_lark_goal_channel_target( + target_path=target_path, + target_name="mew-product", + chat_id=CHAT_ID, + chat_name="Product group", + identity_mode="local_user", + sender_profile="mew", + bot_app_id=APP_ID, + bot_open_id=None, + bot_display_name="mew bot", + cli_bin="fake-lark", + execute=True, + ) + return target_path + + +def test_reconnect_after_upgrade_reuses_legacy_topic_root_without_resend( + tmp_path: Path, +) -> None: + registry = _registry(tmp_path) + registry["goals"][0]["coordination"] = {"registered_agents": ["agent-alpha"]} + binding_path = tmp_path / "binding.json" + target_path = _prep_goal_channel_target(tmp_path) + write_goal_channel_binding( + binding_path, _legacy_v0_binding_payload("om_legacy_root", "agent-alpha") + ) + state: dict[str, Any] = {"sent": {"om_legacy_root": "Old title\nGoal ID: goal-alpha"}} + result = connect_lark_goal_topic( + registry=registry, + goal_id="goal-alpha", + target_path=target_path, + binding_path=binding_path, + app_ref="mew", + chat_id=CHAT_ID, + chat_name="Product group", + agent_id="agent-alpha", + runner=_runner(state), + cli_bin="fake-lark", + ) + assert result["ok"] is True + sends = [call for call in state["calls"] if "+messages-send" in call] + assert sends == [] + assert result["external_write_performed"] is False + assert result["readback_verified"] is True + connection = binding_for_goal( read_goal_channel_binding(binding_path), "goal-alpha", - connection_id=str(connection["connection_id"]), - ) is None + agent_id="agent-alpha", + ) + assert connection is not None + assert connection["enabled"] is True + assert str(connection["topic"]["root_message_id"]) == "om_legacy_root" + + +def test_reconnect_with_mismatched_target_ref_sends_new_topic( + tmp_path: Path, +) -> None: + registry = _registry(tmp_path) + registry["goals"][0]["coordination"] = {"registered_agents": ["agent-alpha"]} + binding_path = tmp_path / "binding.json" + target_path = _prep_goal_channel_target(tmp_path) + write_goal_channel_binding( + binding_path, + _legacy_v0_binding_payload( + "om_legacy_root", "agent-alpha", target_ref="mew-elsewhere" + ), + ) + state: dict[str, Any] = {} + result = connect_lark_goal_topic( + registry=registry, + goal_id="goal-alpha", + target_path=target_path, + binding_path=binding_path, + app_ref="mew", + chat_id=CHAT_ID, + chat_name="Product group", + agent_id="agent-alpha", + runner=_runner(state), + cli_bin="fake-lark", + ) + assert result["ok"] is True + sends = [call for call in state["calls"] if "+messages-send" in call] + assert len(sends) == 1 + connection = binding_for_goal( + read_goal_channel_binding(binding_path), + "goal-alpha", + agent_id="agent-alpha", + ) + assert connection is not None + assert str(connection["topic"]["root_message_id"]) == "om_topic_alpha" + + +def test_set_connection_reconnect_reuses_root_without_resend(tmp_path: Path) -> None: + registry = _registry(tmp_path) + registry["goals"][0]["coordination"] = {"registered_agents": ["agent-alpha"]} + binding_path = tmp_path / "binding.json" + target_path = _prep_goal_channel_target(tmp_path) + state: dict[str, Any] = {} + kwargs = dict( + registry=registry, + goal_id="goal-alpha", + target_path=target_path, + binding_path=binding_path, + app_ref="mew", + chat_id=CHAT_ID, + chat_name="Product group", + agent_id="agent-alpha", + runner=_runner(state), + cli_bin="fake-lark", + ) + assert connect_lark_goal_topic(**kwargs)["ok"] is True + first_sends = [call for call in state["calls"] if "+messages-send" in call] + assert len(first_sends) == 1 + assert connect_lark_goal_topic(**kwargs)["ok"] is True + total_sends = [call for call in state["calls"] if "+messages-send" in call] + assert len(total_sends) == 1 + connection = binding_for_goal( + read_goal_channel_binding(binding_path), + "goal-alpha", + agent_id="agent-alpha", + ) + assert connection is not None + assert str(connection["topic"]["root_message_id"]) == "om_topic_alpha" + assert connection.get("connection_id", "").startswith("lark_") + assert connection.get("receipts"), "a reused root must keep its topic receipt" + + +@pytest.mark.parametrize("failure", ["missing", "wrong_chat", "wrong_goal"]) +def test_reconnect_unverified_root_preserves_binding(tmp_path: Path, failure: str) -> None: + registry = _registry(tmp_path) + registry["goals"][0]["coordination"] = {"registered_agents": ["agent-alpha"]} + target_path = _prep_goal_channel_target(tmp_path) + binding_path = tmp_path / "binding.json" + original = _legacy_v0_binding_payload("om_legacy_root", "agent-alpha") + write_goal_channel_binding(binding_path, original) + state: dict[str, Any] = {} + normal = _runner(state) + + def runner(args, cwd, timeout): + if "+messages-mget" not in args: + return normal(args, cwd, timeout) + state.setdefault("readbacks", []).append(args) + return { + "returncode": 1 if failure == "missing" else 0, + "stdout": json.dumps({"data": {"items": [{ + "message_id": "om_legacy_root", + "chat_id": "oc_other" if failure == "wrong_chat" else CHAT_ID, + "body": {"content": "Goal ID: " + ( + "goal-other" if failure == "wrong_goal" else "goal-alpha" + )}, + }]}}), + "stderr": "", + } + + result = connect_lark_goal_topic( + registry=registry, goal_id="goal-alpha", target_path=target_path, + binding_path=binding_path, app_ref="mew", chat_id=CHAT_ID, + chat_name="Product group", agent_id="agent-alpha", runner=runner, cli_bin="fake-lark", + ) + assert result["ok"] is False + assert result["readback_verified"] is False + assert result["external_write_performed"] is False + assert len(state["readbacks"]) == 1 + assert not any("+messages-send" in args for args in state["calls"]) + assert read_goal_channel_binding(binding_path) == original + + +def test_reconnect_isolates_other_agent_target(tmp_path: Path) -> None: + registry = _registry(tmp_path) + registry["goals"][0]["coordination"] = {"registered_agents": ["agent-alpha", "agent-beta"]} + target_path = _prep_goal_channel_target(tmp_path) + binding_path = tmp_path / "binding.json" + alpha = _legacy_v0_binding_payload("om_existing", "agent-alpha")["bindings"]["goal-alpha"] + beta = _legacy_v0_binding_payload("om_other", "agent-beta", "other-target")["bindings"]["goal-alpha"] + alpha_id = goal_channel_connection_id("goal-alpha", "agent-alpha") + beta_id = goal_channel_connection_id("goal-alpha", "agent-beta") + write_goal_channel_binding(binding_path, { + "schema_version": GOAL_CHANNEL_BINDING_SCHEMA_VERSION, + "bindings": {"goal-alpha": { + "schema_version": GOAL_CHANNEL_CONNECTION_SET_SCHEMA_VERSION, + "connections": {alpha_id: alpha, beta_id: beta}, + }}, + }) + state: dict[str, Any] = {"sent": {"om_existing": "Old title\nGoal ID: goal-alpha"}} + result = connect_lark_goal_topic( + registry=registry, goal_id="goal-alpha", target_path=target_path, + binding_path=binding_path, app_ref="mew", chat_id=CHAT_ID, + chat_name="Product group", agent_id="agent-alpha", runner=_runner(state), cli_bin="fake-lark", + ) + assert result["ok"] is True + assert result["external_write_performed"] is False + assert result["readback_verified"] is True + assert not any("+messages-send" in args for args in state["calls"]) + saved = read_goal_channel_binding(binding_path)["bindings"]["goal-alpha"]["connections"] + assert saved[alpha_id]["topic"]["root_message_id"] == "om_existing" + assert saved[beta_id] == beta