From 61e9a16777922dc16a9503bbb111d262888393fb Mon Sep 17 00:00:00 2001
From: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
Date: Mon, 14 Sep 2026 16:58:50 +0800
Subject: [PATCH 1/6] fix(lark): verify normalized card callbacks
Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
---
loopx/extensions/lark/event_collector.py | 3 +
.../lark/event_collector_runtime.py | 20 ++
.../lark/goal_channel_message_delivery.py | 107 ++++++-
.../extensions/lark/goal_channel_operation.py | 92 +++++-
.../references/repair-patterns.md | 1 +
.../test_lark_event_collector_runtime.py | 1 +
.../test_lark_goal_channel_operation.py | 300 +++++++++++++++++-
7 files changed, 501 insertions(+), 23 deletions(-)
diff --git a/loopx/extensions/lark/event_collector.py b/loopx/extensions/lark/event_collector.py
index 640e0f0700..cc8f733c9d 100644
--- a/loopx/extensions/lark/event_collector.py
+++ b/loopx/extensions/lark/event_collector.py
@@ -735,6 +735,9 @@ def inspect_lark_event_collector(
"operation_callback_last_evidence_at": callback_status.get(
"last_verified_callback_at"
),
+ "operation_callback_last_failure_code": callback_status.get(
+ "last_failure_code"
+ ),
"operation_callback_console_configuration_preflighted": False,
"thread_complete": all(
route["inbox"]["thread_complete"] for route in config["routes"]
diff --git a/loopx/extensions/lark/event_collector_runtime.py b/loopx/extensions/lark/event_collector_runtime.py
index 03e8d9f06a..bfa9dc7453 100644
--- a/loopx/extensions/lark/event_collector_runtime.py
+++ b/loopx/extensions/lark/event_collector_runtime.py
@@ -31,10 +31,27 @@
APP_ID_PATTERN = re.compile(r"cli_[A-Za-z0-9_-]+")
EVENT_READY_PREFIX = "[event] ready "
EVENT_DIAGNOSTIC_PREFIX = "[event] "
+_CALLBACK_FAILURE_CODES = {
+ "operation card delivery was not recorded": "delivery_not_recorded",
+ "operation callback digest drifted": "confirmation_digest_drifted",
+ "recorded operation card digest drifted": "recorded_card_digest_drifted",
+ "operation callback card content drifted": "callback_card_projection_drifted",
+ "operation callback app identity drifted": "callback_app_identity_drifted",
+ "principal is not authorized for this operation": "principal_not_authorized",
+ "operation callback tenant membership is unverified": "membership_unverified",
+ "operation callback does not match the delivered request": "delivery_binding_mismatch",
+ "operation is not awaiting confirmation": "operation_not_awaiting_confirmation",
+ "operation confirmation arrived after expiry": "operation_expired",
+ "operation callback result delivery was not verified": "result_delivery_unverified",
+}
CommandRunner = Callable[..., subprocess.CompletedProcess[str]]
Sleeper = Callable[[float], None]
+def _operation_callback_failure_code(exc: BaseException) -> str:
+ return _CALLBACK_FAILURE_CODES.get(str(exc), "callback_rejected")
+
+
def _run_json(
runner: CommandRunner,
argv: Sequence[str],
@@ -392,6 +409,7 @@ def _write_operation_callback_status(
listener_ready: bool | None = None,
callback_delivery_verified: bool | None = None,
failure_kind: str | None = None,
+ failure_code: str | None = None,
consumer_returncode: int | None = None,
recovered_result_count_delta: int = 0,
result_delivery_failure_count_delta: int = 0,
@@ -444,6 +462,7 @@ def _write_operation_callback_status(
else prior.get("last_verified_callback_at")
),
"last_failure_kind": failure_kind or prior.get("last_failure_kind"),
+ "last_failure_code": failure_code or prior.get("last_failure_code"),
"consumer_returncode": consumer_returncode,
"updated_at": now,
"private_content_returned": False,
@@ -722,6 +741,7 @@ def consume_operation_callbacks() -> None:
listener_active=True,
listener_ready=True,
failure_kind=type(exc).__name__,
+ failure_code=_operation_callback_failure_code(exc),
)
continue
callback_stats["verified"] += 1
diff --git a/loopx/extensions/lark/goal_channel_message_delivery.py b/loopx/extensions/lark/goal_channel_message_delivery.py
index 552e1ed16e..1b6906bf81 100644
--- a/loopx/extensions/lark/goal_channel_message_delivery.py
+++ b/loopx/extensions/lark/goal_channel_message_delivery.py
@@ -130,6 +130,8 @@ def _message_card(value: Mapping[str, Any]) -> Mapping[str, Any] | None:
def _normalized_card_text(card: Mapping[str, Any]) -> str | None:
+ if card.get("schema") == "2.0":
+ return _normalized_card_v2_text(card)
header = card.get("header")
elements = card.get("elements")
if not isinstance(header, Mapping) or not isinstance(elements, list):
@@ -160,12 +162,109 @@ def _normalized_card_text(card: Mapping[str, Any]) -> str | None:
return "\n".join(lines)
-def _message_card_matches(
+def _text_content(value: object) -> str | None:
+ if not isinstance(value, Mapping):
+ return None
+ content = value.get("content")
+ return content if isinstance(content, str) and content else None
+
+
+def _card_v2_element_lines(value: object) -> list[str] | None:
+ if not isinstance(value, Mapping):
+ return None
+ tag = value.get("tag")
+ if tag == "markdown":
+ content = value.get("content")
+ return [content] if isinstance(content, str) and content else None
+ if tag in {"plain_text", "text"}:
+ content = value.get("content") or value.get("text")
+ return [content] if isinstance(content, str) and content else None
+ if tag == "button":
+ label = _text_content(value.get("text"))
+ return [f"[{label}]"] if label else None
+ children: object = None
+ if tag == "column_set":
+ children = value.get("columns")
+ elif tag == "column":
+ children = value.get("elements")
+ if not isinstance(children, list):
+ return None
+ lines: list[str] = []
+ button_labels: list[str] = []
+ for child in children:
+ child_lines = _card_v2_element_lines(child)
+ if child_lines is None:
+ return None
+ if (
+ isinstance(child, Mapping)
+ and child.get("tag") == "column"
+ and all(line.startswith("[") and line.endswith("]") for line in child_lines)
+ ):
+ button_labels.extend(child_lines)
+ else:
+ lines.extend(child_lines)
+ if button_labels:
+ lines.append(" ".join(button_labels))
+ return lines
+
+
+def _normalized_card_v2_text(card: Mapping[str, Any]) -> str | None:
+ header = card.get("header")
+ body = card.get("body")
+ if not isinstance(header, Mapping) or not isinstance(body, Mapping):
+ return None
+ title = _text_content(header.get("title"))
+ subtitle = _text_content(header.get("subtitle"))
+ elements = body.get("elements")
+ tags = header.get("text_tag_list")
+ if not title or not isinstance(elements, list) or not elements:
+ return None
+ attributes = f'title="{title}"'
+ if subtitle:
+ attributes += f' subtitle="{subtitle}"'
+ lines = [f""]
+ if tags is not None:
+ if not isinstance(tags, list):
+ return None
+ for item in tags:
+ if not isinstance(item, Mapping):
+ return None
+ text = _text_content(item.get("text"))
+ if not text:
+ return None
+ lines.append(f"「{text}」")
+ for element in elements:
+ element_lines = _card_v2_element_lines(element)
+ if element_lines is None:
+ return None
+ lines.extend(element_lines)
+ lines.append("")
+ return "\n".join(lines)
+
+
+def card_projection_matches(
+ observed: Mapping[str, Any], expected: Mapping[str, Any]
+) -> bool:
+ """Compare an exact card or its provider-normalized visible projection."""
+
+ if observed == expected:
+ return True
+ observed_text = _normalized_card_text(observed)
+ expected_text = _normalized_card_text(expected)
+ return (
+ observed_text is not None
+ and expected_text is not None
+ and observed_text == expected_text
+ )
+
+
+def message_card_matches(
value: Mapping[str, Any], expected: Mapping[str, Any] | None
) -> bool:
if expected is None:
return False
- if _message_card(value) == expected:
+ observed = _message_card(value)
+ if isinstance(observed, Mapping) and card_projection_matches(observed, expected):
return True
content = value.get("content")
return isinstance(content, str) and content == _normalized_card_text(expected)
@@ -245,7 +344,7 @@ def _existing_message(
and str(message.get("chat_id") or "") == route["chat_id"]
and sender_type == "app"
and sender_app_id == route["bot_app_id"]
- and _message_card_matches(message, card)
+ and message_card_matches(message, card)
):
return str(message["message_id"])
if not _history_is_complete(payload):
@@ -390,7 +489,7 @@ def readback(self, message_id: str) -> Mapping[str, Any]:
result.get("returncode") == 0
and message is not None
and contains_exact_field(message, "chat_id", str(self.route["chat_id"]))
- and _message_card_matches(message, expected_card)
+ and message_card_matches(message, expected_card)
and sender_type == "app"
and sender_app_id == self.route["bot_app_id"]
and auth_verified(
diff --git a/loopx/extensions/lark/goal_channel_operation.py b/loopx/extensions/lark/goal_channel_operation.py
index 4140dcb40b..f8db171d0a 100644
--- a/loopx/extensions/lark/goal_channel_operation.py
+++ b/loopx/extensions/lark/goal_channel_operation.py
@@ -1,5 +1,6 @@
from __future__ import annotations
+from copy import deepcopy
from datetime import datetime, timezone
import hashlib
import html
@@ -24,6 +25,8 @@
)
from .goal_channel_message_delivery import (
GoalChannelMessageDeliverySession,
+ card_projection_matches,
+ message_card_matches,
resolve_bound_goal_channel,
)
from .goal_channel_transport import call, json_payload, lark_args
@@ -295,6 +298,22 @@ def build_goal_channel_operation_card(
}
+def _submitted_confirmation_card(proposal: Mapping[str, Any]) -> dict[str, Any]:
+ """Rebuild the immutable submitted card after the operation has advanced."""
+
+ operation = proposal.get("operation")
+ if not isinstance(operation, Mapping):
+ raise ValueError("typed operation envelope is unavailable")
+ if operation.get("lifecycle_state") == "awaiting_confirmation":
+ return build_goal_channel_operation_card(proposal)
+ replay = deepcopy(dict(proposal))
+ replay_operation = replay.get("operation")
+ if not isinstance(replay_operation, dict):
+ raise ValueError("typed operation envelope is unavailable")
+ replay_operation["lifecycle_state"] = "awaiting_confirmation"
+ return build_goal_channel_operation_card(replay)
+
+
def build_goal_channel_operation_result_card(
proposal: Mapping[str, Any],
) -> dict[str, Any]:
@@ -560,6 +579,56 @@ def _callback_timestamp(value: object) -> str:
)
+def _lark_card_v2_fallback_matches(
+ observed: Mapping[str, Any], expected: Mapping[str, Any]
+) -> bool:
+ """Recognize Lark's message-get fallback for a Card 2.0 payload.
+
+ The provider exposes Card 2.0 through message-get as a title plus an
+ upgrade-client placeholder. ``card.action.trigger`` consumers hydrate
+ ``card_content`` from that endpoint, so its digest cannot equal the
+ submitted Card 2.0 JSON. The exact message, route, app, action digest, and
+ recorded submitted-card digest are checked independently by the caller.
+ """
+
+ if expected.get("schema") != "2.0":
+ return False
+ header = expected.get("header")
+ if not isinstance(header, Mapping):
+ return False
+ title_value = header.get("title")
+ subtitle_value = header.get("subtitle")
+ title = title_value.get("content") if isinstance(title_value, Mapping) else None
+ subtitle = (
+ subtitle_value.get("content") if isinstance(subtitle_value, Mapping) else None
+ )
+ expected_title = "\n".join(
+ item for item in (title, subtitle) if isinstance(item, str) and item
+ )
+ return _lark_card_v2_fallback_matches_title(observed, expected_title)
+
+
+def _lark_card_v2_fallback_matches_title(
+ observed: Mapping[str, Any], expected_title: str
+) -> bool:
+ if set(observed) != {"title", "elements"}:
+ return False
+ elements = observed.get("elements")
+ if observed.get("title") != expected_title or not isinstance(elements, list):
+ return False
+ leaves: list[Mapping[str, Any]] = []
+
+ def collect(value: object) -> bool:
+ if isinstance(value, list):
+ return bool(value) and all(collect(item) for item in value)
+ if not isinstance(value, Mapping) or value.get("tag") not in {"img", "text"}:
+ return False
+ leaves.append(value)
+ return True
+
+ return collect(elements) and any(item.get("tag") == "img" for item in leaves)
+
+
def _operator_membership_verified(
*,
runner: CommandRunner,
@@ -726,16 +795,6 @@ def _find_message(value: object, message_id: str) -> Mapping[str, Any] | None:
return None
-def _message_card(value: Mapping[str, Any]) -> Mapping[str, Any] | None:
- body = value.get("body")
- raw = body.get("content") if isinstance(body, Mapping) else value.get("content")
- try:
- card = json.loads(raw) if isinstance(raw, str) else raw
- except json.JSONDecodeError:
- return None
- return card if isinstance(card, Mapping) else None
-
-
def _result_card_readback_verified(
payload: Mapping[str, Any],
*,
@@ -746,7 +805,6 @@ def _result_card_readback_verified(
) -> bool:
message = _find_message(payload, message_id)
sender = message.get("sender") if isinstance(message, Mapping) else None
- observed_card = _message_card(message) if isinstance(message, Mapping) else None
return bool(
payload.get("ok") is True
and isinstance(message, Mapping)
@@ -754,8 +812,7 @@ def _result_card_readback_verified(
and isinstance(sender, Mapping)
and sender.get("sender_type") == "app"
and sender.get("id") == app_id
- and isinstance(observed_card, Mapping)
- and _digest(observed_card) == _digest(card)
+ and message_card_matches(message, card)
)
@@ -1089,7 +1146,14 @@ def handle_goal_channel_operation_callback(
if action["confirmation_digest"] != operation.get("confirmation_digest"):
raise ActionConflictError("operation callback digest drifted")
if _digest(card) != delivery.get("card_digest"):
- raise ActionConflictError("operation callback card content drifted")
+ expected_card = _submitted_confirmation_card(proposal)
+ if _digest(expected_card) != delivery.get("card_digest"):
+ raise ActionConflictError("recorded operation card digest drifted")
+ if not (
+ card_projection_matches(card, expected_card)
+ or _lark_card_v2_fallback_matches(card, expected_card)
+ ):
+ raise ActionConflictError("operation callback card content drifted")
if profile_app_id != delivery.get("app_id"):
raise ActionConflictError("operation callback app identity drifted")
operator_id = str(event["operator_id"])
diff --git a/skills/loopx-self-repair/references/repair-patterns.md b/skills/loopx-self-repair/references/repair-patterns.md
index ef01f0cf13..ea8f3d469d 100644
--- a/skills/loopx-self-repair/references/repair-patterns.md
+++ b/skills/loopx-self-repair/references/repair-patterns.md
@@ -195,6 +195,7 @@ teaches a reusable control-plane lesson.
| `lark_configured_chat_topic_filter_gap` | A `configured_chat_all` connection can read a message through provider history, but a real-time message from a new topic or topic reply in the same configured chat is absent from the Goal inbox and diagnostics report `topic_mismatch`. | Consuming bot target, configured chat id, incoming chat/root ids, capture scope, route decision reason, and persisted inbox message ids. | The router applied one presentation topic root as an ingress filter before evaluating chat-wide capture scope; profile polling could also route against bindings owned by another bot target in the same chat. | Scope polling to the consuming target first. Prefer an exact topic match; otherwise allow one unambiguous chat-wide route for any root in the configured chat, and fail closed when multiple chat-wide routes remain. Treat the topic root as reply/presentation context rather than an ingress boundary under chat-wide capture, and cover new topics plus multiple bots with focused tests. |
| `lark_inbox_reply_placement_gap` | A top-level chat request is answered with an unnecessary new topic, or structured reply text loses line breaks. | Source `parent_id`/`root_id`, configured placement policy, provider send command, and readback text. | The reply path always forced thread mode and flattened all whitespace before sending. | Resolve placement from the captured source context, retain a compatibility policy for existing configs, preserve line breaks for structured replies, and verify both top-level and in-topic sends with provider readback. |
| `lark_configured_chat_addressing_projection_gap` | `configured_chat_all` reports `reply_due` for a message addressed to another member, Bot-authored prose, or a historical invitation merely because the body also contains the product name; genuine structured Bot mentions may meanwhile lose their addressing evidence after persistence. | Provider-native `mentions`/`mentioned` fields, sender/reply verification summary, canonical event `addressed_to_bot`, content-free urgency counts, and source message settlement. | Provider mention evidence was discarded at the compact inbox boundary, then urgency reconstructed intent from a broad text heuristic such as any `@` plus a product-name token. | Normalize provider-native addressing once into a compact typed flag, persist that flag without raw identities, and make configured-chat urgency accept only the flag or a verified Bot reply. Preserve structured negative mention evidence during normalization, require message-context lookup when the event stream lacks it, and make legacy events without the flag fail closed to material review. Cover another-member mentions, Bot-name prose/self messages, verified direct mentions, and legacy events with public-safe fixtures. |
+| `lark_card_v2_readback_shape_gap` | A Card 2.0 message is visibly delivered and its click reaches the callback consumer, yet delivery remains unverified or the callback fails with a generic action conflict and the card shows no result. | Submitted Card 2.0 payload digest, exact message/route/app readback, CLI-normalized message text, callback `card_content` shape, action digest, authorization checks, and operation lifecycle state. | Delivery compared a rich Card 2.0 payload with the CLI's normalized text, while callback validation compared the same payload digest with Lark's provider-normalized `user_card_content` or its legacy title-and-placeholder fallback. Distinct provider projections of one card were treated as one byte-identical representation. | Give each provider projection a strict semantic verifier: match the complete normalized visible card at delivery and callback, retain the submitted-card digest in the canonical receipt, and accept the bounded legacy fallback only after exact message, route, app, action digest, stored-card digest, and authorized-principal checks. Preserve fail-closed projection-drift tests and expose a stable conflict reason so operators are not asked to click blindly. |
| `repository_delivery_gate_projection_gap` | Git rejects commit or push through an effective global or repository-local guard while quota still exposes only broad `delivery_allowed=true`, so preparation eligibility is mistaken for repository-delivery admission. | Path-free `change-window status` diagnostic, provider verification checks, typed policy decision, interaction-contract repository delivery gate, typed hook registration/result, linked-worktree and separate-clone readback. | Provider discovery exposed only repository-local install state, while Kernel interaction state had no trusted, provider-neutral seam for a capability-owned commit/push decision. | Detect configured external guard surfaces without returning paths or inferring policy; recognize only a bounded legacy signature for preview-first layering. Register a bounded read-only interaction projection hook at the composition root; keep typed validation, slot conflicts, and failure isolation core-owned. Project commit/push admission only from a fully verified typed repository provider, keep preparation/validation distinct, carry `next_eligible_at`, and never grant remote-write authority. |
| `dashboard_verified_mutation_projection_gap` | A preview-locked dashboard mutation reports a successful shared-state readback, but the initiating control still shows its old value; a second click then says the requested setting already exists. | Exact apply receipt, no-change canonical preview, shared-state readback verification, status projection generation/revision, rendered control state, and refresh outcome. | The data adapter verified the canonical write or no-change state, then the UI discarded that receipt and rebound immediately to a separate stale status projection. | Return the verified configuration through apply and no-change preview callbacks and use it as a drawer-scoped read model for the same Goal; refresh the normal projection independently, surface refresh failure without undoing the verified result, and clear the override when the drawer selection changes. Keep preview state visibly pending rather than presenting it as applied. Cover a deliberately stale status response in the browser smoke. |
diff --git a/tests/extensions/test_lark_event_collector_runtime.py b/tests/extensions/test_lark_event_collector_runtime.py
index 29840ccef6..cfca390f20 100644
--- a/tests/extensions/test_lark_event_collector_runtime.py
+++ b/tests/extensions/test_lark_event_collector_runtime.py
@@ -313,6 +313,7 @@ def runner(argv: list[str], **_kwargs: object) -> subprocess.CompletedProcess[st
assert status["callback_delivery_verified"] is True
assert status["listener_ready"] is False
assert status["failed_callback_count"] == 1
+ assert status["last_failure_code"] == "result_delivery_unverified"
assert status["listener_active"] is False
diff --git a/tests/extensions/test_lark_goal_channel_operation.py b/tests/extensions/test_lark_goal_channel_operation.py
index e4eff2e437..69f1f7dfbf 100644
--- a/tests/extensions/test_lark_goal_channel_operation.py
+++ b/tests/extensions/test_lark_goal_channel_operation.py
@@ -172,7 +172,36 @@ def _prepare(store: ChatActionStore, registry_path: Path) -> dict[str, Any]:
)
-def _runner(calls: list[list[str]], sent_cards: dict[str, dict[str, Any]]):
+def _normalized_card_v2(card: Mapping[str, Any]) -> str:
+ header = card["header"]
+ lines = [
+ (
+ f''
+ )
+ ]
+ lines.extend(f"「{item['text']['content']}」" for item in header["text_tag_list"])
+ for element in card["body"]["elements"]:
+ columns = element["columns"]
+ buttons = [
+ column["elements"][0]["text"]["content"]
+ for column in columns
+ if column["elements"][0]["tag"] == "button"
+ ]
+ if buttons:
+ lines.append(" ".join(f"[{label}]" for label in buttons))
+ continue
+ lines.extend(column["elements"][0]["content"] for column in columns)
+ lines.append("")
+ return "\n".join(lines)
+
+
+def _runner(
+ calls: list[list[str]],
+ sent_cards: dict[str, dict[str, Any]],
+ *,
+ normalized_card_v2_readback: bool = False,
+):
def run(
args: list[str], _cwd: Path | None, _timeout: float | None
) -> dict[str, Any]:
@@ -219,7 +248,11 @@ def run(
"chat_id": CHAT_ID,
"sender": {"sender_type": "app", "id": APP_ID},
"deleted": False,
- "body": {"content": json.dumps(card)},
+ **(
+ {"content": _normalized_card_v2(card)}
+ if normalized_card_v2_readback
+ else {"body": {"content": json.dumps(card)}}
+ ),
}
for message_id, card in sent_cards.items()
],
@@ -238,7 +271,15 @@ def run(
"message_id": message_id,
"chat_id": CHAT_ID,
"sender": {"sender_type": "app", "id": APP_ID},
- "body": {"content": json.dumps(sent_cards[message_id])},
+ **(
+ {"content": _normalized_card_v2(sent_cards[message_id])}
+ if normalized_card_v2_readback
+ else {
+ "body": {
+ "content": json.dumps(sent_cards[message_id])
+ }
+ }
+ ),
}
]
},
@@ -281,6 +322,55 @@ def _event(proposal: dict[str, Any], card: dict[str, Any]) -> dict[str, Any]:
}
+def _lark_card_v2_callback_fallback(card: Mapping[str, Any]) -> dict[str, Any]:
+ header = card["header"]
+ return {
+ "title": (f"{header['title']['content']}\n{header['subtitle']['content']}"),
+ "elements": [
+ [
+ {"tag": "img", "image_key": "img_v3_public_fixture"},
+ {"tag": "text", "text": "Upgrade the client to view this card"},
+ {"tag": "text", "text": ""},
+ ]
+ ],
+ }
+
+
+def _lark_card_v2_user_content(card: Mapping[str, Any]) -> dict[str, Any]:
+ """Approximate the Card 2.0 projection returned by user_card_content."""
+
+ normalized = json.loads(json.dumps(card))
+ normalized["config"] = {
+ "enable_forward_interaction": False,
+ "streaming_mode": False,
+ "width_mode": "default",
+ }
+ normalized["header"].pop("icon")
+ sequence = 0
+
+ def visit(value: object) -> None:
+ nonlocal sequence
+ if isinstance(value, list):
+ for item in value:
+ visit(item)
+ return
+ if not isinstance(value, dict):
+ return
+ if value.get("tag") in {"column_set", "column", "markdown", "button"}:
+ sequence += 1
+ value["element_id"] = f"provider_element_{sequence}"
+ if value.get("tag") == "column_set":
+ value["horizontal_align"] = "left"
+ if value.get("tag") == "button":
+ value.pop("behaviors", None)
+ value.pop("confirm", None)
+ for child in value.values():
+ visit(child)
+
+ visit(normalized["body"])
+ return normalized
+
+
def test_card_is_one_bounded_non_forwardable_confirmation_projection(
tmp_path: Path,
) -> None:
@@ -359,8 +449,9 @@ def compile_frame(method: str, params: Mapping[str, Any]) -> dict[str, Any]:
for method, params in calls
)
assert confirmation["header"]["title"]["content"] == "Simulated trade request"
- assert "confirm" not in (
- confirmation["body"]["elements"][3]["columns"][0]["elements"][0]
+ assert (
+ "confirm"
+ not in (confirmation["body"]["elements"][3]["columns"][0]["elements"][0])
)
assert result["header"]["text_tag_list"][0]["text"]["content"] == "模拟完成"
# Lark Card 2.0 rejects `corner_radius` on a column with error 200621,
@@ -621,6 +712,205 @@ def executor(claimed: dict[str, Any]) -> dict[str, Any]:
)
+def test_card_v2_normalized_readback_and_callback_fallback_complete_simulation(
+ tmp_path: Path,
+) -> None:
+ store, registry, runtime, binding, target = _fixture(tmp_path)
+ proposal = _prepare(store, registry)
+ calls: list[list[str]] = []
+ sent_cards: dict[str, dict[str, Any]] = {}
+ runner = _runner(calls, sent_cards, normalized_card_v2_readback=True)
+
+ delivered = deliver_goal_channel_operation_card(
+ proposal_id=proposal["proposal_id"],
+ action_store_root=store.root,
+ runtime_root=runtime,
+ binding_path=binding,
+ target_path=target,
+ execute=True,
+ runner=runner,
+ executor_binding_resolver=lambda _parameters, _runtime: {
+ "revision": "simulator-v0"
+ },
+ )
+
+ assert delivered["status"] == "awaiting_confirmation"
+ durable = store.load(proposal["proposal_id"])
+ assert durable is not None
+ message_id = durable["operation"]["delivery"]["message_id"]
+ card = sent_cards[message_id]
+ event = {
+ **_event(durable, card),
+ "card_content": json.dumps(_lark_card_v2_callback_fallback(card)),
+ }
+ execution_count = 0
+
+ def executor(claimed: dict[str, Any]) -> dict[str, Any]:
+ nonlocal execution_count
+ execution_count += 1
+ operation = claimed["operation"]
+ return {
+ "schema_version": "loopx_operation_outcome_v0",
+ "outcome": "simulated_filled",
+ "projection_verified": True,
+ "operation_id": operation["operation_id"],
+ "payload_digest": operation["payload_digest"],
+ "claim_id": operation["claim"]["claim_id"],
+ "executor_revision": operation["executor_revision"],
+ "summary": "Simulation completed without an external venue write.",
+ "simulation": True,
+ "external_write_performed": False,
+ "observed_at": datetime.now(timezone.utc).isoformat(),
+ }
+
+ receipt = handle_goal_channel_operation_callback(
+ event,
+ runtime_root=runtime,
+ action_store_root=store.root,
+ profile_app_id=APP_ID,
+ cli_bin="lark-cli",
+ profile="operation-bot",
+ runner=runner,
+ executor=executor,
+ )
+ replay = handle_goal_channel_operation_callback(
+ event,
+ runtime_root=runtime,
+ action_store_root=store.root,
+ profile_app_id=APP_ID,
+ cli_bin="lark-cli",
+ profile="operation-bot",
+ runner=runner,
+ executor=executor,
+ )
+
+ assert receipt["outcome"] == "simulated_filled"
+ assert replay["outcome"] == "simulated_filled"
+ assert replay["external_write_performed"] is False
+ assert execution_count == 1
+ assert receipt["card_update_verified"] is True
+ assert (
+ store.load(proposal["proposal_id"])["operation"]["result_delivery"]["transport"]
+ == "callback_update"
+ )
+
+
+def test_provider_normalized_card_v2_callback_and_replay_complete_simulation(
+ tmp_path: Path,
+) -> None:
+ store, registry, runtime, binding, target = _fixture(tmp_path)
+ proposal = _prepare(store, registry)
+ calls: list[list[str]] = []
+ sent_cards: dict[str, dict[str, Any]] = {}
+ runner = _runner(calls, sent_cards, normalized_card_v2_readback=True)
+
+ delivered = deliver_goal_channel_operation_card(
+ proposal_id=proposal["proposal_id"],
+ action_store_root=store.root,
+ runtime_root=runtime,
+ binding_path=binding,
+ target_path=target,
+ execute=True,
+ runner=runner,
+ executor_binding_resolver=lambda _parameters, _runtime: {
+ "revision": "simulator-v0"
+ },
+ )
+ assert delivered["status"] == "awaiting_confirmation"
+ durable = store.load(proposal["proposal_id"])
+ assert durable is not None
+ card = sent_cards[durable["operation"]["delivery"]["message_id"]]
+ event = {
+ **_event(durable, card),
+ "card_content": json.dumps(_lark_card_v2_user_content(card)),
+ }
+ execution_count = 0
+
+ def executor(claimed: dict[str, Any]) -> dict[str, Any]:
+ nonlocal execution_count
+ execution_count += 1
+ operation = claimed["operation"]
+ return {
+ "schema_version": "loopx_operation_outcome_v0",
+ "outcome": "simulated_filled",
+ "projection_verified": True,
+ "operation_id": operation["operation_id"],
+ "payload_digest": operation["payload_digest"],
+ "claim_id": operation["claim"]["claim_id"],
+ "executor_revision": operation["executor_revision"],
+ "summary": "Simulation completed without an external venue write.",
+ "simulation": True,
+ "external_write_performed": False,
+ "observed_at": datetime.now(timezone.utc).isoformat(),
+ }
+
+ first = handle_goal_channel_operation_callback(
+ event,
+ runtime_root=runtime,
+ action_store_root=store.root,
+ profile_app_id=APP_ID,
+ cli_bin="lark-cli",
+ profile="operation-bot",
+ runner=runner,
+ executor=executor,
+ )
+ replay = handle_goal_channel_operation_callback(
+ event,
+ runtime_root=runtime,
+ action_store_root=store.root,
+ profile_app_id=APP_ID,
+ cli_bin="lark-cli",
+ profile="operation-bot",
+ runner=runner,
+ executor=executor,
+ )
+
+ assert first["outcome"] == replay["outcome"] == "simulated_filled"
+ assert first["card_update_verified"] is True
+ assert replay["external_write_performed"] is False
+ assert execution_count == 1
+
+
+def test_card_v2_callback_fallback_rejects_projection_drift(tmp_path: Path) -> None:
+ store, registry, runtime, binding, target = _fixture(tmp_path)
+ proposal = _prepare(store, registry)
+ calls: list[list[str]] = []
+ sent_cards: dict[str, dict[str, Any]] = {}
+ runner = _runner(calls, sent_cards, normalized_card_v2_readback=True)
+ deliver_goal_channel_operation_card(
+ proposal_id=proposal["proposal_id"],
+ action_store_root=store.root,
+ runtime_root=runtime,
+ binding_path=binding,
+ target_path=target,
+ execute=True,
+ runner=runner,
+ executor_binding_resolver=lambda _parameters, _runtime: {
+ "revision": "simulator-v0"
+ },
+ )
+ durable = store.load(proposal["proposal_id"])
+ assert durable is not None
+ card = sent_cards[durable["operation"]["delivery"]["message_id"]]
+ fallback = _lark_card_v2_callback_fallback(card)
+ fallback["title"] = "A different operation\nA different request"
+
+ with pytest.raises(ActionConflictError, match="card content drifted"):
+ handle_goal_channel_operation_callback(
+ {
+ **_event(durable, card),
+ "card_content": json.dumps(fallback),
+ },
+ runtime_root=runtime,
+ action_store_root=store.root,
+ profile_app_id=APP_ID,
+ cli_bin="lark-cli",
+ profile="operation-bot",
+ runner=runner,
+ executor=lambda _proposal: {},
+ )
+
+
def test_concurrent_callback_replay_dispatches_the_claim_once(tmp_path: Path) -> None:
store, registry, runtime, binding, target = _fixture(tmp_path)
proposal = _prepare(store, registry)
From 9b1e64a113a0654e52448b1e966aad7d6cc50d99 Mon Sep 17 00:00:00 2001
From: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
Date: Mon, 14 Sep 2026 17:28:07 +0800
Subject: [PATCH 2/6] fix(lark): recover callback card hydration
Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
---
.../lark/event_collector_runtime.py | 20 ++-
.../lark/goal_channel_message_delivery.py | 9 +-
.../extensions/lark/goal_channel_operation.py | 126 +++++++++++++++---
.../test_lark_event_collector_runtime.py | 3 +-
.../test_lark_goal_channel_operation.py | 43 +++++-
5 files changed, 172 insertions(+), 29 deletions(-)
diff --git a/loopx/extensions/lark/event_collector_runtime.py b/loopx/extensions/lark/event_collector_runtime.py
index bfa9dc7453..3984ca0269 100644
--- a/loopx/extensions/lark/event_collector_runtime.py
+++ b/loopx/extensions/lark/event_collector_runtime.py
@@ -32,6 +32,22 @@
EVENT_READY_PREFIX = "[event] ready "
EVENT_DIAGNOSTIC_PREFIX = "[event] "
_CALLBACK_FAILURE_CODES = {
+ "collector Bot application identity is unverified": "collector_app_identity_unverified",
+ "operation callback event type is unsupported": "callback_event_type_unsupported",
+ "operation callback must come from a button": "callback_action_not_button",
+ "operation callback action_value is invalid": "callback_action_value_invalid",
+ "operation callback action is incomplete": "callback_action_incomplete",
+ "operation callback action schema is unsupported": "callback_action_schema_unsupported",
+ "operation callback decision is unsupported": "callback_decision_unsupported",
+ "operation callback update token is invalid": "callback_update_token_invalid",
+ "operation callback event_id is invalid": "callback_event_id_invalid",
+ "operation callback message_id is invalid": "callback_message_id_invalid",
+ "operation callback chat_id is invalid": "callback_chat_id_invalid",
+ "operation callback operator_id is invalid": "callback_operator_id_invalid",
+ "operation callback host is unsupported": "callback_host_unsupported",
+ "operation callback card content is unavailable": "callback_card_content_unavailable",
+ "operation callback proposal was not found": "callback_proposal_not_found",
+ "typed operation proposal is unavailable": "callback_operation_unavailable",
"operation card delivery was not recorded": "delivery_not_recorded",
"operation callback digest drifted": "confirmation_digest_drifted",
"recorded operation card digest drifted": "recorded_card_digest_drifted",
@@ -993,9 +1009,7 @@ def forward_signal(signum: int, _: object) -> None:
result.update(
{
"operation_callback_listener_started": True,
- "operation_callback_listener_ready": bool(
- callback_stats.get("ready")
- ),
+ "operation_callback_listener_ready": bool(callback_stats.get("ready")),
"operation_callback_received_count": callback_stats["received"],
"operation_callback_verified_count": callback_stats["verified"],
"operation_callback_failure_count": callback_stats["failed"],
diff --git a/loopx/extensions/lark/goal_channel_message_delivery.py b/loopx/extensions/lark/goal_channel_message_delivery.py
index 1b6906bf81..b5e55f3437 100644
--- a/loopx/extensions/lark/goal_channel_message_delivery.py
+++ b/loopx/extensions/lark/goal_channel_message_delivery.py
@@ -129,7 +129,7 @@ def _message_card(value: Mapping[str, Any]) -> Mapping[str, Any] | None:
return parsed if isinstance(parsed, Mapping) else None
-def _normalized_card_text(card: Mapping[str, Any]) -> str | None:
+def normalized_card_text(card: Mapping[str, Any]) -> str | None:
if card.get("schema") == "2.0":
return _normalized_card_v2_text(card)
header = card.get("header")
@@ -249,8 +249,8 @@ def card_projection_matches(
if observed == expected:
return True
- observed_text = _normalized_card_text(observed)
- expected_text = _normalized_card_text(expected)
+ observed_text = normalized_card_text(observed)
+ expected_text = normalized_card_text(expected)
return (
observed_text is not None
and expected_text is not None
@@ -267,7 +267,7 @@ def message_card_matches(
if isinstance(observed, Mapping) and card_projection_matches(observed, expected):
return True
content = value.get("content")
- return isinstance(content, str) and content == _normalized_card_text(expected)
+ return isinstance(content, str) and content == normalized_card_text(expected)
def _message_sender(value: Mapping[str, Any]) -> tuple[str, str]:
@@ -518,6 +518,7 @@ def readback(self, message_id: str) -> Mapping[str, Any]:
__all__ = [
"GoalChannelMessageDeliverySession",
+ "normalized_card_text",
"goal_channel_delivery_route",
"resolve_bound_goal_channel",
]
diff --git a/loopx/extensions/lark/goal_channel_operation.py b/loopx/extensions/lark/goal_channel_operation.py
index f8db171d0a..c93d3bb587 100644
--- a/loopx/extensions/lark/goal_channel_operation.py
+++ b/loopx/extensions/lark/goal_channel_operation.py
@@ -27,6 +27,7 @@
GoalChannelMessageDeliverySession,
card_projection_matches,
message_card_matches,
+ normalized_card_text,
resolve_bound_goal_channel,
)
from .goal_channel_transport import call, json_payload, lark_args
@@ -629,6 +630,91 @@ def collect(value: object) -> bool:
return collect(elements) and any(item.get("tag") == "img" for item in leaves)
+def _callback_card_content_matches(value: object, expected: Mapping[str, Any]) -> bool:
+ """Verify either provider JSON or the documented userDSL callback shape."""
+
+ observed: object = value
+ if isinstance(value, str):
+ if not value:
+ return False
+ try:
+ observed = json.loads(value)
+ except json.JSONDecodeError:
+ return value == normalized_card_text(expected)
+ return bool(
+ isinstance(observed, Mapping)
+ and (
+ card_projection_matches(observed, expected)
+ or _lark_card_v2_fallback_matches(observed, expected)
+ )
+ )
+
+
+def _read_callback_card_content(
+ *,
+ runner: CommandRunner,
+ cli_bin: str,
+ profile: str,
+ message_id: str,
+ chat_id: str,
+ app_id: str,
+) -> object:
+ """Retry the CLI's best-effort callback hydration through exact readback."""
+
+ result = call(
+ runner,
+ lark_args(
+ cli_bin=cli_bin,
+ profile=profile,
+ tail=[
+ "api",
+ "GET",
+ f"/open-apis/im/v1/messages/{message_id}",
+ "--params",
+ json.dumps({"card_msg_content_type": "user_card_content"}),
+ "--as",
+ "bot",
+ ],
+ ),
+ )
+ if result.get("returncode") != 0:
+ return None
+ message = _find_message(json_payload(result), message_id)
+ sender = message.get("sender") if isinstance(message, Mapping) else None
+ if (
+ not isinstance(message, Mapping)
+ or str(message.get("chat_id") or "") != chat_id
+ or not isinstance(sender, Mapping)
+ or sender.get("sender_type") != "app"
+ or sender.get("id") != app_id
+ ):
+ return None
+ body = message.get("body") if isinstance(message, Mapping) else None
+ return body.get("content") if isinstance(body, Mapping) else None
+
+
+def _callback_replays_confirmation(
+ *,
+ confirmation: object,
+ action: Mapping[str, str],
+ event: Mapping[str, Any],
+ operator_principal: str,
+ profile_app_id: str,
+) -> bool:
+ if not isinstance(confirmation, Mapping):
+ return False
+ expected = {
+ "event_id": str(event["event_id"]),
+ "principal": operator_principal,
+ "message_id": str(event["message_id"]),
+ "chat_id": str(event["chat_id"]),
+ "app_id": profile_app_id,
+ "confirmation_digest": action["confirmation_digest"],
+ "decision": action["decision"],
+ }
+ return all(confirmation.get(key) == value for key, value in expected.items())
+
+
def _operator_membership_verified(
*,
runner: CommandRunner,
@@ -1128,13 +1214,6 @@ def handle_goal_channel_operation_callback(
raise ValueError(f"operation callback {field} is invalid")
if str(event.get("host") or "") != "im_message":
raise ValueError("operation callback host is unsupported")
- card_content = event.get("card_content")
- try:
- card = json.loads(card_content) if isinstance(card_content, str) else None
- except json.JSONDecodeError as exc:
- raise ValueError("operation callback card_content is invalid") from exc
- if not isinstance(card, Mapping):
- raise ValueError("operation callback requires exact card_content")
store = ChatActionStore(action_store_root)
proposal = store.load(action["operation_id"])
if proposal is None:
@@ -1145,15 +1224,6 @@ def handle_goal_channel_operation_callback(
raise ActionConflictError("operation card delivery was not recorded")
if action["confirmation_digest"] != operation.get("confirmation_digest"):
raise ActionConflictError("operation callback digest drifted")
- if _digest(card) != delivery.get("card_digest"):
- expected_card = _submitted_confirmation_card(proposal)
- if _digest(expected_card) != delivery.get("card_digest"):
- raise ActionConflictError("recorded operation card digest drifted")
- if not (
- card_projection_matches(card, expected_card)
- or _lark_card_v2_fallback_matches(card, expected_card)
- ):
- raise ActionConflictError("operation callback card content drifted")
if profile_app_id != delivery.get("app_id"):
raise ActionConflictError("operation callback app identity drifted")
operator_id = str(event["operator_id"])
@@ -1161,6 +1231,30 @@ def handle_goal_channel_operation_callback(
chat_id = str(event["chat_id"])
if operator_principal not in set(parameters.get("authorized_principals") or []):
raise ActionConflictError("principal is not authorized for this operation")
+ if not _callback_replays_confirmation(
+ confirmation=operation.get("confirmation"),
+ action=action,
+ event=event,
+ operator_principal=operator_principal,
+ profile_app_id=profile_app_id,
+ ):
+ expected_card = _submitted_confirmation_card(proposal)
+ if _digest(expected_card) != delivery.get("card_digest"):
+ raise ActionConflictError("recorded operation card digest drifted")
+ card_content = event.get("card_content")
+ if card_content is None or card_content == "":
+ card_content = _read_callback_card_content(
+ runner=runner,
+ cli_bin=cli_bin,
+ profile=profile,
+ message_id=str(event["message_id"]),
+ chat_id=chat_id,
+ app_id=profile_app_id,
+ )
+ if card_content is None or card_content == "":
+ raise ValueError("operation callback card content is unavailable")
+ if not _callback_card_content_matches(card_content, expected_card):
+ raise ActionConflictError("operation callback card content drifted")
if not _operator_membership_verified(
runner=runner,
cli_bin=cli_bin,
diff --git a/tests/extensions/test_lark_event_collector_runtime.py b/tests/extensions/test_lark_event_collector_runtime.py
index cfca390f20..11a79fad9b 100644
--- a/tests/extensions/test_lark_event_collector_runtime.py
+++ b/tests/extensions/test_lark_event_collector_runtime.py
@@ -223,8 +223,7 @@ def runner(*_args: object, **_kwargs: object) -> subprocess.CompletedProcess[str
assert status["operation_callback_listener_active"] is True
assert status["operation_callback_listener_ready"] is True
assert (
- status["operation_callback_qualification_state"]
- == "listener_ready_unqualified"
+ status["operation_callback_qualification_state"] == "listener_ready_unqualified"
)
payload = json.loads(callback_status.read_text(encoding="utf-8"))
diff --git a/tests/extensions/test_lark_goal_channel_operation.py b/tests/extensions/test_lark_goal_channel_operation.py
index 69f1f7dfbf..1f7ae34e97 100644
--- a/tests/extensions/test_lark_goal_channel_operation.py
+++ b/tests/extensions/test_lark_goal_channel_operation.py
@@ -284,6 +284,32 @@ def run(
]
},
}
+ elif (
+ "api" in args
+ and "GET" in args
+ and any("/open-apis/im/v1/messages/" in item for item in args)
+ ):
+ endpoint = next(
+ item for item in args if "/open-apis/im/v1/messages/" in item
+ )
+ message_id = endpoint.rsplit("/", 1)[-1]
+ payload = {
+ "ok": True,
+ "data": {
+ "items": [
+ {
+ "message_id": message_id,
+ "chat_id": CHAT_ID,
+ "sender": {"sender_type": "app", "id": APP_ID},
+ "body": {
+ "content": json.dumps(
+ _lark_card_v2_user_content(sent_cards[message_id])
+ )
+ },
+ }
+ ]
+ },
+ }
elif "messages" in args and "patch" in args:
message_id = args[args.index("--message-id") + 1]
update = json.loads(args[args.index("--data") + 1])
@@ -795,8 +821,10 @@ def executor(claimed: dict[str, Any]) -> dict[str, Any]:
)
+@pytest.mark.parametrize("callback_shape", ["provider_json", "userdsl", "empty"])
def test_provider_normalized_card_v2_callback_and_replay_complete_simulation(
tmp_path: Path,
+ callback_shape: str,
) -> None:
store, registry, runtime, binding, target = _fixture(tmp_path)
proposal = _prepare(store, registry)
@@ -820,10 +848,12 @@ def test_provider_normalized_card_v2_callback_and_replay_complete_simulation(
durable = store.load(proposal["proposal_id"])
assert durable is not None
card = sent_cards[durable["operation"]["delivery"]["message_id"]]
- event = {
- **_event(durable, card),
- "card_content": json.dumps(_lark_card_v2_user_content(card)),
- }
+ callback_content = {
+ "provider_json": json.dumps(_lark_card_v2_user_content(card)),
+ "userdsl": _normalized_card_v2(card),
+ "empty": "",
+ }[callback_shape]
+ event = {**_event(durable, card), "card_content": callback_content}
execution_count = 0
def executor(claimed: dict[str, Any]) -> dict[str, Any]:
@@ -869,6 +899,11 @@ def executor(claimed: dict[str, Any]) -> dict[str, Any]:
assert first["card_update_verified"] is True
assert replay["external_write_performed"] is False
assert execution_count == 1
+ if callback_shape == "empty":
+ assert any(
+ "GET" in call and any("/open-apis/im/v1/messages/" in item for item in call)
+ for call in calls
+ )
def test_card_v2_callback_fallback_rejects_projection_drift(tmp_path: Path) -> None:
From f0fefb1628e7ae6ebaa7a6e7a15231919cbca372 Mon Sep 17 00:00:00 2001
From: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
Date: Mon, 14 Sep 2026 17:36:38 +0800
Subject: [PATCH 3/6] fix(lark): retain safe callback failure shape
Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
---
.../lark/event_collector_runtime.py | 68 +++++++++++++++++++
1 file changed, 68 insertions(+)
diff --git a/loopx/extensions/lark/event_collector_runtime.py b/loopx/extensions/lark/event_collector_runtime.py
index 3984ca0269..0b89bf0236 100644
--- a/loopx/extensions/lark/event_collector_runtime.py
+++ b/loopx/extensions/lark/event_collector_runtime.py
@@ -1,6 +1,7 @@
from __future__ import annotations
import json
+import hashlib
import re
import signal
import subprocess
@@ -48,6 +49,17 @@
"operation callback card content is unavailable": "callback_card_content_unavailable",
"operation callback proposal was not found": "callback_proposal_not_found",
"typed operation proposal is unavailable": "callback_operation_unavailable",
+ "typed operation envelope is unavailable": "callback_operation_envelope_unavailable",
+ "operation review plan is unavailable": "callback_review_plan_unavailable",
+ "operation review frame is unavailable": "callback_review_frame_unavailable",
+ "operation projection is unavailable": "callback_projection_unavailable",
+ "operation projection fields are unavailable": "callback_projection_fields_unavailable",
+ "operation callback timestamp is invalid": "callback_timestamp_invalid",
+ "operation confirmation has unsupported or missing fields": "callback_confirmation_invalid",
+ "operation timestamps require a timezone": "callback_timestamp_timezone_missing",
+ "claimed operation disappeared before dispatch": "callback_claim_disappeared",
+ "operation executor outcome does not match the consumed claim": "callback_executor_outcome_invalid",
+ "operation disappeared before result delivery": "callback_operation_disappeared",
"operation card delivery was not recorded": "delivery_not_recorded",
"operation callback digest drifted": "confirmation_digest_drifted",
"recorded operation card digest drifted": "recorded_card_digest_drifted",
@@ -68,6 +80,55 @@ def _operation_callback_failure_code(exc: BaseException) -> str:
return _CALLBACK_FAILURE_CODES.get(str(exc), "callback_rejected")
+def _callback_event_shape(payload: Mapping[str, Any]) -> dict[str, object]:
+ """Return a value-free diagnostic projection for a rejected callback."""
+
+ action_value = payload.get("action_value")
+ try:
+ action = (
+ json.loads(action_value) if isinstance(action_value, str) else action_value
+ )
+ except json.JSONDecodeError:
+ action = None
+ card_content = payload.get("card_content")
+ card_shape = "missing"
+ if isinstance(card_content, str):
+ if not card_content:
+ card_shape = "empty"
+ else:
+ try:
+ parsed_card = json.loads(card_content)
+ except json.JSONDecodeError:
+ card_shape = "text"
+ else:
+ card_shape = (
+ "json_object" if isinstance(parsed_card, Mapping) else "json_other"
+ )
+ elif isinstance(card_content, Mapping):
+ card_shape = "object"
+ elif card_content is not None:
+ card_shape = type(card_content).__name__
+ return {
+ "type_supported": payload.get("type") == "card.action.trigger",
+ "action_is_button": payload.get("action_tag") == "button",
+ "action_is_object": isinstance(action, Mapping),
+ "action_field_count": len(action) if isinstance(action, Mapping) else 0,
+ "event_id_valid": bool(
+ re.fullmatch(r"[A-Za-z0-9._:-]{1,240}", str(payload.get("event_id") or ""))
+ ),
+ "timestamp_is_digits": str(payload.get("timestamp") or "").isdigit(),
+ "operator_id_present": bool(payload.get("operator_id")),
+ "message_id_present": bool(payload.get("message_id")),
+ "chat_id_present": bool(payload.get("chat_id")),
+ "host_supported": payload.get("host") == "im_message",
+ "token_present": bool(payload.get("token")),
+ "card_content_shape": card_shape,
+ "shape_digest": hashlib.sha256(
+ "\0".join(sorted(str(key) for key in payload)).encode()
+ ).hexdigest()[:16],
+ }
+
+
def _run_json(
runner: CommandRunner,
argv: Sequence[str],
@@ -426,6 +487,7 @@ def _write_operation_callback_status(
callback_delivery_verified: bool | None = None,
failure_kind: str | None = None,
failure_code: str | None = None,
+ failure_event_shape: Mapping[str, object] | None = None,
consumer_returncode: int | None = None,
recovered_result_count_delta: int = 0,
result_delivery_failure_count_delta: int = 0,
@@ -479,6 +541,11 @@ def _write_operation_callback_status(
),
"last_failure_kind": failure_kind or prior.get("last_failure_kind"),
"last_failure_code": failure_code or prior.get("last_failure_code"),
+ "last_failure_event_shape": (
+ dict(failure_event_shape)
+ if failure_event_shape is not None
+ else prior.get("last_failure_event_shape")
+ ),
"consumer_returncode": consumer_returncode,
"updated_at": now,
"private_content_returned": False,
@@ -758,6 +825,7 @@ def consume_operation_callbacks() -> None:
listener_ready=True,
failure_kind=type(exc).__name__,
failure_code=_operation_callback_failure_code(exc),
+ failure_event_shape=_callback_event_shape(payload),
)
continue
callback_stats["verified"] += 1
From 8a3ce010c0f9976d99a533d61eb6bd830ac67edc Mon Sep 17 00:00:00 2001
From: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
Date: Mon, 14 Sep 2026 17:43:50 +0800
Subject: [PATCH 4/6] fix(lark): bind callbacks to sent card snapshot
Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
---
loopx/chat_action_store.py | 43 +++++++++++++++++++
.../extensions/lark/goal_channel_operation.py | 11 ++++-
.../test_lark_goal_channel_operation.py | 11 +++++
3 files changed, 64 insertions(+), 1 deletion(-)
diff --git a/loopx/chat_action_store.py b/loopx/chat_action_store.py
index 2968ce8f15..a9a71a38b2 100644
--- a/loopx/chat_action_store.py
+++ b/loopx/chat_action_store.py
@@ -411,6 +411,49 @@ def record_operation_delivery(
self._write(payload)
return proposal
+ def record_operation_delivery_snapshot(
+ self,
+ proposal_id: str,
+ *,
+ submitted_card: Mapping[str, Any],
+ ) -> dict[str, Any]:
+ """Persist the immutable sent card after verifying its recorded digest."""
+
+ safe_card = _safe_json_value(
+ dict(submitted_card), path="operation.delivery.submitted_card"
+ )
+ if not isinstance(safe_card, dict):
+ raise ValueError("operation submitted card must be an object")
+ token = _opaque_id(proposal_id, field="proposal_id")
+ with exclusive_file_lock(
+ self.path,
+ agent_id="loopx-chat",
+ operation="record_operation_delivery_snapshot",
+ ):
+ payload = self._read()
+ proposal = payload["proposals"].get(token)
+ operation = (
+ proposal.get("operation") if isinstance(proposal, dict) else None
+ )
+ delivery = (
+ operation.get("delivery") if isinstance(operation, dict) else None
+ )
+ if not isinstance(delivery, dict):
+ raise ActionConflictError("operation card delivery was not recorded")
+ if _canonical_digest(safe_card) != delivery.get("card_digest"):
+ raise ActionConflictError("operation submitted card digest drifted")
+ existing = delivery.get("submitted_card")
+ if existing is not None:
+ if existing != safe_card:
+ raise ActionConflictError(
+ "operation submitted card snapshot already differs"
+ )
+ return proposal
+ delivery["submitted_card"] = safe_card
+ proposal["updated_at"] = _utc_now()
+ self._write(payload)
+ return proposal
+
def decide_operation(
self,
proposal_id: str,
diff --git a/loopx/extensions/lark/goal_channel_operation.py b/loopx/extensions/lark/goal_channel_operation.py
index c93d3bb587..7e723c0d73 100644
--- a/loopx/extensions/lark/goal_channel_operation.py
+++ b/loopx/extensions/lark/goal_channel_operation.py
@@ -520,6 +520,10 @@ def resolve_current() -> Mapping[str, Any]:
"delivered_at": datetime.now(timezone.utc).isoformat(),
},
)
+ store.record_operation_delivery_snapshot(
+ proposal_id,
+ submitted_card=card,
+ )
return operation_packet(
ok=True,
goal_id=goal_id,
@@ -1238,7 +1242,12 @@ def handle_goal_channel_operation_callback(
operator_principal=operator_principal,
profile_app_id=profile_app_id,
):
- expected_card = _submitted_confirmation_card(proposal)
+ submitted_card = delivery.get("submitted_card")
+ expected_card = (
+ dict(submitted_card)
+ if isinstance(submitted_card, Mapping)
+ else _submitted_confirmation_card(proposal)
+ )
if _digest(expected_card) != delivery.get("card_digest"):
raise ActionConflictError("recorded operation card digest drifted")
card_content = event.get("card_content")
diff --git a/tests/extensions/test_lark_goal_channel_operation.py b/tests/extensions/test_lark_goal_channel_operation.py
index 1f7ae34e97..6f574c0c86 100644
--- a/tests/extensions/test_lark_goal_channel_operation.py
+++ b/tests/extensions/test_lark_goal_channel_operation.py
@@ -825,6 +825,7 @@ def executor(claimed: dict[str, Any]) -> dict[str, Any]:
def test_provider_normalized_card_v2_callback_and_replay_complete_simulation(
tmp_path: Path,
callback_shape: str,
+ monkeypatch: pytest.MonkeyPatch,
) -> None:
store, registry, runtime, binding, target = _fixture(tmp_path)
proposal = _prepare(store, registry)
@@ -848,12 +849,22 @@ def test_provider_normalized_card_v2_callback_and_replay_complete_simulation(
durable = store.load(proposal["proposal_id"])
assert durable is not None
card = sent_cards[durable["operation"]["delivery"]["message_id"]]
+ assert durable["operation"]["delivery"]["submitted_card"] == card
callback_content = {
"provider_json": json.dumps(_lark_card_v2_user_content(card)),
"userdsl": _normalized_card_v2(card),
"empty": "",
}[callback_shape]
event = {**_event(durable, card), "card_content": callback_content}
+
+ def fail_rebuild(_proposal: Mapping[str, Any]) -> dict[str, Any]:
+ raise AssertionError("callback must use the immutable delivery snapshot")
+
+ monkeypatch.setattr(
+ goal_channel_operation,
+ "_submitted_confirmation_card",
+ fail_rebuild,
+ )
execution_count = 0
def executor(claimed: dict[str, Any]) -> dict[str, Any]:
From daccc2c161801946811f04294da7c529aa6356b1 Mon Sep 17 00:00:00 2001
From: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
Date: Mon, 14 Sep 2026 18:26:37 +0800
Subject: [PATCH 5/6] fix(lark): retain callback failure stage
Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
---
.../lark/event_collector_runtime.py | 30 ++++++++++++++++++-
.../test_lark_event_collector_runtime.py | 17 +++++++++++
2 files changed, 46 insertions(+), 1 deletion(-)
diff --git a/loopx/extensions/lark/event_collector_runtime.py b/loopx/extensions/lark/event_collector_runtime.py
index 0b89bf0236..a1d2fa740e 100644
--- a/loopx/extensions/lark/event_collector_runtime.py
+++ b/loopx/extensions/lark/event_collector_runtime.py
@@ -72,6 +72,15 @@
"operation confirmation arrived after expiry": "operation_expired",
"operation callback result delivery was not verified": "result_delivery_unverified",
}
+_CALLBACK_FAILURE_STAGES = {
+ "_callback_action": "parse_action",
+ "_callback_timestamp": "validate_timestamp",
+ "_callback_card_content_matches": "verify_card_content",
+ "_operator_membership_verified": "verify_operator_membership",
+ "decide_operation": "claim_operation",
+ "_execute_claimed_operation": "execute_operation",
+ "_update_callback_card": "deliver_result",
+}
CommandRunner = Callable[..., subprocess.CompletedProcess[str]]
Sleeper = Callable[[float], None]
@@ -80,6 +89,20 @@ def _operation_callback_failure_code(exc: BaseException) -> str:
return _CALLBACK_FAILURE_CODES.get(str(exc), "callback_rejected")
+def _operation_callback_failure_stage(exc: BaseException) -> str:
+ """Return a value-free processing stage for a rejected callback."""
+
+ stage = "handle_callback"
+ traceback = exc.__traceback__
+ while traceback is not None:
+ stage = _CALLBACK_FAILURE_STAGES.get(
+ traceback.tb_frame.f_code.co_name,
+ stage,
+ )
+ traceback = traceback.tb_next
+ return stage
+
+
def _callback_event_shape(payload: Mapping[str, Any]) -> dict[str, object]:
"""Return a value-free diagnostic projection for a rejected callback."""
@@ -108,6 +131,7 @@ def _callback_event_shape(payload: Mapping[str, Any]) -> dict[str, object]:
card_shape = "object"
elif card_content is not None:
card_shape = type(card_content).__name__
+ timestamp = str(payload.get("timestamp") or "")
return {
"type_supported": payload.get("type") == "card.action.trigger",
"action_is_button": payload.get("action_tag") == "button",
@@ -116,7 +140,8 @@ def _callback_event_shape(payload: Mapping[str, Any]) -> dict[str, object]:
"event_id_valid": bool(
re.fullmatch(r"[A-Za-z0-9._:-]{1,240}", str(payload.get("event_id") or ""))
),
- "timestamp_is_digits": str(payload.get("timestamp") or "").isdigit(),
+ "timestamp_is_digits": timestamp.isdigit(),
+ "timestamp_digit_count": len(timestamp),
"operator_id_present": bool(payload.get("operator_id")),
"message_id_present": bool(payload.get("message_id")),
"chat_id_present": bool(payload.get("chat_id")),
@@ -487,6 +512,7 @@ def _write_operation_callback_status(
callback_delivery_verified: bool | None = None,
failure_kind: str | None = None,
failure_code: str | None = None,
+ failure_stage: str | None = None,
failure_event_shape: Mapping[str, object] | None = None,
consumer_returncode: int | None = None,
recovered_result_count_delta: int = 0,
@@ -541,6 +567,7 @@ def _write_operation_callback_status(
),
"last_failure_kind": failure_kind or prior.get("last_failure_kind"),
"last_failure_code": failure_code or prior.get("last_failure_code"),
+ "last_failure_stage": failure_stage or prior.get("last_failure_stage"),
"last_failure_event_shape": (
dict(failure_event_shape)
if failure_event_shape is not None
@@ -825,6 +852,7 @@ def consume_operation_callbacks() -> None:
listener_ready=True,
failure_kind=type(exc).__name__,
failure_code=_operation_callback_failure_code(exc),
+ failure_stage=_operation_callback_failure_stage(exc),
failure_event_shape=_callback_event_shape(payload),
)
continue
diff --git a/tests/extensions/test_lark_event_collector_runtime.py b/tests/extensions/test_lark_event_collector_runtime.py
index 11a79fad9b..cbdb2a7821 100644
--- a/tests/extensions/test_lark_event_collector_runtime.py
+++ b/tests/extensions/test_lark_event_collector_runtime.py
@@ -13,6 +13,7 @@
plan_lark_event_collector,
)
from loopx.extensions.lark.event_collector_runtime import (
+ _callback_event_shape,
_run_json_with_status,
enrich_lark_event_reply_context,
lark_event_requires_reply_context_lookup,
@@ -20,6 +21,21 @@
)
+def test_operation_callback_failure_shape_retains_timestamp_width_without_value() -> None:
+ shape = _callback_event_shape(
+ {
+ "type": "card.action.trigger",
+ "timestamp": "1776409469273",
+ "action_tag": "button",
+ "action_value": "{}",
+ }
+ )
+
+ assert shape["timestamp_is_digits"] is True
+ assert shape["timestamp_digit_count"] == 13
+ assert "1776409469273" not in json.dumps(shape)
+
+
def test_reply_context_lookup_does_not_trust_unrelated_text_mentions() -> None:
bot_name = "Context Bot"
@@ -313,6 +329,7 @@ def runner(argv: list[str], **_kwargs: object) -> subprocess.CompletedProcess[st
assert status["listener_ready"] is False
assert status["failed_callback_count"] == 1
assert status["last_failure_code"] == "result_delivery_unverified"
+ assert status["last_failure_stage"] == "handle_callback"
assert status["listener_active"] is False
From a85dd22cba3d4ceb24bab763147a5f759972ae1d Mon Sep 17 00:00:00 2001
From: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
Date: Mon, 14 Sep 2026 18:44:25 +0800
Subject: [PATCH 6/6] fix(lark): keep actionable card identity strict
Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
---
loopx/extensions/lark/event_collector.py | 6 +++
.../lark/goal_channel_message_delivery.py | 28 ++++++++++-
.../test_lark_event_collector_runtime.py | 14 ++++++
.../test_lark_goal_channel_operation.py | 47 +++++++++++++++++++
4 files changed, 93 insertions(+), 2 deletions(-)
diff --git a/loopx/extensions/lark/event_collector.py b/loopx/extensions/lark/event_collector.py
index cc8f733c9d..b7313e48fc 100644
--- a/loopx/extensions/lark/event_collector.py
+++ b/loopx/extensions/lark/event_collector.py
@@ -738,6 +738,12 @@ def inspect_lark_event_collector(
"operation_callback_last_failure_code": callback_status.get(
"last_failure_code"
),
+ "operation_callback_last_failure_stage": callback_status.get(
+ "last_failure_stage"
+ ),
+ "operation_callback_last_failure_event_shape": callback_status.get(
+ "last_failure_event_shape"
+ ),
"operation_callback_console_configuration_preflighted": False,
"thread_complete": all(
route["inbox"]["thread_complete"] for route in config["routes"]
diff --git a/loopx/extensions/lark/goal_channel_message_delivery.py b/loopx/extensions/lark/goal_channel_message_delivery.py
index b5e55f3437..0dc079b951 100644
--- a/loopx/extensions/lark/goal_channel_message_delivery.py
+++ b/loopx/extensions/lark/goal_channel_message_delivery.py
@@ -258,12 +258,29 @@ def card_projection_matches(
)
+def _has_callback_behavior(value: object) -> bool:
+ if isinstance(value, Mapping):
+ if value.get("type") == "callback":
+ return True
+ return any(_has_callback_behavior(child) for child in value.values())
+ if isinstance(value, list):
+ return any(_has_callback_behavior(child) for child in value)
+ return False
+
+
def message_card_matches(
- value: Mapping[str, Any], expected: Mapping[str, Any] | None
+ value: Mapping[str, Any],
+ expected: Mapping[str, Any] | None,
+ *,
+ allow_normalized: bool = True,
) -> bool:
if expected is None:
return False
observed = _message_card(value)
+ if isinstance(observed, Mapping) and observed == expected:
+ return True
+ if not allow_normalized:
+ return False
if isinstance(observed, Mapping) and card_projection_matches(observed, expected):
return True
content = value.get("content")
@@ -344,7 +361,14 @@ def _existing_message(
and str(message.get("chat_id") or "") == route["chat_id"]
and sender_type == "app"
and sender_app_id == route["bot_app_id"]
- and message_card_matches(message, card)
+ and message_card_matches(
+ message,
+ card,
+ # Provider-normalized Card 2.0 history omits callback
+ # values. Visible equality therefore cannot prove that an
+ # old actionable message carries this operation id/digest.
+ allow_normalized=not _has_callback_behavior(card),
+ )
):
return str(message["message_id"])
if not _history_is_complete(payload):
diff --git a/tests/extensions/test_lark_event_collector_runtime.py b/tests/extensions/test_lark_event_collector_runtime.py
index cbdb2a7821..beeda6dae4 100644
--- a/tests/extensions/test_lark_event_collector_runtime.py
+++ b/tests/extensions/test_lark_event_collector_runtime.py
@@ -218,6 +218,12 @@ def test_operation_callback_status_separates_readiness_from_qualification(
"listener_active": True,
"listener_ready": True,
"callback_delivery_verified": False,
+ "last_failure_code": "callback_timestamp_invalid",
+ "last_failure_stage": "validate_timestamp",
+ "last_failure_event_shape": {
+ "timestamp_is_digits": True,
+ "timestamp_digit_count": 13,
+ },
}
),
encoding="utf-8",
@@ -241,6 +247,14 @@ def runner(*_args: object, **_kwargs: object) -> subprocess.CompletedProcess[str
assert (
status["operation_callback_qualification_state"] == "listener_ready_unqualified"
)
+ assert status["operation_callback_last_failure_code"] == (
+ "callback_timestamp_invalid"
+ )
+ assert status["operation_callback_last_failure_stage"] == "validate_timestamp"
+ assert status["operation_callback_last_failure_event_shape"] == {
+ "timestamp_is_digits": True,
+ "timestamp_digit_count": 13,
+ }
payload = json.loads(callback_status.read_text(encoding="utf-8"))
payload["callback_delivery_verified"] = True
diff --git a/tests/extensions/test_lark_goal_channel_operation.py b/tests/extensions/test_lark_goal_channel_operation.py
index 6f574c0c86..466c6a5bc7 100644
--- a/tests/extensions/test_lark_goal_channel_operation.py
+++ b/tests/extensions/test_lark_goal_channel_operation.py
@@ -17,6 +17,10 @@
GOAL_CHANNEL_BINDING_SCHEMA_VERSION,
write_goal_channel_binding,
)
+from loopx.extensions.lark.goal_channel_message_delivery import (
+ message_card_matches,
+ normalized_card_text,
+)
from loopx.extensions.lark.goal_channel_operation import (
build_goal_channel_operation_card,
build_goal_channel_operation_result_card,
@@ -196,6 +200,49 @@ def _normalized_card_v2(card: Mapping[str, Any]) -> str:
return "\n".join(lines)
+def test_normalized_action_card_cannot_prove_historical_action_identity() -> None:
+ card = {
+ "schema": "2.0",
+ "header": {
+ "title": {"content": "Simulated trade request"},
+ "subtitle": {"content": "Synthetic fixture · no venue call"},
+ "text_tag_list": [],
+ },
+ "body": {
+ "elements": [
+ {
+ "tag": "column_set",
+ "columns": [
+ {
+ "tag": "column",
+ "elements": [
+ {
+ "tag": "button",
+ "text": {"content": "Confirm"},
+ "behaviors": [
+ {
+ "type": "callback",
+ "value": {
+ "operation_id": (
+ "proposal-normalized-binding"
+ )
+ },
+ }
+ ],
+ }
+ ],
+ }
+ ],
+ }
+ ]
+ },
+ }
+ normalized = {"content": normalized_card_text(card)}
+
+ assert message_card_matches(normalized, card)
+ assert not message_card_matches(normalized, card, allow_normalized=False)
+
+
def _runner(
calls: list[list[str]],
sent_cards: dict[str, dict[str, Any]],