From ec58ec474bfe48fefb295d140da4f1354328fa4b Mon Sep 17 00:00:00 2001 From: zichun <26684461+reporkey@users.noreply.github.com> Date: Wed, 29 Jul 2026 01:32:03 +0800 Subject: [PATCH] A: ai_chat_event MySQL ledger (append-only, iSeq/iTurn, Redis hot cache) + token-budget projection (conservative estimator, in-turn demotion, pre-send self-check, prompt_eval_count calibration) --- sql/ai_chat_event.sql | 21 +++++++++++++++++++++ src/main/java/com/xly/agent/AgentIdentity.java | 9 ++++++++- src/main/java/com/xly/agent/EventLogChatMemory.java | 45 ++++++++++++++++++++++++++------------------- src/main/java/com/xly/config/AgentFactory.java | 4 ++-- src/main/java/com/xly/config/TracingChatModelListener.java | 29 +++++++++++++++++++++++++++++ src/main/java/com/xly/service/AuthzService.java | 4 ++-- src/main/java/com/xly/service/EventProjectionService.java | 387 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------------------------------------------------------------------------------------------ src/main/java/com/xly/service/LedgerService.java | 174 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------------------- src/main/java/com/xly/service/TokenEstimator.java | 75 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ src/main/java/com/xly/web/AgentChatController.java | 2 +- src/main/resources/application.yml | 3 +++ src/test/java/com/xly/service/ConversationScopeTest.java | 2 +- src/test/java/com/xly/service/EventLogConcurrencyTest.java | 2 +- src/test/java/com/xly/service/EventProjectionTest.java | 39 ++++++++++++++++++++++++++++++++------- src/test/java/com/xly/service/InMemoryLedger.java | 7 ++++--- 15 files changed, 644 insertions(+), 159 deletions(-) create mode 100644 sql/ai_chat_event.sql create mode 100644 src/main/java/com/xly/service/TokenEstimator.java diff --git a/sql/ai_chat_event.sql b/sql/ai_chat_event.sql new file mode 100644 index 0000000..f66c1b0 --- /dev/null +++ b/sql/ai_chat_event.sql @@ -0,0 +1,21 @@ +-- ai_chat_event:会话事件账本(唯一事实源,append-only,永不删除)。 +-- 每事件一行:一句话 / 一次工具调用 / 一次工具结果 / 一次按钮点击,比"轮"更细。 +-- 字段按 ERP 惯例命名(sMakePerson/tCreateDate/sBrandsId/sSubsidiaryId)。 +-- sPayload = 事件 JSON 全文(含 t/type 与业务字段,与 Redis 热缓存同格式)。 +-- iSeq = 会话内序号 1..n(服务端赋值,UNIQUE 防并发撞号;纯溯源列,不进投影/API)。 +-- iTurn = 轮次 1..m(user/form_submit 开新轮,其余事件继承当前轮;纯溯源列)。 +-- Redis(chat:ledger:{convId},30 天 TTL)只是热缓存:读优先走 Redis,缓存失效回源本表。 +CREATE TABLE IF NOT EXISTS ai_chat_event ( + iId bigint NOT NULL AUTO_INCREMENT PRIMARY KEY, + sConversationId varchar(96) NOT NULL, + iSeq int NOT NULL, + iTurn int NOT NULL DEFAULT 1, + sMakePerson varchar(64) NULL, + sBrandsId varchar(32) NULL, + sSubsidiaryId varchar(32) NULL, + sType varchar(32) NOT NULL, + sPayload mediumtext, + tCreateDate datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, + UNIQUE KEY uk_conv_seq (sConversationId, iSeq), + KEY idx_conv (sConversationId, iId) +); diff --git a/src/main/java/com/xly/agent/AgentIdentity.java b/src/main/java/com/xly/agent/AgentIdentity.java index db1c885..266e025 100644 --- a/src/main/java/com/xly/agent/AgentIdentity.java +++ b/src/main/java/com/xly/agent/AgentIdentity.java @@ -18,12 +18,15 @@ public final class AgentIdentity { private final String token; private final String userId; private final String brandsId; + private final String subsidiaryId; private final Set grantedModuleIds; // null = 全部(管理员) - public AgentIdentity(String token, String userId, String brandsId, Set grantedModuleIds) { + public AgentIdentity(String token, String userId, String brandsId, String subsidiaryId, + Set grantedModuleIds) { this.token = token; this.userId = userId; this.brandsId = brandsId; + this.subsidiaryId = subsidiaryId; this.grantedModuleIds = grantedModuleIds; } @@ -40,6 +43,10 @@ public final class AgentIdentity { return brandsId; } + public String subsidiaryId() { + return subsidiaryId; + } + public boolean canAccessModule(String moduleId) { return grantedModuleIds == null || (moduleId != null && grantedModuleIds.contains(moduleId)); } diff --git a/src/main/java/com/xly/agent/EventLogChatMemory.java b/src/main/java/com/xly/agent/EventLogChatMemory.java index ad7c2d1..16a5faa 100644 --- a/src/main/java/com/xly/agent/EventLogChatMemory.java +++ b/src/main/java/com/xly/agent/EventLogChatMemory.java @@ -5,7 +5,6 @@ import com.xly.service.LedgerService; import dev.langchain4j.data.message.AiMessage; import dev.langchain4j.data.message.ChatMessage; import dev.langchain4j.agent.tool.ToolExecutionRequest; -import dev.langchain4j.data.message.ChatMessageSerializer; import dev.langchain4j.data.message.SystemMessage; import dev.langchain4j.data.message.ToolExecutionResultMessage; import dev.langchain4j.data.message.UserMessage; @@ -17,14 +16,17 @@ import java.util.List; import java.util.Map; /** - * 事件日志上的对话记忆:{@link #add} 把模型环消息逐条 append 成事件(rightPush 原子—— - * 对话流与确认端点并发写互不覆盖,替代旧 chat:mem 整包读改写),{@link #messages()} 读取 - * {@link EventProjectionService} 的四段式投影。日志即唯一事实源,本类不持有任何会话状态。 + * 事件账本上的对话记忆:{@link #add} 把模型环消息逐条 append 成事件(追加原子—— + * 对话流与按钮端点并发写互不覆盖),{@link #messages()} 读取 + * {@link EventProjectionService} 的四段式 token 记账投影。账本即唯一事实源,本类不持有任何会话状态。 * *

写者分工:用户事件由控制器在收到请求时先落账(前端立即可见),本类对 UserMessage 做 * 去重跳过;{@code internalUserTurn=true} 时(反编造护栏重试的注入话术)例外——落一条 * {@code internal} 标记的用户事件,LLM 可见、前端历史不显示。模型消息(tool_call/tool_result/ai) * 只由本类落账。 + * + *

tool_call 事件存**中立格式** {@code calls=[{id,name,args}]}(脱 LangChain4j 序列化耦合), + * 旧版 {@code payload} 载荷投影层仍可读。 */ public class EventLogChatMemory implements ChatMemory { @@ -33,23 +35,23 @@ public class EventLogChatMemory implements ChatMemory { private final String convId; private final LedgerService log; private final EventProjectionService projection; - private final int charBudget; + private final AgentIdentity identity; private final boolean internalUserTurn; - /** system prompt 每次由 provider 提供,不落日志(保持日志纯业务事件)。 */ + /** system prompt 每次由 provider 提供,不落账(保持账本纯业务事件)。 */ private volatile String systemText; /** - * 本轮实际喂给模型的用户文本(原话 + 编排层附加的 grounding/状态后缀)。日志只存原话; + * 本轮实际喂给模型的用户文本(原话 + 编排层附加的后缀)。账本只存原话; * 投影时把当前轮用户消息替换为它。P2 编排层不再加后缀后,此字段恒空。 */ private volatile String currentTurnUserText; public EventLogChatMemory(String convId, LedgerService log, EventProjectionService projection, - int charBudget, boolean internalUserTurn) { + AgentIdentity identity, boolean internalUserTurn) { this.convId = convId; this.log = log; this.projection = projection; - this.charBudget = charBudget; + this.identity = identity; this.internalUserTurn = internalUserTurn; } @@ -75,22 +77,28 @@ public class EventLogChatMemory implements ChatMemory { if (internalUserTurn) { data.put("internal", true); } - log.append(convId, "user", data); + log.append(convId, "user", data, identity); return; } if (m instanceof AiMessage am) { if (am.hasToolExecutionRequests()) { Map data = new LinkedHashMap<>(); - data.put("payload", ChatMessageSerializer.messagesToJson(List.of(am))); data.put("text", am.text() == null ? "" : am.text()); + List> calls = new ArrayList<>(); List tools = new ArrayList<>(); for (ToolExecutionRequest r : am.toolExecutionRequests()) { + Map c = new LinkedHashMap<>(); + c.put("id", r.id() == null ? "" : r.id()); + c.put("name", r.name() == null ? "" : r.name()); + c.put("args", r.arguments() == null ? "" : r.arguments()); + calls.add(c); tools.add(r.name()); } + data.put("calls", calls); data.put("tools", tools); - log.append(convId, "tool_call", data); + log.append(convId, "tool_call", data, identity); } else if (am.text() != null && !am.text().isBlank()) { - log.append(convId, "ai", Map.of("text", am.text())); + log.append(convId, "ai", Map.of("text", am.text()), identity); } return; } @@ -101,13 +109,13 @@ public class EventLogChatMemory implements ChatMemory { data.put("name", tr.toolName() == null ? "" : tr.toolName()); data.put("text", text); data.put("digest", digest(text)); - log.append(convId, "tool_result", data); + log.append(convId, "tool_result", data, identity); } } /** * 控制器已把本轮用户**原话**落账(最近的用户事件是喂给模型文本的前缀、且其后无模型事件) - * → 跳过,避免重复。编排层可能在原话后附加 grounding/状态后缀,故用前缀而非全等判定。 + * → 跳过,避免重复。编排层可能在原话后附加后缀,故用前缀而非全等判定。 */ private boolean alreadyLoggedPrefixOf(String text) { List> tail = log.events(convId, 6); @@ -120,7 +128,7 @@ public class EventLogChatMemory implements ChatMemory { case "form_submit", "ai", "tool_call", "tool_result", "assistant": return false; default: - // confirm/proposal 等按钮事件可能插在中间,继续往前找 + // queued/confirm 等按钮事件可能插在中间,继续往前找 } } return false; @@ -128,8 +136,7 @@ public class EventLogChatMemory implements ChatMemory { @Override public List messages() { - List out = projection.project( - systemText, log.events(convId, LLM_EVENT_WINDOW), charBudget); + List out = projection.project(systemText, log.events(convId, LLM_EVENT_WINDOW)); String full = currentTurnUserText; if (full != null) { // 当前轮用户消息换成实际喂给模型的完整文本(原话+编排后缀) @@ -144,7 +151,7 @@ public class EventLogChatMemory implements ChatMemory { return out; } - /** 日志是唯一事实源,不因模型环异常清史;删除会话走 ConversationService 级联。 */ + /** 账本是唯一事实源,不因模型环异常清史;删除会话走 ConversationService 级联。 */ @Override public void clear() { } diff --git a/src/main/java/com/xly/config/AgentFactory.java b/src/main/java/com/xly/config/AgentFactory.java index e533227..da79c81 100644 --- a/src/main/java/com/xly/config/AgentFactory.java +++ b/src/main/java/com/xly/config/AgentFactory.java @@ -91,9 +91,9 @@ public class AgentFactory { .streamingChatModel(streamingModel) .tools(tools) .maxSequentialToolsInvocations(8) // 循环护栏:防止 askUser/工具无限自我循环 - // 事件日志记忆:逐条 append 原子落账,读取为四段式投影(见 EventLogChatMemory) + // 事件账本记忆:逐条 append 原子落账,读取为四段式 token 记账投影(见 EventLogChatMemory) .chatMemoryProvider(memoryId -> new EventLogChatMemory( - String.valueOf(memoryId), ledger, projection, 6000, internalUserTurn)) + String.valueOf(memoryId), ledger, projection, identity, internalUserTurn)) .systemMessageProvider(memoryId -> systemPromptService.prompt()) .build(); } diff --git a/src/main/java/com/xly/config/TracingChatModelListener.java b/src/main/java/com/xly/config/TracingChatModelListener.java index 48147a1..4138bb1 100644 --- a/src/main/java/com/xly/config/TracingChatModelListener.java +++ b/src/main/java/com/xly/config/TracingChatModelListener.java @@ -53,12 +53,21 @@ public class TracingChatModelListener implements ChatModelListener { this.mapper = mapper; } + @org.springframework.beans.factory.annotation.Value("${llm.context-length:16384}") + private int contextLength; + @Override public void onRequest(ChatModelRequestContext ctx) { ctx.attributes().put("t0", System.nanoTime()); ctx.attributes().put("startTs", Instant.now().toString()); String model = ctx.chatRequest() == null ? null : ctx.chatRequest().modelName(); ctx.attributes().put("model", model == null ? "unknown" : model); + try { // 校准回路:记下本地保守估算,onResponse 时与 prompt_eval_count 对账 + if (ctx.chatRequest() != null && ctx.chatRequest().messages() != null) { + ctx.attributes().put("estIn", com.xly.service.TokenEstimator.estimate(ctx.chatRequest().messages())); + } + } catch (Exception ignore) { + } } @Override @@ -73,10 +82,30 @@ public class TracingChatModelListener implements ChatModelListener { } } catch (Exception ignore) { } + calibrate(ctx.attributes().get("estIn"), in, out); record(String.valueOf(ctx.attributes().get("model")), ctx.attributes().get("t0"), String.valueOf(ctx.attributes().get("startTs")), in, out, null); } + /** + * token 估算校准回路:prompt_eval_count(响应免费自带)vs 本地保守估算。 + * 实际 > 估算 = 估算器不够保守(截头风险);实际+输出逼近 num_ctx = 很可能已被 Ollama 静默截头,都告警。 + */ + private void calibrate(Object est, Integer actualIn, Integer actualOut) { + if (!(est instanceof Integer e) || actualIn == null || actualIn <= 0) { + return; + } + double ratio = actualIn / (double) e; + log.info("LLM token calib est={} actual={} ratio={}", e, actualIn, String.format("%.2f", ratio)); + if (actualIn > e) { + log.warn("LLM token 估算偏低(est={} < actual={})——保守系数不足,存在被 num_ctx 截头的风险", e, actualIn); + } + int total = actualIn + (actualOut == null ? 0 : actualOut); + if (total > contextLength * 0.95) { + log.warn("LLM 上下文逼近 num_ctx={}(in+out={}),Ollama 可能已从最前静默截断——请检查投影预算", contextLength, total); + } + } + @Override public void onError(ChatModelErrorContext ctx) { Throwable e = ctx.error(); diff --git a/src/main/java/com/xly/service/AuthzService.java b/src/main/java/com/xly/service/AuthzService.java index 54d54f7..26f3d37 100644 --- a/src/main/java/com/xly/service/AuthzService.java +++ b/src/main/java/com/xly/service/AuthzService.java @@ -64,7 +64,7 @@ public class AuthzService { String subsidiaryId = w.path("sSubsidiaryId").asText(""); String userType = w.path("sType").asText(""); Set granted = grantedIds(userId, userType, brandsId, subsidiaryId); - return new com.xly.agent.AgentIdentity(token, userId, brandsId, granted); + return new com.xly.agent.AgentIdentity(token, userId, brandsId, subsidiaryId, granted); } return erp.devLoginEnabled() ? devIdentity() : null; } @@ -76,7 +76,7 @@ public class AuthzService { public com.xly.agent.AgentIdentity devIdentity() { String uid = devUserIdOverride != null && !devUserIdOverride.isBlank() ? devUserIdOverride : resolveDevUserId(); Set granted = grantedIds(uid, devUserType, devBrand, devSub); - return new com.xly.agent.AgentIdentity(null, uid, devBrand, granted); + return new com.xly.agent.AgentIdentity(null, uid, devBrand, devSub, granted); } /** null = 全部(管理员);否则 = 有权的 id 集合。 */ diff --git a/src/main/java/com/xly/service/EventProjectionService.java b/src/main/java/com/xly/service/EventProjectionService.java index 4d4577f..c875abb 100644 --- a/src/main/java/com/xly/service/EventProjectionService.java +++ b/src/main/java/com/xly/service/EventProjectionService.java @@ -2,123 +2,302 @@ package com.xly.service; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; +import dev.langchain4j.agent.tool.ToolExecutionRequest; import dev.langchain4j.data.message.AiMessage; import dev.langchain4j.data.message.ChatMessage; import dev.langchain4j.data.message.ChatMessageDeserializer; import dev.langchain4j.data.message.SystemMessage; import dev.langchain4j.data.message.ToolExecutionResultMessage; import dev.langchain4j.data.message.UserMessage; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import java.util.ArrayList; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.function.UnaryOperator; /** - * 事件日志的**唯一**投影/渲染器 —— 前端历史与 LLM 上下文同源于 {@link LedgerService} 的事件流, + * 事件账本的**唯一**投影/渲染器 —— 前端历史与 LLM 上下文同源于 {@link LedgerService} 的事件流, * 文案在这一处生成,不再散落在各控制器里。 * - *

LLM 上下文四段式({@link #project}): + *

LLM 上下文四段式({@link #project}),预算按 token 记账(本地保守估算 + * {@link TokenEstimator},绝对值由 {@code llm.context-length}(镜像 Ollama 的 + * OLLAMA_CONTEXT_LENGTH,经 OpenAI 兼容通道无法逐请求传 num_ctx)导出): *

    - *
  1. 稳定 system prompt(KV 前缀稳定);
  2. - *
  3. 进行中流程卡(激活 skill 全文 + 在办单据状态)——钉在 system 尾部,不参与截断, - * 提议 executed/cancelled 或新 skill 激活时摘下;
  4. - *
  5. 往事摘要区:预算外旧轮 → 确定性一行摘要(零模型调用);
  6. - *
  7. 近期原文区:约 charBudget 字符按**整轮**纳入;旧轮工具结果压 {@value #TOOL_DIGEST_LEN} 字; - * 当前轮(最后一个用户事件起)原样载荷精确重建,工具调用/结果配对不可破坏。
  8. + *
  9. ① 稳定 system prompt + ② 进行中流程卡 —— 永不让位;
  10. + *
  11. ④ 往事摘要区:预算外旧轮 → 确定性一行摘要(零模型调用);
  12. + *
  13. ⑤ 近期原文区:按**整轮**纳入,旧轮工具结果压 {@value #TOOL_DIGEST_LEN} 字;
  14. + *
  15. 当前轮:原样保真、工具配对不可破坏;当前轮自身超限时**轮内让位**——最早的工具结果 + * 降为摘要(保最近 {@value #CURRENT_KEEP_FULL} 个全文),单条结果超硬顶截断。
  16. *
+ * 组装后做**发送前自检**:总 token 超预算 → ⑤ 最旧轮退入 ④、再裁 ④ 最旧行;①②与当前轮不动。 + * 绝不把超限上下文交给 Ollama 静默截头(它从最前截、① 先死)。 */ @Service public class EventProjectionService { + private static final Logger log = LoggerFactory.getLogger(EventProjectionService.class); + public static final int TOOL_DIGEST_LEN = 120; private static final int DIGEST_TURNS_MAX = 40; private static final int DIGEST_LINE_LEN = 90; + /** 模型输出预留(num_ctx 覆盖 prompt+生成)。 */ + private static final int OUTPUT_RESERVE = 1024; + /** 工具 schema(随每次请求发送、不在 messages 里)的估算开销。 */ + private static final int TOOLS_OVERHEAD = 2000; + /** 当前轮增长的乐观小额预留(每次工具往返都会重投影,超出时反应式收缩⑤)。 */ + private static final int TURN_RESERVE = 512; + /** ④ 摘要区 token 上限。 */ + private static final int DIGEST_BUDGET = 800; + /** 当前轮里保全文的最近工具结果个数(轮内让位时更早的降摘要)。 */ + private static final int CURRENT_KEEP_FULL = 2; + /** 当前轮单条工具结果硬顶(token),超出截断。 */ + private static final int SINGLE_RESULT_CAP = 2400; + private final ObjectMapper mapper; + private final int contextLength; - public EventProjectionService(ObjectMapper mapper) { + /** opId → 人类可读状态行(Phase C:流程卡只读展示 ai_op_queue.sStatus)。可空(测试/未装配)。 */ + private volatile UnaryOperator opStatusLookup; + + public EventProjectionService(ObjectMapper mapper, @Value("${llm.context-length:16384}") int contextLength) { this.mapper = mapper; + this.contextLength = contextLength; + } + + public EventProjectionService(ObjectMapper mapper) { + this(mapper, 16384); + } + + public void setOpStatusLookup(UnaryOperator lookup) { + this.opStatusLookup = lookup; + } + + /** 输入 prompt 的总预算(token)。 */ + int promptBudget() { + return Math.max(2048, contextLength - OUTPUT_RESERVE - TOOLS_OVERHEAD); } // ---------------------------------------------------------------- LLM 投影 - /** 四段式 LLM 上下文。systemText 为空时不出 system 消息(如测试)。 */ - public List project(String systemText, List> events, int charBudget) { + /** 四段式 LLM 上下文(token 记账)。systemText 为空时不出 system 消息(如测试)。 */ + public List project(String systemText, List> events) { + int budget = promptBudget(); List>> turns = groupTurns(events); int current = turns.size() - 1; - // ④ 近期原文区:从最新往回按整轮纳入;当前轮必进且不计预算 - int firstVerbatim = current; + // ①② system + 流程卡(永不让位) + SystemMessage sys = null; + int sysEst = 0; + if (systemText != null && !systemText.isBlank()) { + String card = activeCard(events); + sys = SystemMessage.from(card == null ? systemText + : systemText + "\n\n【进行中的流程】\n" + card); + sysEst = TokenEstimator.estimate(sys); + } + + // 当前轮:原样保真;自身超限时轮内让位 + List currentMsgs = new ArrayList<>(); + int curEst = 0; + if (current >= 0) { + int curAllowance = budget - sysEst - TURN_RESERVE; + currentMsgs = renderCurrentTurn(turns.get(current), curAllowance); + curEst = TokenEstimator.estimate(currentMsgs); + } + + // ⑤ 近期原文区:从最新旧轮往回按整轮纳入 + List> oldRendered = new ArrayList<>(); // index 对齐 turns[0..current-1] + int[] oldEst = new int[Math.max(0, current)]; + for (int i = 0; i < current; i++) { + List ms = new ArrayList<>(); + appendTurnMessages(ms, turns.get(i), false); + oldRendered.add(ms); + oldEst[i] = TokenEstimator.estimate(ms); + } + int remaining = budget - sysEst - curEst - TURN_RESERVE; + int firstVerbatim = fillRecent(oldEst, current, remaining, 0); + if (firstVerbatim > 0) { // 有轮进摘要区 → 给 ④ 留预算后重算 + firstVerbatim = fillRecent(oldEst, current, remaining, DIGEST_BUDGET); + } + + // ④ 摘要行 + List digestLines = new ArrayList<>(); + int digestFrom = Math.max(0, firstVerbatim - DIGEST_TURNS_MAX); + for (int i = digestFrom; i < firstVerbatim; i++) { + String line = digestLine(turns.get(i)); + if (!line.isBlank()) { + digestLines.add(line); + } + } + int omitted = digestFrom; + + // 组装 + 发送前自检(⑤ 最旧退 ④ → 裁 ④ 最旧;①② 与当前轮不动) + while (true) { + List out = assemble(sys, digestLines, omitted, oldRendered, firstVerbatim, current, currentMsgs); + int total = TokenEstimator.estimate(out); + if (total <= budget) { + return out; + } + if (firstVerbatim < current) { + String line = digestLine(turns.get(firstVerbatim)); + if (!line.isBlank()) { + digestLines.add(line); + } + firstVerbatim++; + continue; + } + if (!digestLines.isEmpty()) { + digestLines.remove(0); + omitted++; + continue; + } + log.warn("projection over budget even after shrink: total≈{}t budget={}t (当前轮过大)", total, budget); + return out; + } + } + + /** ⑤ 从最新旧轮往回装,返回 firstVerbatim(第一条进原文区的旧轮下标;current 表示一条都装不下)。 */ + private static int fillRecent(int[] oldEst, int current, int remaining, int digestAllowance) { + int budget = remaining - digestAllowance; int used = 0; + int first = current; for (int i = current - 1; i >= 0; i--) { - int size = 0; - for (Map ev : turns.get(i)) { - size += approxLen(ev); - } - if (used + size > charBudget) { + if (used + oldEst[i] > budget) { break; } - used += size; - firstVerbatim = i; + used += oldEst[i]; + first = i; } + return first; + } + private List assemble(SystemMessage sys, List digestLines, int omitted, + List> oldRendered, int firstVerbatim, int current, + List currentMsgs) { List out = new ArrayList<>(); - String card = activeCard(events); - if (systemText != null && !systemText.isBlank()) { - out.add(SystemMessage.from(card == null ? systemText - : systemText + "\n\n【进行中的流程】\n" + card)); + if (sys != null) { + out.add(sys); } - - // ③ 往事摘要区 - if (firstVerbatim > 0) { + if (!digestLines.isEmpty() || omitted > 0) { StringBuilder sb = new StringBuilder("【更早对话摘要(自动截断,仅供参考)】"); - int from = Math.max(0, firstVerbatim - DIGEST_TURNS_MAX); - if (from > 0) { - sb.append("\n(更早 ").append(from).append(" 轮已省略)"); + if (omitted > 0) { + sb.append("\n(更早 ").append(omitted).append(" 轮已省略)"); } - for (int i = from; i < firstVerbatim; i++) { - String line = digestLine(turns.get(i)); - if (!line.isBlank()) { - sb.append("\n- ").append(line); - } + for (String line : digestLines) { + sb.append("\n- ").append(line); } out.add(UserMessage.from(sb.toString())); } + for (int i = firstVerbatim; i < current; i++) { + out.addAll(oldRendered.get(i)); + } + out.addAll(currentMsgs); + return out; + } - for (int i = firstVerbatim; i < turns.size(); i++) { - appendTurnMessages(out, turns.get(i), i == current); + /** + * 当前轮渲染:原样保真;总量超 allowance 时**轮内让位**——最早的工具结果先降 + * {@value #TOOL_DIGEST_LEN} 字摘要(保最近 {@value #CURRENT_KEEP_FULL} 个全文,仍超则只保 1 个)。 + * 工具调用/结果配对永不破坏,单条结果超 {@value #SINGLE_RESULT_CAP}t 截断。 + */ + private List renderCurrentTurn(List> turn, int allowance) { + for (int keepFull = CURRENT_KEEP_FULL; keepFull >= 0; keepFull--) { + List ms = renderCurrentWith(turn, keepFull); + if (TokenEstimator.estimate(ms) <= allowance || keepFull == 0) { + return ms; + } + } + return List.of(); + } + + private List renderCurrentWith(List> turn, int keepFull) { + int results = 0; + for (Map ev : turn) { + if ("tool_result".equals(str(ev.get("type")))) { + results++; + } + } + int demote = Math.max(0, results - keepFull); + List out = new ArrayList<>(); + int seen = 0; + for (Map ev : turn) { + if ("tool_result".equals(str(ev.get("type")))) { + seen++; + String text = seen <= demote ? str(ev.get("digest")) : capped(str(ev.get("text"))); + out.add(ToolExecutionResultMessage.from(str(ev.get("tcId")), str(ev.get("name")), text)); + } else { + appendOneMessage(out, ev, true); + } } return out; } + private static String capped(String text) { + if (TokenEstimator.estimate(text) <= SINGLE_RESULT_CAP) { + return text; + } + // 按估算规则回推一个安全字符数(全 CJK 最坏情形 = 1 字 1 token) + return text.substring(0, Math.min(text.length(), SINGLE_RESULT_CAP)) + "……[超长已截断]"; + } + /** 一轮 → LLM 消息。当前轮工具结果全文保真;旧轮压一行摘要。 */ private void appendTurnMessages(List out, List> turn, boolean isCurrent) { for (Map ev : turn) { - String type = str(ev.get("type")); - switch (type) { - case "user" -> out.add(UserMessage.from(str(ev.get("text")))); - case "form_submit" -> out.add(UserMessage.from(renderFormSubmit(ev))); - case "ai", "assistant", "clarify" -> addAi(out, str(ev.get("text"))); - case "tool_call" -> { - ChatMessage m = fromPayload(str(ev.get("payload"))); - if (m != null) { - out.add(m); - } + if ("tool_result".equals(str(ev.get("type")))) { + String text = isCurrent ? str(ev.get("text")) : str(ev.get("digest")); + out.add(ToolExecutionResultMessage.from(str(ev.get("tcId")), str(ev.get("name")), text)); + } else { + appendOneMessage(out, ev, isCurrent); + } + } + } + + private void appendOneMessage(List out, Map ev, boolean isCurrent) { + String type = str(ev.get("type")); + switch (type) { + case "user" -> out.add(UserMessage.from(str(ev.get("text")))); + case "form_submit" -> out.add(UserMessage.from(renderFormSubmit(ev))); + case "ai", "assistant", "clarify" -> addAi(out, str(ev.get("text"))); + case "tool_call" -> { + ChatMessage m = toolCallMessage(ev); + if (m != null) { + out.add(m); } - case "tool_result" -> { - String text = isCurrent ? str(ev.get("text")) : str(ev.get("digest")); - out.add(ToolExecutionResultMessage.from( - str(ev.get("tcId")), str(ev.get("name")), text)); + } + case "question" -> addAi(out, str(ev.get("question"))); + case "form" -> addAi(out, "已为「" + str(ev.get("entity")) + "」弹出新建表单,等待用户填写提交。"); + case "queued" -> addAi(out, renderQueued(ev)); + case "proposal" -> addAi(out, "已生成待确认提议:" + str(ev.get("summary")) + "(等待用户点确认/取消,尚未执行)"); + case "confirm", "cancel" -> addAi(out, renderOutcome(ev)); + default -> { } // tool(旧版摘要)/skill_active/skill_done 不进正文(skill 由流程卡承载) + } + } + + /** tool_call 事件 → AiMessage:优先中立格式 calls=[{id,name,args}],回退旧版 LangChain4j 序列化载荷。 */ + private ChatMessage toolCallMessage(Map ev) { + Object calls = ev.get("calls"); + if (calls instanceof List list && !list.isEmpty()) { + List reqs = new ArrayList<>(); + for (Object o : list) { + if (o instanceof Map m) { + reqs.add(ToolExecutionRequest.builder() + .id(str(m.get("id"))) + .name(str(m.get("name"))) + .arguments(str(m.get("args"))) + .build()); } - case "question" -> addAi(out, str(ev.get("question"))); - case "form" -> addAi(out, "已为「" + str(ev.get("entity")) + "」弹出新建表单,等待用户填写提交。"); - case "proposal" -> addAi(out, "已生成待确认提议:" + str(ev.get("summary")) + "(等待用户点确认/取消,尚未执行)"); - case "confirm", "cancel" -> addAi(out, renderOutcome(ev)); - default -> { } // tool(旧版摘要)/skill_active/skill_done 不进正文(skill 由流程卡承载) + } + if (!reqs.isEmpty()) { + String text = str(ev.get("text")); + return text.isBlank() ? AiMessage.from(reqs) : AiMessage.from(text, reqs); } } + return fromPayload(str(ev.get("payload"))); } // ---------------------------------------------------------------- 前端历史投影 @@ -138,6 +317,7 @@ public class EventProjectionService { case "ai", "assistant", "clarify" -> addHist(out, str(ev.get("text"))); case "question" -> addHist(out, str(ev.get("question"))); case "form" -> addHist(out, "【表单】新建" + str(ev.get("entity")) + ":" + str(ev.get("message"))); + case "queued" -> addHist(out, "【已提交待办】" + str(ev.get("description"))); case "proposal" -> addHist(out, "【待确认】" + str(ev.get("summary"))); case "confirm" -> addHist(out, ("executed".equals(str(ev.get("status"))) ? "【已执行】" : "【执行失败】") + str(ev.get("description")) + parenOr(str(ev.get("msg")))); @@ -149,10 +329,11 @@ public class EventProjectionService { return out; } - /** agent 路径的 提问/表单/提议 不再单独落显示事件——从工具结果同源推导。 */ + /** agent 路径的 提问/表单/预览 不再单独落显示事件——从工具结果同源推导。 */ private String historyFromToolResult(Map ev) { String name = str(ev.get("name")); - if (!"askUser".equals(name) && !"collectForm".equals(name) && !"proposeWrite".equals(name)) { + if (!"askUser".equals(name) && !"collectForm".equals(name) + && !"previewChange".equals(name) && !"proposeWrite".equals(name)) { return ""; } try { @@ -163,6 +344,9 @@ public class EventProjectionService { if ("collectForm".equals(name) && "form_collect".equals(r.path("type").asText(""))) { return "【表单】新建" + r.path("entity").asText("") + ":" + r.path("message").asText(""); } + if ("previewChange".equals(name) && "change_preview".equals(r.path("type").asText(""))) { + return "【预览】" + r.path("summary").asText(""); + } if ("proposeWrite".equals(name) && !r.path("opId").asText("").isBlank()) { return "【待确认】" + r.path("summary").asText(""); } @@ -175,14 +359,13 @@ public class EventProjectionService { /** * 进行中流程卡:激活 skill 全文(最近一次 {@code skill_active},被 {@code skill_done} 或更新的 - * skill 覆盖前一直钉住)+ 在办单据状态(表单待提交 / 提议待确认;executed/cancelled/failed 即摘下)。 + * skill 覆盖前一直钉住)+ 在办单据状态(表单待提交 / 预览待保存 / 已提交待办的只读进度)。 + * {@code queued} 后技能线摘下(AI 侧流程到写入队列为止),待办行保留到下一个流程事件替换。 * 两者皆无 → null。 */ public String activeCard(List> events) { String skillText = null; - Map lastFlow = null; // form/collectForm/proposal/confirm/cancel/form_submit 中最近的一个 - String proposalOpId = null; - String proposalSummary = null; + Map lastFlow = null; for (Map ev : events) { String type = str(ev.get("type")); @@ -190,34 +373,30 @@ public class EventProjectionService { case "skill_active" -> skillText = str(ev.get("text")); case "skill_done" -> skillText = null; case "form", "form_submit", "proposal" -> lastFlow = ev; - case "confirm", "cancel" -> { // 流程终结:单据线与技能线一起摘下 + case "queued", "confirm", "cancel" -> { // AI 侧流程终结:技能线摘下 lastFlow = ev; skillText = null; } case "tool_result" -> { - String name = str(ev.get("name")); - if ("collectForm".equals(name) || "proposeWrite".equals(name)) { - Map derived = deriveFlowEvent(ev); - if (derived != null) { - lastFlow = derived; - } + Map derived = deriveFlowEvent(ev); + if (derived != null) { + lastFlow = derived; } } default -> { } } - if (lastFlow != null && "proposal".equals(str(lastFlow.get("type")))) { - proposalOpId = str(lastFlow.get("opId")); - proposalSummary = str(lastFlow.get("summary")); - } } String docLine = null; if (lastFlow != null) { switch (str(lastFlow.get("type"))) { case "form" -> docLine = "已为「" + str(lastFlow.get("entity")) + "」弹出新建表单,等待用户填写提交。"; - case "proposal" -> docLine = "待确认提议:" + proposalSummary - + "(opId=" + proposalOpId + ",用户点【确认】才会执行)。"; - default -> { } // form_submit/confirm/cancel = 流程已推进/终结,卡摘下 + case "preview" -> docLine = "已生成写操作预览卡:" + str(lastFlow.get("summary")) + + "(等待用户点卡上按钮提交,**尚未写入**)。"; + case "queued" -> docLine = renderQueuedLine(lastFlow); + case "proposal" -> docLine = "待确认提议:" + str(lastFlow.get("summary")) + + "(opId=" + str(lastFlow.get("opId")) + ",用户点【确认】才会执行)。"; + default -> { } // form_submit/confirm/cancel = 流程已推进/终结 } } @@ -237,17 +416,28 @@ public class EventProjectionService { return sb.toString(); } - /** 工具结果 → 流程事件(collectForm 弹表单 / proposeWrite 出提议),供流程卡推导。 */ + /** 工具结果 → 流程事件(collectForm 弹表单 / previewChange 出预览 / 旧版 proposeWrite),供流程卡推导。 */ private Map deriveFlowEvent(Map ev) { + String name = str(ev.get("name")); + if (!"collectForm".equals(name) && !"previewChange".equals(name) && !"proposeWrite".equals(name)) { + return null; + } try { JsonNode r = mapper.readTree(str(ev.get("text"))); - if ("collectForm".equals(str(ev.get("name"))) && "form_collect".equals(r.path("type").asText(""))) { + if ("collectForm".equals(name) && "form_collect".equals(r.path("type").asText(""))) { Map m = new LinkedHashMap<>(); m.put("type", "form"); m.put("entity", r.path("entity").asText("")); return m; } - if ("proposeWrite".equals(str(ev.get("name"))) && !r.path("opId").asText("").isBlank()) { + if ("previewChange".equals(name) && "change_preview".equals(r.path("type").asText(""))) { + Map m = new LinkedHashMap<>(); + m.put("type", "preview"); + m.put("previewId", r.path("previewId").asText("")); + m.put("summary", r.path("summary").asText("")); + return m; + } + if ("proposeWrite".equals(name) && !r.path("opId").asText("").isBlank()) { Map m = new LinkedHashMap<>(); m.put("type", "proposal"); m.put("opId", r.path("opId").asText("")); @@ -291,6 +481,7 @@ public class EventProjectionService { case "ai", "assistant", "clarify" -> outcome = str(ev.get("text")); case "question" -> outcome = "问:" + str(ev.get("question")); case "form" -> outcome = "弹出「" + str(ev.get("entity")) + "」新建表单"; + case "queued" -> outcome = "已提交待办:" + str(ev.get("description")); case "proposal" -> outcome = "生成待确认提议:" + str(ev.get("summary")); case "confirm" -> outcome = ("executed".equals(str(ev.get("status"))) ? "用户确认,执行成功:" : "确认后执行失败:") + str(ev.get("description")); @@ -323,6 +514,27 @@ public class EventProjectionService { return "提交「" + str(ev.get("entity")) + "」新增表单:" + str(ev.get("fields")); } + private String renderQueued(Map ev) { + return "(系统)用户已保存,操作已提交待办:" + str(ev.get("description")) + + "(AI 侧流程到此为止,是否/何时执行由 ERP 处理;绝不声称已执行成功)"; + } + + /** queued 流程卡行:附 ai_op_queue.sStatus 只读进度(可用时)。 */ + private String renderQueuedLine(Map ev) { + String desc = str(ev.get("description")); + String status = null; + UnaryOperator lookup = opStatusLookup; + Object opIds = ev.get("opIds"); + if (lookup != null && opIds instanceof List ids && !ids.isEmpty()) { + try { + status = lookup.apply(String.valueOf(ids.get(0))); + } catch (Exception ignore) { + } + } + return "已提交待办:" + desc + "(当前状态:" + (status == null || status.isBlank() ? "等待 ERP 处理" : status) + + ";只读进度,勿再重复提交)。"; + } + private String renderOutcome(Map ev) { String desc = str(ev.get("description")); if ("cancel".equals(str(ev.get("type")))) { @@ -335,23 +547,6 @@ public class EventProjectionService { return "(系统)用户确认后执行失败:" + desc + (msg.isBlank() ? "" : ",原因:" + msg); } - private static int approxLen(Map ev) { - String type = str(ev.get("type")); - if ("tool_result".equals(type)) { - return str(ev.get("digest")).length() + 20; // 旧轮只进摘要 - } - if ("tool_call".equals(type)) { - return str(ev.get("payload")).length() / 2 + 20; - } - int n = 0; - for (Object v : ev.values()) { - if (v instanceof String s) { - n += s.length(); - } - } - return Math.max(n, 20); - } - private static void addAi(List out, String text) { if (text != null && !text.isBlank()) { out.add(AiMessage.from(text)); diff --git a/src/main/java/com/xly/service/LedgerService.java b/src/main/java/com/xly/service/LedgerService.java index cbff91f..f607c6c 100644 --- a/src/main/java/com/xly/service/LedgerService.java +++ b/src/main/java/com/xly/service/LedgerService.java @@ -1,33 +1,36 @@ package com.xly.service; import com.fasterxml.jackson.databind.ObjectMapper; +import com.xly.agent.AgentIdentity; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.data.redis.core.StringRedisTemplate; +import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.stereotype.Service; import java.time.Duration; import java.util.ArrayList; +import java.util.Collections; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; /** - * 会话事件日志(append-only)—— 会话的**唯一事实源**:用户话、模型消息(含工具调用/结果的原样载荷)、 - * 表单/提议/确认等按钮事件,全部按发生顺序落一条 Redis LIST。前端历史与 LLM 上下文都是它的投影 - * ({@link EventProjectionService}),不再有 chat:mem 整包读改写与多处手工同步。 + * 会话事件账本(append-only)—— 会话的**唯一事实源**:用户话、模型消息(含工具调用/结果)、 + * 表单/预览/保存等按钮事件,全部按发生顺序落账。前端历史与 LLM 上下文都是它的投影 + * ({@link EventProjectionService})。 * - *

键:{@code chat:ledger:{convId}},元素为事件 JSON {@code {t,type,...}},30 天 TTL; - * rightPush 原子,任意并发写者(对话流、确认端点)互不覆盖。 + *

存储分层:MySQL {@code ai_chat_event} 为权威持久层(每事件一行、永不删除、可溯源, + * 服务端赋 iSeq 会话内序号 + iTurn 轮次),Redis LIST {@code chat:ledger:{convId}}(30 天 TTL) + * 为热缓存加速投影读;缓存失效时回源 MySQL 并回填。两边任一失败都不阻断对话(互为兜底)。 * *

事件类型: * {@code user}(text, internal?)/ {@code ai}(text)/ - * {@code tool_call}(payload=原样序列化的模型消息, text, tools)/ + * {@code tool_call}(text, calls=[{id,name,args}], tools)/ * {@code tool_result}(tcId, name, text, digest)/ - * {@code form_submit}(entity, fields)/ {@code proposal}(opId, summary)/ - * {@code confirm}(opId, status, msg, description)/ {@code cancel}(opId, description)/ + * {@code form_submit}(entity, fields)/ {@code queued}(opIds, description, summaryLines?)/ * {@code skill_active}(name, text)/ {@code skill_done}(name); - * 旧版遗留类型 {@code assistant/clarify/form/question/tool} 仍可读(兼容渲染,自然淘汰)。 + * 旧版遗留类型 {@code assistant/clarify/form/question/tool/proposal/confirm/cancel} 仍可读(兼容渲染,自然淘汰)。 */ @Service public class LedgerService { @@ -35,20 +38,32 @@ public class LedgerService { private static final Logger log = LoggerFactory.getLogger(LedgerService.class); private static final String PREFIX = "chat:ledger:"; private static final Duration TTL = Duration.ofDays(30); + /** MySQL 冷读回源的窗口上限(有界读:超长会话不整包搬运)。 */ + private static final int COLD_READ_MAX = 1000; + /** 开新轮的事件类型(其余事件继承当前轮次)。 */ + private static final java.util.Set TURN_OPENERS = java.util.Set.of("user", "form_submit"); private final StringRedisTemplate redis; private final ObjectMapper mapper; + private final JdbcTemplate jdbc; - public LedgerService(StringRedisTemplate redis, ObjectMapper mapper) { + public LedgerService(StringRedisTemplate redis, ObjectMapper mapper, JdbcTemplate jdbc) { this.redis = redis; this.mapper = mapper; + this.jdbc = jdbc; } - /** 追加一条事件(绝不抛异常——日志失败不能影响对话主流程)。 */ + /** 追加一条事件(无身份上下文的旧签名:sMakePerson 从 convId 前缀推导)。 */ public void append(String convId, String type, Map data) { + append(convId, type, data, null); + } + + /** 追加一条事件(绝不抛异常——账本失败不能影响对话主流程)。 */ + public void append(String convId, String type, Map data, AgentIdentity who) { if (convId == null || convId.isBlank()) { return; } + String json; try { Map ev = new LinkedHashMap<>(); ev.put("t", System.currentTimeMillis()); @@ -56,14 +71,70 @@ public class LedgerService { if (data != null) { ev.putAll(data); } + json = mapper.writeValueAsString(ev); + } catch (Exception e) { + log.warn("ledger serialize failed (conv={}, type={}): {}", convId, type, e.getMessage()); + return; + } + persist(convId, type, json, who); + cache(convId, json); + } + + /** MySQL 权威落库:服务端赋 iSeq/iTurn,唯一索引 + 重试防并发撞号。失败仅告警(Redis 兜底)。 */ + private void persist(String convId, String type, String json, AgentIdentity who) { + if (jdbc == null) { + return; + } + String maker = who != null && who.userId() != null ? who.userId() : userIdOf(convId); + String brands = who == null ? null : who.brandsId(); + String sub = who == null ? null : who.subsidiaryId(); + for (int attempt = 0; attempt < 3; attempt++) { + try { + int seq = 1; + int turn = TURN_OPENERS.contains(type) ? 1 : 0; + List> last = jdbc.queryForList( + "SELECT iSeq, iTurn FROM ai_chat_event WHERE sConversationId=? ORDER BY iId DESC LIMIT 1", + convId); + if (!last.isEmpty()) { + int lastSeq = ((Number) last.get(0).get("iSeq")).intValue(); + int lastTurn = ((Number) last.get(0).get("iTurn")).intValue(); + seq = lastSeq + 1; + turn = TURN_OPENERS.contains(type) ? lastTurn + 1 : lastTurn; + } + if (turn < 1) { + turn = 1; + } + jdbc.update( + "INSERT INTO ai_chat_event(sConversationId,iSeq,iTurn,sMakePerson,sBrandsId,sSubsidiaryId,sType,sPayload,tCreateDate) " + + "VALUES(?,?,?,?,?,?,?,?,NOW())", + convId, seq, turn, maker, brands, sub, type, json); + return; + } catch (org.springframework.dao.DuplicateKeyException dup) { + // 并发撞号:重读末行再试 + } catch (Exception e) { + log.warn("ledger mysql append failed (conv={}, type={}): {}", convId, type, e.getMessage()); + return; + } + } + log.warn("ledger mysql append gave up after retries (conv={}, type={})", convId, type); + } + + private void cache(String convId, String json) { + try { String key = PREFIX + convId; - redis.opsForList().rightPush(key, mapper.writeValueAsString(ev)); + redis.opsForList().rightPush(key, json); redis.expire(key, TTL); } catch (Exception e) { - log.warn("ledger append failed (conv={}, type={}): {}", convId, type, e.getMessage()); + log.warn("ledger redis append failed (conv={}): {}", convId, e.getMessage()); } } + /** convId 形如 {userId}:{local}(ConversationService.scopedId),前缀即用户 id。 */ + private static String userIdOf(String convId) { + int i = convId.indexOf(':'); + return i > 0 ? convId.substring(0, i) : null; + } + /** 全量事件(按发生顺序),供前端历史投影。 */ public List> events(String convId) { return read(convId, 0); @@ -75,28 +146,74 @@ public class LedgerService { } private List> read(String convId, int lastN) { - List> out = new ArrayList<>(); + List raw = null; try { long start = lastN <= 0 ? 0 : -lastN; - List raw = redis.opsForList().range(PREFIX + convId, start, -1); - if (raw == null) { - return out; + if (Boolean.TRUE.equals(redis.hasKey(PREFIX + convId))) { + raw = redis.opsForList().range(PREFIX + convId, start, -1); } - for (String s : raw) { - try { - @SuppressWarnings("unchecked") - Map m = mapper.readValue(s, Map.class); - out.add(m); - } catch (Exception ignore) { + } catch (Exception e) { + log.warn("ledger redis read failed (conv={}): {}", convId, e.getMessage()); + } + if (raw == null) { + raw = coldRead(convId, lastN); + } + List> out = new ArrayList<>(); + for (String s : raw) { + try { + @SuppressWarnings("unchecked") + Map m = mapper.readValue(s, Map.class); + out.add(m); + } catch (Exception ignore) { + } + } + return out; + } + + /** Redis 缓存失效 → 回源 MySQL(有界窗口)并回填缓存。 */ + private List coldRead(String convId, int lastN) { + if (jdbc == null) { + return List.of(); + } + int window = lastN <= 0 ? COLD_READ_MAX : Math.min(lastN, COLD_READ_MAX); + try { + List> rows = jdbc.queryForList( + "SELECT sType, sPayload FROM ai_chat_event WHERE sConversationId=? ORDER BY iId DESC LIMIT " + window, + convId); + if (rows.isEmpty()) { + return List.of(); + } + List out = new ArrayList<>(rows.size()); + for (Map r : rows) { + if ("deleted".equals(String.valueOf(r.get("sType")))) { + break; // 删除标记:更早的事件属于已删除的会话,不再回源 } + Object p = r.get("sPayload"); + if (p != null) { + out.add(p.toString()); + } + } + if (out.isEmpty()) { + return List.of(); } + Collections.reverse(out); + try { // 回填热缓存(尽力而为) + String key = PREFIX + convId; + redis.opsForList().rightPushAll(key, out); + redis.expire(key, TTL); + } catch (Exception ignore) { + } + return out; } catch (Exception e) { - log.warn("ledger read failed (conv={}): {}", convId, e.getMessage()); + log.warn("ledger mysql read failed (conv={}): {}", convId, e.getMessage()); + return List.of(); } - return out; } - /** 日志键是否存在(供会话列表懒清理:内容已过期的会话项从列表剔除)。 */ + /** + * 会话在热缓存或持久层是否有内容(供会话列表懒清理)。MySQL 里的事件永不删除, + * 因此这里只看 Redis 热缓存:30 天没动过的会话从侧栏消失,但账本仍可溯源。 + */ public boolean exists(String convId) { try { return Boolean.TRUE.equals(redis.hasKey(PREFIX + convId)); @@ -105,10 +222,15 @@ public class LedgerService { } } + /** + * 删除会话 = 清热缓存 + 在持久层落一条 {@code deleted} 标记(append-only,旧事件永久保留可溯源, + * 但冷读止于标记——同名会话再启用时不会复活已删除的历史)。 + */ public void delete(String convId) { try { redis.delete(PREFIX + convId); } catch (Exception ignore) { } + persist(convId, "deleted", "{\"type\":\"deleted\"}", null); } } diff --git a/src/main/java/com/xly/service/TokenEstimator.java b/src/main/java/com/xly/service/TokenEstimator.java new file mode 100644 index 0000000..ca6a074 --- /dev/null +++ b/src/main/java/com/xly/service/TokenEstimator.java @@ -0,0 +1,75 @@ +package com.xly.service; + +import dev.langchain4j.data.message.AiMessage; +import dev.langchain4j.data.message.ChatMessage; +import dev.langchain4j.data.message.SystemMessage; +import dev.langchain4j.data.message.ToolExecutionResultMessage; +import dev.langchain4j.data.message.UserMessage; +import dev.langchain4j.agent.tool.ToolExecutionRequest; + +import java.util.List; + +/** + * 本地保守 token 估算器(不依赖任何模型请求)。 + * + *

规则:CJK/全角字符 ≈ 1 token/字;其余字符 ≈ 3 字符/token。**故意高估**—— + * 高估只会少带几轮历史,低估会让 Ollama 从前面静默截头(system prompt 先死),两者代价不对称。 + * 校准回路:每次响应免费自带 prompt_eval_count,{@code TracingChatModelListener} 记录 + * 估算 vs 实际,实际逼近 num_ctx 即告警。 + */ +public final class TokenEstimator { + + /** 每条消息的结构开销(role/分隔符等)。 */ + public static final int MSG_OVERHEAD = 8; + + private TokenEstimator() { + } + + public static int estimate(String s) { + if (s == null || s.isEmpty()) { + return 0; + } + int cjk = 0; + int other = 0; + for (int i = 0; i < s.length(); i++) { + if (s.charAt(i) >= 0x2E80) { + cjk++; + } else { + other++; + } + } + return cjk + (other + 2) / 3; + } + + public static int estimate(ChatMessage m) { + if (m == null) { + return 0; + } + int n = MSG_OVERHEAD; + if (m instanceof SystemMessage sm) { + n += estimate(sm.text()); + } else if (m instanceof UserMessage um) { + n += estimate(um.hasSingleText() ? um.singleText() : String.valueOf(um.contents())); + } else if (m instanceof AiMessage am) { + n += estimate(am.text()); + if (am.hasToolExecutionRequests()) { + for (ToolExecutionRequest r : am.toolExecutionRequests()) { + n += estimate(r.name()) + estimate(r.arguments()) + 6; + } + } + } else if (m instanceof ToolExecutionResultMessage tr) { + n += estimate(tr.text()) + estimate(tr.toolName()); + } else { + n += estimate(String.valueOf(m)); + } + return n; + } + + public static int estimate(List messages) { + int n = 0; + for (ChatMessage m : messages) { + n += estimate(m); + } + return n; + } +} diff --git a/src/main/java/com/xly/web/AgentChatController.java b/src/main/java/com/xly/web/AgentChatController.java index 5450612..7c77840 100644 --- a/src/main/java/com/xly/web/AgentChatController.java +++ b/src/main/java/com/xly/web/AgentChatController.java @@ -95,7 +95,7 @@ public class AgentChatController { final String convId = conversations.scopedId(identity, req.conversationId); conversations.touch(identity.userId(), convId, userInput); - ledger.append(convId, "user", Map.of("text", userInput)); + ledger.append(convId, "user", Map.of("text", userInput), identity); exec.submit(() -> { try { diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index 03a0730..1299283 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -68,6 +68,9 @@ llm: base-url: ${LLM_BASE_URL:http://112.82.245.194:41434/v1} api-key: ${LLM_API_KEY:ollama} # Ollama 不校验(任意非空);云端填真实 key chat-model: ${LLM_CHAT_MODEL:qwen3.6-27b-iq3:latest} + # 镜像 Ollama 侧 OLLAMA_CONTEXT_LENGTH(OpenAI 兼容通道无法逐请求传 num_ctx)。 + # 投影层以它为预算基数;两值漂移由 prompt_eval_count 校准回路兜底(TracingChatModelListener 告警)。 + context-length: ${LLM_CONTEXT_LENGTH:16384} erp: baseurl: ${ERP_BASEURL:http://118.178.19.35:8080/xlyEntry_saas} diff --git a/src/test/java/com/xly/service/ConversationScopeTest.java b/src/test/java/com/xly/service/ConversationScopeTest.java index 4beb00c..d6a81f2 100644 --- a/src/test/java/com/xly/service/ConversationScopeTest.java +++ b/src/test/java/com/xly/service/ConversationScopeTest.java @@ -26,7 +26,7 @@ class ConversationScopeTest { } private static AgentIdentity user(String id) { - return new AgentIdentity("tok", id, "brand", null); + return new AgentIdentity("tok", id, "brand", "sub", null); } @Test diff --git a/src/test/java/com/xly/service/EventLogConcurrencyTest.java b/src/test/java/com/xly/service/EventLogConcurrencyTest.java index 6fe847e..db7c12b 100644 --- a/src/test/java/com/xly/service/EventLogConcurrencyTest.java +++ b/src/test/java/com/xly/service/EventLogConcurrencyTest.java @@ -30,7 +30,7 @@ class EventLogConcurrencyTest { void concurrentStreamAndConfirmLoseNothing() throws Exception { InMemoryLedger log = new InMemoryLedger(); EventLogChatMemory mem = new EventLogChatMemory( - CONV, log, new EventProjectionService(new ObjectMapper()), 6000, false); + CONV, log, new EventProjectionService(new ObjectMapper()), null, false); int writers = 4; int perWriter = 100; diff --git a/src/test/java/com/xly/service/EventProjectionTest.java b/src/test/java/com/xly/service/EventProjectionTest.java index 8132b84..d005fea 100644 --- a/src/test/java/com/xly/service/EventProjectionTest.java +++ b/src/test/java/com/xly/service/EventProjectionTest.java @@ -34,7 +34,7 @@ class EventProjectionTest { } private EventLogChatMemory memory(InMemoryLedger log, boolean internal) { - return new EventLogChatMemory(CONV, log, proj, 6000, internal); + return new EventLogChatMemory(CONV, log, proj, null, internal); } private static AiMessage toolCall(String id, String name, String args) { @@ -92,6 +92,8 @@ class EventProjectionTest { @Test void turnsBeyondBudgetBecomeDeterministicDigest() { + // 小上下文(6000 → promptBudget≈2976t),每轮约 1500t:只装得下 1-2 轮原文,其余进摘要区 + EventProjectionService small = new EventProjectionService(new ObjectMapper(), 6000); InMemoryLedger log = new InMemoryLedger(); for (int i = 1; i <= 10; i++) { log.append(CONV, "user", Map.of("text", "问题" + i + ":" + "长".repeat(1500))); @@ -99,7 +101,7 @@ class EventProjectionTest { } log.append(CONV, "user", Map.of("text", "当前问题")); - List ms = proj.project("SYS", log.events(CONV), 6000); + List ms = small.project("SYS", log.events(CONV)); UserMessage digest = (UserMessage) ms.stream() .filter(m -> m instanceof UserMessage u && u.singleText().startsWith("【更早对话摘要")) .findFirst().orElseThrow(); @@ -107,9 +109,32 @@ class EventProjectionTest { assertTrue(digest.singleText().contains("⇒ 回答1"), "摘要含结果"); long verbatimUsers = ms.stream().filter(m -> m instanceof UserMessage u && !u.singleText().startsWith("【更早对话摘要")).count(); - assertTrue(verbatimUsers <= 5, "近期原文区按 6000 字预算收拢(每轮约 1500 字)"); + assertTrue(verbatimUsers <= 3, "近期原文区按 token 预算收拢(每轮约 1500t)"); assertTrue(ms.get(ms.size() - 1) instanceof UserMessage u && u.singleText().equals("当前问题"), "当前轮永在末尾且完整"); + assertTrue(TokenEstimator.estimate(ms) <= small.promptBudget(), "发送前自检:总量不超预算"); + } + + @Test + void currentTurnDemotesOldestToolResultsWhenOversized() { + // 当前轮 4 个各约 1200t 的工具结果 + 小上下文:最早的降 120 字摘要,最近 2 个保全文,配对不破坏 + EventProjectionService small = new EventProjectionService(new ObjectMapper(), 8000); + InMemoryLedger log = new InMemoryLedger(); + EventLogChatMemory mem = new EventLogChatMemory(CONV, log, small, null, false); + log.append(CONV, "user", Map.of("text", "查四张表")); + for (int i = 1; i <= 4; i++) { + mem.add(toolCall("t" + i, "readFormData", "{\"formId\":\"F" + i + "\"}")); + mem.add(ToolExecutionResultMessage.from("t" + i, "readFormData", "行".repeat(1200))); + } + mem.add(SystemMessage.from("SYS")); + List ms = mem.messages(); + List trs = ms.stream() + .filter(m -> m instanceof ToolExecutionResultMessage) + .map(m -> (ToolExecutionResultMessage) m).toList(); + assertEquals(4, trs.size(), "工具调用/结果配对一个不丢"); + assertTrue(trs.get(0).text().length() <= EventProjectionService.TOOL_DIGEST_LEN + 1, "最早结果降摘要"); + assertEquals(1200, trs.get(3).text().length(), "最近结果保全文"); + assertTrue(TokenEstimator.estimate(ms) <= small.promptBudget(), "轮内让位后不超预算"); } @Test @@ -138,7 +163,7 @@ class EventProjectionTest { log.append(CONV, "ai", Map.of("text", "大约有 1286 个客户。")); memory(log, true).add(UserMessage.from("你上一条回答没有调用任何工具,请先查询真实数据。")); - List ms = proj.project("SYS", log.events(CONV), 6000); + List ms = proj.project("SYS", log.events(CONV)); assertTrue(ms.stream().anyMatch(m -> m instanceof UserMessage u && u.singleText().contains("没有调用任何工具")), "注入话术 LLM 可见"); @@ -160,7 +185,7 @@ class EventProjectionTest { assertTrue(hist.get(1).get("content").startsWith("【待确认】")); assertTrue(hist.get(2).get("content").startsWith("【已执行】")); - List ms = proj.project("SYS", log.events(CONV), 6000); + List ms = proj.project("SYS", log.events(CONV)); assertTrue(ms.stream().anyMatch(m -> m instanceof UserMessage u && u.singleText().contains("提交「报价」"))); assertTrue(ms.stream().anyMatch(m -> m instanceof AiMessage a && a.text().contains("已生成待确认提议"))); assertTrue(ms.stream().anyMatch(m -> m instanceof AiMessage a && a.text().contains("操作执行成功"))); @@ -197,7 +222,7 @@ class EventProjectionTest { "{\"type\":\"form_collect\",\"entity\":\"报价\",\"message\":\"请填写\"}")); String card = proj.activeCard(log.events(CONV)); assertTrue(card != null && card.contains("「报价」"), "弹表单后流程卡钉住"); - assertTrue(proj.project("SYS", log.events(CONV), 6000).get(0).toString().contains("进行中的流程")); + assertTrue(proj.project("SYS", log.events(CONV)).get(0).toString().contains("进行中的流程")); // 中途插入无关查询轮 → 卡不摘 log.append(CONV, "user", Map.of("text", "先查下必胜客电话")); @@ -241,7 +266,7 @@ class EventProjectionTest { assertEquals(5, hist.size(), "tool 摘要事件不进历史,其余照旧"); assertTrue(hist.get(3).get("content").contains("【表单】新建报价")); - List ms = proj.project("SYS", log.events(CONV), 6000); + List ms = proj.project("SYS", log.events(CONV)); assertTrue(ms.stream().anyMatch(m -> m instanceof AiMessage a && a.text().contains("请问改成什么号码"))); assertFalse(ms.stream().anyMatch(m -> m.toString().contains("findForms")), "旧 tool 摘要不进 LLM 正文"); } diff --git a/src/test/java/com/xly/service/InMemoryLedger.java b/src/test/java/com/xly/service/InMemoryLedger.java index fea8a27..3f2802e 100644 --- a/src/test/java/com/xly/service/InMemoryLedger.java +++ b/src/test/java/com/xly/service/InMemoryLedger.java @@ -1,6 +1,7 @@ package com.xly.service; import com.fasterxml.jackson.databind.ObjectMapper; +import com.xly.agent.AgentIdentity; import java.util.ArrayList; import java.util.Collections; @@ -9,7 +10,7 @@ import java.util.List; import java.util.Map; /** - * 测试替身:内存版事件日志。append 与 Redis rightPush 同为「整条原子追加」, + * 测试替身:内存版事件账本。append 与生产(MySQL insert + Redis rightPush)同为「整条原子追加」, * 用于验证投影逻辑与「任意并发写者只追加、绝无整包读改写」的设计性质。 */ class InMemoryLedger extends LedgerService { @@ -18,11 +19,11 @@ class InMemoryLedger extends LedgerService { private final ObjectMapper mapper = new ObjectMapper(); InMemoryLedger() { - super(null, new ObjectMapper()); + super(null, new ObjectMapper(), null); } @Override - public void append(String convId, String type, Map data) { + public void append(String convId, String type, Map data, AgentIdentity who) { Map ev = new LinkedHashMap<>(); ev.put("t", 0L); ev.put("type", type); -- libgit2 0.22.2