From 22a39df58b5ab76d86358187d562fc1f6102e795 Mon Sep 17 00:00:00 2001 From: now-ing Date: Sun, 6 Sep 2026 17:49:29 +0800 Subject: [PATCH 1/2] fix(lark): report inbox cleanup failures as packets during disconnect Disconnect writes the binding removal before unregistering the Agent inbox. When the registry mutation timed out or the registry was briefly unreadable, configure_goal_with_global_sync raised past the caller: the HTTP layer answered invalid_lark_connection while the connection was already gone, and the Agent inbox registration was left orphaned. Catch OSError/ValueError/TimeoutError in _unregister_async_inbox and return the failed result so the existing disconnected_inbox_cleanup_failed packet path reports the durable binding removal faithfully. Signed-off-by: now-ing --- .../extensions/lark/goal_topic_connections.py | 32 ++++---- .../test_lark_goal_topic_connections.py | 78 +++++++++++++++++++ 2 files changed, 94 insertions(+), 16 deletions(-) diff --git a/loopx/extensions/lark/goal_topic_connections.py b/loopx/extensions/lark/goal_topic_connections.py index 6265334fd2..10d6914078 100644 --- a/loopx/extensions/lark/goal_topic_connections.py +++ b/loopx/extensions/lark/goal_topic_connections.py @@ -1406,10 +1406,7 @@ def _without_goal_topic_connection( def _unregister_async_inbox( - *, - removed: Mapping[str, Any] | None, - registry_path: Path | None, - goal_id: str, + *, removed: Mapping[str, Any] | None, registry_path: Path | None, goal_id: str ) -> tuple[dict[str, Any] | None, str]: routing = removed.get("routing") if isinstance(removed, Mapping) else None routing = routing if isinstance(routing, Mapping) else {} @@ -1417,18 +1414,21 @@ def _unregister_async_inbox( if routing.get("ingress_mode") != IngressMode.ASYNC_INBOX.value or not agent_id: return None, agent_id if registry_path is None: - return { - "ok": False, - "error": "source registry path is required to unregister the Agent inbox", - }, agent_id - return configure_goal_with_global_sync( - registry_path=registry_path, - goal_id=goal_id, - runtime_root_override=None, - execute=True, - lark_event_inbox_agent_id=agent_id, - clear_lark_event_inbox_config=True, - ), agent_id + error = "source registry path is required to unregister the Agent inbox" + return {"ok": False, "error": error}, agent_id + try: + return configure_goal_with_global_sync( + registry_path=registry_path, + goal_id=goal_id, + runtime_root_override=None, + execute=True, + lark_event_inbox_agent_id=agent_id, + clear_lark_event_inbox_config=True, + ), agent_id + except (OSError, ValueError, TimeoutError) as exc: + # The binding removal already landed; report the cleanup failure as a + # failed packet instead of raising past the caller mid-disconnect. + return {"ok": False, "error": str(exc)}, agent_id @serialize_goal_binding_mutation diff --git a/tests/extensions/test_lark_goal_topic_connections.py b/tests/extensions/test_lark_goal_topic_connections.py index f5adbc806b..05476c9397 100644 --- a/tests/extensions/test_lark_goal_topic_connections.py +++ b/tests/extensions/test_lark_goal_topic_connections.py @@ -33,6 +33,7 @@ reply_lark_goal_topic, route_lark_topic_event, ) +from loopx.file_lock import LockAcquireTimeoutError from loopx.registry import atomic_write_json APP_ID = "cli_public_fixture" @@ -1477,6 +1478,83 @@ def test_disconnect_reports_agent_inbox_cleanup_failure( ) +@pytest.mark.parametrize( + "failure_factory", + [ + lambda: LockAcquireTimeoutError( + incident={"holder": {"pid": 4242}}, + incident_recorded=False, + incident_channel="test", + ), + lambda: ValueError("goal registry fixture is unreadable"), + ], + ids=["lock-timeout", "registry-error"], +) +def test_disconnect_reports_agent_inbox_cleanup_raise_as_packet( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, + failure_factory: Any, +) -> None: + registry_path = tmp_path / ".loopx" / "registry.json" + registry = _registry(tmp_path) + registry["common_runtime_root"] = str(tmp_path / "runtime") + registry["goals"][0]["coordination"] = {"registered_agents": ["agent-alpha"]} + atomic_write_json(registry_path, registry) + binding_path = tmp_path / "binding.json" + assert connect_lark_goal_topic( + registry=registry, + registry_path=registry_path, + goal_id="goal-alpha", + target_path=tmp_path / "targets.json", + binding_path=binding_path, + app_ref="mew", + chat_id=CHAT_ID, + chat_name="Product group", + agent_id="agent-alpha", + ingress_mode="async_inbox", + runner=_runner({}), + cli_bin="fake-lark", + )["ok"] + connection = binding_for_goal( + read_goal_channel_binding(binding_path), + "goal-alpha", + ) + assert connection is not None + + def _raise(**_kwargs: Any) -> dict[str, Any]: + raise failure_factory() + + monkeypatch.setattr( + "loopx.extensions.lark.goal_topic_connections.configure_goal_with_global_sync", + _raise, + ) + + result = disconnect_lark_goal_topic( + binding_path=binding_path, + registry_path=registry_path, + goal_id="goal-alpha", + connection_id=str(connection["connection_id"]), + ) + + assert result["ok"] is False + assert result["status"] == "disconnected_inbox_cleanup_failed" + assert result["blocker"] == "agent_inbox_unregistration_failed" + assert result["readback_verified"] is False + assert result["details"] == { + "connection_id": connection["connection_id"], + "agent_inbox_unregistered": False, + "agent_id": "agent-alpha", + } + 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]: From 6abf5b28969e5eeb71d7699079e821e957ee77f6 Mon Sep 17 00:00:00 2001 From: now-ing Date: Sun, 6 Sep 2026 18:02:04 +0800 Subject: [PATCH 2/2] fix(lark): align invalid-default fallback between binding readers and writers When the stored default_connection_id no longer matches a live connection, binding_for_goal fell back to the first candidate in JSON insertion order while _without_goal_topic_connection promoted min(connection_id). Notify/setup/automation reads could therefore target a different Agent's topic than the default written back after a disconnect. Select min(connection_id) on the read side so both fallbacks agree. Signed-off-by: now-ing --- .../extensions/lark/goal_channel_contracts.py | 6 ++- .../test_lark_goal_topic_connections.py | 43 +++++++++++++++++++ 2 files changed, 48 insertions(+), 1 deletion(-) diff --git a/loopx/extensions/lark/goal_channel_contracts.py b/loopx/extensions/lark/goal_channel_contracts.py index 531b41fc7f..989b2df6bc 100644 --- a/loopx/extensions/lark/goal_channel_contracts.py +++ b/loopx/extensions/lark/goal_channel_contracts.py @@ -245,7 +245,11 @@ def binding_for_goal( ) if selected is not None: return selected - return candidates[0] if candidates else None + # Keep the invalid-default fallback aligned with the writer in + # _without_goal_topic_connection, which promotes min(connection_id). + if not candidates: + return None + return min(candidates, key=lambda item: str(item.get("connection_id") or "")) def human_gate_auto_notify_enabled(binding: Mapping[str, Any] | None) -> bool: diff --git a/tests/extensions/test_lark_goal_topic_connections.py b/tests/extensions/test_lark_goal_topic_connections.py index 05476c9397..4180a26803 100644 --- a/tests/extensions/test_lark_goal_topic_connections.py +++ b/tests/extensions/test_lark_goal_topic_connections.py @@ -1555,6 +1555,49 @@ def _raise(**_kwargs: Any) -> dict[str, Any]: ) +def test_invalid_default_fallback_selects_same_connection_as_disconnect( + tmp_path: Path, +) -> None: + beta_id = goal_channel_connection_id("goal-alpha", "agent-beta") + zeta_id = goal_channel_connection_id("goal-alpha", "agent-zeta") + omega_id = goal_channel_connection_id("goal-alpha", "agent-omega") + assert min(beta_id, zeta_id, omega_id) == omega_id + binding_path = tmp_path / "binding.json" + write_goal_channel_binding( + binding_path, + { + "schema_version": GOAL_CHANNEL_BINDING_SCHEMA_VERSION, + "bindings": { + "goal-alpha": { + "schema_version": GOAL_CHANNEL_CONNECTION_SET_SCHEMA_VERSION, + "default_connection_id": "lark_missing0000000000000000", + "connections": { + beta_id: {"agent_id": "agent-beta", "enabled": True}, + zeta_id: {"agent_id": "agent-zeta", "enabled": True}, + omega_id: {"agent_id": "agent-omega", "enabled": True}, + }, + } + }, + }, + ) + payload = read_goal_channel_binding(binding_path) + + selected = binding_for_goal(payload, "goal-alpha") + + assert selected is not None + assert str(selected["connection_id"]) == omega_id + + result = disconnect_lark_goal_topic( + binding_path=binding_path, + goal_id="goal-alpha", + connection_id=beta_id, + ) + + assert result["ok"] is True + stored = read_goal_channel_binding(binding_path)["bindings"]["goal-alpha"] + assert stored["default_connection_id"] == omega_id + + def _legacy_v0_binding_payload( root_message_id: str, agent_id: str, target_ref: str = "mew-product" ) -> dict[str, Any]: