From 3afccc8c05893272c1f9c86efc2cf516fcccd851 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sun, 4 Oct 2026 14:47:02 +0800 Subject: [PATCH 01/12] feat(lark): enable native private feedback through shared Inbox lifecycle Signed-off-by: huangruiteng --- loopx/chat_server.py | 6 +- loopx/cli_commands/support_control.py | 1 + loopx/cli_commands/support_control_chat.py | 4 + loopx/extensions/lark/inbox_reactions.py | 363 ++++++++++-------- loopx/extensions/lark/inbox_reply.py | 19 +- .../extensions/lark/private_conversations.py | 43 ++- .../project_registry_io_manifest_v1.json | 4 +- 7 files changed, 261 insertions(+), 179 deletions(-) diff --git a/loopx/chat_server.py b/loopx/chat_server.py index a06d195b86..14a2c83122 100644 --- a/loopx/chat_server.py +++ b/loopx/chat_server.py @@ -1543,6 +1543,7 @@ def serve_chat( verbose: bool = False, enable_goal_subagent_configuration: bool = False, project_workspace_grant: str = "workspace_write", + private_reactions: bool = True, ) -> None: if not is_loopback_host(host): raise ValueError("loopx chat requires a loopback --host such as 127.0.0.1") @@ -1616,7 +1617,8 @@ def serve_chat( observe=lambda profile: observe_lark_conversation_identity(profile=profile, runner=server.lark_runner, cli_bin=server.lark_cli_resolution.command or "lark-cli")) private_transport = LarkPrivateConversations(controller=server.runtime_controller, runtime_root=runtime_root, - runner=server.lark_runner, cli_bin=server.lark_cli_resolution.command or "lark-cli") + runner=server.lark_runner, cli_bin=server.lark_cli_resolution.command or "lark-cli", + reaction_feedback=private_reactions) server.lark_private_conversations = private_transport @@ -1674,7 +1676,7 @@ def _wake_goal_context(session): ).start() url = f"http://{host}:{port}{DEFAULT_CHAT_PATH}" print(f"Serving LoopX Chat at {url}", flush=True) - print("Agent boundary: local adapters, read-only sandbox, approval policy never", flush=True) + print(f"Agent boundary: local adapters, project grant {project_workspace_grant}, approval policy never", flush=True) print("Todo writes: preview-locked on loopback", flush=True) if enable_goal_subagent_configuration: print("Goal sub-agent configuration: preview-locked opt-in enabled", flush=True) diff --git a/loopx/cli_commands/support_control.py b/loopx/cli_commands/support_control.py index 7ed0b19c76..a2e991e79e 100644 --- a/loopx/cli_commands/support_control.py +++ b/loopx/cli_commands/support_control.py @@ -615,6 +615,7 @@ def handle_support_control_command( replace_existing_loopx_chat(args.host, args.port) serve_chat( project_workspace_grant=args.project_workspace_grant, + private_reactions=not getattr(args, "no_private_reactions", False), registry_path=chat_registry_path, runtime_root_override=args.runtime_root, scan_roots=scan_roots, diff --git a/loopx/cli_commands/support_control_chat.py b/loopx/cli_commands/support_control_chat.py index 8af58a7869..9abab64c9e 100644 --- a/loopx/cli_commands/support_control_chat.py +++ b/loopx/cli_commands/support_control_chat.py @@ -156,6 +156,10 @@ def register_chat_and_dashboard_commands( dashboard_parser.add_argument( "--verbose", action="store_true", help="Print HTTP request logs." ) + chat_parser.add_argument( + "--no-private-reactions", action="store_true", + help="Disable received/processing reactions for native Lark private conversations. Enabled by default.", + ) __all__ = ["register_chat_and_dashboard_commands"] diff --git a/loopx/extensions/lark/inbox_reactions.py b/loopx/extensions/lark/inbox_reactions.py index ec692e36f5..a515cd682d 100644 --- a/loopx/extensions/lark/inbox_reactions.py +++ b/loopx/extensions/lark/inbox_reactions.py @@ -25,7 +25,7 @@ RECEIVED_OPERATION_SCHEMA_VERSION = "lark_event_inbox_received_operation_v0" PROCESSED_STATE_FILENAME = "processed.json" REACTION_PHASES = {"received", "processing"} -RECEIVED_OPERATION_PHASES = {"prepared", "created"} +CREATION_OPERATION_PHASES = {"prepared", "created"} REACTION_ID_PATTERN = re.compile(r"[A-Za-z0-9_-]{1,200}") CommandRunner = Callable[[Sequence[str]], Mapping[str, Any]] ReactionCreator = Callable[[str, str], str | None] @@ -96,6 +96,14 @@ def _contains_string_by_key(value: object, key: str, expected: str) -> bool: return False +def _has_more_pages(value: object) -> bool: + if isinstance(value, Mapping): + return value.get("has_more") is True or any(_has_more_pages(child) for child in value.values()) + if isinstance(value, list): + return any(_has_more_pages(child) for child in value) + return False + + def _receipt_path(inbox: Path) -> Path: return inbox / "reactions" / "receipts.json" @@ -109,9 +117,11 @@ def _turn_start_reads_path(inbox: Path) -> Path: return inbox / "reactions" / "turn-start-reads.json" -def _received_operation_path(inbox: Path, message_id: str) -> Path: +def _reaction_creation_operation_path(inbox: Path, message_id: str, reaction_phase: str = "received") -> Path: digest = hashlib.sha256(message_id.encode("utf-8")).hexdigest() - return inbox / "reactions" / "received-operations" / f"{digest}.json" + if reaction_phase not in REACTION_PHASES: + raise ValueError("reaction phase is invalid") + return inbox / "reactions" / f"{reaction_phase}-operations" / f"{digest}.json" def _load_turn_start_reads(inbox: Path) -> set[str]: @@ -194,8 +204,14 @@ def lark_inbox_pending_turn_start_read_message_ids(*, inbox: Path) -> list[str]: ) -def _load_received_operation(*, inbox: Path, message_id: str) -> dict[str, str] | None: - path = _received_operation_path(inbox, message_id) +def _creation_operation_schema(reaction_phase: str) -> str: + # Preserve the existing received journal and reuse its lifecycle for processing. + return (RECEIVED_OPERATION_SCHEMA_VERSION if reaction_phase == "received" + else "lark_event_inbox_processing_operation_v0") + + +def _load_reaction_creation_operation(*, inbox: Path, message_id: str, reaction_phase: str = "received") -> dict[str, str] | None: + path = _reaction_creation_operation_path(inbox, message_id, reaction_phase) if not path.is_file(): return None try: @@ -211,9 +227,9 @@ def _load_received_operation(*, inbox: Path, message_id: str) -> dict[str, str] ) if ( not isinstance(payload, Mapping) - or payload.get("schema_version") != RECEIVED_OPERATION_SCHEMA_VERSION + or payload.get("schema_version") != _creation_operation_schema(reaction_phase) or payload.get("message_id") != message_id - or phase not in RECEIVED_OPERATION_PHASES + or phase not in CREATION_OPERATION_PHASES or not REACTION_EMOJI_PATTERN.fullmatch(emoji_type) or (phase == "created" and not REACTION_ID_PATTERN.fullmatch(reaction_id)) or (phase == "prepared" and reaction_id) @@ -225,16 +241,17 @@ def _load_received_operation(*, inbox: Path, message_id: str) -> dict[str, str] return result -def _write_received_operation( +def _write_reaction_creation_operation( *, inbox: Path, message_id: str, phase: str, emoji_type: str, reaction_id: str = "", + reaction_phase: str = "received", ) -> None: payload = { - "schema_version": RECEIVED_OPERATION_SCHEMA_VERSION, + "schema_version": _creation_operation_schema(reaction_phase), "message_id": message_id, "phase": phase, "emoji_type": emoji_type, @@ -242,11 +259,11 @@ def _write_received_operation( } if reaction_id: payload["reaction_id"] = reaction_id - write_private_json_atomic(_received_operation_path(inbox, message_id), payload) + write_private_json_atomic(_reaction_creation_operation_path(inbox, message_id, reaction_phase), payload) -def _clear_received_operation(*, inbox: Path, message_id: str) -> None: - path = _received_operation_path(inbox, message_id) +def _clear_reaction_creation_operation(*, inbox: Path, message_id: str, reaction_phase: str = "received") -> None: + path = _reaction_creation_operation_path(inbox, message_id, reaction_phase) try: path.unlink() except FileNotFoundError: @@ -499,6 +516,7 @@ def _delete_reaction( payload = _json_object(readback.get("stdout")) return bool( payload.get("ok") is True + and not _has_more_pages(payload) and not _contains_string_by_key(payload, "reaction_id", reaction_id) ) @@ -590,116 +608,67 @@ def ensure_lark_event_inbox_received_reaction_locked( configured=True, captured_pending=False, ) - receipts = lark_inbox_reaction_receipts( - inbox=inbox, - message_id=message_id, - ) - if receipts.get("received") is not None: - _clear_received_operation(inbox=inbox, message_id=message_id) - return _received_reaction_result( - status="already_received", - ok=True, - configured=True, - captured_pending=True, - ) - if receipts.get("processing") is not None: - _clear_received_operation(inbox=inbox, message_id=message_id) + receipts = lark_inbox_reaction_receipts(inbox=inbox, message_id=message_id) + if receipts.get("received") is None and receipts.get("processing") is not None: + _clear_reaction_creation_operation(inbox=inbox, message_id=message_id) return _received_reaction_result( status="already_processing", ok=True, configured=True, captured_pending=True, ) - operation = _load_received_operation(inbox=inbox, message_id=message_id) - if operation and operation["emoji_type"] != emoji_type: - return _received_reaction_result( - status="operation_config_changed", - ok=False, - configured=True, - captured_pending=True, - blocker="lark_inbox_received_reaction_operation_config_changed", - ) - if operation and operation["phase"] == "prepared": - return _received_reaction_result( - status="provider_outcome_uncertain", - ok=False, - configured=True, - captured_pending=True, - blocker="lark_inbox_received_reaction_provider_outcome_uncertain", - ) - if operation and operation["phase"] == "created": - reaction_id = operation["reaction_id"] - try: - record_lark_inbox_reaction( - inbox=inbox, - message_id=message_id, - phase="received", - reaction_id=reaction_id, - emoji_type=emoji_type, - ) - except (OSError, TypeError, ValueError): - return _received_reaction_result( - status="receipt_failed", - ok=False, - configured=True, - captured_pending=True, - external_writes_performed=False, - blocker="lark_inbox_received_reaction_receipt_failed", - ) - _clear_received_operation(inbox=inbox, message_id=message_id) - return _received_reaction_result( - status="receipt_recovered", - ok=True, - configured=True, - captured_pending=True, - ) - _write_received_operation( - inbox=inbox, - message_id=message_id, - phase="prepared", - emoji_type=emoji_type, + return _ensure_reaction_creation_locked( + inbox=inbox, message_id=message_id, emoji_type=emoji_type, reaction_phase="received", + create_reaction=create_reaction, delete_reaction=delete_reaction, ) - created_reaction_id = create_reaction(message_id, emoji_type) - if created_reaction_id is None: - _clear_received_operation(inbox=inbox, message_id=message_id) - return _received_reaction_result( - status="failed", - ok=False, - configured=True, - captured_pending=True, - blocker="lark_inbox_received_reaction_create_failed", - ) - try: - _write_received_operation( - inbox=inbox, - message_id=message_id, - phase="created", - emoji_type=emoji_type, - reaction_id=created_reaction_id, - ) - except (OSError, TypeError, ValueError): - cleaned_up = delete_reaction(message_id, created_reaction_id) - if cleaned_up: - _clear_received_operation(inbox=inbox, message_id=message_id) - return _received_reaction_result( - status="operation_receipt_failed", - ok=False, - configured=True, - captured_pending=True, - created_count=int(not cleaned_up), - external_writes_performed=True, - blocker=( - "lark_inbox_received_reaction_operation_receipt_failed" - if cleaned_up - else "lark_inbox_received_reaction_provider_outcome_uncertain" - ), - ) + + +def _ensure_reaction_creation_locked( + *, inbox: Path, message_id: str, emoji_type: str, reaction_phase: str, + create_reaction: ReactionCreator, delete_reaction: ReactionDeleter, +) -> dict[str, Any]: + """One provider creation lifecycle for received and processing feedback. + + Caller owns the per-source transition/settlement locks. Unknown writes stay + prepared; known provider ids recover their receipt without another create. + """ + receipts = lark_inbox_reaction_receipts( + inbox=inbox, + message_id=message_id, + ) + if receipts.get(reaction_phase) is not None: + _clear_reaction_creation_operation(inbox=inbox, message_id=message_id, reaction_phase=reaction_phase) + return _received_reaction_result( + status=f"already_{reaction_phase}", + ok=True, + configured=True, + captured_pending=True, + ) + operation = _load_reaction_creation_operation(inbox=inbox, message_id=message_id, reaction_phase=reaction_phase) + if operation and operation["emoji_type"] != emoji_type: + return _received_reaction_result( + status="operation_config_changed", + ok=False, + configured=True, + captured_pending=True, + blocker=f"lark_inbox_{reaction_phase}_reaction_operation_config_changed", + ) + if operation and operation["phase"] == "prepared": + return _received_reaction_result( + status="provider_outcome_uncertain", + ok=False, + configured=True, + captured_pending=True, + blocker=f"lark_inbox_{reaction_phase}_reaction_provider_outcome_uncertain", + ) + if operation and operation["phase"] == "created": + reaction_id = operation["reaction_id"] try: record_lark_inbox_reaction( inbox=inbox, message_id=message_id, - phase="received", - reaction_id=created_reaction_id, + phase=reaction_phase, + reaction_id=reaction_id, emoji_type=emoji_type, ) except (OSError, TypeError, ValueError): @@ -708,18 +677,84 @@ def ensure_lark_event_inbox_received_reaction_locked( ok=False, configured=True, captured_pending=True, - created_count=1, - external_writes_performed=True, - blocker="lark_inbox_received_reaction_receipt_failed", + external_writes_performed=False, + blocker=f"lark_inbox_{reaction_phase}_reaction_receipt_failed", ) - _clear_received_operation(inbox=inbox, message_id=message_id) + _clear_reaction_creation_operation(inbox=inbox, message_id=message_id, reaction_phase=reaction_phase) return _received_reaction_result( - status="received", + status="receipt_recovered", ok=True, configured=True, captured_pending=True, + ) + _write_reaction_creation_operation( + inbox=inbox, + message_id=message_id, + phase="prepared", + reaction_phase=reaction_phase, + emoji_type=emoji_type, + ) + created_reaction_id = create_reaction(message_id, emoji_type) + if created_reaction_id is None: + return _received_reaction_result( + status="failed", + ok=False, + configured=True, + captured_pending=True, + blocker=f"lark_inbox_{reaction_phase}_reaction_create_failed", + ) + try: + _write_reaction_creation_operation( + inbox=inbox, + message_id=message_id, + phase="created", + reaction_phase=reaction_phase, + emoji_type=emoji_type, + reaction_id=created_reaction_id, + ) + except (OSError, TypeError, ValueError): + cleaned_up = delete_reaction(message_id, created_reaction_id) + if cleaned_up: + _clear_reaction_creation_operation(inbox=inbox, message_id=message_id, reaction_phase=reaction_phase) + return _received_reaction_result( + status="operation_receipt_failed", + ok=False, + configured=True, + captured_pending=True, + created_count=int(not cleaned_up), + external_writes_performed=True, + blocker=( + f"lark_inbox_{reaction_phase}_reaction_operation_receipt_failed" + if cleaned_up + else f"lark_inbox_{reaction_phase}_reaction_provider_outcome_uncertain" + ), + ) + try: + record_lark_inbox_reaction( + inbox=inbox, + message_id=message_id, + phase=reaction_phase, + reaction_id=created_reaction_id, + emoji_type=emoji_type, + ) + except (OSError, TypeError, ValueError): + return _received_reaction_result( + status="receipt_failed", + ok=False, + configured=True, + captured_pending=True, created_count=1, + external_writes_performed=True, + blocker=f"lark_inbox_{reaction_phase}_reaction_receipt_failed", ) + _clear_reaction_creation_operation(inbox=inbox, message_id=message_id, reaction_phase=reaction_phase) + return _received_reaction_result( + status=reaction_phase, + ok=True, + configured=True, + captured_pending=True, + created_count=1, + ) def ensure_lark_event_inbox_received_reaction( @@ -748,6 +783,21 @@ def ensure_lark_event_inbox_received_reaction( ) +def mark_lark_event_inbox_received( + *, project: str | Path, config_path: str | Path, event: Mapping[str, Any], + runner: CommandRunner = _default_runner, +) -> dict[str, Any]: + config = load_lark_event_inbox_config(project=project, config_path=config_path) + profile = str(config["reply"]["sender_profile"]) + return ensure_lark_event_inbox_received_reaction( + project=project, config_path=config_path, event=event, + create_reaction=lambda message_id, emoji_type: _create_reaction( + runner=runner, profile=profile, message_id=message_id, emoji_type=emoji_type), + delete_reaction=lambda message_id, reaction_id: _delete_reaction( + runner=runner, profile=profile, message_id=message_id, reaction_id=reaction_id), + ) + + def _operation_result( *, operation: str, @@ -798,14 +848,12 @@ def mark_lark_event_inbox_processing( raise ValueError("processing requires a valid Lark message id") inbox = config["inbox_path"] with lark_inbox_reaction_lock(inbox=inbox, message_id=normalized): - if not _captured_pending_message(inbox=inbox, message_id=normalized): - raise ValueError("processing requires a captured pending inbox message") - return _mark_lark_event_inbox_processing_locked( - config=config, - message_id=normalized, - execute=execute, - runner=runner, - ) + with exclusive_file_lock(inbox / ".state" / "settlement", operation="lark_inbox_processing_settlement"): + if not _captured_pending_message(inbox=inbox, message_id=normalized): + raise ValueError("processing requires a captured pending inbox message") + return _mark_lark_event_inbox_processing_locked( + config=config, message_id=normalized, execute=execute, runner=runner, + ) def _mark_lark_event_inbox_processing_locked( @@ -846,50 +894,18 @@ def _mark_lark_event_inbox_processing_locked( profile = str(reply["sender_profile"]) created_count = 0 if processing is None: - reaction_id = _create_reaction( - runner=runner, - profile=profile, - message_id=normalized, - emoji_type=processing_emoji, + creation = _ensure_reaction_creation_locked( + inbox=inbox, message_id=normalized, emoji_type=processing_emoji, reaction_phase="processing", + create_reaction=lambda message_id, emoji_type: _create_reaction( + runner=runner, profile=profile, message_id=message_id, emoji_type=emoji_type), + delete_reaction=lambda message_id, reaction_id: _delete_reaction( + runner=runner, profile=profile, message_id=message_id, reaction_id=reaction_id), ) - if reaction_id is None: - return _operation_result( - operation="processing", - status="failed", - ok=False, - execute=True, - configured=True, - blocker="lark_inbox_processing_reaction_create_failed", - ) - try: - record_lark_inbox_reaction( - inbox=inbox, - message_id=normalized, - phase="processing", - reaction_id=reaction_id, - emoji_type=processing_emoji, - ) - except (OSError, ValueError): - _delete_reaction( - runner=runner, - profile=profile, - message_id=normalized, - reaction_id=reaction_id, - ) - return _operation_result( - operation="processing", - status="failed", - ok=False, - execute=True, - configured=True, - created_count=1, - blocker="lark_inbox_processing_reaction_receipt_failed", - ) - processing = { - "reaction_id": reaction_id, - "emoji_type": processing_emoji, - } - created_count = 1 + if not creation["ok"]: + return _operation_result(operation="processing", status=creation["status"], + ok=False, execute=True, configured=True, created_count=creation["created_count"], + blocker=creation.get("blocker")) + created_count = creation["created_count"] deleted_count = 0 if received is not None and config["reply"].get("received_reaction_policy") != "retain": @@ -960,6 +976,19 @@ def _complete_lark_event_inbox_reactions_locked( ) -> dict[str, Any]: inbox = config["inbox_path"] normalized = message_id + operation = _load_reaction_creation_operation(inbox=inbox, message_id=normalized, reaction_phase="processing") + if operation and execute: + profile = str(config["reply"]["sender_profile"]) + creation = _ensure_reaction_creation_locked( + inbox=inbox, message_id=normalized, emoji_type=operation["emoji_type"], reaction_phase="processing", + create_reaction=lambda *_: None, + delete_reaction=lambda message_id, reaction_id: _delete_reaction( + runner=runner, profile=profile, message_id=message_id, reaction_id=reaction_id), + ) + if not creation["ok"]: + return _operation_result(operation="complete", status=creation["status"], + ok=False, execute=True, configured=True, created_count=creation["created_count"], + blocker=creation.get("blocker")) receipts = lark_inbox_reaction_receipts( inbox=inbox, message_id=normalized, diff --git a/loopx/extensions/lark/inbox_reply.py b/loopx/extensions/lark/inbox_reply.py index a1790f1a03..c213d74082 100644 --- a/loopx/extensions/lark/inbox_reply.py +++ b/loopx/extensions/lark/inbox_reply.py @@ -286,6 +286,7 @@ def _deliver_lark_inbox_outbound( before_send: Callable[[str], Mapping[str, Any]] | None = None, delivery_attempt_recorder: Callable[[Mapping[str, str | None]], None] | None = None, short_message_limit: int | None = DEFAULT_LARK_TEXT_LIMIT, + finalize_reactions: bool = True, ) -> dict[str, Any]: """Deliver through one inbox-configured bot with exact provider readback. @@ -690,15 +691,15 @@ def _deliver_lark_inbox_outbound( execute=True, runner=runner, ) - if verified and source_message_id + if verified and source_message_id and finalize_reactions else {"ok": True} if verified else None ) - reaction_cleanup_verified = bool( + reaction_cleanup_verified = bool(finalize_reactions and reaction_cleanup is not None and reaction_cleanup.get("ok") is True ) - completed = bool(verified and reaction_cleanup_verified) + completed = bool(verified and (not finalize_reactions or reaction_cleanup_verified)) result = _result( status=( "sent_verified" @@ -730,6 +731,7 @@ def _deliver_lark_inbox_outbound( ) if guidance is not None: result["outbound_guidance"] = dict(guidance) + result["reaction_cleanup_deferred"] = not finalize_reactions return result @@ -747,12 +749,16 @@ def reply_lark_event_inbox( before_send: Callable[[str], Mapping[str, Any]] | None = None, delivery_attempt_recorder: Callable[[Mapping[str, str | None]], None] | None = None, short_message_limit: int | None = DEFAULT_LARK_TEXT_LIMIT, + finalize_reactions: bool = True, ) -> dict[str, Any]: """Reply with the explicit inbox-configured bot and placement policy. An answer delivery passes ``short_message_limit=None`` to declare that it is bounded by the provider's request limit rather than by the compact notification length. + + Intermediate admission/progress replies pass ``finalize_reactions=False``; + their verified delivery does not settle the source's processing lifecycle. """ result = _deliver_lark_inbox_outbound( @@ -768,6 +774,7 @@ def reply_lark_event_inbox( before_send=before_send, delivery_attempt_recorder=delivery_attempt_recorder, short_message_limit=short_message_limit, + finalize_reactions=finalize_reactions, ) result.setdefault("content_format", "markdown" if content_format == "markdown" @@ -784,6 +791,7 @@ def verify_lark_inbox_reply( attempt: Mapping[str, Any], runner: CommandRunner = _default_runner, source_membership_verifier: Callable[[], bool] | None = None, + finalize_reactions: bool = True, ) -> dict[str, Any]: """Read back one prior Lark reply without sending another message.""" @@ -937,12 +945,13 @@ def verify_lark_inbox_reply( message_id=message_id, execute=True, runner=runner, - ) + ) if finalize_reactions else {"ok": True} return { "ok": cleanup.get("ok") is True, "verification_performed": True, "reply_verified": True, - "reaction_cleanup_verified": cleanup.get("ok") is True, + "reaction_cleanup_verified": finalize_reactions and cleanup.get("ok") is True, + "reaction_cleanup_deferred": not finalize_reactions, } diff --git a/loopx/extensions/lark/private_conversations.py b/loopx/extensions/lark/private_conversations.py index a5105707eb..acfc9a1db2 100644 --- a/loopx/extensions/lark/private_conversations.py +++ b/loopx/extensions/lark/private_conversations.py @@ -17,13 +17,16 @@ from .event_inbox import acknowledge_lark_event_inbox, ingest_lark_event_inbox from .goal_channel_transport import call, json_payload, lark_args from .inbox_reply import _message, reply_lark_event_inbox, verify_lark_inbox_reply +from .inbox_reactions import mark_lark_event_inbox_processing, mark_lark_event_inbox_received class LarkPrivateConversations: - def __init__(self, *, controller: Any, runtime_root: Path, runner: Any, cli_bin: str) -> None: + def __init__(self, *, controller: Any, runtime_root: Path, runner: Any, cli_bin: str, + reaction_feedback: bool = True) -> None: self.core = ChatExternalConversations(controller) self.bindings = self.core.bindings self.runtime_root, self.runner, self.cli_bin = runtime_root, runner, cli_bin + self.reaction_feedback = reaction_feedback self.root = controller.store.root / "lark-private-deliveries" def profiles(self) -> dict[str, dict[str, str]]: @@ -93,7 +96,8 @@ def _inbox(self, record: dict[str, Any]) -> Path: "inbox_dir": f".loopx/inbox/private-chat/{scope}", "capture_scope": "configured_chat_all", "reply": {"enabled": True, "sender_profile": record["profile"], "sender_identity": "bot", "bot_display_name": observation["bot_display_name"], "chat_id": record["event"]["chat_id"], - "placement_policy": "source_context", "received_reaction_emoji": ""}, + "placement_policy": "source_context", "received_reaction_emoji": "Get" if self.reaction_feedback else "", + "received_reaction_policy": "retain", "processing_reaction_emoji": "OnIt" if self.reaction_feedback else ""}, }) event = record["event"] # Some unsupported attachment events have no rendered content. Keep a @@ -192,6 +196,28 @@ def admit(self, profile: str, event: dict[str, Any]) -> dict[str, Any]: def _reply_runner(self, args: list[str]) -> Any: return self.runner([self.cli_bin, *args[1:]], None, 30) + def _feedback(self, path: Path, record: dict[str, Any], *, processing: bool = False) -> None: + """Render observed Core admission/execution through the shared Inbox owner. + + Presentation failures never reject a persisted Turn. Existing reaction + journals prevent uncertain provider writes from being repeated on replay. + """ + if not self.reaction_feedback: + return + phase = "processing" if processing else "received" + feedback = record.setdefault("feedback", {}) + if feedback.get(phase, {}).get("ok"): + return + config = self._inbox(record) + if processing: + result = mark_lark_event_inbox_processing(project=self.runtime_root, config_path=config, + message_id=record["event"]["message_id"], execute=True, runner=self._reply_runner) + else: + result = mark_lark_event_inbox_received(project=self.runtime_root, config_path=config, + event=record["event"], runner=self._reply_runner) + feedback[phase] = result + _atomic_write_json(path, record) + def _deliver(self, path: Path, record: dict[str, Any], phase: str, text: str) -> bool: if not text: return False @@ -203,6 +229,7 @@ def _deliver(self, path: Path, record: dict[str, Any], phase: str, text: str) -> config = self._inbox(record) kwargs = dict(project=self.runtime_root, config_path=config, message_id=record["event"]["message_id"], text=text, runner=self._reply_runner, + finalize_reactions=phase != "admission" and not (phase == "terminal" and record.get("commission_resources")), source_membership_verifier=lambda: self._source_verified(record)) if phase_state.get("attempt"): result = verify_lark_inbox_reply(**kwargs, attempt=phase_state["attempt"]) @@ -223,7 +250,9 @@ def attempt(value: Any) -> None: result = reply_lark_event_inbox(**kwargs, execute=True, before_send=before_send, delivery_attempt_recorder=attempt, short_message_limit=None) - phase_state["verified"] = result.get("reply_verified") is True + # A verified intermediate reply is not completion. A final reply with + # cleanup owed recovers its original attempt; it must never be resent. + phase_state["verified"] = result.get("reply_verified") is True and result.get("ok") is True phase_state["blocker"] = result.get("blocker") record["deliveries"][phase] = phase_state _atomic_write_json(path, record) @@ -248,6 +277,7 @@ def reconcile(self) -> int: if record["status"] in {"command_queued", "command_completed", "commission_running"}: native = self.core.read_request(record["request_ref"]) if native["status"] == "command_queued": + self._feedback(path, record) self._deliver(path, record, "admission", native["response"]) continue record.update(status=native["status"], response=native.get("response"), @@ -257,16 +287,21 @@ def reconcile(self) -> int: native = self.core.read_request(record["request_ref"]) if native.get("agent_target"): self.bindings.resolve_agent_target(self.bindings.resolve(binding_id=record["binding_id"], **record["source"]), native["agent_target"]) + self._feedback(path, record) # Receipt follows persistent Core admission and is # independent of terminal execution and reply delivery. self._deliver(path, record, "admission", "已持久受理到原 Agent 会话;等待原宿主领取。/status 查看持久队列,/project 返回普通项目对话。实时停止暂不支持,请在原宿主处理。" if native.get("agent_target") else "已持久受理;若已有执行,本条会排队。可发送 /status、/stop 或 /new。") turn = self.core.controller.store.load_turn(record["session_id"], record["turn_id"]) + if turn and turn["status"] in {"starting", "running"}: + self._feedback(path, record, processing=True) if not turn or turn["status"] not in {"completed", "failed", "interrupted", "expired"}: continue response = str((turn.get("response") or {}).get("message") or "") if turn["status"] == "completed" else ( "本次执行已停止。" if turn["status"] == "interrupted" else "本次执行失败或已过期,原会话已保留;请发送 /status 后再决定是否重试。") else: + if record["status"] != "rejected": + self._feedback(path, record) response = _command_text(str(record.get("response_code") or "")) or str(record.get("response") or "") if record.get("status_snapshot"): response = _status_text(record["status_snapshot"], help_requested=native.get("command") == "help") @@ -274,6 +309,8 @@ def reconcile(self) -> int: resources = record.get("commission_resources") or {} if resources.get("session_id") and resources.get("turn_id"): first_turn = self.core.controller.store.load_turn(resources["session_id"], resources["turn_id"]) + if first_turn and first_turn["status"] in {"starting", "running"}: + self._feedback(path, record, processing=True) if not first_turn or first_turn["status"] not in {"completed", "failed", "interrupted", "expired"}: record["status"] = "commission_running" _atomic_write_json(path, record) diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 614af20d82..bfff23a53d 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -471,7 +471,7 @@ }, { "site": "loopx/chat_server.py::.serve_chat::codec_read:load_registry#1", - "line": 1558, + "line": 1559, "column": 16, "kind": "codec_read", "api": "load_registry", @@ -479,7 +479,7 @@ }, { "site": "loopx/chat_server.py::.serve_chat._wake_goal_context::codec_read:load_registry#1", - "line": 1663, + "line": 1665, "column": 20, "kind": "codec_read", "api": "load_registry", From 7efb7bb439a9eb2719f122ceac4b095c21e8e881 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sun, 4 Oct 2026 14:47:02 +0800 Subject: [PATCH 02/12] docs(lark): disclose default reactions and intermediate delivery semantics Signed-off-by: huangruiteng --- .../app-conversation-and-async-inbox-v0.md | 23 +++++++++++++++++++ ...p-conversation-and-async-inbox-v0.zh-CN.md | 15 ++++++++++++ .../extensions/lark/docs/lark-event-inbox.md | 7 ++++++ 3 files changed, 45 insertions(+) diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md index 850c943db4..a2b5046e8b 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -84,6 +84,29 @@ rules while preserving prior text and creating no Goal. It is a host/filesystem result, not live Lark write-workflow, material intake or release qualification. Maintainer review, installed/Lark journeys and broader IM interactions remain open. +## Native private feedback: default and delivery lifecycle + +Native Lark private conversations enable received `Get` and processing `OnIt` +feedback by default. `loopx chat --no-private-reactions` explicitly disables new +feedback writes. Both Apps reuse the existing Inbox provider lifecycle: received +feedback follows durable Core admission; processing requires an observed active +Turn. Queued requests do not appear to be executing. Received feedback is retained +after the result and conveys consumption, never acceptance of the work outcome. + +Intermediate replies and their read-only recovery defer reaction finalization. +Only the final result cleans processing receipts; cleanup failure keeps the +original verified reply recoverable without resending it. Received and processing +creations share a durable prepared/created journal. A known id recovers its receipt; +an uncertain write is not repeated. Incomplete paginated readback cannot establish +that a failed deletion succeeded. Provider permission failure does not cancel +already admitted work or silently escalate host policy. + +This is a bounded provider presentation refactor, not a new Session, queue, +model runner or control-plane owner. Native state drives both direct conversations +and explicit commissions. Synthetic queue/stop/replay/isolation regressions and +real provider canaries are separate evidence; broader incremental cards, media +and permission callbacks remain open acceptance. + ## Bound steward private Chat: explicit new commissions Settings → Lark can now select a steward role independently of ordinary project diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md index 10c6ced05d..690d7912ab 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md @@ -76,6 +76,21 @@ RPC 回执。Lark 专属设置 companion 归 extension;会话、请求、范 合成产品预览:[空管家与项目助手](../../assets/personal-workspace/private-steward-empty.png)、 [窄视口](../../assets/personal-workspace/private-steward-empty-narrow.png)。 +## 原生私聊反馈:默认开启与投递生命周期 + +Lark 原生私聊默认开启收到 `Get` 与处理中 `OnIt`; +`loopx chat --no-private-reactions` 显式关闭新反馈写入。两个 App 复用现有 Inbox +provider 生命周期:收到反馈在 Core 持久受理之后出现,处理中必须有原生 active +Turn 观测。排队不显示为正在执行;收到表情保留到最终结果之后,表示消费而非 +工作验收。中间受理/进度回复及其只读恢复不清理表情;最终结果清理处理中回执。 + +清理失败保留原回复的恢复路径,不重复发送已核验答案。收到与处理中共享 +prepared/created journal:已知 provider id 恢复回执,写入结果不确定时不盲重试; +分页未读完不能证明删除成功。缺少表情权限不取消已受理工作,也不提高宿主策略。 +这是现有 provider 呈现边界的有界重构,不新增 Session、queue、model runner 或 +控制面 authority。合成回归与真实 provider canary 分别记录;增量卡片、媒体和 +权限回调仍需独立验收。 + ## 私聊状态与帮助:授权范围内的观测 `/status` 与 `/help` 复用既有 typed bound-request owner,展示已授权角色、工作区、 diff --git a/loopx/extensions/lark/docs/lark-event-inbox.md b/loopx/extensions/lark/docs/lark-event-inbox.md index 17b9314466..3d3734a7b1 100644 --- a/loopx/extensions/lark/docs/lark-event-inbox.md +++ b/loopx/extensions/lark/docs/lark-event-inbox.md @@ -252,6 +252,13 @@ messages; use `configured_chat_all` for complete collaboration threads: } ``` +Native owner-bound private Chat enables `Get`/`OnIt` by default and retains the +received receipt. Use `loopx chat --no-private-reactions` to disable new private +feedback writes. Native Core admission and active Turn observations drive these +provider effects. Intermediate replies pass `finalize_reactions=False`; their +verified delivery does not clean processing feedback. Final delivery, including +recovery of a prior verified answer, still owns cleanup. + For every reply-enabled Inbox, a missing `reply.received_reaction_emoji` defaults to `Get`. Set it explicitly to the empty string to disable this provider write. The reaction belongs to the same explicit sender boundary as From f578192f8e0b9ec3198f0408e171605cb63a8c67 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sun, 4 Oct 2026 14:47:02 +0800 Subject: [PATCH 03/12] test(lark): qualify feedback queue, replay and cleanup recovery Signed-off-by: huangruiteng --- tests/extensions/test_lark_inbox_feedback.py | 81 ++++++++++++++ tests/extensions/test_lark_inbox_reactions.py | 4 +- tests/test_lark_private_conversations.py | 24 ++++ tests/test_lark_private_feedback.py | 103 ++++++++++++++++++ 4 files changed, 210 insertions(+), 2 deletions(-) create mode 100644 tests/extensions/test_lark_inbox_feedback.py create mode 100644 tests/test_lark_private_feedback.py diff --git a/tests/extensions/test_lark_inbox_feedback.py b/tests/extensions/test_lark_inbox_feedback.py new file mode 100644 index 0000000000..b23c73d260 --- /dev/null +++ b/tests/extensions/test_lark_inbox_feedback.py @@ -0,0 +1,81 @@ +"""Intermediate delivery and durable provider feedback recovery contracts.""" +import subprocess +from pathlib import Path + +from test_lark_inbox_reactions import ReactionRunner, ReplyRunner, _fixture + +from loopx.extensions.lark import inbox_reactions as inbox_reactions_module +from loopx.extensions.lark.inbox_reactions import ( + complete_lark_event_inbox_reactions, lark_inbox_reaction_receipts, + mark_lark_event_inbox_processing, record_lark_inbox_reaction, +) +from loopx.extensions.lark.inbox_reply import reply_lark_event_inbox, verify_lark_inbox_reply + + +def test_intermediate_reply_and_recovery_leave_processing_until_final(tmp_path: Path) -> None: + config, inbox, project = _fixture(tmp_path) + runner = ReplyRunner(readback_text='处理中') + record_lark_inbox_reaction(inbox=inbox, message_id='om_reaction_fixture', phase='processing', + reaction_id='reaction_OnIt', emoji_type='OnIt') + attempts = [] + intermediate = reply_lark_event_inbox(project=project, config_path=config, message_id='om_reaction_fixture', + text='处理中', execute=True, runner=runner, finalize_reactions=False, delivery_attempt_recorder=attempts.append) + assert intermediate['ok'] and intermediate['reply_verified'] + assert intermediate['reaction_cleanup_deferred'] and not intermediate['reaction_cleanup_verified'] + assert 'processing' in lark_inbox_reaction_receipts(inbox=inbox, message_id='om_reaction_fixture') + verified = verify_lark_inbox_reply(project=project, config_path=config, message_id='om_reaction_fixture', + text='处理中', attempt=attempts[0], runner=runner, finalize_reactions=False) + assert verified['ok'] and not verified['reaction_cleanup_verified'] + assert not any('delete' in args for args in runner.calls) + runner.readback_text = '处理完成' + final = reply_lark_event_inbox(project=project, config_path=config, message_id='om_reaction_fixture', + text='处理完成', execute=True, runner=runner) + assert final['ok'] and final['reaction_cleanup_verified'] + assert lark_inbox_reaction_receipts(inbox=inbox, message_id='om_reaction_fixture') == {} + + +def test_processing_known_create_recovers_receipt_before_terminal_cleanup(tmp_path: Path, monkeypatch) -> None: + config, inbox, project = _fixture(tmp_path) + runner = ReactionRunner() + real_record = inbox_reactions_module.record_lark_inbox_reaction + monkeypatch.setattr(inbox_reactions_module, 'record_lark_inbox_reaction', + lambda **_: (_ for _ in ()).throw(OSError('durable receipt unavailable'))) + first = mark_lark_event_inbox_processing(project=project, config_path=config, + message_id='om_reaction_fixture', execute=True, runner=runner) + assert not first['ok'] + assert inbox_reactions_module._load_reaction_creation_operation(inbox=inbox, message_id='om_reaction_fixture', + reaction_phase='processing')['phase'] == 'created' + monkeypatch.setattr(inbox_reactions_module, 'record_lark_inbox_reaction', real_record) + before = len([args for args in runner.calls if 'create' in args]) + completed = complete_lark_event_inbox_reactions(project=project, config_path=config, + message_id='om_reaction_fixture', execute=True, runner=runner) + assert completed['ok'] + assert len([args for args in runner.calls if 'create' in args]) == before + assert lark_inbox_reaction_receipts(inbox=inbox, message_id='om_reaction_fixture') == {} + + +def test_processing_unknown_create_never_repeats_after_recovery(tmp_path: Path) -> None: + config, inbox, project = _fixture(tmp_path) + calls = [] + def uncertain(args): + calls.append(args) + raise subprocess.TimeoutExpired(args, 1) + first = mark_lark_event_inbox_processing(project=project, config_path=config, + message_id='om_reaction_fixture', execute=True, runner=uncertain) + assert not first['ok'] + for operation in (mark_lark_event_inbox_processing, complete_lark_event_inbox_reactions): + recovered = operation(project=project, config_path=config, message_id='om_reaction_fixture', + execute=True, runner=uncertain) + assert recovered['status'] == 'provider_outcome_uncertain' + assert len(calls) == 1 + assert inbox_reactions_module._load_reaction_creation_operation(inbox=inbox, message_id='om_reaction_fixture', + reaction_phase='processing')['phase'] == 'prepared' + + +def test_failed_delete_with_incomplete_readback_is_not_verified() -> None: + def runner(args): + if 'delete' in args: + return {'returncode': 1} + return {'returncode': 0, 'stdout': '{"ok":true,"data":{"items":[],"has_more":true}}'} + assert not inbox_reactions_module._delete_reaction(runner=runner, profile='synthetic-app', + message_id='om_synthetic', reaction_id='reaction_own') diff --git a/tests/extensions/test_lark_inbox_reactions.py b/tests/extensions/test_lark_inbox_reactions.py index e15dbb6444..f47406cf71 100644 --- a/tests/extensions/test_lark_inbox_reactions.py +++ b/tests/extensions/test_lark_inbox_reactions.py @@ -275,7 +275,7 @@ def test_received_reaction_uncertain_operation_never_repeats_provider_create( tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: config, _inbox, project = _fixture(tmp_path) - real_write = inbox_reactions_module._write_received_operation + real_write = inbox_reactions_module._write_reaction_creation_operation write_count = 0 def fail_created_receipt(**kwargs: object) -> None: @@ -287,7 +287,7 @@ def fail_created_receipt(**kwargs: object) -> None: monkeypatch.setattr( inbox_reactions_module, - "_write_received_operation", + "_write_reaction_creation_operation", fail_created_receipt, ) created: list[tuple[str, str]] = [] diff --git a/tests/test_lark_private_conversations.py b/tests/test_lark_private_conversations.py index d1ccdcd84d..853b7d771f 100644 --- a/tests/test_lark_private_conversations.py +++ b/tests/test_lark_private_conversations.py @@ -18,6 +18,10 @@ def __init__(self): self.calls = [] self.writes = [] self.verify_replies = True + self.reactions = {} + self.reaction_creates = [] + self.fail_reaction_create = False + self.fail_reaction_delete = False def event(self, profile, name, text, kind="text"): event = {"schema_version": "lark_event_inbox_event_v0", "event_id": f"event_{name}", @@ -51,6 +55,26 @@ def __call__(self, args, cwd=None, timeout=None): self.writes.append((profile, text)) self.messages[ref] = {"message_id": ref, "body": {"content": content}} data = {"ok": True, "data": {"message_id": ref}} + elif "reactions" in args: + ref = args[args.index("--message-id") + 1] + if "create" in args: + if self.fail_reaction_create: + return {"returncode": 1, "stdout": '{"ok":false}'} + emoji = json.loads(args[args.index("--data") + 1])["reaction_type"]["emoji_type"] + reaction = f"reaction_{len(self.reaction_creates)}" + self.reaction_creates.append((profile, ref, emoji)) + self.reactions[reaction] = (profile, ref, emoji) + data = {"ok": True, "data": {"reaction_id": reaction}} + elif "delete" in args: + if self.fail_reaction_delete: + return {"returncode": 1, "stdout": '{"ok":false}'} + reaction = args[args.index("--reaction-id") + 1] + assert self.reactions[reaction][:2] == (profile, ref) + self.reactions.pop(reaction) + data = {"ok": True} + else: + data = {"ok": True, "data": {"items": [{"reaction_id": key} for key, row in self.reactions.items() + if row[:2] == (profile, ref)], "has_more": False}} else: pytest.fail(f"unnecessary provider operation: {args[2:5]}") return {"returncode": 0, "stdout": json.dumps(data), "stderr": ""} diff --git a/tests/test_lark_private_feedback.py b/tests/test_lark_private_feedback.py new file mode 100644 index 0000000000..c40d287316 --- /dev/null +++ b/tests/test_lark_private_feedback.py @@ -0,0 +1,103 @@ +"""Native state drives presentation; provider receipts never own admission.""" +import json +import time + +import pytest +from test_chat_ordinary_project import ordinary # noqa: F401 +from test_lark_private_conversations import connect + +from loopx.extensions.lark.private_conversations import LarkPrivateConversations + + +def active(store, row): + deadline = time.monotonic() + 10 + while store.load_session(row['session_id']).get('active_turn_id') != row['turn_id']: + assert time.monotonic() < deadline + time.sleep(.01) + + +def test_default_feedback_tracks_queue_execution_stop_and_app_isolation(ordinary): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + try: + slow = provider.event('notes-app', 'feedback_slow', 'wait for interrupt') + transport.admit('notes-app', slow) + first = transport.core.pending()[0] + active(store, first) + queued = provider.event('notes-app', 'feedback_queued', 'follow-up') + transport.admit('notes-app', queued) + transport.reconcile() + assert ('notes-app', slow['message_id'], 'OnIt') in provider.reaction_creates + assert ('notes-app', queued['message_id'], 'Get') in provider.reaction_creates + assert ('notes-app', queued['message_id'], 'OnIt') not in provider.reaction_creates + assert any(row[1:] == (slow['message_id'], 'OnIt') for row in provider.reactions.values()) + # A new provider instance reuses durable receipts, not in-memory emoji state. + replay = LarkPrivateConversations(controller=runtime, runtime_root=transport.runtime_root, + runner=provider, cli_bin='lark-cli') + before = list(provider.reaction_creates) + replay.admit('notes-app', {**slow, 'event_id': 'replayed'}) + replay.reconcile() + assert provider.reaction_creates == before + other = provider.event('steward-app', 'feedback_other', '/status') + replay.admit('steward-app', other) + replay.reconcile() + assert ('steward-app', other['message_id'], 'Get') in provider.reaction_creates + assert not any(profile == 'steward-app' and ref == slow['message_id'] + for profile, ref, _ in provider.reaction_creates) + replay.admit('notes-app', provider.event('notes-app', 'feedback_stop', '/stop')) + follow = next(row for row in replay.core.pending() if row['message'] == 'follow-up') + runtime.wait_for_turn(session_id=follow['session_id'], turn_id=follow['turn_id'], timeout_sec=10) + replay.reconcile() + assert not any(row[1:] == (slow['message_id'], 'OnIt') for row in provider.reactions.values()) + assert ('notes-app', slow['message_id'], 'Get') in provider.reactions.values() + assert all(row['goal_id'] is None for row in store.list_sessions()) + finally: + runtime.close() + + +def test_terminal_cleanup_recovers_verified_answer_without_resend(ordinary): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + try: + event = provider.event('notes-app', 'cleanup_slow', 'wait for interrupt') + transport.admit('notes-app', event) + row = transport.core.pending()[0] + active(store, row) + transport.reconcile() + transport.admit('notes-app', provider.event('notes-app', 'cleanup_stop', '/stop')) + runtime.wait_for_turn(session_id=row['session_id'], turn_id=row['turn_id'], timeout_sec=10) + provider.fail_reaction_delete = True + transport.reconcile() + record_path = transport.root / f"{row['request_ref']}.json" + record = json.loads(record_path.read_text()) + assert record['status'] != 'delivered' + assert record['deliveries']['terminal']['verified'] is False + assert record['deliveries']['terminal']['attempt'] + before = list(provider.writes) + transport.reconcile() + assert provider.writes == before + provider.fail_reaction_delete = False + transport.reconcile() + assert provider.writes == before + assert json.loads(record_path.read_text())['status'] == 'delivered' + finally: + runtime.close() + + +@pytest.mark.parametrize('enabled', [False, True]) +def test_opt_out_or_missing_reaction_permission_preserves_real_admission(ordinary, enabled): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + transport.reaction_feedback = enabled + provider.fail_reaction_create = True + try: + rejected = provider.event('notes-app', 'wrong_feedback', '/status') + assert transport.admit('steward-app', rejected)['status'] == 'audience_rejected' + event = provider.event('notes-app', 'no_scope_feedback', 'plain request') + assert transport.admit('notes-app', event)['status'] == 'durably_accepted' + row = transport.core.pending()[0] + runtime.wait_for_turn(session_id=row['session_id'], turn_id=row['turn_id'], timeout_sec=10) + assert transport.reconcile() == 1 + assert provider.reaction_creates == [] + assert bool(any('reactions' in call for call in provider.calls)) == enabled + assert any(text == 'Runtime response.' for _, text in provider.writes) + assert len(store.list_sessions()) == 1 + finally: + runtime.close() From a626e417ecabb9c7891b4ccd1fa9a9fbe8bc6010 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sun, 4 Oct 2026 15:11:44 +0800 Subject: [PATCH 04/12] Preserve native project defaults and render general private replies Signed-off-by: huangruiteng --- loopx/chat.py | 84 ++++++++++++++++++- loopx/chat_agent.py | 26 +++++- .../extensions/lark/private_conversations.py | 3 +- 3 files changed, 105 insertions(+), 8 deletions(-) diff --git a/loopx/chat.py b/loopx/chat.py index 619380bfa6..488825fda0 100644 --- a/loopx/chat.py +++ b/loopx/chat.py @@ -5,6 +5,7 @@ import re from pathlib import Path from typing import Any, Iterable, Mapping +from urllib.parse import unquote from .todos import add_goal_todo from .public_safe_text import LOCAL_PATH_SURFACE_PATTERN @@ -122,6 +123,83 @@ def replace_absolute_path(match: re.Match[str]) -> str: return _local_path_pattern(replacements).sub(replace_absolute_path, redacted) +def redact_response_markdown(text: str, *, protected_paths: Iterable[Path | str] = ()) -> str: + """Keep local inline-link labels without publishing unusable destinations. + + This bounded display repair handles balanced inline links, not reference + resolution or action admission. Code stays opaque to the link repair and + all output still passes through the existing path privacy owner. + """ + protected = tuple(protected_paths) + lines: list[str] = [] + fence: tuple[str, int] | None = None + for line in str(text or "").splitlines(keepends=True): + marker = re.match(r"^ {0,3}(`{3,}|~{3,})", line) + if fence is not None: + lines.append(line) + if marker and marker[1][0] == fence[0] and len(marker[1]) >= fence[1] and not line[marker.end():].strip(): + fence = None + continue + if marker: + fence = (marker[1][0], len(marker[1])) + lines.append(line) + continue + edits: list[tuple[int, int, str]] = [] + cursor, ticks = 0, 0 + while cursor < len(line): + if line[cursor] == "\\" and not ticks: + cursor += 2 + continue + if line[cursor] == "`": + end = cursor + 1 + while end < len(line) and line[end] == "`": + end += 1 + size = end - cursor + ticks = size if not ticks else 0 if size == ticks else ticks + cursor = end + continue + if ticks or line[cursor] != "[": + cursor += 1 + continue + start, end, depth = cursor, cursor + 1, 1 + while end < len(line) and depth: + if line[end] == "\\": + end += 2 + continue + depth += (line[end] == "[") - (line[end] == "]") + end += 1 + if depth or line[end:end + 1] != "(": + cursor = end + continue + destination_start, close, depth = end + 1, end + 1, 1 + while close < len(line) and depth: + if line[close] == "\\": + close += 2 + continue + depth += (line[close] == "(") - (line[close] == ")") + close += 1 + if depth: + cursor = close + continue + destination = line[destination_start:close - 1] + decoded = unquote(destination) + local = decoded.lstrip("< ").lower().startswith("file:") or any( + redact_local_paths(value, protected_paths=protected) != value + for value in (destination, decoded) + ) + if local: + image_start = start - 1 if start and line[start - 1] == "!" else start + edits.append((image_start, close, line[start + 1:end - 1])) + cursor = close + cursor, chunks = 0, [] + for start, end, label in edits: + chunks.extend((line[cursor:start], label)) + cursor = end + chunks.append(line[cursor:]) + lines.append("".join(chunks)) + return redact_local_paths("".join(lines), protected_paths=protected) + + class VisibleResponseStreamFilter: """Stream safe operator text while withholding the structured review envelope.""" @@ -415,7 +493,7 @@ def normalize_agent_response( from .capabilities.manager_context import normalize_request handoff = normalize_request(payload.get("context_handoff")) protected = tuple(protected_paths) - message = redact_local_paths( + message = redact_response_markdown( str(payload.get("message") or ""), protected_paths=protected, ).strip() @@ -468,7 +546,7 @@ def parse_agent_response( if isinstance(salvaged, str) and salvaged.strip(): return { "schema_version": CHAT_AGENT_RESPONSE_SCHEMA_VERSION, - "message": redact_local_paths(salvaged, protected_paths=protected).strip(), + "message": redact_response_markdown(salvaged, protected_paths=protected).strip(), "proposals": [], "protected_action": None, "gate": None, @@ -501,7 +579,7 @@ def parse_agent_response( raw_text = visible or salvaged_message return { "schema_version": CHAT_AGENT_RESPONSE_SCHEMA_VERSION, - "message": redact_local_paths(raw_text, protected_paths=protected).strip(), + "message": redact_response_markdown(raw_text, protected_paths=protected).strip(), "proposals": [], "protected_action": None, "gate": None, diff --git a/loopx/chat_agent.py b/loopx/chat_agent.py index 902b8e9bc9..6f5f25a95d 100644 --- a/loopx/chat_agent.py +++ b/loopx/chat_agent.py @@ -413,6 +413,8 @@ def _turn_prompt( "Do not expose chain-of-thought, tool narration, intended steps, or scratch work. " "First write the complete operator-facing answer as safe Markdown text. Give a simple question a direct sourced answer; for a complex task, lead with the judgment and then explain the material evidence, comparisons, decisions and limitations at useful depth. " "Use short sentences or lines so the answer can stream. Avoid gratuitous headings, boilerplate, raw ID inventories and more than five actionable items. " + "For completed work, explain the useful result and any material limitation; keep routine tool logs, test commands and implementation details out of the default reply unless they help the operator decide or were requested. " + "Use readable lists, emphasis, quotes and fenced code when they clarify the answer. Link only to real accessible resources; present local deliverables as workspace-relative inline code with a descriptive label, never a fabricated web link or an absolute machine path. " "Do not emit executable HTML. The complete answer must stay in this conversation, even when a separate report artifact also exists. " "Then append exactly one machine-readable envelope whose message field repeats that complete answer. This envelope is hidden protocol metadata and is required even for ordinary questions or exact-wording replies; user formatting instructions govern the visible answer, not omission of this metadata. " "protected_action must be null or an object shaped as " @@ -614,6 +616,22 @@ def start( request_id=1, ) session._notify("initialized", {}) + thread_request_id = 2 + if project_context is not None: + # Codex owns project/default configuration. Resume otherwise + # retains the old thread's effort even after an owner edits it. + configured = session._request( + "config/read", {"cwd": str(root), "includeLayers": False}, + request_id=thread_request_id, + ).get("config", {}) + if not isinstance(configured, dict): + raise session._runtime_error("Codex project configuration is unavailable.") + if any(configured.get(key) is not None and not isinstance(configured[key], str) + for key in ("model", "model_reasoning_effort")): + raise session._runtime_error("Codex project model configuration is invalid.") + model = model or configured.get("model") + reasoning_effort = reasoning_effort or configured.get("model_reasoning_effort") + thread_request_id += 1 thread_result = session._request( "thread/resume" if resume_thread_id else "thread/start", { @@ -646,18 +664,18 @@ def start( else {} ), }, - request_id=2, + request_id=thread_request_id, ) if model and thread_result.get("model") not in {None, model}: raise session._runtime_error( - "Codex did not apply the requested manager model." + "Codex did not apply the requested model." ) if reasoning_effort and thread_result.get("reasoningEffort") not in { None, reasoning_effort, }: raise session._runtime_error( - "Codex did not apply the requested manager reasoning effort." + "Codex did not apply the requested reasoning effort." ) session.model = thread_result.get("model") or model session.reasoning_effort = thread_result.get("reasoningEffort") or reasoning_effort @@ -674,7 +692,7 @@ def start( # that public-safe context in each Turn prompt. Codex Goal mode is reserved # for autonomous execution; enabling it here causes conversational messages # to be treated as continuation ticks instead of the current user task. - session.next_request_id = 3 + session.next_request_id = thread_request_id + 1 return session except _LegacyModelCatalogSchemaError as exc: session.close() diff --git a/loopx/extensions/lark/private_conversations.py b/loopx/extensions/lark/private_conversations.py index acfc9a1db2..8f21d94168 100644 --- a/loopx/extensions/lark/private_conversations.py +++ b/loopx/extensions/lark/private_conversations.py @@ -249,7 +249,8 @@ def attempt(value: Any) -> None: _atomic_write_json(path, record) result = reply_lark_event_inbox(**kwargs, execute=True, before_send=before_send, - delivery_attempt_recorder=attempt, short_message_limit=None) + delivery_attempt_recorder=attempt, short_message_limit=None, + content_format="markdown") # A verified intermediate reply is not completion. A final reply with # cleanup owed recovers its original attempt; it must never be resent. phase_state["verified"] = result.get("reply_verified") is True and result.get("ok") is True From 670db0f2d31a4b60c79af65050c82050f223653a Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sun, 4 Oct 2026 15:11:44 +0800 Subject: [PATCH 05/12] Document private presentation and exact project resume boundaries Signed-off-by: huangruiteng --- .../rfcs/app-conversation-and-async-inbox-v0.md | 10 ++++++++++ .../rfcs/app-conversation-and-async-inbox-v0.zh-CN.md | 6 ++++++ loopx/extensions/lark/docs/lark-event-inbox.md | 10 ++++++++++ 3 files changed, 26 insertions(+) diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md index a2b5046e8b..1eb2efa77f 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -247,6 +247,16 @@ execution. Stopping a conversation does not silently stop delegated work. ## Product expression: what to borrow and what remains unproven +Native private replies now share Markdown post presentation for lists, quotations, +public links and code. Final response redaction preserves local inline-link labels +without publishing fake destinations; it is not local artifact delivery. Existing +plain-text attempts recover their original verified format without another send. +Ordinary project Codex adapters resolve current host project model/effort defaults +on start and exact resume; an explicit executor choice takes precedence. An idle +service restart reattaches the adapter after configuration edits, preserving its +Session/thread and grants. Media, incremental presentation and permission journeys +retain their separate acceptance gaps. + The [Lorca release post](https://x.com/localhost_4173/status/2103454978220470708) was inspected on 2026-09-25. It contains a **static screenshot**, not a verified interactive or recovery demonstration. The screenshot shows a named conversation diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md index 690d7912ab..4a5d27fd11 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md @@ -78,6 +78,12 @@ RPC 回执。Lark 专属设置 companion 归 extension;会话、请求、范 ## 原生私聊反馈:默认开启与投递生命周期 +通用个人助手的最终回复复用 Markdown post,保留列表、引用、公开链接和代码; +本地链接脱敏时保留可读标签,不生成指向 `[project]` 的假链接。旧纯文本投递 +仍按原 attempt 读回恢复,不因新默认而重发。图片/文件交付与真实增量仍需独立验收。 +普通项目 Codex adapter 在启动和确切 resume 时读取宿主当前项目模型/effort 默认; +显式选择优先,配置修改后空闲重启服务以重新接续 adapter,保留 Session/thread 和 grants。 + Lark 原生私聊默认开启收到 `Get` 与处理中 `OnIt`; `loopx chat --no-private-reactions` 显式关闭新反馈写入。两个 App 复用现有 Inbox provider 生命周期:收到反馈在 Core 持久受理之后出现,处理中必须有原生 active diff --git a/loopx/extensions/lark/docs/lark-event-inbox.md b/loopx/extensions/lark/docs/lark-event-inbox.md index 3d3734a7b1..1c3bee69c3 100644 --- a/loopx/extensions/lark/docs/lark-event-inbox.md +++ b/loopx/extensions/lark/docs/lark-event-inbox.md @@ -259,6 +259,16 @@ provider effects. Intermediate replies pass `finalize_reactions=False`; their verified delivery does not clean processing feedback. Final delivery, including recovery of a prior verified answer, still owns cleanup. +Native private replies use the shared Markdown `post` renderer, preserving lists, +links, quotes and code. Recovery verifies the original attempt's format, including +plain-text replies sent before this default changed, without another send. Local +inline links keep their readable labels rather than linking to redaction tokens; +machine paths remain private. This does not deliver local files or images. +Ordinary project Codex start/resume resolves the host's current project model and +effort defaults before attaching the exact thread; explicit executor choices win. +Restart an idle Chat service after changing host defaults to reattach its adapter. +No new thread, audience, sandbox grant or manager configuration is implied. + For every reply-enabled Inbox, a missing `reply.received_reaction_emoji` defaults to `Get`. Set it explicitly to the empty string to disable this provider write. The reaction belongs to the same explicit sender boundary as From 09498450436e68fbfa3b941a4e4379fb7b5aae06 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sun, 4 Oct 2026 15:11:44 +0800 Subject: [PATCH 06/12] Cover project default adoption and native rich reply recovery Signed-off-by: huangruiteng --- tests/test_chat_ordinary_project.py | 32 +++++++++++++++++ tests/test_chat_path_redaction.py | 25 +++++++++++++- tests/test_lark_private_conversations.py | 9 ++--- tests/test_lark_private_feedback.py | 44 ++++++++++++++++++++++++ 4 files changed, 105 insertions(+), 5 deletions(-) diff --git a/tests/test_chat_ordinary_project.py b/tests/test_chat_ordinary_project.py index 36a168dd0d..850b44f384 100644 --- a/tests/test_chat_ordinary_project.py +++ b/tests/test_chat_ordinary_project.py @@ -150,3 +150,35 @@ def test_default_project_host_grant_is_write_and_read_only_launch_is_enforced(or bindings.configure(transport_ref="notes-app", project_ref=projects["projects"][0]["project_ref"], executor_endpoint_id="codex", project_grant="workspace_write") assert bindings.read()["bindings"][0] == binding runtime.close() + + +def test_project_codex_defaults_apply_on_exact_resume_without_replacing_context(ordinary): + store, runtime, contexts, request, capture, fake, workspace = ordinary + settings = workspace / "host-model-fixture.json" + settings.write_text(json.dumps({"model": "gpt-6.1-sol", "model_reasoning_effort": "medium"})) + fake.write_text(fake.read_text().replace( + ' elif method in {"thread/start", "thread/resume"}:', + f' elif method == "config/read":\n result = {{"config": json.load(open({str(settings)!r}))}}\n' + ' elif method in {"thread/start", "thread/resume"}:')) + ref = contexts.available()[0]["project_ref"] + session, resumed = runtime.open_session(goal_id=None, agent_id="codex", work_dir=workspace, + objective="ordinary conversation", project_ref=ref, mode="resume_latest") + assert not resumed + original = runtime.adapters[session["session_id"]].session + assert (original.model, original.reasoning_effort) == ("gpt-6.1-sol", "medium") + settings.write_text(json.dumps({"model": "gpt-6.1-sol", "model_reasoning_effort": "high"})) + runtime.adapters.pop(session["session_id"]).close_session() + restored, resumed = runtime.open_session(goal_id=None, agent_id="codex", work_dir=workspace, + objective="ordinary conversation", project_ref=ref, mode="resume_latest") + assert resumed and restored["session_id"] == session["session_id"] + assert restored["upstream_thread_id"] == session["upstream_thread_id"] + adapter = runtime.adapters[session["session_id"]].session + assert (adapter.model, adapter.reasoning_effort) == ("gpt-6.1-sol", "high") + requests = [json.loads(line) for line in capture.read_text().splitlines()] + resume = next(row for row in requests if row["method"] == "thread/resume") + assert resume["params"]["model"] == "gpt-6.1-sol" + assert resume["params"]["config"]["model_reasoning_effort"] == "high" + assert resume["params"]["threadId"] == session["upstream_thread_id"] + assert resume["params"]["sandbox"] == "read-only" and resume["params"]["approvalPolicy"] == "never" + assert all(row["goal_id"] is None for row in store.list_sessions()) + runtime.close() diff --git a/tests/test_chat_path_redaction.py b/tests/test_chat_path_redaction.py index c2c4118da2..10ecefe10c 100644 --- a/tests/test_chat_path_redaction.py +++ b/tests/test_chat_path_redaction.py @@ -8,7 +8,7 @@ import pytest -from loopx.chat import VisibleResponseStreamFilter, parse_agent_response, redact_local_paths +from loopx.chat import VisibleResponseStreamFilter, parse_agent_response, redact_local_paths, redact_response_markdown from loopx.chat_status_api import ChatStatusRequestMixin @@ -35,6 +35,29 @@ def test_canonical_private_path_shapes_and_public_urls(): assert redact_local_paths(public) == public +@pytest.mark.parametrize("destination", [ + "/custom-volume/project/report.md#result", "/custom-volume/project/a(b).md:12", + ' "title"', + "file:///custom-volume/project/report.md", "%2Fcustom-volume%2Fproject%2Freport.md", + r"Q:\private state\report.md", +]) +def test_response_local_links_keep_labels_without_broken_or_private_destinations(destination): + text = f"- [报告 [结果]]({destination});[公开来源](https://example.org/a(b))." + result = parse_agent_response(text, protected_paths=["/custom-volume/project", r"Q:\private state"])["message"] + assert result == "- 报告 [结果];[公开来源](https://example.org/a(b))." + + +def test_response_link_repair_preserves_code_and_does_not_rewrite_status_json(): + link = "[label](/custom-volume/project/report.md)" + code = f"`{link}`\n```md\n{link}\n```\n" + assert redact_response_markdown(code, protected_paths=["/custom-volume/project"]) == code.replace( + "/custom-volume/project/report.md", "[project]") + assert json.loads(redact_local_paths(json.dumps({"message": link}), protected_paths=["/custom-volume/project"])) == { + "message": "[label]([project])"} + raw = '' + json.dumps({"message": link}) + '' + assert parse_agent_response(raw, protected_paths=["/custom-volume/project"])["message"] == "label" + + @pytest.mark.parametrize("root", [ "/custom-volume/private-state", "/custom-volume/private state", "/custom-volume/" + "r" * 190, diff --git a/tests/test_lark_private_conversations.py b/tests/test_lark_private_conversations.py index 853b7d771f..e842fed523 100644 --- a/tests/test_lark_private_conversations.py +++ b/tests/test_lark_private_conversations.py @@ -46,14 +46,15 @@ def __call__(self, args, cwd=None, timeout=None): message = {**message, "body": {"content": json.dumps({"text": "not yet visible"})}} data = {"ok": True, "data": {"items": [message]}} elif "+messages-send" in args: - text = args[args.index("--text") + 1] - content = json.dumps({"text": text}) + kind = "post" if "--content" in args else "text" + content = args[args.index("--content") + 1] if kind == "post" else json.dumps({"text": args[args.index("--text") + 1]}) + text = json.loads(content)["zh_cn"]["content"][0][0]["text"] if kind == "post" else json.loads(content)["text"] if "--dry-run" in args: - data = {"ok": True, "api": [{"body": {"content": content}}]} + data = {"ok": True, "api": [{"body": {"content": content, "msg_type": kind}}]} else: ref = f"om_out_{len(self.writes)}" self.writes.append((profile, text)) - self.messages[ref] = {"message_id": ref, "body": {"content": content}} + self.messages[ref] = {"message_id": ref, "msg_type": kind, "body": {"content": content}} data = {"ok": True, "data": {"message_id": ref}} elif "reactions" in args: ref = args[args.index("--message-id") + 1] diff --git a/tests/test_lark_private_feedback.py b/tests/test_lark_private_feedback.py index c40d287316..81079f3450 100644 --- a/tests/test_lark_private_feedback.py +++ b/tests/test_lark_private_feedback.py @@ -101,3 +101,47 @@ def test_opt_out_or_missing_reaction_permission_preserves_real_admission(ordinar assert len(store.list_sessions()) == 1 finally: runtime.close() + + +def test_native_private_default_post_preserves_general_result_structure_and_safe_links(ordinary): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + fake, workspace = ordinary[-2:] + text = (f"**结果**\n\n- [本地报告]({workspace}/report.md)\n" + "- [公开来源](https://example.org/source)\n\n> 待确认限制\n\n```python\nprint('ok')\n```") + fake.write_text(fake.read_text().replace('"message": "Runtime response.",', f'"message": {text!r},')) + try: + event = provider.event('notes-app', 'rich_result', 'ordinary work') + transport.admit('notes-app', event) + row = transport.core.pending()[0] + runtime.wait_for_turn(session_id=row['session_id'], turn_id=row['turn_id'], timeout_sec=10) + assert transport.reconcile() == 1 + final = provider.messages[f'om_out_{len(provider.writes) - 1}'] + assert final['msg_type'] == 'post' + visible = json.loads(final['body']['content'])['zh_cn']['content'][0][0]['text'] + assert visible == text.replace(f'[本地报告]({workspace}/report.md)', '本地报告') + assert '[project]' not in visible and str(workspace) not in visible + assert transport.core.read_request(row['request_ref'])['delivery_verified'] is True + finally: + runtime.close() + + +def test_native_post_default_recovers_an_existing_plain_text_attempt_without_resend(ordinary, monkeypatch): # noqa: F811 + import loopx.extensions.lark.private_conversations as native + store, runtime, provider, transport = connect(ordinary) + send = native.reply_lark_event_inbox + def old_text_send(**kwargs): + return send(**{**kwargs, 'content_format': 'text'}) + try: + event = provider.event('notes-app', 'old_text_receipt', '/status') + transport.admit('notes-app', event) + provider.verify_replies = False + with monkeypatch.context() as patch: + patch.setattr(native, 'reply_lark_event_inbox', old_text_send) + assert transport.reconcile() == 0 + assert len(provider.writes) == 1 + provider.verify_replies = True + assert transport.reconcile() == 1 + assert len(provider.writes) == 1 + assert provider.messages['om_out_0']['msg_type'] == 'text' + finally: + runtime.close() From cadecda932334a3597bdfec8f1179ee929376fff Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sun, 4 Oct 2026 15:25:43 +0800 Subject: [PATCH 07/12] Avoid duplicate admission messages when native feedback succeeds Signed-off-by: huangruiteng --- loopx/extensions/lark/private_conversations.py | 14 +++++++++++++- 1 file changed, 13 insertions(+), 1 deletion(-) diff --git a/loopx/extensions/lark/private_conversations.py b/loopx/extensions/lark/private_conversations.py index 8f21d94168..61dd80e94b 100644 --- a/loopx/extensions/lark/private_conversations.py +++ b/loopx/extensions/lark/private_conversations.py @@ -291,8 +291,20 @@ def reconcile(self) -> int: self._feedback(path, record) # Receipt follows persistent Core admission and is # independent of terminal execution and reply delivery. - self._deliver(path, record, "admission", "已持久受理到原 Agent 会话;等待原宿主领取。/status 查看持久队列,/project 返回普通项目对话。实时停止暂不支持,请在原宿主处理。" if native.get("agent_target") else "已持久受理;若已有执行,本条会排队。可发送 /status、/stop 或 /new。") turn = self.core.controller.store.load_turn(record["session_id"], record["turn_id"]) + # Get already acknowledges admission. Reserve a text + # notification for a real wait or unavailable feedback; + # resume any old attempt with its original exact text. + admission = (record["deliveries"].get("admission") or {}).get("text") + if not admission: + if native.get("agent_target"): + admission = "已收到,等待原 Agent 宿主处理。/status 查看状态,/project 返回项目对话。" + elif turn and turn["status"] == "queued": + admission = "正在处理前一条,这条已排队。" + elif (record.get("feedback", {}).get("received") or {}).get("ok") is not True: + admission = "已收到,正在处理。" + if admission: + self._deliver(path, record, "admission", admission) if turn and turn["status"] in {"starting", "running"}: self._feedback(path, record, processing=True) if not turn or turn["status"] not in {"completed", "failed", "interrupted", "expired"}: From 02a1dcc0a981abe0183917a8445b686dea74cd67 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sun, 4 Oct 2026 15:25:43 +0800 Subject: [PATCH 08/12] Explain concise queue and reaction fallback presentation Signed-off-by: huangruiteng --- loopx/extensions/lark/docs/lark-event-inbox.md | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/loopx/extensions/lark/docs/lark-event-inbox.md b/loopx/extensions/lark/docs/lark-event-inbox.md index 1c3bee69c3..2b393899b2 100644 --- a/loopx/extensions/lark/docs/lark-event-inbox.md +++ b/loopx/extensions/lark/docs/lark-event-inbox.md @@ -255,7 +255,10 @@ messages; use `configured_chat_all` for complete collaboration threads: Native owner-bound private Chat enables `Get`/`OnIt` by default and retains the received receipt. Use `loopx chat --no-private-reactions` to disable new private feedback writes. Native Core admission and active Turn observations drive these -provider effects. Intermediate replies pass `finalize_reactions=False`; their +provider effects. Successful received feedback avoids a duplicate admission +message; an observed queued Turn still gets a concise waiting notice. Missing +reaction permission or explicit opt-out retains a text receipt. Existing attempts +recover their original text. Intermediate replies pass `finalize_reactions=False`; their verified delivery does not clean processing feedback. Final delivery, including recovery of a prior verified answer, still owns cleanup. From d042e2a330b6f8bfed7d2cf0cd9f9b71f128c140 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sun, 4 Oct 2026 15:25:43 +0800 Subject: [PATCH 09/12] Qualify queue notices and text fallback without duplicate receipts Signed-off-by: huangruiteng --- tests/test_lark_private_conversations.py | 2 +- tests/test_lark_private_feedback.py | 3 +++ 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/tests/test_lark_private_conversations.py b/tests/test_lark_private_conversations.py index e842fed523..fdedbe3235 100644 --- a/tests/test_lark_private_conversations.py +++ b/tests/test_lark_private_conversations.py @@ -117,7 +117,7 @@ def test_native_private_admission_queue_other_app_stop_and_verified_delivery(ord assert runtime.wait_for_turn(session_id=other_row["session_id"], turn_id=other_row["turn_id"], timeout_sec=10)["status"] == "completed" transport.reconcile() assert any(profile == "steward-app" and text == "Runtime response." for profile, text in provider.writes) - assert any(profile == "notes-app" and "持久受理" in text for profile, text in provider.writes) + assert any(profile == "notes-app" and "已排队" in text for profile, text in provider.writes) stop = provider.event("notes-app", "stop", "/stop") assert transport.admit("notes-app", stop)["status"] == "command_recorded" queued_row = next(row for row in transport.core.pending() if row["message"] == "follow-up") diff --git a/tests/test_lark_private_feedback.py b/tests/test_lark_private_feedback.py index 81079f3450..c3b2d343a2 100644 --- a/tests/test_lark_private_feedback.py +++ b/tests/test_lark_private_feedback.py @@ -29,6 +29,7 @@ def test_default_feedback_tracks_queue_execution_stop_and_app_isolation(ordinary assert ('notes-app', slow['message_id'], 'OnIt') in provider.reaction_creates assert ('notes-app', queued['message_id'], 'Get') in provider.reaction_creates assert ('notes-app', queued['message_id'], 'OnIt') not in provider.reaction_creates + assert provider.writes == [('notes-app', '正在处理前一条,这条已排队。')] assert any(row[1:] == (slow['message_id'], 'OnIt') for row in provider.reactions.values()) # A new provider instance reuses durable receipts, not in-memory emoji state. replay = LarkPrivateConversations(controller=runtime, runtime_root=transport.runtime_root, @@ -98,6 +99,7 @@ def test_opt_out_or_missing_reaction_permission_preserves_real_admission(ordinar assert provider.reaction_creates == [] assert bool(any('reactions' in call for call in provider.calls)) == enabled assert any(text == 'Runtime response.' for _, text in provider.writes) + assert any(text == '已收到,正在处理。' for _, text in provider.writes) assert len(store.list_sessions()) == 1 finally: runtime.close() @@ -117,6 +119,7 @@ def test_native_private_default_post_preserves_general_result_structure_and_safe assert transport.reconcile() == 1 final = provider.messages[f'om_out_{len(provider.writes) - 1}'] assert final['msg_type'] == 'post' + assert len(provider.writes) == 1 # emoji admission does not duplicate the answer visible = json.loads(final['body']['content'])['zh_cn']['content'][0][0]['text'] assert visible == text.replace(f'[本地报告]({workspace}/report.md)', '本地报告') assert '[project]' not in visible and str(workspace) not in visible From 4640e9963b9374f51c711c732c3933a135b706ed Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sun, 4 Oct 2026 18:44:57 +0800 Subject: [PATCH 10/12] Prevent slow private reply readbacks from holding control feedback Signed-off-by: huangruiteng --- .../app-conversation-and-async-inbox-v0.md | 7 + ...p-conversation-and-async-inbox-v0.zh-CN.md | 5 + .../native_chat/external_conversations.py | 5 +- .../lark/goal_topic_runtime_service.py | 38 +++- .../extensions/lark/private_conversations.py | 177 +++++++++--------- tests/test_lark_private_conversations.py | 104 ++++++++++ 6 files changed, 243 insertions(+), 93 deletions(-) diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md index 1eb2efa77f..b59152ab7d 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -101,6 +101,13 @@ an uncertain write is not repeated. Incomplete paginated readback cannot establi that a failed deletion succeeded. Provider permission failure does not cancel already admitted work or silently escalate host policy. +Private delivery reconciles at most four persisted requests concurrently, with +no executor backlog. A blocked reply readback does not hold an independent queue +notice or control response; per-request locks and journals retain no-resend +recovery. Core recovers each exact request under the existing Session fence. +Provider verification cost and saturated-worker latency remain separate from +model concurrency and are not certified by a synthetic blocked-readback test. + This is a bounded provider presentation refactor, not a new Session, queue, model runner or control-plane owner. Native state drives both direct conversations and explicit commissions. Synthetic queue/stop/replay/isolation regressions and diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md index 4a5d27fd11..b50629817b 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md @@ -93,6 +93,11 @@ Turn 观测。排队不显示为正在执行;收到表情保留到最终结果 清理失败保留原回复的恢复路径,不重复发送已核验答案。收到与处理中共享 prepared/created journal:已知 provider id 恢复回执,写入结果不确定时不盲重试; 分页未读完不能证明删除成功。缺少表情权限不取消已受理工作,也不提高宿主策略。 +私聊投递最多并行处理 4 条持久请求,没有 executor 内存排队。一个回复读回阻塞 +不再挡住独立的排队提示或控制反馈;每条请求的锁与 journal 保留恢复不重发语义。 +Core 只恢复确切请求,沿用原 Session fence。Provider 核验成本与 worker 饱和时延 +仍需独立测量,合成阻塞测试不证明实时 SLO,也不增加模型并发授权。 + 这是现有 provider 呈现边界的有界重构,不新增 Session、queue、model runner 或 控制面 authority。合成回归与真实 provider canary 分别记录;增量卡片、媒体和 权限回调仍需独立验收。 diff --git a/loopx/capabilities/native_chat/external_conversations.py b/loopx/capabilities/native_chat/external_conversations.py index 3f50d3646a..2a98ae962d 100644 --- a/loopx/capabilities/native_chat/external_conversations.py +++ b/loopx/capabilities/native_chat/external_conversations.py @@ -305,8 +305,9 @@ def record_delivery(self, request_ref: str, *, session_id: str | None, turn_id: row["delivery_verified"] = True _atomic_write_json(path, row) - def recover(self) -> None: - for row in self.pending(): + def recover(self, *, request_ref: str | None = None) -> None: + rows = self.pending() if request_ref is None else [self.read_request(request_ref)] + for row in rows: try: if row.get("delivery_verified"): continue diff --git a/loopx/extensions/lark/goal_topic_runtime_service.py b/loopx/extensions/lark/goal_topic_runtime_service.py index 9f4d007403..8e543444d5 100644 --- a/loopx/extensions/lark/goal_topic_runtime_service.py +++ b/loopx/extensions/lark/goal_topic_runtime_service.py @@ -3,6 +3,7 @@ from __future__ import annotations from collections.abc import Callable, Mapping +from concurrent.futures import Future, ThreadPoolExecutor from datetime import datetime, timezone import hashlib import logging @@ -73,12 +74,37 @@ def start(self) -> None: self._delivery_thread.start() def _reconcile_private(self) -> None: - while not self._closed.is_set(): - try: - self.private_conversations.reconcile() - except Exception: - logging.getLogger(__name__).warning("Private Chat delivery reconciliation is pending") - self._closed.wait(1) + # Keep slow provider writes/readbacks off the scan loop. At most four + # persisted requests run, with no executor backlog or extra listener. + # The per-request journal lock retains ambiguous-write/no-resend safety. + transport = self.private_conversations + if transport is None: + return + pending: dict[Path, Future[int]] = {} + attempted: dict[Path, int] = {} + with ThreadPoolExecutor(max_workers=4, thread_name_prefix="loopx-lark-private-reply") as pool: + while not self._closed.is_set(): + for path, future in list(pending.items()): + if not future.done(): + continue + pending.pop(path) + try: + future.result() + except Exception: + logging.getLogger(__name__).warning("Private Chat delivery reconciliation is pending") + try: + paths = transport.pending_delivery_paths() + attempted = {path: count for path, count in attempted.items() if path in paths} + for path in sorted(paths, key=lambda path: attempted.get(path, 0)): + if len(pending) >= 4 or self._closed.is_set(): + break + if path in pending: + continue + pending[path] = pool.submit(transport.reconcile_request, path) + attempted[path] = attempted.get(path, 0) + 1 + except Exception: + logging.getLogger(__name__).warning("Private Chat delivery discovery is pending") + self._closed.wait(1) def _refresh_on_start(self) -> None: while not self._closed.is_set(): diff --git a/loopx/extensions/lark/private_conversations.py b/loopx/extensions/lark/private_conversations.py index 61dd80e94b..b8730b9a9c 100644 --- a/loopx/extensions/lark/private_conversations.py +++ b/loopx/extensions/lark/private_conversations.py @@ -6,7 +6,7 @@ from __future__ import annotations import json -from collections.abc import Mapping +from collections.abc import Mapping, Sequence from pathlib import Path from typing import Any @@ -193,7 +193,7 @@ def admit(self, profile: str, event: dict[str, Any]) -> dict[str, Any]: "queue_rejected" if admitted["status"] == "rejected" else "command_recorded") return {"status": status} - def _reply_runner(self, args: list[str]) -> Any: + def _reply_runner(self, args: Sequence[str]) -> Any: return self.runner([self.cli_bin, *args[1:]], None, 30) def _feedback(self, path: Path, record: dict[str, Any], *, processing: bool = False) -> None: @@ -259,93 +259,100 @@ def attempt(value: Any) -> None: _atomic_write_json(path, record) return bool(phase_state["verified"]) + def pending_delivery_paths(self) -> list[Path]: + """The durable transport store is the queue; no in-memory admission.""" + return [path for path in sorted(self.root.glob("*.json")) + if _read_json(path)["status"] != "delivered"] + def reconcile(self) -> int: - self.core.recover() - # A crash after Core admission but before the transport correlation - # record is repaired by replaying the same message identity. - for pending in sorted(self.root.glob("*.json")): - row = _read_json(pending) - if row["status"] == "captured": - self.admit(row["profile"], row["event"]) - delivered = 0 - for path in sorted(self.root.glob("*.json")): - with exclusive_file_lock(path, operation="deliver_private_chat_request"): - record = _read_json(path) - if record["status"] in {"captured", "delivered"}: - continue - try: - self.bindings.resolve(binding_id=record["binding_id"], **record["source"]) - if record["status"] in {"command_queued", "command_completed", "commission_running"}: - native = self.core.read_request(record["request_ref"]) - if native["status"] == "command_queued": - self._feedback(path, record) - self._deliver(path, record, "admission", native["response"]) - continue - record.update(status=native["status"], response=native.get("response"), - commission_resources=native.get("commission_resources"), - status_snapshot=native.get("status_snapshot")) - if record["status"] == "accepted": - native = self.core.read_request(record["request_ref"]) + return sum(self.reconcile_request(path) for path in self.pending_delivery_paths()) + + def reconcile_request(self, path: Path) -> int: + # Repair only this persisted source. Core still owns admission, recovery + # and the shared Session queue/fence; transport workers never run models. + row = _read_json(path) + if row["status"] == "captured": + self.admit(row["profile"], row["event"]) + native_path = self.core.root / f"{row['request_ref']}.json" + if native_path.exists(): + self.core.recover(request_ref=row["request_ref"]) + with exclusive_file_lock(path, operation="deliver_private_chat_request"): + record = _read_json(path) + if record["status"] in {"captured", "delivered"}: + return 0 + try: + self.bindings.resolve(binding_id=record["binding_id"], **record["source"]) + if record["status"] in {"command_queued", "command_completed", "commission_running"}: + native = self.core.read_request(record["request_ref"]) + if native["status"] == "command_queued": + self._feedback(path, record) + self._deliver(path, record, "admission", native["response"]) + return 0 + record.update(status=native["status"], response=native.get("response"), + commission_resources=native.get("commission_resources"), + status_snapshot=native.get("status_snapshot")) + if record["status"] == "accepted": + native = self.core.read_request(record["request_ref"]) + if native.get("agent_target"): + self.bindings.resolve_agent_target(self.bindings.resolve(binding_id=record["binding_id"], **record["source"]), native["agent_target"]) + self._feedback(path, record) + # Receipt follows persistent Core admission and is + # independent of terminal execution and reply delivery. + turn = self.core.controller.store.load_turn(record["session_id"], record["turn_id"]) + # Get already acknowledges admission. Reserve a text + # notification for a real wait or unavailable feedback; + # resume any old attempt with its original exact text. + admission = (record["deliveries"].get("admission") or {}).get("text") + if not admission: if native.get("agent_target"): - self.bindings.resolve_agent_target(self.bindings.resolve(binding_id=record["binding_id"], **record["source"]), native["agent_target"]) + admission = "已收到,等待原 Agent 宿主处理。/status 查看状态,/project 返回项目对话。" + elif turn and turn["status"] == "queued": + admission = "正在处理前一条,这条已排队。" + elif (record.get("feedback", {}).get("received") or {}).get("ok") is not True: + admission = "已收到,正在处理。" + if admission: + self._deliver(path, record, "admission", admission) + if turn and turn["status"] in {"starting", "running"}: + self._feedback(path, record, processing=True) + if not turn or turn["status"] not in {"completed", "failed", "interrupted", "expired"}: + return 0 + response = str((turn.get("response") or {}).get("message") or "") if turn["status"] == "completed" else ( + "本次执行已停止。" if turn["status"] == "interrupted" else + "本次执行失败或已过期,原会话已保留;请发送 /status 后再决定是否重试。") + else: + if record["status"] != "rejected": self._feedback(path, record) - # Receipt follows persistent Core admission and is - # independent of terminal execution and reply delivery. - turn = self.core.controller.store.load_turn(record["session_id"], record["turn_id"]) - # Get already acknowledges admission. Reserve a text - # notification for a real wait or unavailable feedback; - # resume any old attempt with its original exact text. - admission = (record["deliveries"].get("admission") or {}).get("text") - if not admission: - if native.get("agent_target"): - admission = "已收到,等待原 Agent 宿主处理。/status 查看状态,/project 返回项目对话。" - elif turn and turn["status"] == "queued": - admission = "正在处理前一条,这条已排队。" - elif (record.get("feedback", {}).get("received") or {}).get("ok") is not True: - admission = "已收到,正在处理。" - if admission: - self._deliver(path, record, "admission", admission) - if turn and turn["status"] in {"starting", "running"}: + response = _command_text(str(record.get("response_code") or "")) or str(record.get("response") or "") + if record.get("status_snapshot"): + response = _status_text(record["status_snapshot"], help_requested=native.get("command") == "help") + if self._deliver(path, record, "terminal", response): + resources = record.get("commission_resources") or {} + if resources.get("session_id") and resources.get("turn_id"): + first_turn = self.core.controller.store.load_turn(resources["session_id"], resources["turn_id"]) + if first_turn and first_turn["status"] in {"starting", "running"}: self._feedback(path, record, processing=True) - if not turn or turn["status"] not in {"completed", "failed", "interrupted", "expired"}: - continue - response = str((turn.get("response") or {}).get("message") or "") if turn["status"] == "completed" else ( - "本次执行已停止。" if turn["status"] == "interrupted" else - "本次执行失败或已过期,原会话已保留;请发送 /status 后再决定是否重试。") - else: - if record["status"] != "rejected": - self._feedback(path, record) - response = _command_text(str(record.get("response_code") or "")) or str(record.get("response") or "") - if record.get("status_snapshot"): - response = _status_text(record["status_snapshot"], help_requested=native.get("command") == "help") - if self._deliver(path, record, "terminal", response): - resources = record.get("commission_resources") or {} - if resources.get("session_id") and resources.get("turn_id"): - first_turn = self.core.controller.store.load_turn(resources["session_id"], resources["turn_id"]) - if first_turn and first_turn["status"] in {"starting", "running"}: - self._feedback(path, record, processing=True) - if not first_turn or first_turn["status"] not in {"completed", "failed", "interrupted", "expired"}: - record["status"] = "commission_running" - _atomic_write_json(path, record) - continue - result_text = str((first_turn.get("response") or {}).get("message") or "") - result_text = "委托执行结果:\n" + result_text if first_turn["status"] == "completed" else "委托首轮执行未完成;原 Goal 和回执已保留,请查看状态后决定恢复。" - proposal_id = native.get("proposal_id") - if proposal_id: - result_text += f"\n如需恢复暂停或额度受限的原执行:/resume-commission {proposal_id} --tokens N(N 为包含历史用量的总上限,须大于已用量;不会重开线程)。" - if not self._deliver(path, record, "commission_result", result_text): - continue - config = self._inbox(record) - acknowledge_lark_event_inbox(project=self.runtime_root, config_path=config, - message_ids=[record["event"]["message_id"]], execute=True) - self.core.record_delivery(record["request_ref"], session_id=record.get("session_id"), turn_id=record.get("turn_id")) - record["status"] = "delivered" - _atomic_write_json(path, record) - delivered += 1 - except (KeyError, ValueError, OSError): - # Keep pending receipts available for actionable recovery. - continue - return delivered + if not first_turn or first_turn["status"] not in {"completed", "failed", "interrupted", "expired"}: + record["status"] = "commission_running" + _atomic_write_json(path, record) + return 0 + result_text = str((first_turn.get("response") or {}).get("message") or "") + result_text = "委托执行结果:\n" + result_text if first_turn["status"] == "completed" else "委托首轮执行未完成;原 Goal 和回执已保留,请查看状态后决定恢复。" + proposal_id = native.get("proposal_id") + if proposal_id: + result_text += f"\n如需恢复暂停或额度受限的原执行:/resume-commission {proposal_id} --tokens N(N 为包含历史用量的总上限,须大于已用量;不会重开线程)。" + if not self._deliver(path, record, "commission_result", result_text): + return 0 + config = self._inbox(record) + acknowledge_lark_event_inbox(project=self.runtime_root, config_path=config, + message_ids=[record["event"]["message_id"]], execute=True) + self.core.record_delivery(record["request_ref"], session_id=record.get("session_id"), turn_id=record.get("turn_id")) + record["status"] = "delivered" + _atomic_write_json(path, record) + return 1 + except (KeyError, ValueError, OSError): + # Keep pending receipts available for actionable recovery. + return 0 + return 0 def _command_text(code: str) -> str: diff --git a/tests/test_lark_private_conversations.py b/tests/test_lark_private_conversations.py index fdedbe3235..8e4c174285 100644 --- a/tests/test_lark_private_conversations.py +++ b/tests/test_lark_private_conversations.py @@ -284,3 +284,107 @@ def observe(profile): assert transport.admit("steward-app", rejected)["status"] == "audience_rejected" finally: runtime.close() + + +def test_slow_reply_readback_does_not_hold_independent_stop_or_other_app(ordinary, monkeypatch): # noqa: F811 + import threading + from loopx.extensions.lark.goal_topic_runtime_service import LarkGoalTopicRuntimeService + + store, runtime, provider, transport = connect(ordinary) + blocked, release = threading.Event(), threading.Event() + original = Provider.__call__ + + def delayed(self, args, cwd=None, timeout=None): + if "+messages-mget" in args and args[args.index("--message-ids") + 1] == "om_out_0": + blocked.set() + assert release.wait(10) + return original(self, args, cwd, timeout) + + monkeypatch.setattr(Provider, "__call__", delayed) + transport.admit("notes-app", provider.event("notes-app", "old-status", "/status")) + service = LarkGoalTopicRuntimeService(snapshot_provider=lambda: {}, runtime_root=store.root.parent, + runtime_controller=runtime, private_conversations=transport) + worker = threading.Thread(target=service._reconcile_private) + worker.start() + try: + assert blocked.wait(10) + transport.admit("notes-app", provider.event("notes-app", "urgent-stop", "/stop")) + transport.admit("steward-app", provider.event("steward-app", "other-status", "/status")) + deadline = time.monotonic() + 5 + while len(provider.writes) < 3: + assert time.monotonic() < deadline + time.sleep(.01) + # The old provider read is still blocked, yet both exact audiences have + # received feedback. No Session/model was needed for these controls. + assert not release.is_set() + assert ("notes-app", "当前没有正在执行的消息。") in provider.writes + assert any(profile == "steward-app" and "角色:普通项目对话" in text + for profile, text in provider.writes) + assert store.list_sessions() == [] + finally: + release.set() + service._closed.set() + worker.join(10) + runtime.close() + assert not worker.is_alive() + writes = list(provider.writes) + transport.reconcile() + assert provider.writes == writes # Recovery does not resend an attempted reply. + + +def test_private_reply_workers_have_no_executor_backlog_and_stop_scheduling(tmp_path): + import threading + from loopx.extensions.lark.goal_topic_runtime_service import LarkGoalTopicRuntimeService + + release, full = threading.Event(), threading.Event() + lock = threading.Lock() + calls = [] + + class Pending: + def pending_delivery_paths(self): + return [tmp_path / str(index) for index in range(9)] + + def reconcile_request(self, path): + with lock: + calls.append(path) + if len(calls) == 4: + full.set() + assert release.wait(10) + return 0 + + service = LarkGoalTopicRuntimeService(snapshot_provider=lambda: {}, runtime_root=tmp_path, + runtime_controller=None, private_conversations=Pending()) + worker = threading.Thread(target=service._reconcile_private) + worker.start() + try: + assert full.wait(5) + assert len(calls) == 4 and len(set(calls)) == 4 + service._closed.set() + release.set() + worker.join(5) + assert not worker.is_alive() + assert len(calls) == 4 # Five durable requests remain; none were queued in memory. + finally: + service._closed.set() + release.set() + worker.join(10) + + +def test_scoped_recovery_does_not_probe_unrelated_app_and_revocation_blocks_reply(ordinary): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + try: + for profile in ["notes-app", "steward-app"]: + transport.admit(profile, provider.event(profile, profile, "/status")) + first = next(row for row in transport.core.pending() + if row["binding_id"] == transport.bindings.read()["bindings"][0]["binding_id"]) + observed = [] + original = transport.bindings.observe + transport.bindings.observe = lambda profile: (observed.append(profile), original(profile))[1] + transport.core.recover(request_ref=first["request_ref"]) + assert observed == ["notes-app"] + binding = transport.bindings.read() + transport.bindings.disconnect(first["binding_id"], expected_revision=binding["revision"]) + assert transport.reconcile_request(transport.root / f"{first['request_ref']}.json") == 0 + assert provider.writes == [] and store.list_sessions() == [] + finally: + runtime.close() From cbd7c00301264aa7be3346f8554601151480e706 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sun, 4 Oct 2026 19:06:46 +0800 Subject: [PATCH 11/12] Reuse operation-local private inbox preparation without weakening grants Signed-off-by: huangruiteng --- .../app-conversation-and-async-inbox-v0.md | 4 ++ ...p-conversation-and-async-inbox-v0.zh-CN.md | 4 ++ .../native_chat/external_conversations.py | 1 - .../extensions/lark/private_conversations.py | 49 ++++++++++++------- tests/test_lark_private_conversations.py | 22 +++++++++ 5 files changed, 60 insertions(+), 20 deletions(-) diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md index b59152ab7d..2a1106cbc8 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -107,6 +107,10 @@ notice or control response; per-request locks and journals retain no-resend recovery. Core recovers each exact request under the existing Session fence. Provider verification cost and saturated-worker latency remain separate from model concurrency and are not certified by a synthetic blocked-readback test. +Admission reuses its initial App observation for source reading, then rechecks +authority inside the Core source fence before creating work. Each locked delivery +reuses its Inbox configuration; provider writes still verify current identity, +grants and the exact original message. No observation is cached across requests. This is a bounded provider presentation refactor, not a new Session, queue, model runner or control-plane owner. Native state drives both direct conversations diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md index b50629817b..3a7a7c2881 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md @@ -78,6 +78,10 @@ RPC 回执。Lark 专属设置 companion 归 extension;会话、请求、范 ## 原生私聊反馈:默认开启与投递生命周期 +受理复用本次初始 App 观测读取原消息,Core 在来源 fence 内再次核验授权后才创建 +执行。同一加锁投递复用 Inbox 配置;发送前仍核验当前身份、授权与确切原消息。 +观测不跨请求缓存。 + 通用个人助手的最终回复复用 Markdown post,保留列表、引用、公开链接和代码; 本地链接脱敏时保留可读标签,不生成指向 `[project]` 的假链接。旧纯文本投递 仍按原 attempt 读回恢复,不因新默认而重发。图片/文件交付与真实增量仍需独立验收。 diff --git a/loopx/capabilities/native_chat/external_conversations.py b/loopx/capabilities/native_chat/external_conversations.py index 2a98ae962d..77cf27ca78 100644 --- a/loopx/capabilities/native_chat/external_conversations.py +++ b/loopx/capabilities/native_chat/external_conversations.py @@ -30,7 +30,6 @@ def admit(self, *, binding_id: str, source: dict[str, Any], request_ref: str, raise ValueError("invalid external request reference") if command not in {None, "agents", "select_agent", "select_project", "status", "help", "new", "stop", "unsupported", "commission", "confirm_commission", "cancel_commission", "stop_commission", "resume_commission"}: raise ValueError("unsupported external conversation command") - selected = self.bindings.resolve(binding_id=binding_id, **source) path = self.root / f"{request_ref}.json" with exclusive_file_lock(self.root / "source-fences" / f"{binding_id}.{source['source_ref']}.json", operation="route_external_chat_request"), exclusive_file_lock(path, operation="admit_external_chat_request"): selected = self.bindings.resolve(binding_id=binding_id, **source) diff --git a/loopx/extensions/lark/private_conversations.py b/loopx/extensions/lark/private_conversations.py index b8730b9a9c..fe89504a61 100644 --- a/loopx/extensions/lark/private_conversations.py +++ b/loopx/extensions/lark/private_conversations.py @@ -6,7 +6,7 @@ from __future__ import annotations import json -from collections.abc import Mapping, Sequence +from collections.abc import Callable, Mapping, Sequence from pathlib import Path from typing import Any @@ -58,7 +58,7 @@ def health(self) -> dict[str, dict[str, int]]: def _binding(self, profile: str) -> dict[str, Any]: return next(row for row in self.bindings.read()["bindings"] if row["transport_ref"] == profile) - def _source_message(self, record: dict[str, Any]) -> Mapping[str, Any] | None: + def _source_message(self, record: dict[str, Any], *, selected: dict[str, Any] | None = None) -> Mapping[str, Any] | None: """Read the exact source under this App, then recheck the Core audience. A p2p source read proves the destination without requiring group-member @@ -66,7 +66,8 @@ def _source_message(self, record: dict[str, Any]) -> Mapping[str, Any] | None: """ event = record["event"] try: - selected = self.bindings.resolve(binding_id=record["binding_id"], **record["source"]) + if selected is None: + selected = self.bindings.resolve(binding_id=record["binding_id"], **record["source"]) native_path = self.core.root / f"{record['request_ref']}.json" native = _read_json(native_path) if native_path.exists() else {} if native.get("agent_target"): @@ -113,7 +114,7 @@ def admit(self, profile: str, event: dict[str, Any]) -> dict[str, Any]: try: binding = self._binding(profile) source = lark_private_source(provider_ref=binding["provider_ref"], event=event) - self.bindings.resolve(binding_id=binding["binding_id"], **source) + selected = self.bindings.resolve(binding_id=binding["binding_id"], **source) except (KeyError, StopIteration, ValueError): return {"status": "audience_rejected"} request = identity_ref(binding["provider_ref"], event["message_id"]) @@ -131,7 +132,9 @@ def admit(self, profile: str, event: dict[str, Any]) -> dict[str, Any]: "profile": profile, "binding_id": binding["binding_id"], "source": source, "event": event, "deliveries": {}, "status": "captured"} _atomic_write_json(path, record) - source_message = self._source_message(record) + # Reuse this admission preflight only; Core rechecks under its + # source fence, and outbound writes always perform a fresh check. + source_message = self._source_message(record, selected=selected) if source_message is None: return {"status": "source_verification_failed"} message_type = str(event.get("message_type") or "") @@ -196,7 +199,7 @@ def admit(self, profile: str, event: dict[str, Any]) -> dict[str, Any]: def _reply_runner(self, args: Sequence[str]) -> Any: return self.runner([self.cli_bin, *args[1:]], None, 30) - def _feedback(self, path: Path, record: dict[str, Any], *, processing: bool = False) -> None: + def _feedback(self, path: Path, record: dict[str, Any], *, inbox: Callable[[], Path], processing: bool = False) -> None: """Render observed Core admission/execution through the shared Inbox owner. Presentation failures never reject a persisted Turn. Existing reaction @@ -208,7 +211,7 @@ def _feedback(self, path: Path, record: dict[str, Any], *, processing: bool = Fa feedback = record.setdefault("feedback", {}) if feedback.get(phase, {}).get("ok"): return - config = self._inbox(record) + config = inbox() if processing: result = mark_lark_event_inbox_processing(project=self.runtime_root, config_path=config, message_id=record["event"]["message_id"], execute=True, runner=self._reply_runner) @@ -218,7 +221,7 @@ def _feedback(self, path: Path, record: dict[str, Any], *, processing: bool = Fa feedback[phase] = result _atomic_write_json(path, record) - def _deliver(self, path: Path, record: dict[str, Any], phase: str, text: str) -> bool: + def _deliver(self, path: Path, record: dict[str, Any], phase: str, text: str, *, inbox: Callable[[], Path]) -> bool: if not text: return False phase_state = record["deliveries"].get(phase, {}) @@ -226,7 +229,7 @@ def _deliver(self, path: Path, record: dict[str, Any], phase: str, text: str) -> return True if phase_state.get("text") not in (None, text): raise ValueError("delivery content changed after an attempt") - config = self._inbox(record) + config = inbox() kwargs = dict(project=self.runtime_root, config_path=config, message_id=record["event"]["message_id"], text=text, runner=self._reply_runner, finalize_reactions=phase != "admission" and not (phase == "terminal" and record.get("commission_resources")), @@ -280,13 +283,21 @@ def reconcile_request(self, path: Path) -> int: record = _read_json(path) if record["status"] in {"captured", "delivered"}: return 0 + config: Path | None = None + + def inbox() -> Path: + nonlocal config + if config is None: + config = self._inbox(record) + return config + try: self.bindings.resolve(binding_id=record["binding_id"], **record["source"]) if record["status"] in {"command_queued", "command_completed", "commission_running"}: native = self.core.read_request(record["request_ref"]) if native["status"] == "command_queued": - self._feedback(path, record) - self._deliver(path, record, "admission", native["response"]) + self._feedback(path, record, inbox=inbox) + self._deliver(path, record, "admission", native["response"], inbox=inbox) return 0 record.update(status=native["status"], response=native.get("response"), commission_resources=native.get("commission_resources"), @@ -295,7 +306,7 @@ def reconcile_request(self, path: Path) -> int: native = self.core.read_request(record["request_ref"]) if native.get("agent_target"): self.bindings.resolve_agent_target(self.bindings.resolve(binding_id=record["binding_id"], **record["source"]), native["agent_target"]) - self._feedback(path, record) + self._feedback(path, record, inbox=inbox) # Receipt follows persistent Core admission and is # independent of terminal execution and reply delivery. turn = self.core.controller.store.load_turn(record["session_id"], record["turn_id"]) @@ -311,9 +322,9 @@ def reconcile_request(self, path: Path) -> int: elif (record.get("feedback", {}).get("received") or {}).get("ok") is not True: admission = "已收到,正在处理。" if admission: - self._deliver(path, record, "admission", admission) + self._deliver(path, record, "admission", admission, inbox=inbox) if turn and turn["status"] in {"starting", "running"}: - self._feedback(path, record, processing=True) + self._feedback(path, record, inbox=inbox, processing=True) if not turn or turn["status"] not in {"completed", "failed", "interrupted", "expired"}: return 0 response = str((turn.get("response") or {}).get("message") or "") if turn["status"] == "completed" else ( @@ -321,16 +332,16 @@ def reconcile_request(self, path: Path) -> int: "本次执行失败或已过期,原会话已保留;请发送 /status 后再决定是否重试。") else: if record["status"] != "rejected": - self._feedback(path, record) + self._feedback(path, record, inbox=inbox) response = _command_text(str(record.get("response_code") or "")) or str(record.get("response") or "") if record.get("status_snapshot"): response = _status_text(record["status_snapshot"], help_requested=native.get("command") == "help") - if self._deliver(path, record, "terminal", response): + if self._deliver(path, record, "terminal", response, inbox=inbox): resources = record.get("commission_resources") or {} if resources.get("session_id") and resources.get("turn_id"): first_turn = self.core.controller.store.load_turn(resources["session_id"], resources["turn_id"]) if first_turn and first_turn["status"] in {"starting", "running"}: - self._feedback(path, record, processing=True) + self._feedback(path, record, inbox=inbox, processing=True) if not first_turn or first_turn["status"] not in {"completed", "failed", "interrupted", "expired"}: record["status"] = "commission_running" _atomic_write_json(path, record) @@ -340,9 +351,9 @@ def reconcile_request(self, path: Path) -> int: proposal_id = native.get("proposal_id") if proposal_id: result_text += f"\n如需恢复暂停或额度受限的原执行:/resume-commission {proposal_id} --tokens N(N 为包含历史用量的总上限,须大于已用量;不会重开线程)。" - if not self._deliver(path, record, "commission_result", result_text): + if not self._deliver(path, record, "commission_result", result_text, inbox=inbox): return 0 - config = self._inbox(record) + config = inbox() acknowledge_lark_event_inbox(project=self.runtime_root, config_path=config, message_ids=[record["event"]["message_id"]], execute=True) self.core.record_delivery(record["request_ref"], session_id=record.get("session_id"), turn_id=record.get("turn_id")) diff --git a/tests/test_lark_private_conversations.py b/tests/test_lark_private_conversations.py index 8e4c174285..db00a37adf 100644 --- a/tests/test_lark_private_conversations.py +++ b/tests/test_lark_private_conversations.py @@ -167,6 +167,28 @@ def test_existing_consumer_dispatches_private_admission_without_waiting_for_mode assert admitted == [("notes-app", event)] +def test_revocation_during_source_read_prevents_native_admission(ordinary): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + try: + event = provider.event("notes-app", "revoked", "must not execute") + binding = transport.bindings.read()["bindings"][0] + + def revoke_after_read(args, cwd=None, timeout=None): + result = provider(args, cwd, timeout) + if "+messages-mget" in args: + transport.bindings.disconnect(binding["binding_id"], + expected_revision=transport.bindings.read()["revision"]) + return result + + transport.runner = revoke_after_read + assert transport.admit("notes-app", event)["status"] == "command_rejected" + assert transport.core.pending() == [] + assert store.list_sessions() == [] + assert provider.writes == [] + finally: + runtime.close() + + def test_new_session_replay_cannot_close_a_later_session(ordinary): # noqa: F811 store, runtime, provider, transport = connect(ordinary) try: From d1015be969a51232691eaceccd4c000d8ce86e5c Mon Sep 17 00:00:00 2001 From: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> Date: Sun, 4 Oct 2026 19:56:45 +0800 Subject: [PATCH 12/12] Separate private listener discovery from message authorization Signed-off-by: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> --- .../app-conversation-and-async-inbox-v0.md | 10 ++ ...p-conversation-and-async-inbox-v0.zh-CN.md | 7 ++ .../extensions/lark/private_conversations.py | 37 ++++-- tests/test_lark_private_conversations.py | 118 +++++++++++++++++- 4 files changed, 159 insertions(+), 13 deletions(-) diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md index 2a1106cbc8..46eb930f54 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -112,6 +112,16 @@ authority inside the Core source fence before creating work. Each locked deliver reuses its Inbox configuration; provider writes still verify current identity, grants and the exact original message. No observation is cached across requests. +Listener discovery reads the CLI's local profile inventory once per snapshot, +compares the current App with its Core binding, and derives the existing +machine/App lease from that App. Network authorization failures are not durable +disconnects. An unreadable inventory preserves a running stream; an explicit +binding removal, missing profile or replacement App stops its old route. This +does not authorize a message: admission and outbound delivery still freshly +verify the App, owner, source and grants. Actual stream fault/recovery and +disconnect regressions qualify this boundary; sustained live availability and +timely end-to-end feedback remain separate acceptance. + This is a bounded provider presentation refactor, not a new Session, queue, model runner or control-plane owner. Native state drives both direct conversations and explicit commissions. Synthetic queue/stop/replay/isolation regressions and diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md index 3a7a7c2881..2b4e27516c 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md @@ -102,6 +102,13 @@ prepared/created journal:已知 provider id 恢复回执,写入结果不确 Core 只恢复确切请求,沿用原 Session fence。Provider 核验成本与 worker 饱和时延 仍需独立测量,合成阻塞测试不证明实时 SLO,也不增加模型并发授权。 +Listener 每次快照只读一次 CLI 本地 profile 清单,核对当前 App 与 Core binding, +由实际 App 推导现有的机器/App 消费者锁。网络授权探测失败不再被当成永久解绑。 +本地清单暂时不可读时保留运行中的 stream;明确解绑、profile 移除或替换 App +则停止旧路由。这不授予消息权限:入站受理和出站投递仍重新核验 App、本人、来源 +与 grants。真实 stream 的故障/恢复及解绑回归覆盖这一边界;持续 live 可用性与 +端到端及时反馈仍需独立验收。 + 这是现有 provider 呈现边界的有界重构,不新增 Session、queue、model runner 或 控制面 authority。合成回归与真实 provider canary 分别记录;增量卡片、媒体和 权限回调仍需独立验收。 diff --git a/loopx/extensions/lark/private_conversations.py b/loopx/extensions/lark/private_conversations.py index fe89504a61..c6e4b71dd5 100644 --- a/loopx/extensions/lark/private_conversations.py +++ b/loopx/extensions/lark/private_conversations.py @@ -5,6 +5,7 @@ """ from __future__ import annotations +import hashlib import json from collections.abc import Callable, Mapping, Sequence from pathlib import Path @@ -15,7 +16,7 @@ from ...file_lock import exclusive_file_lock from .conversation_identity import identity_ref, lark_private_source from .event_inbox import acknowledge_lark_event_inbox, ingest_lark_event_inbox -from .goal_channel_transport import call, json_payload, lark_args +from .goal_channel_transport import APP_ID_PATTERN, call, json_payload, lark_args from .inbox_reply import _message, reply_lark_event_inbox, verify_lark_inbox_reply from .inbox_reactions import mark_lark_event_inbox_processing, mark_lark_event_inbox_received @@ -30,18 +31,34 @@ def __init__(self, *, controller: Any, runtime_root: Path, runner: Any, cli_bin: self.root = controller.store.root / "lark-private-deliveries" def profiles(self) -> dict[str, dict[str, str]]: + """Discover current App configuration, never authorize a conversation. + + A network/token probe failing is not a durable disconnect. The local + profile inventory supplies the actual App lease identity; Core still + freshly verifies App/owner/source/grants at admission and delivery. + """ + bindings = self.bindings.read()["bindings"] + if not bindings: + return {} + result = call(self.runner, [self.cli_bin, "profile", "list"]) + try: + configured = json.loads(str(result.get("stdout") or "")) + except json.JSONDecodeError as exc: + raise OSError("Lark profile configuration could not be read") from exc + if result.get("returncode") != 0 or not isinstance(configured, list): + # An unreadable inventory is unknown, not an empty configuration. + # The existing stream watcher retains its route on read failures. + raise OSError("Lark profile configuration could not be read") + apps = {str(row.get("name") or ""): str(row.get("appId") or "") + for row in configured if isinstance(row, Mapping)} profiles = {} - for row in self.bindings.read()["bindings"]: - try: - observation = self.bindings.observe(row["transport_ref"]) - except (ValueError, OSError): - # One App's expired login must not block the independently - # verified App from acquiring its own listener lease. - continue - if any(observation.get(field) != row[field] for field in ["provider_ref", "operator_ref"]): + for row in bindings: + app_id = apps.get(row["transport_ref"], "") + if (row.get("enabled") is not True or not APP_ID_PATTERN.fullmatch(app_id) + or identity_ref(app_id) != row["provider_ref"]): continue profiles[row["transport_ref"]] = {"cli_bin": self.cli_bin, "provider_ref": row["provider_ref"], - "binding_id": row["binding_id"], "consumer_ref": observation["consumer_ref"]} + "binding_id": row["binding_id"], "consumer_ref": hashlib.sha256(app_id.encode("utf-8")).hexdigest()[:32]} return profiles def health(self) -> dict[str, dict[str, int]]: diff --git a/tests/test_lark_private_conversations.py b/tests/test_lark_private_conversations.py index db00a37adf..77e3228ae0 100644 --- a/tests/test_lark_private_conversations.py +++ b/tests/test_lark_private_conversations.py @@ -22,6 +22,8 @@ def __init__(self): self.reaction_creates = [] self.fail_reaction_create = False self.fail_reaction_delete = False + self.profile_apps = {profile: f"cli_{profile.replace('-', '_')}" + for profile in ["notes-app", "steward-app"]} def event(self, profile, name, text, kind="text"): event = {"schema_version": "lark_event_inbox_event_v0", "event_id": f"event_{name}", @@ -34,6 +36,10 @@ def event(self, profile, name, text, kind="text"): def __call__(self, args, cwd=None, timeout=None): self.calls.append(list(args)) + if args[1:] == ["profile", "list"]: + return {"returncode": 0, "stdout": json.dumps([ + {"name": profile, "appId": app_id} + for profile, app_id in self.profile_apps.items()]), "stderr": ""} profile = args[args.index("--profile") + 1] if "auth" in args: data = {"ok": True, "appId": f"cli_{profile.replace('-', '_')}", "identities": { @@ -288,7 +294,7 @@ def unverified(*_args, **_kwargs): _app_identity_for_private_guard("notes-app", unverified, "lark-cli") -def test_one_expired_app_identity_does_not_hide_the_other_listener(ordinary): # noqa: F811 +def test_listener_discovery_is_separate_from_current_owner_authorization(ordinary): # noqa: F811 _, runtime, _, transport = connect(ordinary) original = transport.bindings.observe def observe(profile): @@ -297,17 +303,123 @@ def observe(profile): return original(profile) transport.bindings.observe = observe try: - assert set(transport.profiles()) == {"notes-app"} + provider = transport.runner + provider.calls.clear() + assert set(transport.profiles()) == {"notes-app", "steward-app"} + assert provider.calls == [["lark-cli", "profile", "list"]] transport.bindings.observe = lambda profile: {**original(profile), "operator_ref": "f" * 24} - assert transport.profiles() == {} + assert set(transport.profiles()) == {"notes-app", "steward-app"} transport.bindings.observe = observe rejected = {"message_id": "om_unknown", "chat_id": "oc_steward_app", "sender_id": "ou_steward_app", "chat_type": "p2p", "sender_type": "user", "message_type": "text", "content": "unverified"} assert transport.admit("steward-app", rejected)["status"] == "audience_rejected" + assert transport.admit("notes-app", provider.event("notes-app", "healthy", "/status"))["status"] == "command_recorded" finally: runtime.close() +def test_listener_lease_uses_current_app_and_detects_removal_or_retarget(ordinary): # noqa: F811 + import hashlib + + _, runtime, provider, transport = connect(ordinary) + try: + profiles = transport.profiles() + assert profiles["notes-app"]["consumer_ref"] == hashlib.sha256(b"cli_notes_app").hexdigest()[:32] + # A new App under the old profile must not run under the old App's lease. + provider.profile_apps["notes-app"] = "cli_replacement" + assert set(transport.profiles()) == {"steward-app"} + del provider.profile_apps["notes-app"] + assert set(transport.profiles()) == {"steward-app"} + provider.profile_apps["notes-app"] = "cli_notes_app" + binding = transport._binding("notes-app") + transport.bindings.disconnect(binding["binding_id"], expected_revision=transport.bindings.read()["revision"]) + assert set(transport.profiles()) == {"steward-app"} + finally: + runtime.close() + + +@pytest.mark.parametrize("result", [ + {"returncode": 1, "stdout": "", "timed_out": True}, + {"returncode": 0, "stdout": "not JSON"}, + {"returncode": 0, "stdout": "{}"}, +]) +def test_unreadable_local_profile_inventory_is_unknown_not_removed(ordinary, result): # noqa: F811 + _, runtime, _, transport = connect(ordinary) + try: + transport.runner = lambda *_args: result + with pytest.raises(OSError, match="configuration could not be read"): + transport.profiles() + finally: + runtime.close() + + +def test_actual_stream_survives_probe_and_inventory_fault_then_stops_on_disconnect(ordinary): # noqa: F811 + import threading + from loopx.extensions.lark.goal_topic_runtime import stream_lark_goal_topic_profile + + _, runtime, provider, transport = connect(ordinary) + released, listening, fault_seen, recovered = (threading.Event() for _ in range(4)) + result = {} + fail_inventory = False + + def snapshot(): + if fail_inventory and not fault_seen.is_set(): + saved = transport.runner + transport.runner = lambda *_args: {"returncode": 1, "stdout": ""} + try: + return {"private_profiles": transport.profiles()} + finally: + transport.runner = saved + fault_seen.set() + profiles = transport.profiles() + if fault_seen.is_set(): + recovered.set() + return {"private_profiles": profiles} + + class Lines: + def __iter__(self): + yield "[event] ready event_key=im.message.receive_v1\n" + assert released.wait(10) + + class Consumer: + stdout = Lines() + def poll(self): + return 0 if released.is_set() else None + def wait(self, timeout=None): + assert released.wait(timeout or 10) + return 0 + def terminate(self): + released.set() + kill = terminate + + stop = threading.Event() + def run(): + result.update(stream_lark_goal_topic_profile(profile="notes-app", snapshot_provider=snapshot, + stop=stop, runtime_root=transport.runtime_root, answer=lambda *_args: "unused", + private_admitter=transport.admit, process_factory=lambda *_args: Consumer(), + health_sink=lambda update: listening.set() if update.get("status") == "listening" else None)) + + worker = threading.Thread(target=run) + try: + worker.start() + assert listening.wait(3) + transport.bindings.observe = lambda _profile: (_ for _ in ()).throw(ValueError("verification unavailable")) + fail_inventory = True + assert fault_seen.wait(3) and recovered.wait(3) + assert worker.is_alive() and not released.is_set() and not stop.is_set() + assert transport.admit("notes-app", provider.event("notes-app", "unverified", "/status"))["status"] == "audience_rejected" + binding = transport._binding("notes-app") + transport.bindings.disconnect(binding["binding_id"], expected_revision=transport.bindings.read()["revision"]) + worker.join(4) + assert not worker.is_alive() and stop.is_set() + assert result["status"] == "configuration_removed" + finally: + stop.set() + released.set() + worker.join(4) + runtime.close() + + def test_slow_reply_readback_does_not_hold_independent_stop_or_other_app(ordinary, monkeypatch): # noqa: F811 import threading from loopx.extensions.lark.goal_topic_runtime_service import LarkGoalTopicRuntimeService