diff --git a/.claude/agents/ai-chat.md b/.claude/agents/ai-chat.md index ebfc8be39..473bd4112 100644 --- a/.claude/agents/ai-chat.md +++ b/.claude/agents/ai-chat.md @@ -104,13 +104,16 @@ description: AI 对话编排领域。任务涉及编排器 AgentOrchestrator、T **可靠性层(2026-08 harness 加固,治"跑一半停了")** - LLM timeout 600s(application.yml open-router.timeout;0.36 的单值=OkHttp callTimeout 整通墙钟上限,不是空闲超时)。 -- **流式模型必须 `logResponses(false)`(`ChatModelFactory.streamingBuilder`,两个流式通道共用的唯一构建口径)**。这不是调优是可靠性契约,改回 true 会让**整个传输层错误处理静默失效**:openai4j 0.23 的 `StreamingRequestExecutor$2.onFailure` 在该开关打开时先调 `ResponseLoggingInterceptor.log(response)` 再走 errorHandler,而 okhttp-sse 的 `RealEventSource.onFailure(call, e)` 在「连接失败/被断、压根没拿到响应」这条路径上传的 response **恒为 null**(另一条 `processResponse` 失败分支才是 t==null/response!=null,两者互斥,所以 response 为 null 时 t 必非 null),于是 `response.code()` 抛 NPE;`onFailure` 只 catch IOException,NPE 掀掉 OkHttp Dispatcher 线程,**紧随其后的 errorHandler 那一行永远走不到**。表现:本轮既不 onComplete 也不 onError,SSE 零字节,只能等看门狗兜底;后端日志里唯一痕迹是 `Exception in thread "OkHttp Dispatcher" ... Cannot invoke "okhttp3.Response.code()" because "response" is null`。**排障陷阱**:真正的 IOException 在这条路上被彻底销毁(`LOGGER.debug("onFailure()", t)` 那行本身也在开关内、且全仓无 logging 级别配置停在 INFO 不打印),所以「日志里只有 NPE、看不到网络错误」不代表网络没问题——修好这个开关才拿得到底层异常。回归守护 `StreamingTransportFailureTest`(连不上必须回调 onError;它走工厂那份真实 builder,用例里自己拼 builder 就永远是绿的)。非流式 `OpenAiChatModel` 走 SyncRequestExecutor 没这条路径,**所以「辅助模型秒回成功」不能用来证明流式通道的网络正常**(不同 executor、不同 OkHttpClient/连接池、且不带工具定义)。 +- **流式通道是自有的 `service/ai/OpenRouterStreamingChatModel`,不再是 langchain4j 0.36 的 `OpenAiStreamingChatModel`(2026-09-02,dev-board#364)**。唯一构建口径 `ChatModelFactory.streamingModel(apiKey, baseUrl, modelId, timeout)`,平台通道与 BYOK 两个流式路径都走它;Ollama 流式仍是 langchain4j 的。换实现的直接原因是**思考型模型**:OpenRouter 对 Kimi K3(`moonshotai/kimi-k3`)这类模型从第 4 秒起就流式返回 `delta.reasoning`(真机探测:每个思考 chunk 是 `content:""` + `reasoning:"…"` + `reasoning_details:[…]`,前面夹 `: OPENROUTER PROCESSING` 注释保活),而 openai4j 0.23 的 `Delta` 只有 role/content/toolCalls/functionCall 四个字段,reasoning 在反序列化那一刻就丢了、注释行被 okhttp-sse 静默吞掉,langchain4j 只对非 null 的 content 调 `onNext`——于是几百秒的思考期间编排器收到的全是 `onNext("")`(**恰好把看门狗喂活、又一个字节都不往前端发**),用户看到的就是「思考中 281 秒、什么都没有、分不清死机还是在想」。这条流在 langchain4j 那一层没有任何钩子能拿到 reasoning,所以只能自己读 HTTP/SSE。**刻意复用不重写**:`InternalOpenAiHelper.toOpenAiMessages/toTools`(含 ImageContent 编组)、openai4j 的 `Json`(请求体,snake_case + NON_NULL + INDENT_OUTPUT——**断言请求体时先去空白**)、`OpenAiStreamingResponseBuilder`(tool_calls 按 index 拼装、usage、finish_reason),本类只管 HTTP + SSE 行协议 + 多转发两条通道。与旧实现对齐的请求参数:`stream=true`、`stream_options.include_usage=true`、`temperature=0.7`;错误语义对齐:非 2xx 抛 `OpenAiHttpException(code, body)`(`LlmErrorClassifier` 按状态码分类),IOException 原样 onError,**HTTP 200 里用 data 事件送来的 `{"error":{...}}` 也当错误**(旧实现会按空回复静默收尾)。护栏 `OpenRouterStreamingChatModelTest`(假服务端回放真机抓到的片段形状)+ `StreamingTransportFailureTest`(连不上必须 onError;它走工厂那份真实口径)。**旧的 `logResponses(false)` 地雷随之消失**(openai4j 的 `StreamingRequestExecutor$2.onFailure` 在 response==null 时先调 `ResponseLoggingInterceptor.log` 抛 NPE、errorHandler 永远走不到),但非流式 `OpenAiChatModel` 仍必须 `logRequests(false)`(请求体物化,理由在 `streamingModel` 的 javadoc)。**「辅助模型秒回成功」仍不能用来证明流式通道的网络正常**(不同 HTTP 客户端/连接池、且不带工具定义)。 + - **`ReasoningStreamingHandler`**(extends `StreamingResponseHandler`)多两个 default 方法:`onReasoning(delta)` 与 `onKeepAlive()`。客户端只对 `instanceof` 这个接口的 handler 转发,回放评测与各测试的脚本模型按老接口写不受影响。`AgentStreamHandler` 实现它:reasoning → SSE `reasoning_delta`(**不进 fullContentBuilder、不进编辑器流、不过标签解析**:思考文本不是正文,不落库、不回喂模型——契约 D);两者都刷新看门狗的 `lastActivityNanos`。**`streamedAnyReasoning` 与 `streamedAnyToken` 刻意分开**:看门狗选时限时任一为真都算「流已开始」(思考几分钟是正常的,改用 180s 停滞时限),而编排器的「可安全重放」判定仍只看正文——思考卡重放一遍无害,正文重放才会让用户看到重复内容。护栏 `AgentStreamHandlerReasoningTest`。 + - **看门狗首字节 60s 保持不变**:真正的零字节死流仍在 60s 被掐;思考型模型靠 reasoning 增量 + OpenRouter 保活注释刷新活动时间,不会再被误杀。**注意 K3 的思考也是按输出单价计费的**($15/M),思考 281 秒的那一轮反复被掐重放会成倍烧钱——这就是首字节时限不能靠「调大」而必须靠「认得出模型还活着」来解决的原因。 + - **前端**:`useAgentStream.handleEvent` 认 `reasoning_delta` → `appendReasoning()`——没有过程卡时写顶层 `bubble.thinking.content`(ghost 态的 ThinkingCard 实时滚动显示),已有工具过程后挂到最后一个过程卡的 thinking 条目(与 `` 标签的落点同口径,否则第二轮起的思考会把首轮顶层卡的时长越算越长)。**不过 `processTextStream`**:思考文本里出现 `` 字样只是模型自言自语。等待首 token 的活性计数本来就有(`sendMessage` 起算 `thinking.startTime`,ThinkingCard 按 `chat.thinkingLive` 读秒);新增的是 **SSE 链路状态 `linkStatus`**(`{state:'live'|'reconnecting', attempt}`,`scheduleReconnect` 置 reconnecting、建连成功与 `resetSSE` 回 live),ChatInterface 输入区据此渲染 `chat.linkReconnecting` 提示条——之前断线重连只写 console.warn,用户看到的是计时器一直走、分不清模型在想还是连接死了。**前端判死阈值 `HEARTBEAT_STALE_MS=45000` = 后端 `SseEmitterService.HEARTBEAT_INTERVAL_SECONDS=15` 的 3 倍**,两边任一改动都要同步(`reasoning-stream.test.mjs` 与 `SseEmitterServiceTest.heartbeatSweepReachesEveryLiveConnection` 各守一侧)。Office 插件的 `sse.js` 对未知事件名直接忽略,`reasoning_delta` 不影响任务窗格。 - **「AI 全线连不上」优先怀疑 JVM 里冻住的代理端口,不要先怀疑密钥或网络**(2026-08-16 实证,两个 e2e home + 用户真机三处复现)。macOS 上**任何 JVM 启动时都会把系统代理设置自动灌进** `http(s).proxyHost/Port` 系统属性——**不需要任何 `-D`、不需要 `JAVA_TOOL_OPTIONS`**(裸 `java Foo.java` 就已经有 `https.proxyHost=127.0.0.1`),OkHttp 走 `ProxySelector.getDefault()` 于是全部 AI 流量被送去本地代理端口。桌面后端是**长命 JVM**(开 app 起、连跑数天),启动那刻把端口**冻住**;用户的代理工具换端口或重启后(实测 1235 → 8234),后端仍在拨旧端口,**每一个 AI 请求都 `ConnectException: Connection refused`**。 - 判定三件套:`jcmd <后端PID> VM.system_properties | grep proxy` 拿 JVM 冻住的端口 → `scutil --proxy` 拿系统当前端口 → `nc -z 127.0.0.1 <旧端口>` 确认旧端口已死。两者不一致就是它。 - **已自愈**:`service/SystemProxyRefresher.java` 每 60s 对齐一次(`scutil --proxy` → `System.setProperty`),开关 `network.proxy.auto-refresh`(默认 true)。成立前提是 `DefaultProxySelector` 每次 `select()` 都重读系统属性、运行期 `setProperty` 立即生效(由 `SystemProxyRefresherTest` 的端到端用例守住);**运行期打开 `java.net.useSystemProxies` 无效**(类初始化时固化,返 DIRECT),所以只能自己读 OS 再写属性。启用条件刻意收窄成「macOS + 启动时继承到回环代理」:非回环的企业代理端口稳定,动它只有风险。**启动时系统没开代理的情况不接管**(没有被冻住的旧端口,不存在要治的病),那种情况仍靠重启后端。老版本(≤ v0.16.0)没有这层自愈,临时解仍是重启 app。 - **表现极具迷惑性,两个假信号**:① 修复前流式路撞上文那个 NPE 被吞、静默 180s,日志里只有 NPE 看不到 ConnectException;② **同步路(辅助模型起标题/记忆/分类器)会「秒回」**——但那是 RetryUtils 重试 3 次约 1.4s 全败后写入的**兜底字面量「新对话」**,不是成功。**排障时先看标题是不是字面量「新对话」**,别拿它当"通道正常"的证据。 - **找日志别找错地方**:`-Duser.home=` 会整体改写 `~/.aiworkdeck` 的位置,e2e 后端的日志在 `/run/backend.log`。在真实 `~/.aiworkdeck/logs/backend.log` 里翻 e2e 的证据只会得出「什么都没有」的错误结论。 -- `AgentStreamHandler`:终态幂等(AtomicBoolean terminated)+ **流看门狗** armInactivityWatchdog(**首字节 60s / 停滞 180s**,5s 轮询)——两条时限刻意分开:停滞时限要照顾「生成长工具参数时中途静默几十秒」所以必须给足,而「从头到尾零字节」没有这种正当理由,合成一个值就是让用户干等三分钟。首字节这条只在 `streamedAnyToken == false` 时生效,而这恰好就是编排器判定「可安全重放」的条件,所以误杀代价上限是白跑一轮、不会让用户看到重复或半截内容。守护 `AgentStreamWatchdogTest`。 +- `AgentStreamHandler`:终态幂等(AtomicBoolean terminated)+ **流看门狗** armInactivityWatchdog(**首字节 60s / 停滞 180s**,5s 轮询)——两条时限刻意分开:停滞时限要照顾「生成长工具参数时中途静默几十秒」所以必须给足,而「从头到尾零字节」没有这种正当理由,合成一个值就是让用户干等三分钟。首字节这条只在 `streamedAnyToken == false && streamedAnyReasoning == false` 时生效(前者恰好是编排器判定「可安全重放」的条件),所以误杀代价上限是白跑一轮、不会让用户看到重复或半截内容;思考增量与 OpenRouter 保活注释(`onReasoning` / `onKeepAlive`)都刷新活动时间,思考型模型静默几分钟不会被首字节时限掐掉。守护 `AgentStreamWatchdogTest` + `AgentStreamHandlerReasoningTest`。 - `AgentOrchestrator.setOnError`:失败按 `LlmErrorClassifier.Kind` 分类(**七类**:RATE_LIMITED / TRANSIENT / MODEL_UNAVAILABLE / REGION_BLOCKED / **QUOTA_EXHAUSTED** / **CONTEXT_OVERFLOW** / FATAL,OpenAiHttpException 的结构化状态码优先于文本匹配),且**零 token 已流出**才允许重放。限流退避 30/60s ×2(限流窗口按分钟计,用 8/16/32 会在同一窗口连撞三次白烧预算),瞬时 8/16/32s ×3(RunGuard.llmRetries,成功轮与切模型后清零);用户文案两套,限流说「限流等待中」不说「服务不可用」。 - **QUOTA_EXHAUSTED = 配额/余额耗尽**(2026-08 对标 dsh):402、或 4xx + 配额语义(insufficient credits/quota/balance、quota exceeded、余额不足…)。**判定先于 429**——余额耗尽很多服务商也回 429,但它是终局:不退避(重试白烧)、不换模型(同一账户换哪个都没钱)。SSE error 载荷带 `AI_QUOTA_EXHAUSTED` 标记(`LlmErrorClassifier.QUOTA_EXHAUSTED_MARKER`),前端 useAgentStream includes 命中换中文引导(自备 Key 去服务商充值 / 平台通道去官网查额度分配)。 - **CONTEXT_OVERFLOW = 上下文超窗**(400 + 上下文语义,先于通用 400→FATAL 判定):不退避(原样重发必撞同一个 400)、不走故障转移链,走**专用恢复通道**——`RunLoopCompactor.forceCompact`(跳过阈值判断)强制压缩后同 depth 重放一次。**重试凭证 = compact 返回了新实例(确实缩小了)**,压不动直接终态(载荷带 `AI_CONTEXT_OVERFLOW` 标记换中文引导)。预算 `RunGuard.overflowCompactions` 1 次/轮,成功轮清零(长任务「涨→压→涨→压」合法)。存在意义:主动 compaction 靠 chars/token=2 估算,中文语料系统性低估,服务商的 400 是最后的事实来源。 @@ -158,7 +161,7 @@ ChatInterface.handleSubmit(~:927)→ useAgentStream.sendMessage(确保 SSE ## SSE 事件名清单 -connected / bubble_start / text_delta / artifact / token_usage / bubble_end(status: finished|paused|awaiting_approval|awaiting_input)/ error / cancelled / file_change / client_action / title_update / doc_stream_data(旧名 wps_stream_data 双轨待摘)/ state_recovery(断线重连快照)/ run_state / plan_update / **skill_update** / background_task_start|complete / task_progress / heartbeat / subtask_progress。前端分派均在 useAgentStream.handleEvent。超限 paused 契约见 PR#172。 +connected / bubble_start / text_delta / **reasoning_delta**(思考型模型的 reasoning 增量,`{"content":"…"}`,只进思考卡、不进正文与历史;state_recovery 快照不含它,重连后思考文本不回放)/ artifact / token_usage / bubble_end(status: finished|paused|awaiting_approval|awaiting_input)/ error / cancelled / file_change / client_action / title_update / doc_stream_data(旧名 wps_stream_data 双轨待摘)/ state_recovery(断线重连快照)/ run_state / plan_update / **skill_update** / background_task_start|complete / task_progress / heartbeat / subtask_progress。前端分派均在 useAgentStream.handleEvent。超限 paused 契约见 PR#172。 **`skill_update`(本轮生效的 skill 清单)**:载荷 `{"skills":[{"id","name","source"}]}`,source ∈ `auto`(触发词自动命中)/ `manual`(用户在面板里主动选的,含旧字段 pinnedSkillId);`name` 已由 `SkillRouter.displayName` 按应用语言解析(en 优先 name_en)。生产者只有 `AgentOrchestrator`,紧跟 `activateForTurn` 之后发一次。 - **每轮必发、空也发**:前端拿它做整表覆写(`useAgentStream.activeSkills`),漏发一次上一轮的 chip 就一直挂着,用户以为某个技能还生效着。 diff --git a/backend/src/main/java/com/checkba/service/SystemProxyRefresher.java b/backend/src/main/java/com/checkba/service/SystemProxyRefresher.java index 19ac8624c..5c8344114 100644 --- a/backend/src/main/java/com/checkba/service/SystemProxyRefresher.java +++ b/backend/src/main/java/com/checkba/service/SystemProxyRefresher.java @@ -17,7 +17,8 @@ * 于是全部出站流量被送去本地代理端口。而桌面后端是**长命 JVM**(开 app 起、连跑数天), * 把端口冻在启动那一刻;用户的代理软件换端口或重启后(实测 1235 → 8234), * 后端一直在拨那个已经没人监听的旧端口,**每一个 AI 请求都 ConnectException**。 - * 修复前这个失败还会被 openai4j 吞掉(见 {@code ChatModelFactory.streamingBuilder}), + * 修复前这个失败还会被 openai4j 吞掉(旧流式通道的 logResponses NPE 地雷, + * 见 {@code ChatModelFactory.streamingModel} 的 javadoc;流式通道已换成自有实现), * 用户看到的是"点了发送三分钟没反应"。 * *

为什么必须改属性而不是开 {@code java.net.useSystemProxies}:那个开关在 diff --git a/backend/src/main/java/com/checkba/service/ai/AgentStreamHandler.java b/backend/src/main/java/com/checkba/service/ai/AgentStreamHandler.java index 6db80b51b..e42586afe 100644 --- a/backend/src/main/java/com/checkba/service/ai/AgentStreamHandler.java +++ b/backend/src/main/java/com/checkba/service/ai/AgentStreamHandler.java @@ -1,7 +1,6 @@ package com.checkba.service.ai; import dev.langchain4j.model.output.Response; -import dev.langchain4j.model.StreamingResponseHandler; import dev.langchain4j.data.message.AiMessage; import java.util.UUID; @@ -10,7 +9,7 @@ * 负责将 LLM 的流式回调转换为前端 SSE 协议事件。 * 并收集最终完整的回复用于存储和计费。 */ -public class AgentStreamHandler implements StreamingResponseHandler { +public class AgentStreamHandler implements ReasoningStreamingHandler { private static final org.slf4j.Logger log = org.slf4j.LoggerFactory.getLogger(AgentStreamHandler.class); @@ -38,7 +37,12 @@ public class AgentStreamHandler implements StreamingResponseHandler { new java.util.concurrent.atomic.AtomicBoolean(false); // 是否已有 token 流出(重试决策依据:零 token 的失败轮可安全重放,不会给用户看重复内容) private volatile boolean streamedAnyToken = false; - // 最近一次流活动时间(onNext 刷新),看门狗据此判定"流停滞" + // 是否已有思考增量流出(思考型模型)。刻意与 streamedAnyToken 分开: + // - 看门狗选时限时两者任一为真都算「流已开始」,改用停滞时限(思考几分钟是正常的); + // - 编排器判「可安全重放」仍只看 streamedAnyToken——思考文本重放一遍用户只是再看一次 + // 思考卡,正文重放才会出现重复内容。 + private volatile boolean streamedAnyReasoning = false; + // 最近一次流活动时间(onNext / onReasoning / onKeepAlive 刷新),看门狗据此判定"流停滞" private volatile long lastActivityNanos = System.nanoTime(); private volatile java.util.concurrent.ScheduledFuture watchdogFuture; @@ -75,7 +79,7 @@ public void armInactivityWatchdog(int firstTokenSeconds, int inactivitySeconds) watchdogFuture = WATCHDOG.scheduleWithFixedDelay(() -> { if (terminated.get()) return; long idleSec = (System.nanoTime() - lastActivityNanos) / 1_000_000_000L; - boolean started = streamedAnyToken; + boolean started = streamedAnyToken || streamedAnyReasoning; int limitSec = started ? inactivitySeconds : firstTokenSeconds; if (idleSec >= limitSec) { log.warn("Stream {} for {}s (limit {}s) for {}, terminating round via watchdog", @@ -138,6 +142,32 @@ public void onNext(String token) { } } + /** + * 思考增量(dev-board#364):原样转发成 SSE {@code reasoning_delta},前端实时渲染进思考卡。 + * + *

刻意不进 {@link #fullContentBuilder}、不进编辑器流、不过标签解析:思考文本不是模型正文, + * 不落库、不回喂模型(契约 D:模型只看 content),也不该被写进文档。 + */ + @Override + public void onReasoning(String reasoningDelta) { + if (terminated.get() || reasoningDelta == null || reasoningDelta.isEmpty()) return; + lastActivityNanos = System.nanoTime(); + streamedAnyReasoning = true; + sseEmitterService.send(conversationId, "reasoning_delta", + "{\"content\":\"" + escapeJson(reasoningDelta) + "\"}"); + } + + /** 传输层保活注释:只刷新看门狗,不产生任何事件。 */ + @Override + public void onKeepAlive() { + if (terminated.get()) return; + lastActivityNanos = System.nanoTime(); + } + + public boolean hasStreamedReasoning() { + return streamedAnyReasoning; + } + // ==================== Editor Stream Filtering Logic(过滤后实时写入编辑器文档;SSE 事件名双轨 doc_stream_data/wps_stream_data,见 AgentOrchestrator) ==================== // Buffer for editor stream parser to handle split tags diff --git a/backend/src/main/java/com/checkba/service/ai/ChatModelFactory.java b/backend/src/main/java/com/checkba/service/ai/ChatModelFactory.java index af9d02cad..e088a22b0 100644 --- a/backend/src/main/java/com/checkba/service/ai/ChatModelFactory.java +++ b/backend/src/main/java/com/checkba/service/ai/ChatModelFactory.java @@ -330,7 +330,7 @@ private ChatLanguageModel getOrCreateOpenRouterModel(String modelId) { .baseUrl(baseUrl) .modelName(modelId) .timeout(config.getTimeout()) - // logRequests 必须为 false,理由见 streamingBuilder 的 javadoc(请求体物化) + // logRequests 必须为 false,理由见 streamingModel 的 javadoc(请求体物化) .logRequests(false) .logResponses(true) // Custom Headers for OpenRouter @@ -426,7 +426,7 @@ private ChatLanguageModel getOrCreatePlatformModel(String modelId) { .baseUrl(config.getBaseUrl()) .modelName(modelId) .timeout(config.getTimeout()) - // logRequests 必须为 false,理由见 streamingBuilder 的 javadoc(请求体物化) + // logRequests 必须为 false,理由见 streamingModel 的 javadoc(请求体物化) .logRequests(false) .logResponses(true) .build(); @@ -434,44 +434,25 @@ private ChatLanguageModel getOrCreatePlatformModel(String modelId) { } /** - * 两个流式通道(平台通道 / OpenRouter BYOK)共用的构建口径。 + * 两个流式通道(平台通道 / OpenRouter BYOK)共用的构建口径:自有的 + * {@link OpenRouterStreamingChatModel},不再是 langchain4j 0.36 的 OpenAiStreamingChatModel。 + * 换实现的原因(思考增量被 openai4j 的 Delta 丢掉、保活注释被吞)写在那个类的 javadoc 里, + * dev-board#364。 * - *

logResponses 必须为 false,这是可靠性契约不是调优。openai4j 0.23 的 - * {@code StreamingRequestExecutor$2.onFailure} 在该开关打开时,会先对 response 调 - * {@code ResponseLoggingInterceptor.log(...)},之后才走 errorHandler;而 okhttp-sse 的 - * {@code RealEventSource.onFailure(call, e)} 在「连接失败/被断、压根没拿到响应」这条路径上 - * 传的 response 恒为 null,于是 log() 里的 {@code response.code()} 抛 NPE, - * 而 onFailure 只 catch IOException —— 异常掀掉 OkHttp Dispatcher 线程, - * 紧随其后的 errorHandler 那一行永远走不到。 - * 表现是本轮既不 onComplete 也不 onError:传输层错误被整条吞掉, - * 只能等 {@link AgentStreamHandler} 的看门狗兜底(AGENT 模式下用户干等三分钟)。 - * 关掉后 onFailure 直接走 errorHandler,错误正常传到编排器的分类重试(IOException → TRANSIENT)。 - * - *

顺带一提这两行 DEBUG 日志本来也没人看:全仓没有任何 logging 级别配置, - * {@code dev.ai4j.openai4j} 停在 Spring 默认的 INFO,一行都不会打印; - * 但 slf4j 的参数是提前求值的,所以 NPE 照抛。 - * - *

非流式的 {@code OpenAiChatModel} 走 SyncRequestExecutor,没有这条路径,故不受影响。 - * - *

logRequests 也必须为 false(2026-08-29,随图片多模态一起改的)。 - * openai4j 0.23 的 {@code RequestLoggingInterceptor.logDebug} 在把参数交给 - * {@code Logger.debug} 之前先执行 {@code getBody(request)}(字节码实证: - * {@code invokestatic getBody} 在 {@code invokeinterface Logger.debug} 之前), - * 而这个版本的 getBody 不截断 base64(截断是 1.x 之后才加的)。 - * 也就是说日志级别停在 INFO 一行都不打印,却每次请求都把整个请求体物化成一个 String。 - * 纯文本时代这只是浪费;接上视觉后一张 5MB 的图 base64 后约 6.7MB, - * 一轮工具循环最多 30 轮,等于每次对话白造几百 MB 的一次性字符串垃圾。 - * 同理由三处构建口径(本方法 + 两个非流式 OpenAiChatModel)一起关掉。 + *

历史地雷备忘(仍适用于下面两个非流式 {@code OpenAiChatModel}): + * logRequests 必须为 false(2026-08-29,随图片多模态一起改的)——openai4j 0.23 的 + * {@code RequestLoggingInterceptor.logDebug} 在把参数交给 {@code Logger.debug} 之前 + * 先执行 {@code getBody(request)}(字节码实证),且这个版本不截断 base64。日志级别停在 INFO + * 一行都不打印,却每次请求都把整个请求体物化成一个 String;一张 5MB 的图 base64 后约 6.7MB, + * 一轮工具循环最多 30 轮。 + * 旧的流式通道还有一条 logResponses 必须为 false 的地雷({@code StreamingRequestExecutor$2.onFailure} + * 在 response==null 时先调 {@code ResponseLoggingInterceptor.log} 抛 NPE,errorHandler 永远走不到, + * 传输层错误被整条吞掉只能等看门狗)——自有实现持有 HTTP 层后这条路径不存在了, + * 但 {@code StreamingTransportFailureTest} 仍守着「连不上必须回调 onError」这条契约。 */ - static dev.langchain4j.model.openai.OpenAiStreamingChatModel.OpenAiStreamingChatModelBuilder - streamingBuilder(String apiKey, String baseUrl, String modelId, java.time.Duration timeout) { - return dev.langchain4j.model.openai.OpenAiStreamingChatModel.builder() - .apiKey(apiKey) - .baseUrl(baseUrl) - .modelName(modelId) - .timeout(timeout) - .logRequests(false) - .logResponses(false); + static dev.langchain4j.model.chat.StreamingChatLanguageModel + streamingModel(String apiKey, String baseUrl, String modelId, java.time.Duration timeout) { + return new OpenRouterStreamingChatModel(apiKey, baseUrl, modelId, timeout); } private dev.langchain4j.model.chat.StreamingChatLanguageModel getOrCreatePlatformStreamingModel(String modelId) { @@ -481,8 +462,7 @@ private dev.langchain4j.model.chat.StreamingChatLanguageModel getOrCreatePlatfor return streamingModelCache.computeIfAbsent(cacheKey, k -> { log.info("Creating new AWD Cloud StreamingChatModel for: {}", modelId); AiModelProperties.OpenRouter config = aiModelProperties.getOpenRouter(); - return streamingBuilder(apiKey, config.getBaseUrl(), modelId, config.getTimeout()) - .build(); + return streamingModel(apiKey, config.getBaseUrl(), modelId, config.getTimeout()); }); } @@ -526,12 +506,7 @@ private dev.langchain4j.model.chat.StreamingChatLanguageModel getOrCreateOpenRou String apiKey = resolveOpenRouterApiKey(); String baseUrl = resolveOpenRouterBaseUrl(); - return streamingBuilder(apiKey, baseUrl, modelId, config.getTimeout()) - // .defaultRequestProperties(Map.of( - // "HTTP-Referer", "https://checkba.com", - // "X-Title", "Checkba AI WorkDeck" - // )) - .build(); + return streamingModel(apiKey, baseUrl, modelId, config.getTimeout()); }); } diff --git a/backend/src/main/java/com/checkba/service/ai/OpenRouterStreamingChatModel.java b/backend/src/main/java/com/checkba/service/ai/OpenRouterStreamingChatModel.java new file mode 100644 index 000000000..297919994 --- /dev/null +++ b/backend/src/main/java/com/checkba/service/ai/OpenRouterStreamingChatModel.java @@ -0,0 +1,273 @@ +package com.checkba.service.ai; + +import com.fasterxml.jackson.databind.DeserializationFeature; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import dev.ai4j.openai4j.Json; +import dev.ai4j.openai4j.OpenAiHttpException; +import dev.ai4j.openai4j.chat.ChatCompletionChoice; +import dev.ai4j.openai4j.chat.ChatCompletionRequest; +import dev.ai4j.openai4j.chat.ChatCompletionResponse; +import dev.ai4j.openai4j.chat.Delta; +import dev.ai4j.openai4j.shared.StreamOptions; +import dev.langchain4j.agent.tool.ToolSpecification; +import dev.langchain4j.data.message.AiMessage; +import dev.langchain4j.data.message.ChatMessage; +import dev.langchain4j.model.StreamingResponseHandler; +import dev.langchain4j.model.chat.StreamingChatLanguageModel; +import dev.langchain4j.model.openai.InternalOpenAiHelper; +import dev.langchain4j.model.openai.OpenAiStreamingResponseBuilder; +import okhttp3.Call; +import okhttp3.Callback; +import okhttp3.MediaType; +import okhttp3.OkHttpClient; +import okhttp3.Request; +import okhttp3.RequestBody; +import okhttp3.Response; +import okhttp3.ResponseBody; +import okio.BufferedSource; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.IOException; +import java.time.Duration; +import java.util.List; +import java.util.concurrent.atomic.AtomicBoolean; + +/** + * OpenAI 兼容流式通道(OpenRouter BYOK / 平台通道,两者都打 OpenRouter)的自有实现。 + * + *

为什么不再用 langchain4j 0.36 的 {@code OpenAiStreamingChatModel}(dev-board#364): + *

    + *
  1. 思考型模型的 {@code delta.reasoning} 在 openai4j 0.23 的 {@code Delta} 里没有对应字段, + * 反序列化那一刻就丢了,链路上游没有任何钩子能拿到它。要给用户实时看到思考过程、 + * 要让看门狗知道模型还活着,只能自己读这条流。
  2. + *
  3. 它的 okhttp-sse 解析器会静默吞掉 SSE 注释行,而 OpenRouter 恰恰用 + * {@code : OPENROUTER PROCESSING} 注释在模型静默期间保活——这是「上游还在跑」的唯一信号。
  4. + *
  5. 它的传输层错误处理有一条 NPE 路径({@code logResponses=true} 时 + * {@code ResponseLoggingInterceptor.log(null response)}),历史上靠关掉日志开关绕过; + * 自己持有 HTTP 层之后这条地雷自然消失。
  6. + *
+ * + *

刻意复用、不重写的部分:消息与工具定义到 OpenAI 报文的转换 + * ({@link InternalOpenAiHelper#toOpenAiMessages} / {@link InternalOpenAiHelper#toTools}, + * 含 ImageContent 多模态编组)、请求体序列化(openai4j 的 {@link Json},snake_case + NON_NULL)、 + * 流式片段到 {@code Response} 的组装({@link OpenAiStreamingResponseBuilder}, + * 含 tool_calls 按 index 拼接、usage、finish_reason)。这些是协议细节最多、最容易写错的地方, + * 本类只负责 HTTP + SSE 行协议 + 把 reasoning/注释行多转发两条通道。 + * + *

与 langchain4j 原实现保持一致的请求参数:{@code stream=true}、 + * {@code stream_options.include_usage=true}、{@code temperature=0.7}(0.36 的默认值)。 + * 错误语义也对齐:非 2xx 抛 {@link OpenAiHttpException}({@code LlmErrorClassifier} 按状态码分类), + * 连接失败等 IOException 原样回调 {@code onError}。 + * + *

logRequests/logResponses 这类请求体物化的坑在本类不存在:请求体只序列化一次直接发出, + * 响应按行消费不落整段字符串。 + */ +public final class OpenRouterStreamingChatModel implements StreamingChatLanguageModel { + + private static final Logger log = LoggerFactory.getLogger(OpenRouterStreamingChatModel.class); + private static final MediaType JSON = MediaType.get("application/json; charset=utf-8"); + /** langchain4j 0.36 OpenAiStreamingChatModel 的默认温度,保持行为不变。 */ + private static final double DEFAULT_TEMPERATURE = 0.7; + /** 只用于读 reasoning 与 error 两个 openai4j 不认识的字段;DTO 本身仍交给 openai4j 的注解解析。 */ + private static final ObjectMapper LENIENT = new ObjectMapper() + .configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false); + + private final OkHttpClient client; + private final String endpoint; + private final String apiKey; + private final String modelName; + + public OpenRouterStreamingChatModel(String apiKey, String baseUrl, String modelName, Duration timeout) { + this.apiKey = apiKey; + this.modelName = modelName; + String base = baseUrl == null ? "" : baseUrl; + this.endpoint = (base.endsWith("/") ? base.substring(0, base.length() - 1) : base) + "/chat/completions"; + Duration t = timeout == null ? Duration.ofSeconds(60) : timeout; + // 四个超时同值,与 openai4j 0.23 的 OpenAiClient 口径一致:callTimeout 是整通墙钟上限, + // readTimeout 靠 OpenRouter 的保活注释刷新,不会在模型静默思考时误触发 + this.client = new OkHttpClient.Builder() + .callTimeout(t) + .connectTimeout(t) + .readTimeout(t) + .writeTimeout(t) + .build(); + } + + public String modelName() { + return modelName; + } + + @Override + public void generate(List messages, StreamingResponseHandler handler) { + generate(messages, (List) null, handler); + } + + @Override + public void generate(List messages, ToolSpecification toolSpecification, + StreamingResponseHandler handler) { + // 单工具强制调用形态本仓不用;按 langchain4j 的语义退化成「只提供这一个工具」 + generate(messages, toolSpecification == null ? null : List.of(toolSpecification), handler); + } + + @Override + public void generate(List messages, List toolSpecifications, + StreamingResponseHandler handler) { + ChatCompletionRequest.Builder rb = ChatCompletionRequest.builder() + .stream(true) + .streamOptions(StreamOptions.builder().includeUsage(true).build()) + .model(modelName) + .messages(InternalOpenAiHelper.toOpenAiMessages(messages)) + .temperature(DEFAULT_TEMPERATURE); + if (toolSpecifications != null && !toolSpecifications.isEmpty()) { + rb.tools(InternalOpenAiHelper.toTools(toolSpecifications, false)); + } + String body = Json.toJson(rb.build()); + + Request request = new Request.Builder() + .url(endpoint) + .header("Authorization", "Bearer " + (apiKey == null ? "" : apiKey)) + .header("Accept", "text/event-stream") + .header("User-Agent", "AI-WorkDeck") + .post(RequestBody.create(body, JSON)) + .build(); + + StreamSession session = new StreamSession(handler); + client.newCall(request).enqueue(new Callback() { + @Override + public void onFailure(Call call, IOException e) { + session.fail(e); + } + + @Override + public void onResponse(Call call, Response response) { + try (Response r = response) { + ResponseBody rb0 = r.body(); + if (!r.isSuccessful()) { + String errBody = rb0 == null ? "" : rb0.string(); + session.fail(new OpenAiHttpException(r.code(), errBody)); + return; + } + if (rb0 == null) { + session.fail(new IOException("empty response body from " + endpoint)); + return; + } + session.consume(rb0.source()); + } catch (Throwable t) { + session.fail(t); + } + } + }); + } + + /** 一次调用的流状态:SSE 行协议 + 终态幂等。 */ + private static final class StreamSession { + private final StreamingResponseHandler handler; + private final ReasoningStreamingHandler reasoningHandler; + private final OpenAiStreamingResponseBuilder builder = new OpenAiStreamingResponseBuilder(); + private final AtomicBoolean settled = new AtomicBoolean(false); + + StreamSession(StreamingResponseHandler handler) { + this.handler = handler; + this.reasoningHandler = handler instanceof ReasoningStreamingHandler rh ? rh : null; + } + + void consume(BufferedSource source) throws IOException { + StringBuilder data = new StringBuilder(); + String line; + while ((line = source.readUtf8Line()) != null) { + if (settled.get()) return; + if (line.isEmpty()) { + if (data.length() > 0) { + boolean done = dispatch(data.toString()); + data.setLength(0); + if (done) return; + } + continue; + } + if (line.charAt(0) == ':') { + // SSE 注释行:OpenRouter 的 ": OPENROUTER PROCESSING" 保活 + if (reasoningHandler != null) reasoningHandler.onKeepAlive(); + continue; + } + if (line.startsWith("data:")) { + if (data.length() > 0) data.append('\n'); + data.append(line.substring(5).trim()); + } + // event:/id:/retry: 行在 chat completions 流里不出现,忽略 + } + if (data.length() > 0 && dispatch(data.toString())) return; + // 没等到 [DONE] 就 EOF:与 okhttp-sse 的 onClosed 行为一致,按已收到的内容收尾 + complete(); + } + + /** 返回 true 表示流已终结([DONE] 或错误),调用方停止读取。 */ + private boolean dispatch(String payload) { + if ("[DONE]".equals(payload)) { + complete(); + return true; + } + JsonNode root; + try { + root = LENIENT.readTree(payload); + } catch (IOException e) { + log.warn("Unparseable SSE chunk ignored: {}", abbreviate(payload)); + return false; + } + // OpenRouter 会在 HTTP 200 之后用 data 事件送上游错误({"error":{"code":429,...}}), + // 不认这一形态的话本轮会按空回复静默收尾 + JsonNode error = root.get("error"); + if (error != null && error.isObject() && !root.has("choices")) { + int code = error.path("code").isInt() ? error.path("code").asInt() : 500; + fail(new OpenAiHttpException(code, payload)); + return true; + } + ChatCompletionResponse chunk; + try { + chunk = LENIENT.readValue(payload, ChatCompletionResponse.class); + } catch (IOException e) { + log.warn("SSE chunk does not fit ChatCompletionResponse, ignored: {}", abbreviate(payload)); + return false; + } + builder.append(chunk); + List choices = chunk.choices(); + if (choices != null && !choices.isEmpty()) { + Delta delta = choices.get(0).delta(); + if (delta != null && delta.content() != null) { + handler.onNext(delta.content()); + } + } + if (reasoningHandler != null) { + String reasoning = reasoningDeltaOf(root); + if (reasoning != null && !reasoning.isEmpty()) { + reasoningHandler.onReasoning(reasoning); + } + } + return false; + } + + /** OpenRouter 统一成 {@code delta.reasoning};各家原生兼容端点多用 {@code reasoning_content}。 */ + static String reasoningDeltaOf(JsonNode root) { + JsonNode delta = root.path("choices").path(0).path("delta"); + if (delta.isMissingNode()) return null; + JsonNode r = delta.get("reasoning"); + if (r == null || r.isNull()) r = delta.get("reasoning_content"); + return r == null || !r.isTextual() ? null : r.asText(); + } + + private void complete() { + if (!settled.compareAndSet(false, true)) return; + handler.onComplete(builder.build()); + } + + void fail(Throwable t) { + if (!settled.compareAndSet(false, true)) return; + handler.onError(t); + } + + private static String abbreviate(String s) { + return s.length() > 200 ? s.substring(0, 200) + "..." : s; + } + } +} diff --git a/backend/src/main/java/com/checkba/service/ai/ReasoningStreamingHandler.java b/backend/src/main/java/com/checkba/service/ai/ReasoningStreamingHandler.java new file mode 100644 index 000000000..064157cbe --- /dev/null +++ b/backend/src/main/java/com/checkba/service/ai/ReasoningStreamingHandler.java @@ -0,0 +1,31 @@ +package com.checkba.service.ai; + +import dev.langchain4j.data.message.AiMessage; +import dev.langchain4j.model.StreamingResponseHandler; + +/** + * 在 langchain4j 0.36 的 {@link StreamingResponseHandler} 之上补两条通道,给思考型模型用。 + * + *

背景(dev-board#364):OpenRouter 对 Kimi K3 这类思考型模型从第 4 秒起就在流式返回 + * {@code delta.reasoning},同时 {@code delta.content} 恒为空串,直到思考结束才开始吐正文。 + * openai4j 0.23 的 {@code Delta} 只有 role/content/toolCalls/functionCall 四个字段, + * reasoning 在反序列化那一刻就被丢掉;langchain4j 又只对非 null 的 content 调 {@code onNext}, + * 于是几百秒的思考期间编排器收到的全是 {@code onNext("")}——看门狗被空 token 喂活、 + * 前端一个字节都收不到,用户分不清是死机还是在想。 + * + *

两个方法都是 default 空实现:回放评测与各测试里的脚本模型仍按老接口写,不受影响。 + */ +public interface ReasoningStreamingHandler extends StreamingResponseHandler { + + /** 思考增量(OpenRouter 的 {@code delta.reasoning} / 各家原生的 {@code reasoning_content})。 */ + default void onReasoning(String reasoningDelta) { + } + + /** + * 传输层有字节但不是内容:OpenRouter 在模型生成期间每隔几秒发一行 + * {@code : OPENROUTER PROCESSING} 注释保活。它证明「连接活着、上游还在跑」, + * 看门狗据此刷新活动时间,不再把静默思考的模型当成死流。 + */ + default void onKeepAlive() { + } +} diff --git a/backend/src/main/java/com/checkba/service/ai/SseEmitterService.java b/backend/src/main/java/com/checkba/service/ai/SseEmitterService.java index a9b83dc37..005ed814a 100644 --- a/backend/src/main/java/com/checkba/service/ai/SseEmitterService.java +++ b/backend/src/main/java/com/checkba/service/ai/SseEmitterService.java @@ -82,24 +82,39 @@ private static boolean bufferable(String eventName) { @jakarta.annotation.PostConstruct void startHeartbeat() { - heartbeatScheduler.scheduleWithFixedDelay(() -> { - // 整个任务体必须吞掉 Throwable:scheduleWithFixedDelay 的语义是"任务抛出即 - // 永久取消后续执行"。心跳一旦停摆,全体在线客户端 40 秒后同时判死连接、 - // 齐刷刷重连,而且直到进程重启都好不了——单点故障放大成全局故障。 - // send() 内部只捕获 Exception,Error(OOM/StackOverflow)会漏出来。 - try { - for (String id : emitters.keySet()) { - try { - send(id, "heartbeat", "{\"ts\":" + System.currentTimeMillis() + "}"); - } catch (Throwable t) { - // 单个会话的心跳失败不许连累其余会话 - log.debug("Heartbeat failed for {}", id); - } + heartbeatScheduler.scheduleWithFixedDelay(this::heartbeatSweep, + HEARTBEAT_INTERVAL_SECONDS, HEARTBEAT_INTERVAL_SECONDS, java.util.concurrent.TimeUnit.SECONDS); + } + + /** + * 一轮心跳广播;返回本轮推送到的连接数(回归用例的观察口)。 + * 从调度 lambda 里提出来是为了能不等 15 秒就测「每个在线连接都收到 heartbeat」。 + */ + int heartbeatSweep() { + // 整个任务体必须吞掉 Throwable:scheduleWithFixedDelay 的语义是"任务抛出即 + // 永久取消后续执行"。心跳一旦停摆,全体在线客户端 40 秒后同时判死连接、 + // 齐刷刷重连,而且直到进程重启都好不了——单点故障放大成全局故障。 + // send() 内部只捕获 Exception,Error(OOM/StackOverflow)会漏出来。 + int sent = 0; + try { + for (String id : emitters.keySet()) { + try { + send(id, "heartbeat", "{\"ts\":" + System.currentTimeMillis() + "}"); + sent++; + } catch (Throwable t) { + // 单个会话的心跳失败不许连累其余会话 + log.debug("Heartbeat failed for {}", id); } - } catch (Throwable t) { - log.warn("Heartbeat sweep failed, continuing", t); } - }, HEARTBEAT_INTERVAL_SECONDS, HEARTBEAT_INTERVAL_SECONDS, java.util.concurrent.TimeUnit.SECONDS); + } catch (Throwable t) { + log.warn("Heartbeat sweep failed, continuing", t); + } + return sent; + } + + /** 心跳间隔(秒)。前端 useAgentStream 的 HEARTBEAT_STALE_MS 按它的 3 倍设,改这里要同步那边。 */ + static long heartbeatIntervalSeconds() { + return HEARTBEAT_INTERVAL_SECONDS; } @jakarta.annotation.PreDestroy diff --git a/backend/src/test/java/com/checkba/service/ai/AgentStreamHandlerReasoningTest.java b/backend/src/test/java/com/checkba/service/ai/AgentStreamHandlerReasoningTest.java new file mode 100644 index 000000000..8e889f8ea --- /dev/null +++ b/backend/src/test/java/com/checkba/service/ai/AgentStreamHandlerReasoningTest.java @@ -0,0 +1,89 @@ +package com.checkba.service.ai; + +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; + +/** + * 思考增量在编排层的两条契约(dev-board#364): + *

    + *
  1. 转发成 SSE {@code reasoning_delta},且不进正文(不落库、不回喂模型);
  2. + *
  3. 思考型模型静默几分钟是正常的——思考增量与保活注释都要刷新看门狗, + * 否则 60 秒首字节时限会把正在思考的模型掐掉、按瞬时错误重放,白烧一轮思考 token。
  4. + *
+ */ +class AgentStreamHandlerReasoningTest { + + private static AgentStreamHandler handler(SseEmitterService sse, AtomicReference sink, CountDownLatch fired) { + AgentStreamHandler h = new AgentStreamHandler(sse, "conv-reasoning-test", + mock(TokenUsageService.class), "1", 1L, "moonshotai/kimi-k3", 0L); + h.setOnError(err -> { + sink.set(err); + fired.countDown(); + }); + return h; + } + + @Test + @DisplayName("思考增量转发成 reasoning_delta 事件,不进 text_delta、不算作已流出正文") + void reasoningIsForwardedAsItsOwnEvent() { + SseEmitterService sse = mock(SseEmitterService.class); + AgentStreamHandler h = handler(sse, new AtomicReference<>(), new CountDownLatch(1)); + + h.onReasoning("先核对\"第三条\""); + + ArgumentCaptor payload = ArgumentCaptor.forClass(Object.class); + verify(sse).send(eq("conv-reasoning-test"), eq("reasoning_delta"), payload.capture()); + assertEquals("{\"content\":\"先核对\\\"第三条\\\"\"}", String.valueOf(payload.getValue()), + "思考文本要经过 JSON 转义,引号/换行不能把事件打坏"); + verify(sse, never()).send(anyString(), eq("text_delta"), anyString()); + assertTrue(h.hasStreamedReasoning()); + assertFalse(h.hasStreamedTokens(), + "可重放判定只看正文:思考卡重放一遍无害,正文重放才会让用户看到重复内容"); + } + + @Test + @DisplayName("只有思考增量、还没有正文时,看门狗改用停滞时限而不是首字节时限") + void reasoningKeepsWatchdogFromFiringOnFirstTokenBudget() throws Exception { + AtomicReference err = new AtomicReference<>(); + CountDownLatch fired = new CountDownLatch(1); + AgentStreamHandler h = handler(mock(SseEmitterService.class), err, fired); + + h.armInactivityWatchdog(1, 3600); + h.onReasoning("模型在想"); + + assertFalse(fired.await(8, TimeUnit.SECONDS), + "思考型模型正在吐 reasoning 却按「零字节」被掐——这就是 K3 被 60 秒看门狗反复打断的病灶"); + } + + @Test + @DisplayName("保活注释刷新看门狗:上游只发 keep-alive 也不许按死流终止") + void keepAliveRefreshesWatchdog() throws Exception { + AtomicReference err = new AtomicReference<>(); + CountDownLatch fired = new CountDownLatch(1); + AgentStreamHandler h = handler(mock(SseEmitterService.class), err, fired); + + // 首字节时限 2s,但每 500ms 来一次保活:idle 永远到不了 2s + h.armInactivityWatchdog(2, 3600); + long until = System.currentTimeMillis() + 6500; + while (System.currentTimeMillis() < until) { + h.onKeepAlive(); + if (fired.await(500, TimeUnit.MILLISECONDS)) break; + } + assertFalse(fired.getCount() == 0, + "OpenRouter 的 \": OPENROUTER PROCESSING\" 是上游还在跑的证据,看门狗必须认它"); + } +} diff --git a/backend/src/test/java/com/checkba/service/ai/ChatModelFactoryTest.java b/backend/src/test/java/com/checkba/service/ai/ChatModelFactoryTest.java index 0c42dc06a..42eb96582 100644 --- a/backend/src/test/java/com/checkba/service/ai/ChatModelFactoryTest.java +++ b/backend/src/test/java/com/checkba/service/ai/ChatModelFactoryTest.java @@ -7,7 +7,6 @@ import dev.langchain4j.model.ollama.OllamaChatModel; import dev.langchain4j.model.ollama.OllamaStreamingChatModel; import dev.langchain4j.model.openai.OpenAiChatModel; -import dev.langchain4j.model.openai.OpenAiStreamingChatModel; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; @@ -74,7 +73,7 @@ void openRouterProviderWithNullModelUsesOpenRouter() { "供应商切到 OPENROUTER 后空 modelId 不应回退本地 Ollama"); StreamingChatLanguageModel streaming = factory.getStreamingChatModel(null); - assertInstanceOf(OpenAiStreamingChatModel.class, streaming); + assertInstanceOf(OpenRouterStreamingChatModel.class, streaming); } @Test @@ -86,7 +85,7 @@ void openRouterProviderWithUnknownModelFallsBackToOpenRouterDefault() { assertInstanceOf(OpenAiChatModel.class, model); StreamingChatLanguageModel streaming = factory.getStreamingChatModel("vendor/some-unlisted-model"); - assertInstanceOf(OpenAiStreamingChatModel.class, streaming); + assertInstanceOf(OpenRouterStreamingChatModel.class, streaming); } @Test @@ -221,7 +220,7 @@ void resolveTargetAgreesWithStreamingDispatch() { properties.setProvider(AiModelProperties.Provider.OPENROUTER); assertEquals(AiModelProperties.Provider.OPENROUTER, factory.resolveTarget("moonshotai/kimi-k3", false).channel()); - assertInstanceOf(OpenAiStreamingChatModel.class, factory.getStreamingChatModel("moonshotai/kimi-k3")); + assertInstanceOf(OpenRouterStreamingChatModel.class, factory.getStreamingChatModel("moonshotai/kimi-k3")); when(systemSettingService.get(eq("ai.activeProvider"), any())).thenReturn("OLLAMA"); factory.clearCache(); @@ -342,7 +341,7 @@ void platformChannelTakesPrecedenceOverAllowlistShortcut() { when(platformAiChannel.keyFingerprint()).thenReturn("abc123"); assertInstanceOf(OpenAiChatModel.class, factory.getChatModel("anthropic/claude-sonnet-5")); - assertInstanceOf(OpenAiStreamingChatModel.class, factory.getStreamingChatModel("anthropic/claude-sonnet-5")); + assertInstanceOf(OpenRouterStreamingChatModel.class, factory.getStreamingChatModel("anthropic/claude-sonnet-5")); // 白名单短路分支绝不能先命中——那条路用的是 BYOK 的 key verify(platformAiChannel, atLeastOnce()).apiKey(); verify(systemSettingService, never()).get(eq("external.openrouter.apiKey"), any()); @@ -372,7 +371,7 @@ void failoverCandidateStaysOnPlatformChannel() { when(platformAiChannel.keyFingerprint()).thenReturn("abc123"); // 编排器故障转移就是拿备选 modelId 再调一次工厂——通道由 provider 决定,与 modelId 无关 - assertInstanceOf(OpenAiStreamingChatModel.class, + assertInstanceOf(OpenRouterStreamingChatModel.class, factory.getStreamingChatModel("qwen/qwen3.7-flash")); verify(platformAiChannel, atLeastOnce()).apiKey(); verify(systemSettingService, never()).get(eq("external.openrouter.apiKey"), any()); diff --git a/backend/src/test/java/com/checkba/service/ai/OpenRouterStreamingChatModelTest.java b/backend/src/test/java/com/checkba/service/ai/OpenRouterStreamingChatModelTest.java new file mode 100644 index 000000000..255dafdc3 --- /dev/null +++ b/backend/src/test/java/com/checkba/service/ai/OpenRouterStreamingChatModelTest.java @@ -0,0 +1,235 @@ +package com.checkba.service.ai; + +import com.sun.net.httpserver.HttpServer; +import dev.ai4j.openai4j.OpenAiHttpException; +import dev.langchain4j.agent.tool.ToolSpecification; +import dev.langchain4j.data.message.AiMessage; +import dev.langchain4j.data.message.UserMessage; +import dev.langchain4j.model.StreamingResponseHandler; +import dev.langchain4j.model.output.FinishReason; +import dev.langchain4j.model.output.Response; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import java.io.IOException; +import java.io.OutputStream; +import java.net.InetSocketAddress; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * 自有流式客户端的协议契约(dev-board#364)。 + * + *

假服务端回放的是 2026-09-02 对 OpenRouter {@code moonshotai/kimi-k3} 真机探测抓到的 + * 片段形状:思考期间每个 chunk 都是 {@code content:""} + {@code reasoning:"…"} + + * {@code reasoning_details:[…]},前面夹着 {@code : OPENROUTER PROCESSING} 注释行; + * 正文开始后 reasoning 变 null;最后一个 chunk 只带 usage、choices 为空。 + * 旧实现(langchain4j 0.36 / openai4j 0.23)对这条流的行为是:思考全部丢掉、 + * 注释行吞掉、只往 handler 灌一串空字符串——这就是「思考 281 秒前端零提示」的成因。 + */ +class OpenRouterStreamingChatModelTest { + + private HttpServer server; + private volatile String lastRequestBody; + private volatile String lastAuthHeader; + + @BeforeEach + void startServer() throws IOException { + server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + server.start(); + } + + @AfterEach + void stopServer() { + server.stop(0); + } + + private String baseUrl() { + return "http://127.0.0.1:" + server.getAddress().getPort() + "/api/v1"; + } + + /** 注册一条 200 + text/event-stream 的回放;body 里的 "\n" 原样写出。 */ + private void serveStream(String sseBody) { + serve(200, "text/event-stream", sseBody); + } + + private void serve(int status, String contentType, String body) { + server.createContext("/api/v1/chat/completions", ex -> { + // openai4j 的 Json 开了 INDENT_OUTPUT("stream" : true),断言前去掉空白 + lastRequestBody = new String(ex.getRequestBody().readAllBytes(), StandardCharsets.UTF_8).replaceAll("\\s+", ""); + lastAuthHeader = ex.getRequestHeaders().getFirst("Authorization"); + byte[] bytes = body.getBytes(StandardCharsets.UTF_8); + ex.getResponseHeaders().add("Content-Type", contentType); + ex.sendResponseHeaders(status, bytes.length); + try (OutputStream os = ex.getResponseBody()) { + os.write(bytes); + } + }); + } + + private static String chunk(String deltaJson, String finish) { + return "data: {\"id\":\"gen-1\",\"object\":\"chat.completion.chunk\",\"choices\":[{\"index\":0," + + "\"delta\":" + deltaJson + ",\"finish_reason\":" + (finish == null ? "null" : "\"" + finish + "\"") + + "}]}\n\n"; + } + + /** 收集三条通道的测试 handler。 */ + private static final class Collector implements ReasoningStreamingHandler { + final StringBuilder content = new StringBuilder(); + final StringBuilder reasoning = new StringBuilder(); + final AtomicInteger keepAlives = new AtomicInteger(); + final AtomicInteger emptyTokens = new AtomicInteger(); + final AtomicReference> done = new AtomicReference<>(); + final AtomicReference error = new AtomicReference<>(); + final CountDownLatch settled = new CountDownLatch(1); + + @Override public void onNext(String token) { + if (token.isEmpty()) emptyTokens.incrementAndGet(); + content.append(token); + } + @Override public void onReasoning(String d) { reasoning.append(d); } + @Override public void onKeepAlive() { keepAlives.incrementAndGet(); } + @Override public void onComplete(Response r) { done.set(r); settled.countDown(); } + @Override public void onError(Throwable t) { error.set(t); settled.countDown(); } + + void await() throws InterruptedException { + assertTrue(settled.await(10, TimeUnit.SECONDS), "流必须在 10 秒内到达终态"); + } + } + + @Test + @DisplayName("思考增量走 onReasoning、正文走 onNext、保活注释走 onKeepAlive,三条通道互不串") + void reasoningContentAndKeepAliveAreSeparated() throws Exception { + serveStream("" + + ": OPENROUTER PROCESSING\n\n" + + ": OPENROUTER PROCESSING\n\n" + + chunk("{\"role\":\"assistant\",\"content\":\"\",\"reasoning\":\"先看\"," + + "\"reasoning_details\":[{\"type\":\"reasoning.text\",\"text\":\"先看\"}]}", null) + + chunk("{\"role\":\"assistant\",\"content\":\"\",\"reasoning\":\"条款\"," + + "\"reasoning_details\":[{\"type\":\"reasoning.text\",\"text\":\"条款\"}]}", null) + + chunk("{\"role\":\"assistant\",\"content\":\"修订\",\"reasoning\":null}", null) + + chunk("{\"role\":\"assistant\",\"content\":\"完成\"}", null) + + chunk("{\"role\":\"assistant\",\"content\":\"\"}", "stop") + + "data: {\"id\":\"gen-1\",\"choices\":[],\"usage\":{\"prompt_tokens\":10,\"completion_tokens\":5,\"total_tokens\":15}}\n\n" + + "data: [DONE]\n\n"); + + OpenRouterStreamingChatModel model = new OpenRouterStreamingChatModel( + "sk-test", baseUrl(), "moonshotai/kimi-k3", Duration.ofSeconds(5)); + Collector c = new Collector(); + model.generate(List.of(UserMessage.from("修订这份合同")), c); + c.await(); + + assertNull(c.error.get(), () -> "不该出错:" + c.error.get()); + assertEquals("先看条款", c.reasoning.toString(), "reasoning 字段必须逐段转发,旧实现在这里是空串"); + assertEquals("修订完成", c.content.toString(), "正文只含 content,思考文本一个字都不许混进去"); + assertTrue(c.keepAlives.get() >= 2, "SSE 注释行是「上游还在跑」的信号,必须转发给看门狗"); + Response r = c.done.get(); + assertNotNull(r); + assertEquals("修订完成", r.content().text()); + assertEquals(FinishReason.STOP, r.finishReason()); + assertNotNull(r.tokenUsage()); + assertEquals(10, r.tokenUsage().inputTokenCount()); + assertEquals(5, r.tokenUsage().outputTokenCount()); + + assertEquals("Bearer sk-test", lastAuthHeader); + assertTrue(lastRequestBody.contains("\"stream\":true"), lastRequestBody); + assertTrue(lastRequestBody.contains("\"stream_options\""), "与 langchain4j 一致:要 usage 就得带 stream_options"); + assertTrue(lastRequestBody.contains("\"include_usage\":true"), lastRequestBody); + assertTrue(lastRequestBody.contains("\"model\":\"moonshotai/kimi-k3\""), lastRequestBody); + assertFalse(lastRequestBody.contains("\"tools\""), "没传工具就不该出现 tools 字段"); + } + + @Test + @DisplayName("工具调用增量按 index 拼装成 ToolExecutionRequest,工具定义随请求下发") + void toolCallsAreAssembled() throws Exception { + serveStream("" + + chunk("{\"role\":\"assistant\",\"content\":null,\"tool_calls\":[{\"index\":0,\"id\":\"call_1\"," + + "\"type\":\"function\",\"function\":{\"name\":\"read_file\",\"arguments\":\"{\\\"pa\"}}]}", null) + + chunk("{\"tool_calls\":[{\"index\":0,\"function\":{\"arguments\":\"th\\\":\\\"a.docx\\\"}\"}}]}", null) + + chunk("{}", "tool_calls") + + "data: [DONE]\n\n"); + + OpenRouterStreamingChatModel model = new OpenRouterStreamingChatModel( + "sk-test", baseUrl(), "deepseek/deepseek-v4-flash", Duration.ofSeconds(5)); + Collector c = new Collector(); + ToolSpecification spec = ToolSpecification.builder().name("read_file").description("read a file").build(); + model.generate(List.of(UserMessage.from("读 a.docx")), List.of(spec), c); + c.await(); + + assertNull(c.error.get(), () -> "不该出错:" + c.error.get()); + Response r = c.done.get(); + assertTrue(r.content().hasToolExecutionRequests()); + assertEquals("read_file", r.content().toolExecutionRequests().get(0).name()); + assertEquals("{\"path\":\"a.docx\"}", r.content().toolExecutionRequests().get(0).arguments()); + assertEquals(FinishReason.TOOL_EXECUTION, r.finishReason()); + assertTrue(lastRequestBody.contains("\"tools\":[{\"type\":\"function\""), lastRequestBody); + assertTrue(lastRequestBody.contains("\"name\":\"read_file\""), lastRequestBody); + } + + @Test + @DisplayName("不实现 ReasoningStreamingHandler 的旧 handler 照常拿正文,思考增量静默跳过") + void plainHandlerStillGetsContent() throws Exception { + serveStream("" + + chunk("{\"content\":\"\",\"reasoning\":\"想一下\"}", null) + + chunk("{\"content\":\"答案\"}", "stop") + + "data: [DONE]\n\n"); + OpenRouterStreamingChatModel model = new OpenRouterStreamingChatModel( + "sk-test", baseUrl(), "deepseek/deepseek-v4-flash", Duration.ofSeconds(5)); + StringBuilder content = new StringBuilder(); + CountDownLatch settled = new CountDownLatch(1); + AtomicReference err = new AtomicReference<>(); + model.generate(List.of(UserMessage.from("hi")), new StreamingResponseHandler() { + @Override public void onNext(String token) { content.append(token); } + @Override public void onComplete(Response response) { settled.countDown(); } + @Override public void onError(Throwable error) { err.set(error); settled.countDown(); } + }); + assertTrue(settled.await(10, TimeUnit.SECONDS)); + assertNull(err.get()); + assertEquals("答案", content.toString()); + } + + @Test + @DisplayName("非 2xx 映射成 OpenAiHttpException 并带状态码——LlmErrorClassifier 靠它把 429 归入限流") + void httpErrorKeepsStatusCode() throws Exception { + serve(429, "application/json", "{\"error\":{\"message\":\"Rate limit exceeded\",\"code\":429}}"); + OpenRouterStreamingChatModel model = new OpenRouterStreamingChatModel( + "sk-test", baseUrl(), "deepseek/deepseek-v4-flash", Duration.ofSeconds(5)); + Collector c = new Collector(); + model.generate(List.of(UserMessage.from("hi")), c); + c.await(); + + assertNull(c.done.get(), "HTTP 错误绝不能按成功收尾"); + assertTrue(c.error.get() instanceof OpenAiHttpException, String.valueOf(c.error.get())); + assertEquals(429, ((OpenAiHttpException) c.error.get()).code()); + assertEquals(LlmErrorClassifier.Kind.RATE_LIMITED, LlmErrorClassifier.classify(c.error.get())); + } + + @Test + @DisplayName("HTTP 200 里用 data 事件送来的上游错误也要当错误,不许按空回复静默收尾") + void inStreamErrorObjectIsSurfaced() throws Exception { + serveStream("data: {\"error\":{\"message\":\"Provider returned error\",\"code\":502}}\n\n"); + OpenRouterStreamingChatModel model = new OpenRouterStreamingChatModel( + "sk-test", baseUrl(), "deepseek/deepseek-v4-flash", Duration.ofSeconds(5)); + Collector c = new Collector(); + model.generate(List.of(UserMessage.from("hi")), c); + c.await(); + + assertNull(c.done.get()); + assertTrue(c.error.get() instanceof OpenAiHttpException, String.valueOf(c.error.get())); + assertEquals(502, ((OpenAiHttpException) c.error.get()).code()); + assertEquals(LlmErrorClassifier.Kind.TRANSIENT, LlmErrorClassifier.classify(c.error.get())); + } +} diff --git a/backend/src/test/java/com/checkba/service/ai/SseEmitterServiceTest.java b/backend/src/test/java/com/checkba/service/ai/SseEmitterServiceTest.java index 7a76f8c31..2eb388d10 100644 --- a/backend/src/test/java/com/checkba/service/ai/SseEmitterServiceTest.java +++ b/backend/src/test/java/com/checkba/service/ai/SseEmitterServiceTest.java @@ -108,6 +108,18 @@ void heartbeatsAreNotBuffered() { assertEquals(java.util.List.of("text_delta"), svc.bufferedEventNamesSince(id, 0)); } + @Test + void heartbeatSweepReachesEveryLiveConnection() { + // 心跳是前端区分「连接活着、模型还在想」与「连接断了」的唯一依据(dev-board#364)。 + // 调度器每 15s 调一次 heartbeatSweep;这里直接调,不等真实间隔。 + SseEmitterService svc = new SseEmitterService(); + svc.createConnection("conv-hb-1", "paneA"); + svc.createConnection("conv-hb-2", "paneA"); + assertEquals(2, svc.heartbeatSweep(), "每个在线连接都要收到一次 heartbeat"); + // 前端 useAgentStream 的 HEARTBEAT_STALE_MS=45000 按「3 个心跳周期」设:改间隔要同步那边 + assertEquals(15L, SseEmitterService.heartbeatIntervalSeconds()); + } + @Test void noLastEventIdReplaysNothing() { SseEmitterService svc = new SseEmitterService(); diff --git a/backend/src/test/java/com/checkba/service/ai/StreamingTransportFailureTest.java b/backend/src/test/java/com/checkba/service/ai/StreamingTransportFailureTest.java index 19eee5a7b..04a85d050 100644 --- a/backend/src/test/java/com/checkba/service/ai/StreamingTransportFailureTest.java +++ b/backend/src/test/java/com/checkba/service/ai/StreamingTransportFailureTest.java @@ -38,8 +38,10 @@ * at okhttp3.internal.sse.RealEventSource.onFailure(RealEventSource.kt:91) * * - *

本用例走 {@link ChatModelFactory#streamingBuilder} —— 必须是工厂那份真实口径, - * 用例里自己拼一个 builder 就永远是绿的,起不到守护作用。 + *

本用例走 {@link ChatModelFactory#streamingModel} —— 必须是工厂那份真实口径, + * 用例里自己拼一个模型就永远是绿的,起不到守护作用。 + * 流式通道 2026-09 换成自有的 {@link OpenRouterStreamingChatModel}(dev-board#364)之后, + * 上面那条 NPE 路径已不存在,但「连不上必须 onError」的契约不变,本用例继续守着它。 */ class StreamingTransportFailureTest { @@ -53,11 +55,11 @@ private static int portWithNothingListening() throws IOException { @Test @DisplayName("连不上时必须回调 onError,而不是把异常吞在 OkHttp Dispatcher 上让本轮静默挂死") void transportFailureReachesOnError() throws Exception { - StreamingChatLanguageModel model = ChatModelFactory.streamingBuilder( + StreamingChatLanguageModel model = ChatModelFactory.streamingModel( "test-key", "http://127.0.0.1:" + portWithNothingListening() + "/api/v1", "deepseek/deepseek-v4-flash", - Duration.ofSeconds(5)).build(); + Duration.ofSeconds(5)); CountDownLatch settled = new CountDownLatch(1); AtomicReference captured = new AtomicReference<>(); diff --git a/frontend/src/components/ChatInterface.vue b/frontend/src/components/ChatInterface.vue index 94e805465..b95f00356 100644 --- a/frontend/src/components/ChatInterface.vue +++ b/frontend/src/components/ChatInterface.vue @@ -467,6 +467,11 @@ {{ continueHint }} {{ $t('chat.continueRun') }} + + + {{ $t('chat.linkReconnecting', { attempt: linkStatus.attempt }) }} + @@ -719,6 +724,7 @@ export default { fileChanges, agentPaused, agentRunStatus, + linkStatus, activeSkills, skillNotice, reattachSSE, @@ -2531,6 +2537,7 @@ export default { cleanTitle, recentDotClass, agentRunStatus, + linkStatus, addFile, removeContextFile, removePastedImage, @@ -4337,6 +4344,22 @@ export default { background: var(--awd-accent-hover); } +/* SSE 断连提示条:外形对齐 continue-bar,只有一行文字、没有按钮(重连是自动的) */ +.link-bar { + display: flex; + align-items: center; + margin: 0 12px 6px; + padding: 6px 12px; + background: var(--awd-bg); + border: 1px solid var(--awd-warning); + border-radius: 8px; +} + +.link-hint { + font-size: 11px; + color: var(--awd-warning-text); +} + /* 后台任务控制条(停止):外形对齐 continue-bar,但用中性底色—— 这不是「需要你处理」的黄色警示,只是一个随时可用的控制 */ .task-control-bar { diff --git a/frontend/src/composables/useAgentStream.js b/frontend/src/composables/useAgentStream.js index 0a829aa8a..55bb745dd 100644 --- a/frontend/src/composables/useAgentStream.js +++ b/frontend/src/composables/useAgentStream.js @@ -28,6 +28,11 @@ export function useAgentStream() { const bubbles = ref([]) const isConnected = ref(false) const isStreaming = ref(false) + // SSE 链路状态(dev-board#364):'live' = 连接活着(心跳在跳,等模型是正常的); + // 'reconnecting' = 心跳超时/流断了,正在退避重连(attempt 是第几次)。 + // 之前断线重连只写 console.warn,用户看到的是计时器一直走、什么都不发生—— + // 分不清「模型在想」和「连接死了」。这个状态给输入区一条明确的提示条。 + const linkStatus = ref({ state: 'live', attempt: 0 }) const error = ref(null) const currentConversationId = ref(null) // STATE: Token Usage Tracking (Session Cumulative) @@ -115,6 +120,7 @@ export function useAgentStream() { if (!currentConversationId.value) return const delay = Math.min(30000, 1000 * Math.pow(2, reconnectAttempts)) reconnectAttempts++ + linkStatus.value = { state: 'reconnecting', attempt: reconnectAttempts } console.warn(`[AgentStream] SSE 断开(${reason}),${delay}ms 后自动重连(第 ${reconnectAttempts} 次)`) reconnectTimer = setTimeout(async () => { reconnectTimer = null @@ -202,6 +208,7 @@ export function useAgentStream() { // Reset connection states isConnected.value = false isStreaming.value = false + linkStatus.value = { state: 'live', attempt: 0 } // Reset parser state resetParser() // Reset event parser state @@ -360,6 +367,7 @@ export function useAgentStream() { isConnected.value = true reconnectAttempts = 0 + linkStatus.value = { state: 'live', attempt: 0 } lastSseActivityAt = Date.now() startHeartbeatMonitor() // 网络恢复/回前台时经模块级单例回调触发本实例重连 @@ -831,7 +839,17 @@ export function useAgentStream() { return } - if (evt === 'text_delta') { + if (evt === 'reasoning_delta') { + // 思考型模型的 reasoning 增量(dev-board#364):后端按 OpenRouter 的 delta.reasoning + // 原样转发,这里直接写进思考卡,**不过标签解析器**——思考文本不是协议正文, + // 里面出现 / 字样只是模型在自言自语,不能当成标签处理。 + try { + const d = JSON.parse(dataStr) + appendReasoning(d.content || '') + } catch (e) { + appendReasoning(dataStr) + } + } else if (evt === 'text_delta') { try { const d = JSON.parse(dataStr) // 调试日志:显示 text_delta 内容 @@ -1149,6 +1167,30 @@ export function useAgentStream() { } } + // 思考增量的落点与 标签同一套:还没有过程卡时写顶层思考卡(ghost 态实时 + // 滚动显示),已经有工具过程后(多轮工具循环中间的再思考)挂到最后一个过程卡的 + // 思考条目上——与 flushContent 的 thinking 分支同口径,否则第二轮起的思考会被 + // 记到首轮的顶层卡上、把首轮的时长越算越长。 + const appendReasoning = (text) => { + const bubble = currentAssistantBubble.value + if (!bubble || !text) return + if (bubble.processes.length > 0) { + const lastProc = bubble.processes[bubble.processes.length - 1] + const lastItem = lastProc.items.length > 0 ? lastProc.items[lastProc.items.length - 1] : null + if (!lastItem || lastItem.type !== 'thinking' || lastItem.status === 'done') { + lastProc.items.push({ type: 'thinking', status: 'thinking', content: text, startTime: Date.now(), fromReasoning: true }) + } else { + lastItem.content += text + } + return + } + if (bubble.thinking.status !== 'thinking') { + bubble.thinking.status = 'thinking' + if (!bubble.thinking.startTime) bubble.thinking.startTime = Date.now() + } + bubble.thinking.content += text + } + const flushContent = (text) => { const bubble = currentAssistantBubble.value if (!bubble || !text) return @@ -1273,6 +1315,18 @@ export function useAgentStream() { th.endTime = Date.now() th.duration = th.startTime ? (th.endTime - th.startTime) / 1000 : 0 } + // reasoning_delta 在过程卡里建的思考条目没有 来收尾:正文/下一个标签 + // 一到就算想完了,否则那张过程卡会一直显示「运行中」到整轮结束 + const procs = bubble && bubble.processes + if (procs && procs.length > 0) { + const items = procs[procs.length - 1].items || [] + const last = items[items.length - 1] + if (last && last.type === 'thinking' && last.status === 'thinking' && last.fromReasoning) { + last.status = 'done' + last.endTime = Date.now() + last.duration = last.startTime ? (last.endTime - last.startTime) / 1000 : 0 + } + } } const handleTag = (tagName, isClose, attrs, fullTag) => { @@ -1654,6 +1708,8 @@ export function useAgentStream() { agentPaused, agentRunStatus, agentAwaitingInput, + // SSE 链路状态:'live' / 'reconnecting'(含第几次),输入区据此显示断连提示条 + linkStatus, // 本轮生效的 Skill 与「刚自动加载了一个技能」的轻提示 activeSkills, skillNotice, diff --git a/frontend/src/locales/en-US/chat.js b/frontend/src/locales/en-US/chat.js index 0dd108dd1..368fb1c90 100644 --- a/frontend/src/locales/en-US/chat.js +++ b/frontend/src/locales/en-US/chat.js @@ -105,6 +105,8 @@ export default { busyToast: 'The AI is running another task. Please try again later.', abortToast: 'Stopping. In-flight calls may still return briefly.', continueRun: 'Resume', + // SSE reconnect bar (dev-board#364): heartbeat timed out or the stream ended unexpectedly + linkReconnecting: 'Connection to the server was lost; reconnecting automatically (attempt {attempt}). The AI keeps running in the background and will resume once reconnected.', continueHintInterrupted: 'The app was closed during the last run. The task was interrupted.', continueHintPaused: 'Step limit reached for this run. The task is paused.', continuePrompt: 'Continue', diff --git a/frontend/src/locales/zh-CN/chat.js b/frontend/src/locales/zh-CN/chat.js index 75e4a6540..472f028ba 100644 --- a/frontend/src/locales/zh-CN/chat.js +++ b/frontend/src/locales/zh-CN/chat.js @@ -105,6 +105,8 @@ export default { busyToast: 'AI 正在执行其他任务,请稍后再试', abortToast: '正在停止,在途的调用可能还会回一小段', continueRun: '继续执行', + // SSE 断连提示条(dev-board#364):心跳超时或流意外结束,后台在自动重连 + linkReconnecting: '与服务器的连接已断开,正在自动重连(第 {attempt} 次)。AI 仍在后台运行,重连后会续上。', continueHintInterrupted: '上次任务执行中应用被关闭,任务已中断', continueHintPaused: '已达单轮执行步数上限,任务已暂停', continuePrompt: '继续', diff --git a/frontend/tests/project-home/reasoning-stream.test.mjs b/frontend/tests/project-home/reasoning-stream.test.mjs new file mode 100644 index 000000000..9393823c0 --- /dev/null +++ b/frontend/tests/project-home/reasoning-stream.test.mjs @@ -0,0 +1,135 @@ +// dev-board#364:思考型模型(Kimi K3 等)思考几百秒期间前端零提示。 +// +// 三条契约: +// 1. 等待首 token 的活性计时——ThinkingCard 从「发送」那一刻起读秒(startTime 由 +// useAgentStream 在 sendMessage 里写入),文案走 chat.thinkingLive 两套 locale; +// 2. reasoning_delta 事件实时进思考卡,不过标签解析器; +// 3. 心跳超时/流断了要给用户看得见的提示(linkStatus → chat.linkReconnecting), +// 且前端判死阈值与后端心跳间隔(15s)保持 3 倍关系。 +// +// useAgentStream.js 带 @/ 别名与 uni 全局,node 直接 import 不进来,那一侧做源码级 +// 契约断言;ThinkingCard 的读秒用真实 Vue 响应式跑(与 thinking-card-ghost-collapse +// 同一套挂载法)。 +import { test } from 'node:test' +import assert from 'node:assert/strict' +import { readFileSync } from 'node:fs' +import { ref, watch, computed, onMounted, onUnmounted } from 'vue' + +const read = (p) => readFileSync(new URL(p, import.meta.url), 'utf8') +const STREAM = read('../../src/composables/useAgentStream.js') +const CHAT_UI = read('../../src/components/ChatInterface.vue') +const CARD = read('../../src/components/AgentMessage/ThinkingCard.vue') +const ZH = read('../../src/locales/zh-CN/chat.js') +const EN = read('../../src/locales/en-US/chat.js') +const SSE_JAVA = read('../../../backend/src/main/java/com/checkba/service/ai/SseEmitterService.java') + +const stripComments = (s) => s.replace(/\/\*[\s\S]*?\*\//g, '').replace(/^\s*\/\/.*$/gm, '') +const CODE = stripComments(STREAM) + +// 与 thinking-card-ghost-collapse.test.mjs 同款:真实 setInterval 会让进程退不出去 +const realSetInterval = globalThis.setInterval +globalThis.setInterval = (fn, ms, ...args) => { + const t = realSetInterval(fn, ms, ...args) + if (t && typeof t.unref === 'function') t.unref() + return t +} + +function mountThinkingCard(props) { + const body = CARD.match(/