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
This commit is contained in:
yong 2026-06-28 11:45:52 +08:00 committed by GitHub
parent b32ee26454
commit 7ea72087bf
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
3 changed files with 36 additions and 24 deletions

View File

@ -46,9 +46,9 @@ class FeishuChannel(Channel):
Message flow: Message flow:
1. User sends a message bot adds "OK" emoji reaction 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 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 5. Bot adds "DONE" emoji reaction to the original message
""" """
@ -269,7 +269,7 @@ class FeishuChannel(Channel):
content = json.dumps({"file_key": file_key}) content = json.dumps({"file_key": file_key})
if msg.thread_ts: 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) await asyncio.to_thread(self._api_client.im.v1.message.reply, request)
else: 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() 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 return None
content = self._build_card_content(text) 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 = await asyncio.to_thread(self._api_client.im.v1.message.reply, request)
response_data = getattr(response, "data", None) response_data = getattr(response, "data", None)
return getattr(response_data, "message_id", 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) 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 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.""" """Start running-card creation once per source message."""
running_card_id = self._running_card_ids.get(source_message_id) running_card_id = self._running_card_ids.get(source_message_id)
if running_card_id: if running_card_id:
@ -522,7 +522,7 @@ class FeishuChannel(Channel):
self._running_card_tasks.pop(source_message_id, None) self._running_card_tasks.pop(source_message_id, None)
self._log_task_error(task, "create_running_card", source_message_id) 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.""" """Ensure the in-thread running card exists and track its message ID."""
running_card_id = self._running_card_ids.get(source_message_id) running_card_id = self._running_card_ids.get(source_message_id)
if running_card_id: if running_card_id:
@ -779,9 +779,8 @@ class FeishuChannel(Channel):
msg_id = message.message_id msg_id = message.message_id
sender_id = event.event.sender.sender_id.open_id sender_id = event.event.sender.sender_id.open_id
# root_id is set when the message is a reply within a Feishu thread. root_id = getattr(message, "root_id", None) or None
# Use it as topic_id so all replies share the same DeerFlow thread. chat_type = getattr(message, "chat_type", None)
root_id = self._non_empty_str(getattr(message, "root_id", None))
parent_id = self._non_empty_str(getattr(message, "parent_id", None)) parent_id = self._non_empty_str(getattr(message, "parent_id", None))
feishu_thread_id = self._non_empty_str(getattr(message, "thread_id", None)) feishu_thread_id = self._non_empty_str(getattr(message, "thread_id", None))
@ -844,12 +843,13 @@ class FeishuChannel(Channel):
text = text.strip() text = text.strip()
logger.info( 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, chat_id,
msg_id, msg_id,
root_id, root_id,
parent_id, parent_id,
feishu_thread_id, feishu_thread_id,
chat_type,
sender_id, sender_id,
len(text or ""), len(text or ""),
) )
@ -882,9 +882,9 @@ class FeishuChannel(Channel):
else: else:
msg_type = InboundMessageType.CHAT msg_type = InboundMessageType.CHAT
# Prefer any platform message id that already maps to a DeerFlow # topic_id determines which LangGraph thread the message maps to.
# thread. This keeps replies to bot clarification cards in the # P2P chats: topic_id=None so all messages share one thread (like Telegram DMs).
# original conversation even when Feishu reports the card as root. # But check stored mappings first for backward compatibility with pre-upgrade P2P threads.
topic_id, resolved_from_stored_mapping = self._resolve_topic_id( topic_id, resolved_from_stored_mapping = self._resolve_topic_id(
chat_id, chat_id,
msg_id, msg_id,
@ -892,6 +892,8 @@ class FeishuChannel(Channel):
parent_id=parent_id, parent_id=parent_id,
thread_id=feishu_thread_id, thread_id=feishu_thread_id,
) )
if chat_type == "p2p" and not resolved_from_stored_mapping:
topic_id = None
resolved_from_pending = False resolved_from_pending = False
if msg_type == InboundMessageType.CHAT and not resolved_from_stored_mapping: if msg_type == InboundMessageType.CHAT and not resolved_from_stored_mapping:
pending = self._consume_pending_clarification(chat_id, sender_id) pending = self._consume_pending_clarification(chat_id, sender_id)

View File

@ -54,7 +54,8 @@ DEFAULT_RUN_CONTEXT: dict[str, Any] = {
"is_plan_mode": False, "is_plan_mode": False,
"subagent_enabled": 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 # Stream modes requested from the runtime, and the SSE event names under which
# the message-tuple stream may arrive: the embedded runtime (and LangGraph # the message-tuple stream may arrive: the embedded runtime (and LangGraph
# Platform) deliver the requested "messages-tuple" mode as event "messages". # Platform) deliver the requested "messages-tuple" mode as event "messages".
@ -1403,6 +1404,7 @@ class ChannelManager:
current_message_id: str | None = None current_message_id: str | None = None
latest_text = "" latest_text = ""
last_published_text = "" last_published_text = ""
last_published_len = 0
last_publish_at = 0.0 last_publish_at = 0.0
stream_error: BaseException | None = None stream_error: BaseException | None = None
stream_kwargs: dict[str, Any] = { stream_kwargs: dict[str, Any] = {
@ -1430,23 +1432,30 @@ class ChannelManager:
latest_text = accumulated_text latest_text = accumulated_text
elif event == "values" and isinstance(data, (dict, list)): elif event == "values" and isinstance(data, (dict, list)):
last_values = data last_values = data
snapshot_text = _extract_response_text(data) # Clarification text is only in the values snapshot;
if snapshot_text: # publish it so the user sees the question mid-stream.
latest_text = snapshot_text 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: if not latest_text or latest_text == last_published_text:
continue continue
now = time.monotonic() now = time.monotonic()
if last_published_text and now - last_publish_at < STREAM_UPDATE_MIN_INTERVAL_SECONDS: new_chars = len(latest_text) - last_published_len
continue # 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( await self.bus.publish_outbound(
OutboundMessage( OutboundMessage(
channel_name=msg.channel_name, channel_name=msg.channel_name,
chat_id=msg.chat_id, chat_id=msg.chat_id,
thread_id=thread_id, thread_id=thread_id,
text=latest_text, text=display_text,
is_final=False, is_final=False,
thread_ts=msg.thread_ts, thread_ts=msg.thread_ts,
connection_id=msg.connection_id, connection_id=msg.connection_id,
@ -1455,6 +1464,7 @@ class ChannelManager:
) )
) )
last_published_text = latest_text last_published_text = latest_text
last_published_len = len(latest_text)
last_publish_at = now last_publish_at = now
except Exception as exc: except Exception as exc:
stream_error = exc stream_error = exc

View File

@ -1571,7 +1571,7 @@ class TestChannelManager:
await manager.stop() await manager.stop()
mock_client.runs.stream.assert_called_once() 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 [msg.is_final for msg in outbound_received] == [False, False, True]
assert all(msg.thread_ts == "om-source-1" for msg in outbound_received) assert all(msg.thread_ts == "om-source-1" for msg in outbound_received)
@ -1642,7 +1642,7 @@ class TestChannelManager:
await manager.stop() await manager.stop()
mock_client.runs.stream.assert_called_once() 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 [msg.is_final for msg in outbound_received] == [False, False, True]
_run(go()) _run(go())
@ -1700,8 +1700,8 @@ class TestChannelManager:
await manager.stop() await manager.stop()
assert [msg.is_final for msg in outbound_received] == [False, False, True] assert [msg.is_final for msg in outbound_received] == [False, False, True]
assert outbound_received[0].text == "Thinking" assert outbound_received[0].text == "Thinking"
assert outbound_received[1].text == "Which environment?" assert outbound_received[1].text == "Which environment?"
assert outbound_received[2].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 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 assert outbound_received[-1].metadata[PENDING_CLARIFICATION_METADATA_KEY] is True