diff --git a/src/langbot/libs/wecom_ai_bot_api/ws_client.py b/src/langbot/libs/wecom_ai_bot_api/ws_client.py index eef41531f..3f4cb0504 100644 --- a/src/langbot/libs/wecom_ai_bot_api/ws_client.py +++ b/src/langbot/libs/wecom_ai_bot_api/ws_client.py @@ -691,6 +691,18 @@ class WecomBotWsClient: if not _re.sub(r'[\s​‌‍]', '', next_content): return True + # A blank *final* snapshot would close the stream with an empty + # bubble and strand the real answer in a separate reply_text + # message. Keep the session open so a following non-blank chunk can + # finalize it. The non-final branch above already skips blank + # snapshots; final snapshots with earlier content are non-blank + # here because ``next_content`` falls back to the previous content. + if is_final: + import re as _re + + if not _re.sub(r'[\s\u200b\u200c\u200d\ufeff]', '', next_content): + return True + # Generate feedback_id for final chunk feedback_id = '' if is_final: diff --git a/src/langbot/pkg/pipeline/wrapper/wrapper.py b/src/langbot/pkg/pipeline/wrapper/wrapper.py index eeb2daec2..9ea275e1d 100644 --- a/src/langbot/pkg/pipeline/wrapper/wrapper.py +++ b/src/langbot/pkg/pipeline/wrapper/wrapper.py @@ -25,6 +25,50 @@ class ResponseWrapper(stage.PipelineStage): async def initialize(self, pipeline_config: dict): pass + def _is_final_assistant_message(self, result) -> bool: + """Whether *result* is the agent's final, tool-call-free answer. + + Intermediate streaming chunks and tool-call rounds must NOT trigger + outbound attachment collection — only the terminal assistant message. + """ + if getattr(result, 'role', None) != 'assistant': + return False + if result.tool_calls: + return False + if isinstance(result, provider_message.MessageChunk): + return bool(result.is_final) + return True + + async def _append_outbound_attachments( + self, + query: pipeline_query.Query, + message_chain: platform_message.MessageChain, + ) -> None: + """Collect sandbox outbox files and append them to *message_chain*. + + Runs at most once per query (guarded by a query variable) and never + raises into the pipeline — attachment delivery is best-effort. + """ + if query.variables.get('_sandbox_outbound_collected'): + return + box_service = getattr(self.ap, 'box_service', None) + if box_service is None or not getattr(box_service, 'available', False): + return + query.variables['_sandbox_outbound_collected'] = True + try: + attachments = await box_service.collect_outbound_attachments(query) + except Exception as e: + self.ap.logger.warning(f'Outbound attachment collection failed: {e}') + return + for att in attachments: + att_type = att.get('type') + if att_type == 'Image': + message_chain.append(platform_message.Image(base64=att['base64'])) + elif att_type == 'Voice': + message_chain.append(platform_message.Voice(base64=att['base64'])) + else: + message_chain.append(platform_message.File(name=att.get('name', 'file'), base64=att['base64'])) + async def process( self, query: pipeline_query.Query, @@ -61,7 +105,7 @@ class ResponseWrapper(stage.PipelineStage): reply_text = '' - if result.content or result.attachments: # 有内容 + if result.content: # 有内容 reply_text = str(result.get_content_platform_message_chain()) # ============= 触发插件事件 =============== @@ -96,6 +140,11 @@ class ResponseWrapper(stage.PipelineStage): reply_chain = result.get_content_platform_message_chain() is_plugin_reply = False + # Only collect Box outbox files for a terminal assistant + # response; also accept provider final chunks here. + if self._is_final_assistant_message(result): + await self._append_outbound_attachments(query, reply_chain) + query.resp_message_chain.append(reply_chain) if is_plugin_reply: plugin_diagnostics.record_last_plugin_response_source( @@ -112,25 +161,31 @@ class ResponseWrapper(stage.PipelineStage): elif ( isinstance(result, provider_message.MessageChunk) and result.is_final and not result.tool_calls ): - # Final streaming chunk with no text content but - # no implicit file collection. + # A blank final chunk may carry Box outbox files, but do not + # send an empty chain: streaming adapters treat it as final. reply_chain = platform_message.MessageChain([]) - query.resp_message_chain.append(reply_chain) - yield entities.StageProcessResult( - result_type=entities.ResultType.CONTINUE, - new_query=query, - ) + await self._append_outbound_attachments(query, reply_chain) + if len(reply_chain) > 0: + query.resp_message_chain.append(reply_chain) + yield entities.StageProcessResult( + result_type=entities.ResultType.CONTINUE, + new_query=query, + ) if result.tool_calls is not None and len(result.tool_calls) > 0: # 有函数调用 function_names = [tc.function.name for tc in result.tool_calls] reply_text = f'Call {".".join(function_names)}...' - query.resp_message_chain.append( - platform_message.MessageChain([platform_message.Plain(text=reply_text)]) - ) - + # Only surface the tool-call notice when the pipeline option + # is enabled. Emitting it unconditionally would become the + # final streaming chunk, closing the stream early and pushing + # the real answer into a separate message. if query.pipeline_config['output']['misc']['track-function-calls']: + query.resp_message_chain.append( + platform_message.MessageChain([platform_message.Plain(text=reply_text)]) + ) + event = events.NormalMessageResponded( launcher_type=query.launcher_type.value, launcher_id=query.launcher_id, diff --git a/tests/unit_tests/platform/test_wecombot_stream_empty_final.py b/tests/unit_tests/platform/test_wecombot_stream_empty_final.py new file mode 100644 index 000000000..e21d163b6 --- /dev/null +++ b/tests/unit_tests/platform/test_wecombot_stream_empty_final.py @@ -0,0 +1,119 @@ +"""Regression tests for WeCom AI Bot streaming chunks with blank content. + +A blank *final* snapshot used to be forwarded to the client as-is. On WeCom the +stream frame replaces the displayed bubble, so this closed the stream with an +empty bubble and stranded the model's real answer in a separate ``reply_text`` +message (observed as "empty message, then the real reply"). + +These tests lock in the fixed behaviour: + +* blank snapshots are never sent, final or not; +* a blank final snapshot does not close the stream session, so a following + non-blank chunk can still finalize it; +* a blank final snapshot with earlier content still finalizes with that + content (the previous-content fallback keeps it non-blank). +""" + +import pytest + +import langbot.pkg.core.app # noqa: F401 +from langbot.libs.wecom_ai_bot_api.ws_client import WecomBotWsClient + + +class Logger: + def __init__(self): + self.warnings = [] + self.errors = [] + + async def warning(self, message): + self.warnings.append(message) + + async def error(self, message): + self.errors.append(message) + + async def info(self, message): + return None + + +class RecordingClient(WecomBotWsClient): + """A client that records reply frames instead of touching a socket.""" + + def __init__(self): + super().__init__(bot_id='bot', secret='secret', logger=Logger()) + self.replies = [] + + async def _send_reply(self, req_id, body, cmd='aibot_respond_msg'): + self.replies.append((cmd, req_id, body)) + return {'errcode': 0} + + def seed_stream(self, msg_id='msg-1'): + self._stream_ids[msg_id] = 'req-1|stream-1' + self._stream_sessions[msg_id] = { + 'req_id': 'req-1', + 'stream_id': 'stream-1', + 'msg_id': msg_id, + 'user_id': 'user-1', + 'chat_id': 'chat-1', + 'chat_type': 'single', + } + + def stream_still_open(self, msg_id='msg-1'): + return msg_id in self._stream_ids + + def stream_frames(self): + return [body for _cmd, _req, body in self.replies if body.get('msgtype') == 'stream'] + + +@pytest.mark.asyncio +async def test_blank_final_chunk_does_not_close_stream(): + client = RecordingClient() + client.seed_stream() + + ok = await client.push_stream_chunk('msg-1', '', is_final=True) + + assert ok is True + assert client.stream_frames() == [] + assert client.stream_still_open() + + +@pytest.mark.asyncio +async def test_blank_non_final_chunk_is_skipped(): + client = RecordingClient() + client.seed_stream() + + ok = await client.push_stream_chunk('msg-1', ' \u200b', is_final=False) + + assert ok is True + assert client.stream_frames() == [] + assert client.stream_still_open() + + +@pytest.mark.asyncio +async def test_blank_final_after_content_finalizes_with_previous_content(): + client = RecordingClient() + client.seed_stream() + + await client.push_stream_chunk('msg-1', 'hello', is_final=False) + ok = await client.push_stream_chunk('msg-1', '', is_final=True) + + frames = client.stream_frames() + assert ok is True + assert len(frames) == 2 + assert frames[-1]['stream']['content'] == 'hello' + assert frames[-1]['stream']['finish'] is True + assert not client.stream_still_open() + + +@pytest.mark.asyncio +async def test_real_final_chunk_is_sent_and_closes_stream(): + client = RecordingClient() + client.seed_stream() + + ok = await client.push_stream_chunk('msg-1', 'the answer', is_final=True) + + frames = client.stream_frames() + assert ok is True + assert len(frames) == 1 + assert frames[-1]['stream']['content'] == 'the answer' + assert frames[-1]['stream']['finish'] is True + assert not client.stream_still_open() diff --git a/web/tests/unit/oss-cloud-ui-privacy.test.mjs b/web/tests/unit/oss-cloud-ui-privacy.test.mjs index 1e83634a0..4481aa18f 100644 --- a/web/tests/unit/oss-cloud-ui-privacy.test.mjs +++ b/web/tests/unit/oss-cloud-ui-privacy.test.mjs @@ -75,6 +75,13 @@ test('does not expose storage analysis in Cloud settings or via a deep link', () ); }); +test('renders storage analysis with the storage icon in the user menu', () => { + assert.match( + homeSidebarSource, + /canViewStorageAnalysis\s*&&\s*\(\s*[\s\S]*?\{t\('storageAnalysis\.title'\)\}/, + ); +}); + test('loads plugin pages through the authenticated Workspace-scoped asset route', () => { assert.match(pluginPageSource, /useAuthenticatedPluginAsset/); assert.match(