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 #project}),预算按 token 记账(本地保守估算 * {@link TokenEstimator},绝对值由 {@code llm.context-length}(镜像 Ollama 的 * OLLAMA_CONTEXT_LENGTH,经 OpenAI 兼容通道无法逐请求传 num_ctx)导出): *

    *
  1. ① 稳定 system prompt + ② 进行中流程卡 —— 永不让位;
  2. *
  3. ④ 往事摘要区:预算外旧轮 → 确定性一行摘要(零模型调用);
  4. *
  5. ⑤ 近期原文区:按**整轮**纳入,旧轮工具结果压 {@value #TOOL_DIGEST_LEN} 字;
  6. *
  7. 当前轮:原样保真、工具配对不可破坏;当前轮自身超限时**轮内让位**——最早的工具结果 * 降为摘要(保最近 {@value #CURRENT_KEEP_FULL} 个全文),单条结果超硬顶截断。
  8. *
* 组装后做**发送前自检**:总 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; /** 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 上下文(token 记账)。systemText 为空时不出 system 消息(如测试)。 */ public List project(String systemText, List> events) { int budget = promptBudget(); List>> turns = groupTurns(events); int current = turns.size() - 1; // ①② 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--) { if (used + oldEst[i] > budget) { break; } 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<>(); if (sys != null) { out.add(sys); } if (!digestLines.isEmpty() || omitted > 0) { StringBuilder sb = new StringBuilder("【更早对话摘要(自动截断,仅供参考)】"); if (omitted > 0) { sb.append("\n(更早 ").append(omitted).append(" 轮已省略)"); } 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; } /** * 当前轮渲染:原样保真;总量超 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) { 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 "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()); } } 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"))); } // ---------------------------------------------------------------- 前端历史投影 /** 前端历史(API 形状不变:{role:user|ai, content})。internal 用户事件(护栏重试话术)不显示。 */ public List> historyView(List> events) { List> out = new ArrayList<>(); for (Map ev : events) { String type = str(ev.get("type")); switch (type) { case "user" -> { if (!Boolean.TRUE.equals(ev.get("internal"))) { out.add(Map.of("role", "user", "content", str(ev.get("text")))); } } case "form_submit" -> out.add(Map.of("role", "user", "content", renderFormSubmit(ev))); 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")))); case "cancel" -> addHist(out, "【已取消】" + str(ev.get("description"))); case "tool_result" -> addHist(out, historyFromToolResult(ev)); default -> { } } } return out; } /** agent 路径的 提问/表单/预览 不再单独落显示事件——从工具结果同源推导。 */ private String historyFromToolResult(Map ev) { String name = str(ev.get("name")); if (!"askUser".equals(name) && !"collectForm".equals(name) && !"previewChange".equals(name) && !"proposeWrite".equals(name)) { return ""; } try { JsonNode r = mapper.readTree(str(ev.get("text"))); if ("askUser".equals(name) && "question".equals(r.path("type").asText(""))) { return r.path("question").asText(""); } 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(""); } } catch (Exception ignore) { } return ""; } // ---------------------------------------------------------------- 流程卡(② 段) /** * 进行中流程卡:激活 skill 全文(最近一次 {@code skill_active},被 {@code skill_done} 或更新的 * skill 覆盖前一直钉住)+ 在办单据状态(表单待提交 / 预览待保存 / 已提交待办的只读进度)。 * {@code queued} 后技能线摘下(AI 侧流程到写入队列为止),待办行保留到下一个流程事件替换。 * 两者皆无 → null。 */ public String activeCard(List> events) { String skillText = null; Map lastFlow = null; for (Map ev : events) { String type = str(ev.get("type")); switch (type) { case "skill_active" -> skillText = str(ev.get("text")); case "skill_done" -> skillText = null; case "form", "form_submit", "proposal" -> lastFlow = ev; case "queued", "confirm", "cancel" -> { // AI 侧流程终结:技能线摘下 lastFlow = ev; skillText = null; } case "tool_result" -> { Map derived = deriveFlowEvent(ev); if (derived != null) { lastFlow = derived; } } default -> { } } } String docLine = null; if (lastFlow != null) { switch (str(lastFlow.get("type"))) { case "form" -> docLine = "已为「" + str(lastFlow.get("entity")) + "」弹出新建表单,等待用户填写提交。"; 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 = 流程已推进/终结 } } if (skillText == null && docLine == null) { return null; } StringBuilder sb = new StringBuilder(); if (skillText != null) { sb.append(skillText); } if (docLine != null) { if (sb.length() > 0) { sb.append("\n"); } sb.append("在办单据:").append(docLine); } return sb.toString(); } /** 工具结果 → 流程事件(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(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 ("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("")); m.put("summary", r.path("summary").asText("")); return m; } } catch (Exception ignore) { } return null; } // ---------------------------------------------------------------- 内部 /** 按用户事件(user/form_submit)分轮;轮前遗留事件并入首轮。 */ private static List>> groupTurns(List> events) { List>> turns = new ArrayList<>(); List> cur = new ArrayList<>(); for (Map ev : events) { String type = str(ev.get("type")); if (("user".equals(type) || "form_submit".equals(type)) && !cur.isEmpty()) { turns.add(cur); cur = new ArrayList<>(); } cur.add(ev); } if (!cur.isEmpty()) { turns.add(cur); } return turns; } /** 一轮 → 一行确定性摘要:用户说了什么 ⇒ 结果是什么。 */ private String digestLine(List> turn) { String userText = ""; String outcome = ""; for (Map ev : turn) { String type = str(ev.get("type")); switch (type) { case "user" -> userText = str(ev.get("text")); case "form_submit" -> userText = renderFormSubmit(ev); 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")); case "cancel" -> outcome = "用户取消:" + str(ev.get("description")); case "tool_result" -> { String h = historyFromToolResult(ev); if (!h.isBlank()) { outcome = h; } } default -> { } } } if (userText.isBlank() && outcome.isBlank()) { return ""; } return clip(userText, 40) + (outcome.isBlank() ? "" : " ⇒ " + clip(outcome, DIGEST_LINE_LEN)); } private ChatMessage fromPayload(String payload) { try { List ms = ChatMessageDeserializer.messagesFromJson(payload); return ms.isEmpty() ? null : ms.get(0); } catch (Exception e) { return null; } } private String renderFormSubmit(Map ev) { 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")))) { return "(系统)用户取消了该操作:" + desc; } if ("executed".equals(str(ev.get("status")))) { return "(系统)用户已确认,操作执行成功:" + desc; } String msg = str(ev.get("msg")); return "(系统)用户确认后执行失败:" + desc + (msg.isBlank() ? "" : ",原因:" + msg); } private static void addAi(List out, String text) { if (text != null && !text.isBlank()) { out.add(AiMessage.from(text)); } } private static void addHist(List> out, String text) { if (text != null && !text.isBlank()) { out.add(Map.of("role", "ai", "content", text)); } } private static String clip(String s, int max) { String t = s == null ? "" : s.replace('\n', ' ').trim(); return t.length() > max ? t.substring(0, max) + "…" : t; } private static String parenOr(String msg) { return msg == null || msg.isBlank() ? "" : ("(" + msg + ")"); } private static String str(Object o) { return o == null ? "" : o.toString(); } }