From 7ea72087bfa3080ea73c65c46ba82177ae07329e Mon Sep 17 00:00:00 2001 From: yong <103179041+yong326@users.noreply.github.com> Date: Sun, 28 Jun 2026 11:45:52 +0800 Subject: [PATCH] fix(feishu): stop creating thread topics and throttle card updates (#3810) * fix(feishu): stop creating thread topics and throttle card updates - Remove reply_in_thread(True) so replies appear as normal messages (#1332) - Use chat_type=p2p for shared LangGraph thread in P2P chats - Add two-level stream throttle (1.0s interval / 60 char buffer) with OR logic - Add cursor indicator for streaming - Filter values from overwriting delta text during streaming Closes #3801 * fix(feishu): stop creating thread topics, throttle card updates, preserve clarification - Remove reply_in_thread(True) so replies appear as normal messages (#1332) - Use chat_type=p2p for shared LangGraph thread in P2P chats, with stored mapping fallback for backward compatibility with pre-upgrade threads - Add two-level stream throttle (1.0s interval / 60 char buffer) with OR logic and cursor indicator for non-final outbounds - Re-publish clarification text from values snapshots mid-stream - Update streaming tests to the new contract (cursor glyph, clarification) Closes #3801 --- backend/app/channels/feishu.py | 28 +++++++++++++++------------- backend/app/channels/manager.py | 24 +++++++++++++++++------- backend/tests/test_channels.py | 8 ++++---- 3 files changed, 36 insertions(+), 24 deletions(-) diff --git a/backend/app/channels/feishu.py b/backend/app/channels/feishu.py index ec53eb131..48cc8e8ba 100644 --- a/backend/app/channels/feishu.py +++ b/backend/app/channels/feishu.py @@ -46,9 +46,9 @@ class FeishuChannel(Channel): Message flow: 1. User sends a message → bot adds "OK" emoji reaction - 2. Bot replies in thread: "Working on it......" + 2. Bot replies with a card: "Working on it......" 3. Agent processes the message and returns a result - 4. Bot replies in thread with the result + 4. Bot updates the card with the result 5. Bot adds "DONE" emoji reaction to the original message """ @@ -269,7 +269,7 @@ class FeishuChannel(Channel): content = json.dumps({"file_key": file_key}) if msg.thread_ts: - request = self._ReplyMessageRequest.builder().message_id(msg.thread_ts).request_body(self._ReplyMessageRequestBody.builder().msg_type(msg_type).content(content).reply_in_thread(True).build()).build() + request = self._ReplyMessageRequest.builder().message_id(msg.thread_ts).request_body(self._ReplyMessageRequestBody.builder().msg_type(msg_type).content(content).build()).build() await asyncio.to_thread(self._api_client.im.v1.message.reply, request) else: request = self._CreateMessageRequest.builder().receive_id_type("chat_id").request_body(self._CreateMessageRequestBody.builder().receive_id(msg.chat_id).msg_type(msg_type).content(content).build()).build() @@ -460,7 +460,7 @@ class FeishuChannel(Channel): return None content = self._build_card_content(text) - request = self._ReplyMessageRequest.builder().message_id(message_id).request_body(self._ReplyMessageRequestBody.builder().msg_type("interactive").content(content).reply_in_thread(True).build()).build() + request = self._ReplyMessageRequest.builder().message_id(message_id).request_body(self._ReplyMessageRequestBody.builder().msg_type("interactive").content(content).build()).build() response = await asyncio.to_thread(self._api_client.im.v1.message.reply, request) response_data = getattr(response, "data", None) return getattr(response_data, "message_id", None) @@ -502,7 +502,7 @@ class FeishuChannel(Channel): logger.warning("[Feishu] running card creation returned no message_id for source=%s, subsequent updates will fall back to new replies", source_message_id) return running_card_id - def _ensure_running_card_started(self, source_message_id: str, text: str = "Working on it...") -> asyncio.Task | None: + def _ensure_running_card_started(self, source_message_id: str, text: str = "thinking...") -> asyncio.Task | None: """Start running-card creation once per source message.""" running_card_id = self._running_card_ids.get(source_message_id) if running_card_id: @@ -522,7 +522,7 @@ class FeishuChannel(Channel): self._running_card_tasks.pop(source_message_id, None) self._log_task_error(task, "create_running_card", source_message_id) - async def _ensure_running_card(self, source_message_id: str, text: str = "Working on it...") -> str | None: + async def _ensure_running_card(self, source_message_id: str, text: str = "thinking...") -> str | None: """Ensure the in-thread running card exists and track its message ID.""" running_card_id = self._running_card_ids.get(source_message_id) if running_card_id: @@ -779,9 +779,8 @@ class FeishuChannel(Channel): msg_id = message.message_id sender_id = event.event.sender.sender_id.open_id - # root_id is set when the message is a reply within a Feishu thread. - # Use it as topic_id so all replies share the same DeerFlow thread. - root_id = self._non_empty_str(getattr(message, "root_id", None)) + root_id = getattr(message, "root_id", None) or None + chat_type = getattr(message, "chat_type", None) parent_id = self._non_empty_str(getattr(message, "parent_id", None)) feishu_thread_id = self._non_empty_str(getattr(message, "thread_id", None)) @@ -844,12 +843,13 @@ class FeishuChannel(Channel): text = text.strip() logger.info( - "[Feishu] parsed message: chat_id=%s, msg_id=%s, root_id=%s, parent_id=%s, thread_id=%s, sender=%s, text_len=%d", + "[Feishu] parsed message: chat_id=%s, msg_id=%s, root_id=%s, parent_id=%s, thread_id=%s, chat_type=%s, sender=%s, text_len=%d", chat_id, msg_id, root_id, parent_id, feishu_thread_id, + chat_type, sender_id, len(text or ""), ) @@ -882,9 +882,9 @@ class FeishuChannel(Channel): else: msg_type = InboundMessageType.CHAT - # Prefer any platform message id that already maps to a DeerFlow - # thread. This keeps replies to bot clarification cards in the - # original conversation even when Feishu reports the card as root. + # topic_id determines which LangGraph thread the message maps to. + # P2P chats: topic_id=None so all messages share one thread (like Telegram DMs). + # But check stored mappings first for backward compatibility with pre-upgrade P2P threads. topic_id, resolved_from_stored_mapping = self._resolve_topic_id( chat_id, msg_id, @@ -892,6 +892,8 @@ class FeishuChannel(Channel): parent_id=parent_id, thread_id=feishu_thread_id, ) + if chat_type == "p2p" and not resolved_from_stored_mapping: + topic_id = None resolved_from_pending = False if msg_type == InboundMessageType.CHAT and not resolved_from_stored_mapping: pending = self._consume_pending_clarification(chat_id, sender_id) diff --git a/backend/app/channels/manager.py b/backend/app/channels/manager.py index 235ecce52..519c28235 100644 --- a/backend/app/channels/manager.py +++ b/backend/app/channels/manager.py @@ -54,7 +54,8 @@ DEFAULT_RUN_CONTEXT: dict[str, Any] = { "is_plan_mode": False, "subagent_enabled": False, } -STREAM_UPDATE_MIN_INTERVAL_SECONDS = 0.35 +STREAM_UPDATE_MIN_INTERVAL_SECONDS = 1.0 +STREAM_UPDATE_MIN_CHARS = 60 # flush immediately when this many chars accumulate # Stream modes requested from the runtime, and the SSE event names under which # the message-tuple stream may arrive: the embedded runtime (and LangGraph # Platform) deliver the requested "messages-tuple" mode as event "messages". @@ -1403,6 +1404,7 @@ class ChannelManager: current_message_id: str | None = None latest_text = "" last_published_text = "" + last_published_len = 0 last_publish_at = 0.0 stream_error: BaseException | None = None stream_kwargs: dict[str, Any] = { @@ -1430,23 +1432,30 @@ class ChannelManager: latest_text = accumulated_text elif event == "values" and isinstance(data, (dict, list)): last_values = data - snapshot_text = _extract_response_text(data) - if snapshot_text: - latest_text = snapshot_text + # Clarification text is only in the values snapshot; + # publish it so the user sees the question mid-stream. + if _has_current_turn_clarification(data): + clarification_text = _extract_response_text(data) + if clarification_text and clarification_text != latest_text: + latest_text = clarification_text if not latest_text or latest_text == last_published_text: continue now = time.monotonic() - if last_published_text and now - last_publish_at < STREAM_UPDATE_MIN_INTERVAL_SECONDS: - continue + new_chars = len(latest_text) - last_published_len + # OR logic: flush when interval elapsed OR enough chars accumulated + if last_published_text: + if now - last_publish_at < STREAM_UPDATE_MIN_INTERVAL_SECONDS and new_chars < STREAM_UPDATE_MIN_CHARS: + continue + display_text = latest_text + " ▉" await self.bus.publish_outbound( OutboundMessage( channel_name=msg.channel_name, chat_id=msg.chat_id, thread_id=thread_id, - text=latest_text, + text=display_text, is_final=False, thread_ts=msg.thread_ts, connection_id=msg.connection_id, @@ -1455,6 +1464,7 @@ class ChannelManager: ) ) last_published_text = latest_text + last_published_len = len(latest_text) last_publish_at = now except Exception as exc: stream_error = exc diff --git a/backend/tests/test_channels.py b/backend/tests/test_channels.py index 7cddb4551..bf7312d46 100644 --- a/backend/tests/test_channels.py +++ b/backend/tests/test_channels.py @@ -1571,7 +1571,7 @@ class TestChannelManager: await manager.stop() mock_client.runs.stream.assert_called_once() - assert [msg.text for msg in outbound_received] == ["Hello", "Hello world", "Hello world"] + assert [msg.text for msg in outbound_received] == ["Hello ▉", "Hello world ▉", "Hello world"] assert [msg.is_final for msg in outbound_received] == [False, False, True] assert all(msg.thread_ts == "om-source-1" for msg in outbound_received) @@ -1642,7 +1642,7 @@ class TestChannelManager: await manager.stop() mock_client.runs.stream.assert_called_once() - assert [msg.text for msg in outbound_received] == ["Hello", "Hello world", "Hello world"] + assert [msg.text for msg in outbound_received] == ["Hello ▉", "Hello world ▉", "Hello world"] assert [msg.is_final for msg in outbound_received] == [False, False, True] _run(go()) @@ -1700,8 +1700,8 @@ class TestChannelManager: await manager.stop() assert [msg.is_final for msg in outbound_received] == [False, False, True] - assert outbound_received[0].text == "Thinking" - assert outbound_received[1].text == "Which environment?" + assert outbound_received[0].text == "Thinking ▉" + assert outbound_received[1].text == "Which environment? ▉" assert outbound_received[2].text == "Which environment?" assert all(PENDING_CLARIFICATION_METADATA_KEY not in msg.metadata for msg in outbound_received[:-1]) assert outbound_received[-1].metadata[PENDING_CLARIFICATION_METADATA_KEY] is True