From 8749c915cd5700c0e77315c2f614fa9fd8934215 Mon Sep 17 00:00:00 2001 From: matevip Date: Wed, 22 Jul 2026 13:40:22 +0800 Subject: [PATCH] feat(wecom): roll progress bubble per stage narration --- .../channel/wecom/WeComChannelAdapter.java | 97 ++++++++++++++++++- .../channel/wecom/WeComProcessStreamTest.java | 35 +++++++ 2 files changed, 129 insertions(+), 3 deletions(-) diff --git a/mateclaw-server/src/main/java/vip/mate/channel/wecom/WeComChannelAdapter.java b/mateclaw-server/src/main/java/vip/mate/channel/wecom/WeComChannelAdapter.java index 45b87083..eda57adb 100644 --- a/mateclaw-server/src/main/java/vip/mate/channel/wecom/WeComChannelAdapter.java +++ b/mateclaw-server/src/main/java/vip/mate/channel/wecom/WeComChannelAdapter.java @@ -36,6 +36,7 @@ import java.util.*; import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; /** * 企业微信智能机器人渠道适配器 — WebSocket 长连接模式 @@ -1425,9 +1426,9 @@ public class WeComChannelAdapter extends AbstractChannelAdapter implements Strea return finalContent; } - /** Event-driven progress rendering into the existing processing-stream bubble. */ + /** Event-driven progress rendering into the processing-stream bubble, with per-stage bubble rolling. */ private void consumeWithProgress(Flux stream, ChannelMessage message, - String replyTarget, WeComReplyContext ctx, + String replyTarget, WeComReplyContext initialCtx, StringBuilder contentAccumulator) { boolean showThinking = !getConfigBoolean("filter_thinking", true); boolean standaloneToolMessages = !getConfigBoolean("filter_tool_messages", true); @@ -1438,9 +1439,15 @@ public class WeComChannelAdapter extends AbstractChannelAdapter implements Strea if (keepaliveScheduler != null) { // Silent stretches (long LLM calls with no events) keep showing a // fresh elapsed-time snapshot instead of the static placeholder. - keepaliveScheduler.attachTextSupplier(ctx.processingStreamId(), progress::snapshot); + keepaliveScheduler.attachTextSupplier(initialCtx.processingStreamId(), progress::snapshot); } + // The live progress bubble rolls forward on every stage narration: + // the current stream is finished with the narration text (making it + // a permanent bubble in place) and a fresh stream id opens below it + // as the new progress bubble, so chat chronology stays intact and + // the final answer always lands in the newest bubble. + AtomicReference liveCtx = new AtomicReference<>(initialCtx); final long[] lastFlushAt = {0L}; stream.doOnNext(delta -> { boolean flushNow = false; @@ -1449,6 +1456,25 @@ public class WeComChannelAdapter extends AbstractChannelAdapter implements Strea if (standaloneToolMessages) { maybeSendToolEventMessage(replyTarget, delta.eventType(), delta.eventData()); } + } else if (delta.segmentOnly()) { + // Per-stage narration ("我来查一下…"), emitted once per agent + // loop iteration. Each becomes its own permanent bubble and is + // excluded from the final answer: glued together they read as + // a wall of text, and persisted they pollute the next turn's + // LLM history with unanswered chain-of-thought. + String narration = delta.content() != null ? delta.content().trim() : ""; + if (!narration.isEmpty()) { + WeComReplyContext ctx = liveCtx.get(); + if (replyContexts.get(replyTarget) == ctx) { + liveCtx.set(rollProgressBubble(replyTarget, ctx, narration, progress)); + } else { + // Bubble already force-finished (180s ceiling) — the + // narration still goes out as a plain message. + sendMessage(replyTarget, narration); + } + lastFlushAt[0] = System.currentTimeMillis(); + } + return; } else { if (delta.thinking() != null) { progress.onThinkingDelta(delta.thinking()); @@ -1462,6 +1488,7 @@ public class WeComChannelAdapter extends AbstractChannelAdapter implements Strea if (!flushNow && now - lastFlushAt[0] < minIntervalMs) { return; } + WeComReplyContext ctx = liveCtx.get(); // The keepalive force-finish (180s ceiling) evicts the reply // context; once that happens the stream slot is closed and // further overwrites would be silently rejected — stop pushing. @@ -1481,6 +1508,70 @@ public class WeComChannelAdapter extends AbstractChannelAdapter implements Strea }).blockLast(Duration.ofMinutes(10)); } + /** + * Finalize the current progress bubble with a stage narration and open a + * fresh stream as the next progress bubble. + *

+ * The narration goes through the channel renderer (thinking/tool-tag + * filters, table formatting, length split): the first segment overwrites + * the current bubble with {@code finish=true}, overflow segments ride + * plain messages. The replacement context is registered in + * {@link #replyContexts} so {@code renderAndSend} / approval cards keep + * working against the newest bubble, and keepalive restarts on the new + * stream with the live progress snapshot. + */ + private WeComReplyContext rollProgressBubble(String replyTarget, WeComReplyContext ctx, + String narration, WeComProgressRenderer progress) { + boolean filterThinking = getConfigBoolean("filter_thinking", true); + boolean filterToolMessages = getConfigBoolean("filter_tool_messages", true); + String format = getConfigString("message_format", "auto"); + int maxLen = vip.mate.channel.ChannelMessageRenderer.PLATFORM_LIMITS + .getOrDefault(getChannelType(), 2048); + List segments = vip.mate.channel.ChannelMessageRenderer.renderForChannel( + narration, filterThinking, filterToolMessages, format, maxLen); + if (segments.isEmpty()) { + // Narration entirely filtered away — keep the current bubble. + return ctx; + } + if (keepaliveScheduler != null) { + // Stop before the finish chunk so a refresh tick can't race it + // on the same stream. + keepaliveScheduler.stop(ctx.processingStreamId()); + } + boolean first = true; + for (String rawSegment : segments) { + String segment = formatMarkdownTables(rawSegment); + if (first) { + first = false; + try { + replyStream(ctx.frameReqId(), ctx.processingStreamId(), segment, true); + } catch (Exception e) { + log.debug("[wecom] stage bubble finalize failed: {}", e.getMessage()); + } + } else { + sendMessage(replyTarget, segment); + } + } + + String nextStreamId = generateReqId("stream"); + WeComReplyContext next = new WeComReplyContext(ctx.frameReqId(), nextStreamId); + replyContexts.put(replyTarget, next); + try { + replyStream(ctx.frameReqId(), nextStreamId, progress.snapshot(), false); + } catch (Exception e) { + log.debug("[wecom] next progress bubble open failed: {}", e.getMessage()); + } + if (keepaliveScheduler != null) { + try { + keepaliveScheduler.start(this, ctx.frameReqId(), nextStreamId, replyTarget); + keepaliveScheduler.attachTextSupplier(nextStreamId, progress::snapshot); + } catch (Exception e) { + log.debug("[wecom] keepalive restart failed: {}", e.getMessage()); + } + } + return next; + } + /** * Standalone tool-call trace messages, sent only when the channel's * {@code filter_tool_messages} toggle is off: the user opted into seeing diff --git a/mateclaw-server/src/test/java/vip/mate/channel/wecom/WeComProcessStreamTest.java b/mateclaw-server/src/test/java/vip/mate/channel/wecom/WeComProcessStreamTest.java index feb57b27..8cbd943e 100644 --- a/mateclaw-server/src/test/java/vip/mate/channel/wecom/WeComProcessStreamTest.java +++ b/mateclaw-server/src/test/java/vip/mate/channel/wecom/WeComProcessStreamTest.java @@ -68,6 +68,41 @@ class WeComProcessStreamTest { assertTrue(String.valueOf(finalChunk.get("content")).contains("现在是下午三点。")); } + @Test + @DisplayName("stage narrations roll the bubble: each stage finishes its own bubble, final answer excludes them") + void stageNarrationsRollBubbles() throws Exception { + TestableAdapter adapter = newAdapter("{\"progress_interval_ms\": 0}"); + seedReplyContext(adapter, "alice", "req-1", "stream-1"); + + Flux stream = Flux.just( + StreamDelta.segmentOnly("我先查一下当前时间:", null), + StreamDelta.event("tool_call_started", + Map.of("toolCallId", "c1", "toolName", "get_time")), + StreamDelta.event("tool_call_completed", + Map.of("toolCallId", "c1", "toolName", "get_time", "success", true)), + StreamDelta.segmentOnly("时间拿到了,再查会议室:", null), + new StreamDelta("1 号会议室空闲,已预约。", null)); + + String result = adapter.processStream(stream, inbound("alice"), "wecom:alice"); + + // Narrations are excluded from the returned (persisted) final answer. + assertEquals("1 号会议室空闲,已预约。", result); + + List> streamBodies = streamBodies(adapter.drainFrames()); + List> finished = streamBodies.stream() + .filter(s -> Boolean.TRUE.equals(s.get("finish"))).toList(); + // Three finished bubbles in chronological order: narration #1, + // narration #2, final answer — each on its own stream id. + assertEquals(3, finished.size(), "each stage plus the final answer closes one bubble"); + assertTrue(String.valueOf(finished.get(0).get("content")).contains("我先查一下当前时间")); + assertTrue(String.valueOf(finished.get(1).get("content")).contains("再查会议室")); + assertTrue(String.valueOf(finished.get(2).get("content")).contains("已预约")); + assertEquals(3, finished.stream().map(s -> s.get("id")).distinct().count(), + "each finished bubble must ride its own stream id"); + // The first narration finalizes the original placeholder stream. + assertEquals("stream-1", finished.get(0).get("id")); + } + @Test @DisplayName("stream_progress=false degrades to accumulate-then-send with no interim overwrites") void progressDisabledDegrades() throws Exception {