Skip to content
1 change: 1 addition & 0 deletions examples/capability-extension-registry-smoke.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
16 changes: 11 additions & 5 deletions examples/lark-goal-topic-connection-smoke.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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, "")},
}
]
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down
50 changes: 50 additions & 0 deletions loopx/extensions/lark/goal_channel_contracts.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 ""
5 changes: 5 additions & 0 deletions loopx/extensions/lark/goal_channel_transport.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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)
)
)


Expand Down
99 changes: 55 additions & 44 deletions loopx/extensions/lark/goal_topic_connections.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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,
)

Expand Down Expand Up @@ -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={
Expand Down
Loading