diff --git a/src/main/java/com/xly/agent/EventLogChatMemory.java b/src/main/java/com/xly/agent/EventLogChatMemory.java new file mode 100644 index 0000000..ad7c2d1 --- /dev/null +++ b/src/main/java/com/xly/agent/EventLogChatMemory.java @@ -0,0 +1,157 @@ +package com.xly.agent; + +import com.xly.service.EventProjectionService; +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; +import dev.langchain4j.memory.ChatMemory; + +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; + +/** + * 事件日志上的对话记忆:{@link #add} 把模型环消息逐条 append 成事件(rightPush 原子—— + * 对话流与确认端点并发写互不覆盖,替代旧 chat:mem 整包读改写),{@link #messages()} 读取 + * {@link EventProjectionService} 的四段式投影。日志即唯一事实源,本类不持有任何会话状态。 + * + *

写者分工:用户事件由控制器在收到请求时先落账(前端立即可见),本类对 UserMessage 做 + * 去重跳过;{@code internalUserTurn=true} 时(反编造护栏重试的注入话术)例外——落一条 + * {@code internal} 标记的用户事件,LLM 可见、前端历史不显示。模型消息(tool_call/tool_result/ai) + * 只由本类落账。 + */ +public class EventLogChatMemory implements ChatMemory { + + private static final int LLM_EVENT_WINDOW = 400; + + private final String convId; + private final LedgerService log; + private final EventProjectionService projection; + private final int charBudget; + private final boolean internalUserTurn; + + /** 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) { + this.convId = convId; + this.log = log; + this.projection = projection; + this.charBudget = charBudget; + this.internalUserTurn = internalUserTurn; + } + + @Override + public Object id() { + return convId; + } + + @Override + public void add(ChatMessage m) { + if (m instanceof SystemMessage sm) { + systemText = sm.text(); + return; + } + if (m instanceof UserMessage um) { + String text = um.hasSingleText() ? um.singleText() : String.valueOf(um.contents()); + if (!internalUserTurn && alreadyLoggedPrefixOf(text)) { + currentTurnUserText = text; + return; + } + Map data = new LinkedHashMap<>(); + data.put("text", text); + if (internalUserTurn) { + data.put("internal", true); + } + log.append(convId, "user", data); + 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 tools = new ArrayList<>(); + for (ToolExecutionRequest r : am.toolExecutionRequests()) { + tools.add(r.name()); + } + data.put("tools", tools); + log.append(convId, "tool_call", data); + } else if (am.text() != null && !am.text().isBlank()) { + log.append(convId, "ai", Map.of("text", am.text())); + } + return; + } + if (m instanceof ToolExecutionResultMessage tr) { + String text = tr.text() == null ? "" : tr.text(); + Map data = new LinkedHashMap<>(); + data.put("tcId", tr.id() == null ? "" : tr.id()); + data.put("name", tr.toolName() == null ? "" : tr.toolName()); + data.put("text", text); + data.put("digest", digest(text)); + log.append(convId, "tool_result", data); + } + } + + /** + * 控制器已把本轮用户**原话**落账(最近的用户事件是喂给模型文本的前缀、且其后无模型事件) + * → 跳过,避免重复。编排层可能在原话后附加 grounding/状态后缀,故用前缀而非全等判定。 + */ + private boolean alreadyLoggedPrefixOf(String text) { + List> tail = log.events(convId, 6); + for (int i = tail.size() - 1; i >= 0; i--) { + String type = String.valueOf(tail.get(i).get("type")); + switch (type) { + case "user": + String logged = String.valueOf(tail.get(i).get("text")); + return !logged.isEmpty() && text.startsWith(logged); + case "form_submit", "ai", "tool_call", "tool_result", "assistant": + return false; + default: + // confirm/proposal 等按钮事件可能插在中间,继续往前找 + } + } + return false; + } + + @Override + public List messages() { + List out = projection.project( + systemText, log.events(convId, LLM_EVENT_WINDOW), charBudget); + String full = currentTurnUserText; + if (full != null) { + // 当前轮用户消息换成实际喂给模型的完整文本(原话+编排后缀) + for (int i = out.size() - 1; i >= 0; i--) { + if (out.get(i) instanceof UserMessage u && u.hasSingleText() + && full.startsWith(u.singleText())) { + out.set(i, UserMessage.from(full)); + break; + } + } + } + return out; + } + + /** 日志是唯一事实源,不因模型环异常清史;删除会话走 ConversationService 级联。 */ + @Override + public void clear() { + } + + private static String digest(String text) { + String t = text.replace('\n', ' ').trim(); + return t.length() > EventProjectionService.TOOL_DIGEST_LEN + ? t.substring(0, EventProjectionService.TOOL_DIGEST_LEN) + "…" : t; + } +} diff --git a/src/main/java/com/xly/agent/ProjectedChatMemory.java b/src/main/java/com/xly/agent/ProjectedChatMemory.java deleted file mode 100644 index cc6c8e6..0000000 --- a/src/main/java/com/xly/agent/ProjectedChatMemory.java +++ /dev/null @@ -1,173 +0,0 @@ -package com.xly.agent; - -import dev.langchain4j.agent.tool.ToolExecutionRequest; -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.memory.ChatMemory; -import dev.langchain4j.store.memory.chat.ChatMemoryStore; - -import java.util.ArrayList; -import java.util.List; - -/** - * 带**投影**的对话记忆:存储层保留完整消息({@link ChatMemoryStore},硬上限按整轮裁剪), - * 读取层({@link #messages()})做 token 预算投影——替换按条数计窗的 MessageWindowChatMemory - * (工具消息占条数导致真实轮次只有 5-8 轮,且长短不均)。 - * - *

投影规则: - *

- */ -public class ProjectedChatMemory implements ChatMemory { - - private static final int TOOL_DIGEST_LEN = 120; - - private final Object id; - private final ChatMemoryStore store; - private final int charBudget; - private final int hardCapMessages; - - public ProjectedChatMemory(Object id, ChatMemoryStore store, int charBudget, int hardCapMessages) { - this.id = id; - this.store = store; - this.charBudget = charBudget; - this.hardCapMessages = hardCapMessages; - } - - @Override - public Object id() { - return id; - } - - @Override - public void add(ChatMessage m) { - List full = new ArrayList<>(store.getMessages(id)); - if (m instanceof SystemMessage sm) { - if (!full.isEmpty() && full.get(0) instanceof SystemMessage cur) { - if (cur.text().equals(sm.text())) { - return; - } - full.set(0, sm); - } else { - full.add(0, sm); - } - } else { - full.add(m); - trimToCap(full); - } - store.updateMessages(id, full); - } - - @Override - public List messages() { - return project(new ArrayList<>(store.getMessages(id))); - } - - @Override - public void clear() { - store.deleteMessages(id); - } - - private List project(List full) { - if (full.isEmpty()) { - return full; - } - SystemMessage sys = full.get(0) instanceof SystemMessage s ? s : null; - List body = full.subList(sys == null ? 0 : 1, full.size()); - - int lastUser = 0; - for (int i = body.size() - 1; i >= 0; i--) { - if (body.get(i) instanceof UserMessage) { - lastUser = i; - break; - } - } - List tail = new ArrayList<>(body.subList(lastUser, body.size())); - - List head = new ArrayList<>(); - int used = 0; - int turnEnd = lastUser; - for (int i = lastUser - 1; i >= 0 && used < charBudget; i--) { - if (!(body.get(i) instanceof UserMessage)) { - continue; - } - List turn = new ArrayList<>(); - int size = 0; - for (int k = i; k < turnEnd; k++) { - ChatMessage c = collapse(body.get(k)); - turn.add(c); - size += approxLen(c); - } - if (used + size > charBudget && !head.isEmpty()) { - break; - } - head.addAll(0, turn); - used += size; - turnEnd = i; - } - - List out = new ArrayList<>(); - if (sys != null) { - out.add(sys); - } - out.addAll(head); - out.addAll(tail); - return out; - } - - /** 历史轮的工具结果压成一行摘要(当前轮不经过此路径,配对结构完整)。 */ - private static ChatMessage collapse(ChatMessage m) { - if (m instanceof ToolExecutionResultMessage t) { - String txt = t.text() == null ? "" : t.text().replace('\n', ' ').trim(); - if (txt.length() > TOOL_DIGEST_LEN) { - txt = txt.substring(0, TOOL_DIGEST_LEN) + "…"; - } - return ToolExecutionResultMessage.from(t.id(), t.toolName(), txt); - } - return m; - } - - private static int approxLen(ChatMessage m) { - if (m instanceof UserMessage u && u.hasSingleText()) { - return u.singleText().length(); - } - if (m instanceof AiMessage a) { - int n = a.text() == null ? 0 : a.text().length(); - if (a.hasToolExecutionRequests()) { - for (ToolExecutionRequest r : a.toolExecutionRequests()) { - n += (r.arguments() == null ? 0 : r.arguments().length()) + 20; - } - } - return n; - } - if (m instanceof ToolExecutionResultMessage t) { - return t.text() == null ? 0 : t.text().length(); - } - return 50; - } - - /** 存储硬上限:超限时从最旧的整轮开始删(system 保留)。 */ - private void trimToCap(List full) { - int start = !full.isEmpty() && full.get(0) instanceof SystemMessage ? 1 : 0; - while (full.size() > hardCapMessages) { - int next = -1; - for (int i = start + 1; i < full.size(); i++) { - if (full.get(i) instanceof UserMessage) { - next = i; - break; - } - } - if (next < 0) { - break; - } - full.subList(start, next).clear(); - } - } -} diff --git a/src/main/java/com/xly/config/AgentFactory.java b/src/main/java/com/xly/config/AgentFactory.java index 1726b6a..85ba392 100644 --- a/src/main/java/com/xly/config/AgentFactory.java +++ b/src/main/java/com/xly/config/AgentFactory.java @@ -2,9 +2,12 @@ package com.xly.config; import com.fasterxml.jackson.databind.ObjectMapper; import com.xly.agent.AgentIdentity; +import com.xly.agent.EventLogChatMemory; import com.xly.agent.ReActAgent; import com.xly.service.ErpClient; +import com.xly.service.EventProjectionService; import com.xly.service.FormResolverService; +import com.xly.service.LedgerService; import com.xly.service.OpService; import com.xly.service.SystemPromptService; import com.xly.tool.ErpReadTool; @@ -12,7 +15,6 @@ import com.xly.tool.FormCollectTool; import com.xly.tool.InteractionTool; import com.xly.tool.KgQueryTool; import com.xly.tool.ProposeWriteTool; -import com.xly.agent.ProjectedChatMemory; import dev.langchain4j.model.chat.StreamingChatModel; import dev.langchain4j.service.AiServices; import org.springframework.beans.factory.annotation.Qualifier; @@ -34,7 +36,8 @@ import org.springframework.stereotype.Component; public class AgentFactory { private final StreamingChatModel streamingModel; - private final RedisChatMemoryStore memoryStore; + private final LedgerService ledger; + private final EventProjectionService projection; private final SystemPromptService systemPromptService; private final ErpClient erp; @@ -47,12 +50,13 @@ public class AgentFactory { private final InteractionTool interactionTool; public AgentFactory(@Qualifier("agentStreamingModel") StreamingChatModel streamingModel, - RedisChatMemoryStore memoryStore, + LedgerService ledger, EventProjectionService projection, SystemPromptService systemPromptService, ErpClient erp, JdbcTemplate jdbc, FormResolverService resolver, OpService ops, ObjectMapper mapper, KgQueryTool kgQueryTool, InteractionTool interactionTool) { this.streamingModel = streamingModel; - this.memoryStore = memoryStore; + this.ledger = ledger; + this.projection = projection; this.systemPromptService = systemPromptService; this.erp = erp; this.jdbc = jdbc; @@ -63,11 +67,18 @@ public class AgentFactory { this.interactionTool = interactionTool; } + public ReActAgent build(AgentIdentity identity) { + return build(identity, false); + } + /** * 组装 ReAct agent:固定 6 工具 + 唯一 system prompt。{@code maxSequentialToolsInvocations} * 作为循环护栏,杜绝「反复追问同一问题」这类失控 ReAct 循环。 + * + * @param internalUserTurn true=本次的用户消息是系统注入话术(护栏重试),落账时带 internal 标记, + * 前端历史不显示 */ - public ReActAgent build(AgentIdentity identity) { + public ReActAgent build(AgentIdentity identity, boolean internalUserTurn) { Object[] tools = new Object[]{kgQueryTool, interactionTool, new ErpReadTool(erp, resolver, identity), new ProposeWriteTool(erp, jdbc, ops, mapper, identity, resolver), @@ -77,8 +88,9 @@ public class AgentFactory { .streamingChatModel(streamingModel) .tools(tools) .maxSequentialToolsInvocations(8) // 循环护栏:防止 askUser/工具无限自我循环 - // token 预算投影记忆:存储保完整、读取按预算收拢,旧工具结果压一行(见 ProjectedChatMemory) - .chatMemoryProvider(memoryId -> new ProjectedChatMemory(memoryId, memoryStore, 6000, 80)) + // 事件日志记忆:逐条 append 原子落账,读取为四段式投影(见 EventLogChatMemory) + .chatMemoryProvider(memoryId -> new EventLogChatMemory( + String.valueOf(memoryId), ledger, projection, 6000, internalUserTurn)) .systemMessageProvider(memoryId -> systemPromptService.prompt()) .build(); } diff --git a/src/main/java/com/xly/config/RedisChatMemoryStore.java b/src/main/java/com/xly/config/RedisChatMemoryStore.java index 50afe22..d798c7a 100644 --- a/src/main/java/com/xly/config/RedisChatMemoryStore.java +++ b/src/main/java/com/xly/config/RedisChatMemoryStore.java @@ -1,29 +1,23 @@ package com.xly.config; -import dev.langchain4j.data.message.AiMessage; import dev.langchain4j.data.message.ChatMessage; import dev.langchain4j.data.message.ChatMessageDeserializer; -import dev.langchain4j.data.message.ChatMessageSerializer; -import dev.langchain4j.data.message.UserMessage; -import dev.langchain4j.store.memory.chat.ChatMemoryStore; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.stereotype.Component; -import java.time.Duration; import java.util.ArrayList; import java.util.List; /** - * Redis 持久化的 ChatMemoryStore —— 会话记忆按 conversationId 存到 Redis,重启/换实例不丢, - * 也让「重登录后接回上一次会话」成立(会话状态按稳定身份+conversationId 存,而非按 token)。 + * 旧版消息记忆({@code chat:mem:{convId}})的**遗留只读**兼容层。 * - *

键:{@code chat:mem:{conversationId}};值:LangChain4j 序列化后的消息 JSON。30 天 TTL。 + *

事实源已迁移到事件日志({@link com.xly.service.LedgerService},{@code chat:ledger:}), + * 本类不再写入:只保留旧会话的历史兼容读取与删除级联,30 天 TTL 到期后自然淘汰,届时可整体删除。 */ @Component -public class RedisChatMemoryStore implements ChatMemoryStore { +public class RedisChatMemoryStore { private static final String PREFIX = "chat:mem:"; - private static final Duration TTL = Duration.ofDays(30); private final StringRedisTemplate redis; @@ -31,7 +25,6 @@ public class RedisChatMemoryStore implements ChatMemoryStore { this.redis = redis; } - @Override public List getMessages(Object memoryId) { String json = redis.opsForValue().get(PREFIX + memoryId); if (json == null || json.isBlank()) { @@ -40,25 +33,15 @@ public class RedisChatMemoryStore implements ChatMemoryStore { return ChatMessageDeserializer.messagesFromJson(json); } - @Override - public void updateMessages(Object memoryId, List messages) { - redis.opsForValue().set(PREFIX + memoryId, ChatMessageSerializer.messagesToJson(messages), TTL); + public boolean exists(Object memoryId) { + try { + return Boolean.TRUE.equals(redis.hasKey(PREFIX + memoryId)); + } catch (Exception e) { + return true; // 读失败时不误判为「可清理」 + } } - @Override public void deleteMessages(Object memoryId) { redis.delete(PREFIX + memoryId); } - - /** 确定性路径(表单/澄清/提议/确认)不经过 LLM 记忆——用它把该轮补进存储,修补记忆空洞。 */ - public void appendTurn(Object memoryId, String userText, String aiText) { - List full = new ArrayList<>(getMessages(memoryId)); - if (userText != null && !userText.isBlank()) { - full.add(UserMessage.from(userText)); - } - if (aiText != null && !aiText.isBlank()) { - full.add(AiMessage.from(aiText)); - } - updateMessages(memoryId, full); - } } diff --git a/src/main/java/com/xly/service/ConversationService.java b/src/main/java/com/xly/service/ConversationService.java index 1472346..36616df 100644 --- a/src/main/java/com/xly/service/ConversationService.java +++ b/src/main/java/com/xly/service/ConversationService.java @@ -23,20 +23,25 @@ import java.util.Map; public class ConversationService { private static final String CONVS_KEY = "chat:convs:"; + /** 懒清理宽限:新建会话可能短暂没有任何事件,updatedAt 在此窗口内的不清。 */ + private static final long CLEANUP_GRACE_MS = 3L * 24 * 3600 * 1000; private final StringRedisTemplate redis; private final RedisChatMemoryStore memoryStore; private final ObjectMapper mapper; private final LedgerService ledger; private final StateService state; + private final EventProjectionService projection; public ConversationService(StringRedisTemplate redis, RedisChatMemoryStore memoryStore, - ObjectMapper mapper, LedgerService ledger, StateService state) { + ObjectMapper mapper, LedgerService ledger, StateService state, + EventProjectionService projection) { this.redis = redis; this.memoryStore = memoryStore; this.mapper = mapper; this.ledger = ledger; this.state = state; + this.projection = projection; } /** @@ -105,14 +110,21 @@ public class ConversationService { } } - /** 该用户的会话列表,按最近更新倒序。 */ + /** 该用户的会话列表,按最近更新倒序。内容键已过期(30 天 TTL)的项懒清理剔除。 */ public List> list(String userId) { Map all = redis.opsForHash().entries(CONVS_KEY + userId); List> out = new ArrayList<>(); + long now = System.currentTimeMillis(); for (Object v : all.values()) { try { @SuppressWarnings("unchecked") Map m = mapper.readValue(v.toString(), Map.class); + String convId = String.valueOf(m.get("id")); + if (now - num(m.get("updatedAt")) > CLEANUP_GRACE_MS + && !ledger.exists(convId) && !memoryStore.exists(convId)) { + redis.opsForHash().delete(CONVS_KEY + userId, convId); + continue; + } out.add(m); } catch (Exception ignore) { } @@ -129,28 +141,13 @@ public class ConversationService { } /** - * 会话历史:优先从**会话账本**重放(含确定性路径的表单/澄清/提议/确认结果——消息记忆里没有这些); - * 无账本的旧会话退回消息记忆。映射为 {role:user|ai, content}。 + * 会话历史 = 事件日志的前端投影(与 LLM 上下文同源,见 {@link EventProjectionService}); + * 无日志的旧会话退回消息记忆兼容读取。映射为 {role:user|ai, content}。 */ public List> history(String convId) { List> events = ledger.events(convId); if (!events.isEmpty()) { - List> out = new ArrayList<>(); - for (Map e : events) { - String type = String.valueOf(e.get("type")); - switch (type) { - case "user" -> out.add(Map.of("role", "user", "content", s(e.get("text")))); - case "assistant", "clarify" -> addAi(out, s(e.get("text"))); - case "question" -> addAi(out, s(e.get("question"))); - case "form" -> addAi(out, "【表单】新建" + s(e.get("entity")) + ":" + s(e.get("message"))); - case "proposal" -> addAi(out, "【待确认】" + s(e.get("summary"))); - case "confirm" -> addAi(out, ("executed".equals(s(e.get("status"))) ? "【已执行】" : "【执行失败】") - + s(e.get("description")) + blankOr(s(e.get("msg")))); - case "cancel" -> addAi(out, "【已取消】" + s(e.get("description"))); - default -> { } // tool 等内部事件不进历史 - } - } - return out; + return projection.historyView(events); } List> out = new ArrayList<>(); for (ChatMessage m : memoryStore.getMessages(convId)) { @@ -164,20 +161,6 @@ public class ConversationService { return out; } - private static void addAi(List> out, String text) { - if (text != null && !text.isBlank()) { - out.add(Map.of("role", "ai", "content", text)); - } - } - - private static String s(Object o) { - return o == null ? "" : o.toString(); - } - - private static String blankOr(String msg) { - return msg == null || msg.isBlank() ? "" : ("(" + msg + ")"); - } - private String deriveTitle(String s) { if (s == null) { return "新会话"; diff --git a/src/main/java/com/xly/service/EventProjectionService.java b/src/main/java/com/xly/service/EventProjectionService.java new file mode 100644 index 0000000..96d789a --- /dev/null +++ b/src/main/java/com/xly/service/EventProjectionService.java @@ -0,0 +1,375 @@ +package com.xly.service; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +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.springframework.stereotype.Service; + +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; + +/** + * 事件日志的**唯一**投影/渲染器 —— 前端历史与 LLM 上下文同源于 {@link LedgerService} 的事件流, + * 文案在这一处生成,不再散落在各控制器里。 + * + *

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

    + *
  1. 稳定 system prompt(KV 前缀稳定);
  2. + *
  3. 进行中流程卡(激活 skill 全文 + 在办单据状态)——钉在 system 尾部,不参与截断, + * 提议 executed/cancelled 或新 skill 激活时摘下;
  4. + *
  5. 往事摘要区:预算外旧轮 → 确定性一行摘要(零模型调用);
  6. + *
  7. 近期原文区:约 charBudget 字符按**整轮**纳入;旧轮工具结果压 {@value #TOOL_DIGEST_LEN} 字; + * 当前轮(最后一个用户事件起)原样载荷精确重建,工具调用/结果配对不可破坏。
  8. + *
+ */ +@Service +public class EventProjectionService { + + public static final int TOOL_DIGEST_LEN = 120; + private static final int DIGEST_TURNS_MAX = 40; + private static final int DIGEST_LINE_LEN = 90; + + private final ObjectMapper mapper; + + public EventProjectionService(ObjectMapper mapper) { + this.mapper = mapper; + } + + // ---------------------------------------------------------------- LLM 投影 + + /** 四段式 LLM 上下文。systemText 为空时不出 system 消息(如测试)。 */ + public List project(String systemText, List> events, int charBudget) { + List>> turns = groupTurns(events); + int current = turns.size() - 1; + + // ④ 近期原文区:从最新往回按整轮纳入;当前轮必进且不计预算 + int firstVerbatim = current; + int used = 0; + for (int i = current - 1; i >= 0; i--) { + int size = 0; + for (Map ev : turns.get(i)) { + size += approxLen(ev); + } + if (used + size > charBudget) { + break; + } + used += size; + firstVerbatim = i; + } + + 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 (firstVerbatim > 0) { + StringBuilder sb = new StringBuilder("【更早对话摘要(自动截断,仅供参考)】"); + int from = Math.max(0, firstVerbatim - DIGEST_TURNS_MAX); + if (from > 0) { + sb.append("\n(更早 ").append(from).append(" 轮已省略)"); + } + for (int i = from; i < firstVerbatim; i++) { + String line = digestLine(turns.get(i)); + if (!line.isBlank()) { + sb.append("\n- ").append(line); + } + } + out.add(UserMessage.from(sb.toString())); + } + + for (int i = firstVerbatim; i < turns.size(); i++) { + appendTurnMessages(out, turns.get(i), i == current); + } + return out; + } + + /** 一轮 → 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); + } + } + 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 "proposal" -> addAi(out, "已生成待确认提议:" + str(ev.get("summary")) + "(等待用户点确认/取消,尚未执行)"); + case "confirm", "cancel" -> addAi(out, renderOutcome(ev)); + default -> { } // tool(旧版摘要)/skill_active/skill_done 不进正文(skill 由流程卡承载) + } + } + } + + // ---------------------------------------------------------------- 前端历史投影 + + /** 前端历史(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 "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) && !"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 ("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 覆盖前一直钉住)+ 在办单据状态(表单待提交 / 提议待确认;executed/cancelled/failed 即摘下)。 + * 两者皆无 → 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; + + 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", "confirm", "cancel" -> lastFlow = ev; + 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; + } + } + } + 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 = 流程已推进/终结,卡摘下 + } + } + + 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 弹表单 / proposeWrite 出提议),供流程卡推导。 */ + private Map deriveFlowEvent(Map ev) { + try { + JsonNode r = mapper.readTree(str(ev.get("text"))); + if ("collectForm".equals(str(ev.get("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()) { + 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 "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 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 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)); + } + } + + 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(); + } +} diff --git a/src/main/java/com/xly/service/LedgerService.java b/src/main/java/com/xly/service/LedgerService.java index 1aaced9..cbff91f 100644 --- a/src/main/java/com/xly/service/LedgerService.java +++ b/src/main/java/com/xly/service/LedgerService.java @@ -13,13 +13,21 @@ import java.util.List; import java.util.Map; /** - * 会话账本(append-only 事件流)—— 会话里发生过的**一切**按序落账,包括不经过 LLM 的确定性路径 - * (表单弹出/澄清/写提议/确认结果),修补「确定性路径不进对话记忆」的记忆空洞;前端历史从账本重放。 + * 会话事件日志(append-only)—— 会话的**唯一事实源**:用户话、模型消息(含工具调用/结果的原样载荷)、 + * 表单/提议/确认等按钮事件,全部按发生顺序落一条 Redis LIST。前端历史与 LLM 上下文都是它的投影 + * ({@link EventProjectionService}),不再有 chat:mem 整包读改写与多处手工同步。 * - *

键:Redis LIST {@code chat:ledger:{convId}},元素为事件 JSON {@code {t,type,...}},30 天 TTL。 - * 事件类型:{@code user}(text)/ {@code assistant}(text)/ {@code clarify}(text)/ - * {@code form}(entity,message)/ {@code question}(question,options)/ {@code tool}(name,digest)/ - * {@code proposal}(opId,summary)/ {@code confirm}(opId,status,msg)/ {@code cancel}(opId,description)。 + *

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

事件类型: + * {@code user}(text, internal?)/ {@code ai}(text)/ + * {@code tool_call}(payload=原样序列化的模型消息, text, 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 skill_active}(name, text)/ {@code skill_done}(name); + * 旧版遗留类型 {@code assistant/clarify/form/question/tool} 仍可读(兼容渲染,自然淘汰)。 */ @Service public class LedgerService { @@ -36,7 +44,7 @@ public class LedgerService { this.mapper = mapper; } - /** 追加一条事件(绝不抛异常——账本失败不能影响对话主流程)。 */ + /** 追加一条事件(绝不抛异常——日志失败不能影响对话主流程)。 */ public void append(String convId, String type, Map data) { if (convId == null || convId.isBlank()) { return; @@ -56,11 +64,21 @@ public class LedgerService { } } - /** 全量事件(按发生顺序),供前端历史重放。 */ + /** 全量事件(按发生顺序),供前端历史投影。 */ public List> events(String convId) { + return read(convId, 0); + } + + /** 最近 lastN 条事件,供 LLM 上下文投影(超长会话不整包搬运)。 */ + public List> events(String convId, int lastN) { + return read(convId, lastN); + } + + private List> read(String convId, int lastN) { List> out = new ArrayList<>(); try { - List raw = redis.opsForList().range(PREFIX + convId, 0, -1); + long start = lastN <= 0 ? 0 : -lastN; + List raw = redis.opsForList().range(PREFIX + convId, start, -1); if (raw == null) { return out; } @@ -78,6 +96,15 @@ public class LedgerService { return out; } + /** 日志键是否存在(供会话列表懒清理:内容已过期的会话项从列表剔除)。 */ + public boolean exists(String convId) { + try { + return Boolean.TRUE.equals(redis.hasKey(PREFIX + convId)); + } catch (Exception e) { + return true; // 读失败时不误删列表项 + } + } + public void delete(String convId) { try { redis.delete(PREFIX + convId); diff --git a/src/main/java/com/xly/web/AgentChatController.java b/src/main/java/com/xly/web/AgentChatController.java index 189a567..e218253 100644 --- a/src/main/java/com/xly/web/AgentChatController.java +++ b/src/main/java/com/xly/web/AgentChatController.java @@ -6,7 +6,6 @@ import com.xly.agent.AgentIdentity; import com.xly.agent.Intent; import com.xly.agent.ReActAgent; import com.xly.config.AgentFactory; -import com.xly.config.RedisChatMemoryStore; import com.xly.service.AuthzService; import com.xly.service.ConversationService; import com.xly.service.FormResolverService; @@ -74,14 +73,13 @@ public class AgentChatController { private final FormResolverService resolver; private final LedgerService ledger; private final StateService state; - private final RedisChatMemoryStore memoryStore; private final ExecutorService exec = Executors.newCachedThreadPool(); public AgentChatController(AgentFactory agentFactory, AuthzService authz, ObjectMapper mapper, ConversationService conversations, OpService ops, IntentService intentService, SlotFillService slotFill, FormResolverService resolver, LedgerService ledger, - StateService state, RedisChatMemoryStore memoryStore) { + StateService state) { this.agentFactory = agentFactory; this.authz = authz; this.mapper = mapper; @@ -92,7 +90,6 @@ public class AgentChatController { this.resolver = resolver; this.ledger = ledger; this.state = state; - this.memoryStore = memoryStore; } public static class ChatReq { @@ -168,7 +165,7 @@ public class AgentChatController { }); String userText = "提交「" + entity + "」新增表单:" + parts; conversations.touch(identity.userId(), convId, userText); - ledger.append(convId, "user", Map.of("text", userText)); + ledger.append(convId, "form_submit", Map.of("entity", entity, "fields", parts.toString())); String fieldsJson; try { @@ -187,14 +184,12 @@ public class AgentChatController { String summary = r.path("summary").asText(""); ledger.append(convId, "proposal", Map.of("opId", opId, "summary", summary)); state.setActiveDoc(convId, entity, "", opId, "proposed"); - appendMemoryTurn(convId, userText, "已生成待确认提议:" + summary + "(等待用户点确认/取消)"); out.put("opId", opId); out.put("summary", summary); out.put("message", r.path("message").asText("已生成待确认操作,请点【确认】。")); } else { String err = r.path("error").asText("无法完成该新增。"); ledger.append(convId, "assistant", Map.of("text", err)); - appendMemoryTurn(convId, userText, err); out.put("error", err); } } catch (Exception e) { @@ -270,7 +265,6 @@ public class AgentChatController { send(emitter, "token", hint); ledger.append(convId, "form", Map.of("entity", entity, "message", hint)); state.setActiveDoc(convId, entity, "", "", "collecting"); - appendMemoryTurn(convId, userInput, "已为「" + entity + "」弹出新建表单,等待用户填写提交。"); send(emitter, "done", ""); emitter.complete(); return true; @@ -280,7 +274,6 @@ public class AgentChatController { if (!err.isBlank()) { send(emitter, "token", err); ledger.append(convId, "assistant", Map.of("text", err)); - appendMemoryTurn(convId, userInput, err); send(emitter, "done", ""); emitter.complete(); return true; @@ -291,15 +284,6 @@ public class AgentChatController { return false; } - /** 确定性路径不经过 LLM 记忆——把这轮 用户话+系统答复 补进对话记忆,修补记忆空洞。 */ - private void appendMemoryTurn(String convId, String userText, String aiText) { - try { - memoryStore.appendTurn(convId, userText, aiText); - } catch (Exception e) { - log.warn("append memory turn failed (conv={}): {}", convId, e.getMessage()); - } - } - /** 运行 ReAct agent,把流式回调转成 SSE。queryGuard=true 时启用查询反编造护栏。 */ private void runAgent(SseEmitter emitter, String convId, AgentIdentity identity, String text, boolean queryGuard) { @@ -314,7 +298,8 @@ public class AgentChatController { private void runAgentAttempt(SseEmitter emitter, String convId, AgentIdentity identity, String text, boolean queryGuard, boolean allowRetry) { try { - ReActAgent agent = agentFactory.build(identity); + // 重试轮的注入话术以 internal 用户事件落账:LLM 可见、前端历史不显示 + ReActAgent agent = agentFactory.build(identity, !allowRetry); AtomicInteger toolCalls = new AtomicInteger(); AtomicBoolean proposed = new AtomicBoolean(false); TokenStream ts = agent.chat(convId, text); @@ -340,9 +325,7 @@ public class AgentChatController { if (!proposed.get() && WRITE_CLAIM.matcher(answer).find()) { send(emitter, "token", "\n\n⚠️ 系统提示:本条回复没有真正生成待确认操作(没有出现提议卡片即未生效),请重新描述一次您要做的操作。"); } - if (!answer.isBlank()) { - ledger.append(convId, "assistant", Map.of("text", answer)); - } + // 最终答复由 EventLogChatMemory 落账(ai 事件),这里不再重复写 send(emitter, "done", ""); emitter.complete(); }) @@ -444,12 +427,10 @@ public class AgentChatController { send(emitter, "token", r.path("message").asText("已生成待确认操作,请点【确认】。")); ledger.append(convId, "proposal", Map.of("opId", opId, "summary", summary)); state.setActiveDoc(convId, ent, record, opId, "proposed"); - appendMemoryTurn(convId, userInput, "已生成待确认提议:" + summary + "(等待用户点确认/取消)"); } else { String err = r.path("error").asText("无法完成该操作。"); send(emitter, "token", err); ledger.append(convId, "assistant", Map.of("text", err)); - appendMemoryTurn(convId, userInput, err); } } catch (Exception e) { send(emitter, "error", "服务异常:" + e.getMessage()); @@ -497,7 +478,6 @@ public class AgentChatController { + "。请一起告诉我,我再为你生成待确认的操作。"; send(emitter, "token", text); ledger.append(convId, "clarify", Map.of("text", text)); - appendMemoryTurn(convId, userInput, text); send(emitter, "done", ""); emitter.complete(); } @@ -511,7 +491,10 @@ public class AgentChatController { return b; } - /** 工具执行回调:清掉工具前旁白(reset);按工具类型推对应卡片/控件事件;同步落账本/状态槽。 */ + /** + * 工具执行回调:清掉工具前旁白(reset);按工具类型推对应卡片/控件事件。 + * 落账由 EventLogChatMemory 统一完成(tool_call/tool_result 事件),这里只做 SSE 与 op 绑定。 + */ private void handleToolExecuted(SseEmitter emitter, String convId, ToolExecution te, AtomicBoolean proposed) { send(emitter, "reset", ""); try { @@ -531,7 +514,6 @@ public class AgentChatController { card.put("summary", summary); sendEvent(emitter, card); proposed.set(true); - ledger.append(convId, "proposal", Map.of("opId", opId, "summary", summary)); state.setActiveDoc(convId, "", summary, opId, "proposed"); } } else if ("askUser".equals(toolName) || "collectForm".equals(toolName)) { @@ -539,22 +521,10 @@ public class AgentChatController { String type = r.path("type").asText(""); if ("question".equals(type)) { sendEvent(emitter, mapper.convertValue(r, Map.class)); - ledger.append(convId, "question", Map.of( - "question", r.path("question").asText(""), - "options", mapper.convertValue(r.path("options"), List.class))); } else if ("form_collect".equals(type)) { sendEvent(emitter, mapper.convertValue(r, Map.class)); - String entity = r.path("entity").asText(""); - ledger.append(convId, "form", Map.of("entity", entity, - "message", r.path("message").asText(""))); - state.setActiveDoc(convId, entity, "", "", "collecting"); - } - } else { - String digest = te.result().replace('\n', ' ').trim(); - if (digest.length() > 100) { - digest = digest.substring(0, 100) + "…"; + state.setActiveDoc(convId, r.path("entity").asText(""), "", "", "collecting"); } - ledger.append(convId, "tool", Map.of("name", toolName, "digest", digest)); } } catch (Exception e) { log.warn("handle tool result failed ({})", te.request() == null ? "?" : te.request().name(), e); diff --git a/src/main/java/com/xly/web/OpController.java b/src/main/java/com/xly/web/OpController.java index d3c3ad8..e0232a0 100644 --- a/src/main/java/com/xly/web/OpController.java +++ b/src/main/java/com/xly/web/OpController.java @@ -3,7 +3,6 @@ package com.xly.web; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import com.xly.agent.AgentIdentity; -import com.xly.config.RedisChatMemoryStore; import com.xly.service.AuditService; import com.xly.service.AuthzService; import com.xly.service.ConversationService; @@ -53,7 +52,6 @@ public class OpController { private final ObjectMapper mapper; private final LedgerService ledger; private final StateService state; - private final RedisChatMemoryStore memoryStore; private final AuthzService authz; private final ConversationService conversations; private final FormResolverService resolver; @@ -63,7 +61,7 @@ public class OpController { private boolean execStagingEnabled; public OpController(OpService ops, ErpClient erp, AuditService audit, ObjectMapper mapper, - LedgerService ledger, StateService state, RedisChatMemoryStore memoryStore, + LedgerService ledger, StateService state, AuthzService authz, ConversationService conversations, FormResolverService resolver) { this.ops = ops; @@ -72,13 +70,12 @@ public class OpController { this.mapper = mapper; this.ledger = ledger; this.state = state; - this.memoryStore = memoryStore; this.authz = authz; this.conversations = conversations; this.resolver = resolver; } - /** 确认/取消的结果落账本+状态槽+对话记忆(记忆空洞修补:LLM 下轮就知道这单已执行/失败/取消)。 */ + /** 确认/取消的结果落事件日志(rightPush 原子,流式输出中途点确认也不覆盖对话事件)——LLM 下轮经投影即知这单已执行/失败/取消。 */ private void recordOutcome(String conv, String opId, String status, String msg, String description) { if (conv == null || conv.isBlank() || "null".equals(conv)) { return; @@ -90,15 +87,6 @@ public class OpController { "msg", msg == null ? "" : msg, "description", description == null ? "" : description)); state.updateDocStage(conv, opId, status); - String note; - if ("executed".equals(status)) { - note = "(系统)用户已确认,操作执行成功:" + description; - } else if ("cancelled".equals(status)) { - note = "(系统)用户取消了该操作:" + description; - } else { - note = "(系统)用户确认后执行失败:" + description + (msg == null || msg.isBlank() ? "" : (",原因:" + msg)); - } - memoryStore.appendTurn(conv, null, note); } catch (Exception e) { log.warn("record op outcome failed (conv={}, op={}): {}", conv, opId, e.getMessage()); } diff --git a/src/test/java/com/xly/service/ConversationScopeTest.java b/src/test/java/com/xly/service/ConversationScopeTest.java index 1be863d..644c0e8 100644 --- a/src/test/java/com/xly/service/ConversationScopeTest.java +++ b/src/test/java/com/xly/service/ConversationScopeTest.java @@ -22,7 +22,8 @@ class ConversationScopeTest { private ConversationService service(StringRedisTemplate redis) { return new ConversationService(redis, mock(RedisChatMemoryStore.class), new ObjectMapper(), - mock(LedgerService.class), mock(StateService.class)); + mock(LedgerService.class), mock(StateService.class), + new EventProjectionService(new ObjectMapper())); } private static AgentIdentity user(String id) { diff --git a/src/test/java/com/xly/service/EventLogConcurrencyTest.java b/src/test/java/com/xly/service/EventLogConcurrencyTest.java new file mode 100644 index 0000000..6fe847e --- /dev/null +++ b/src/test/java/com/xly/service/EventLogConcurrencyTest.java @@ -0,0 +1,93 @@ +package com.xly.service; + +import com.xly.agent.EventLogChatMemory; +import com.fasterxml.jackson.databind.ObjectMapper; +import dev.langchain4j.data.message.AiMessage; +import dev.langchain4j.data.message.ToolExecutionResultMessage; +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * 「流式输出中途点【确认】」的并发不变量:对话流(EventLogChatMemory 落模型消息)与确认端点 + * (recordOutcome 落 confirm 事件)同时写同一会话时,**所有事件都在、各写者内部有序**。 + * 旧 chat:mem 整包读改写在此场景必丢更新;事件日志逐条追加从机制上消除该竞态。 + */ +class EventLogConcurrencyTest { + + private static final String CONV = "u1:c-conc"; + + @Test + void concurrentStreamAndConfirmLoseNothing() throws Exception { + InMemoryLedger log = new InMemoryLedger(); + EventLogChatMemory mem = new EventLogChatMemory( + CONV, log, new EventProjectionService(new ObjectMapper()), 6000, false); + + int writers = 4; + int perWriter = 100; + ExecutorService pool = Executors.newFixedThreadPool(writers); + CountDownLatch start = new CountDownLatch(1); + CountDownLatch done = new CountDownLatch(writers); + + for (int w = 0; w < writers; w++) { + final int id = w; + pool.submit(() -> { + try { + start.await(); + for (int i = 0; i < perWriter; i++) { + if (id % 2 == 0) { // 对话流:模型消息经记忆适配器落账 + if (i % 2 == 0) { + mem.add(AiMessage.from("w" + id + "-answer-" + i)); + } else { + mem.add(ToolExecutionResultMessage.from( + "t" + id + "-" + i, "readFormData", "w" + id + "-result-" + i)); + } + } else { // 确认端点:按钮事件直接落账 + log.append(CONV, "confirm", Map.of( + "opId", "w" + id + "-" + i, "status", "executed", + "msg", "", "description", "op w" + id + "-" + i)); + } + } + } catch (InterruptedException ignored) { + } finally { + done.countDown(); + } + }); + } + start.countDown(); + assertTrue(done.await(30, TimeUnit.SECONDS)); + pool.shutdown(); + + List> events = log.events(CONV); + assertEquals(writers * perWriter, events.size(), "并发写者的事件一条不丢"); + + for (int w = 0; w < writers; w++) { + List seq = new ArrayList<>(); + for (Map ev : events) { + String marker = String.valueOf( + ev.containsKey("opId") ? ev.get("opId") + : ev.containsKey("tcId") ? ev.get("tcId") : ev.get("text")); + String prefix = marker.startsWith("t" + w + "-") ? "t" + w + "-" : "w" + w + "-"; + if (marker.startsWith(prefix)) { + String tail = marker.substring(marker.lastIndexOf('-') + 1); + try { + seq.add(Integer.parseInt(tail)); + } catch (NumberFormatException ignore) { + } + } + } + for (int i = 1; i < seq.size(); i++) { + assertTrue(seq.get(i) > seq.get(i - 1), "写者 " + w + " 自身顺序保持"); + } + } + } +} diff --git a/src/test/java/com/xly/service/EventProjectionTest.java b/src/test/java/com/xly/service/EventProjectionTest.java new file mode 100644 index 0000000..8132b84 --- /dev/null +++ b/src/test/java/com/xly/service/EventProjectionTest.java @@ -0,0 +1,248 @@ +package com.xly.service; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.xly.agent.EventLogChatMemory; +import dev.langchain4j.data.message.AiMessage; +import dev.langchain4j.data.message.ChatMessage; +import dev.langchain4j.agent.tool.ToolExecutionRequest; +import dev.langchain4j.data.message.SystemMessage; +import dev.langchain4j.data.message.ToolExecutionResultMessage; +import dev.langchain4j.data.message.UserMessage; +import org.junit.jupiter.api.Test; + +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * 事件日志 → LLM 四段投影 / 前端历史 的正确性:当前轮原样载荷保真(工具配对不可破坏)、 + * 旧轮工具结果压摘要、预算按整轮收拢、超预算旧轮进确定性摘要区、internal 用户事件 + * LLM 可见但历史不显示、进行中流程卡的钉住与摘下。 + */ +class EventProjectionTest { + + private static final String CONV = "u1:c-test"; + + private final EventProjectionService proj = new EventProjectionService(new ObjectMapper()); + + private EventLogChatMemory memory(InMemoryLedger log) { + return memory(log, false); + } + + private EventLogChatMemory memory(InMemoryLedger log, boolean internal) { + return new EventLogChatMemory(CONV, log, proj, 6000, internal); + } + + private static AiMessage toolCall(String id, String name, String args) { + return AiMessage.from(ToolExecutionRequest.builder().id(id).name(name).arguments(args).build()); + } + + @Test + void currentTurnIsReconstructedExactly() { + InMemoryLedger log = new InMemoryLedger(); + EventLogChatMemory mem = memory(log); + + log.append(CONV, "user", Map.of("text", "查一下必胜客")); + int before = log.size(); + mem.add(UserMessage.from("查一下必胜客")); // 控制器已落账 → 去重跳过 + assertEquals(before, log.size(), "同文用户消息不能重复落账"); + + String longResult = "必胜客:地址上海市静安区南京西路1266号,".repeat(10); // >120 字 + mem.add(toolCall("t1", "lookupRecord", "{\"entityKeyword\":\"客户\",\"recordKeyword\":\"必胜客\"}")); + mem.add(ToolExecutionResultMessage.from("t1", "lookupRecord", longResult)); + mem.add(SystemMessage.from("SYS")); + + List ms = mem.messages(); + assertTrue(ms.get(0) instanceof SystemMessage s && s.text().startsWith("SYS")); + AiMessage call = (AiMessage) ms.stream().filter(m -> m instanceof AiMessage a && a.hasToolExecutionRequests()) + .findFirst().orElseThrow(); + assertEquals("t1", call.toolExecutionRequests().get(0).id()); + assertTrue(call.toolExecutionRequests().get(0).arguments().contains("必胜客"), "工具参数原样保真"); + ToolExecutionResultMessage tr = (ToolExecutionResultMessage) ms.stream() + .filter(m -> m instanceof ToolExecutionResultMessage).findFirst().orElseThrow(); + assertEquals(longResult, tr.text(), "当前轮工具结果不截断"); + assertEquals("t1", tr.id(), "工具调用/结果配对保持"); + } + + @Test + void previousTurnToolResultsAreCollapsed() { + InMemoryLedger log = new InMemoryLedger(); + EventLogChatMemory mem = memory(log); + + log.append(CONV, "user", Map.of("text", "查一下必胜客")); + String longResult = "共 6 条记录:".repeat(40); + mem.add(toolCall("t1", "readFormData", "{}")); + mem.add(ToolExecutionResultMessage.from("t1", "readFormData", longResult)); + mem.add(AiMessage.from("共 6 个客户。")); + log.append(CONV, "user", Map.of("text", "下一页")); + + mem.add(SystemMessage.from("SYS")); + List ms = mem.messages(); + ToolExecutionResultMessage tr = (ToolExecutionResultMessage) ms.stream() + .filter(m -> m instanceof ToolExecutionResultMessage).findFirst().orElseThrow(); + assertTrue(tr.text().length() <= EventProjectionService.TOOL_DIGEST_LEN + 1, "旧轮工具结果压一行摘要"); + assertTrue(ms.stream().anyMatch(m -> m instanceof AiMessage a && "共 6 个客户。".equals(a.text())), + "旧轮最终答复保留"); + assertTrue(ms.get(ms.size() - 1) instanceof UserMessage u && u.singleText().equals("下一页")); + } + + @Test + void turnsBeyondBudgetBecomeDeterministicDigest() { + InMemoryLedger log = new InMemoryLedger(); + for (int i = 1; i <= 10; i++) { + log.append(CONV, "user", Map.of("text", "问题" + i + ":" + "长".repeat(1500))); + log.append(CONV, "ai", Map.of("text", "回答" + i)); + } + log.append(CONV, "user", Map.of("text", "当前问题")); + + List ms = proj.project("SYS", log.events(CONV), 6000); + UserMessage digest = (UserMessage) ms.stream() + .filter(m -> m instanceof UserMessage u && u.singleText().startsWith("【更早对话摘要")) + .findFirst().orElseThrow(); + assertTrue(digest.singleText().contains("问题1:"), "最早轮进摘要区"); + 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(ms.get(ms.size() - 1) instanceof UserMessage u && u.singleText().equals("当前问题"), + "当前轮永在末尾且完整"); + } + + @Test + void groundedSuffixDoesNotDuplicateUserEventButFeedsLlm() { + InMemoryLedger log = new InMemoryLedger(); + EventLogChatMemory mem = memory(log); + + log.append(CONV, "user", Map.of("text", "有哪些客户是上海的?")); + int before = log.size(); + String grounded = "有哪些客户是上海的?\n\n(意图分析,仅供参考:意图=查询)\n\n(会话状态,仅供参考:上轮意图=查询)"; + mem.add(UserMessage.from(grounded)); + assertEquals(before, log.size(), "加注后缀的同轮用户消息不能落成第二条事件"); + + mem.add(SystemMessage.from("SYS")); + List ms = mem.messages(); + assertTrue(ms.get(ms.size() - 1) instanceof UserMessage u && u.singleText().equals(grounded), + "LLM 看到的是完整加注文本"); + List> hist = proj.historyView(log.events(CONV)); + assertEquals("有哪些客户是上海的?", hist.get(0).get("content"), "前端历史只显示原话"); + } + + @Test + void internalUserTurnVisibleToLlmButHiddenFromHistory() { + InMemoryLedger log = new InMemoryLedger(); + log.append(CONV, "user", Map.of("text", "有多少个客户?")); + log.append(CONV, "ai", Map.of("text", "大约有 1286 个客户。")); + memory(log, true).add(UserMessage.from("你上一条回答没有调用任何工具,请先查询真实数据。")); + + List ms = proj.project("SYS", log.events(CONV), 6000); + assertTrue(ms.stream().anyMatch(m -> m instanceof UserMessage u + && u.singleText().contains("没有调用任何工具")), "注入话术 LLM 可见"); + + List> hist = proj.historyView(log.events(CONV)); + assertEquals(2, hist.size(), "internal 事件不进前端历史"); + assertEquals("有多少个客户?", hist.get(0).get("content")); + } + + @Test + void buttonPathEventsRenderInBothProjections() { + InMemoryLedger log = new InMemoryLedger(); + log.append(CONV, "form_submit", Map.of("entity", "报价", "fields", "客户=必胜客,产品=纸盒,数量=5000")); + log.append(CONV, "proposal", Map.of("opId", "OP1", "summary", "新增报价:必胜客/纸盒 5000 个")); + log.append(CONV, "confirm", Map.of("opId", "OP1", "status", "executed", "msg", "", "description", "新增报价单")); + + List> hist = proj.historyView(log.events(CONV)); + assertEquals("user", hist.get(0).get("role")); + assertTrue(hist.get(0).get("content").contains("提交「报价」新增表单")); + assertTrue(hist.get(1).get("content").startsWith("【待确认】")); + assertTrue(hist.get(2).get("content").startsWith("【已执行】")); + + List ms = proj.project("SYS", log.events(CONV), 6000); + 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("操作执行成功"))); + } + + @Test + void agentPathQuestionAndProposalDeriveFromToolResults() { + InMemoryLedger log = new InMemoryLedger(); + EventLogChatMemory mem = memory(log); + log.append(CONV, "user", Map.of("text", "作废那张送货单")); + mem.add(toolCall("t1", "askUser", "{}")); + mem.add(ToolExecutionResultMessage.from("t1", "askUser", + "{\"type\":\"question\",\"question\":\"请问是哪一张送货单?\",\"options\":[\"DH1\",\"DH2\"]}")); + log.append(CONV, "user", Map.of("text", "第一张")); + mem.add(toolCall("t2", "proposeWrite", "{\"action\":\"invalid\"}")); + mem.add(ToolExecutionResultMessage.from("t2", "proposeWrite", + "{\"opId\":\"OP2\",\"summary\":\"作废 送货单 DH1\",\"message\":\"请点确认\"}")); + + List> hist = proj.historyView(log.events(CONV)); + assertTrue(hist.stream().anyMatch(h -> h.get("content").equals("请问是哪一张送货单?")), + "askUser 的问题从工具结果同源推导进历史"); + assertTrue(hist.stream().anyMatch(h -> h.get("content").equals("【待确认】作废 送货单 DH1"))); + } + + @Test + void activeFlowCardPinsAndUnpins() { + InMemoryLedger log = new InMemoryLedger(); + EventLogChatMemory mem = memory(log); + + // 表单弹出(agent 路径)→ 卡钉住 + log.append(CONV, "user", Map.of("text", "报价纸盒")); + mem.add(toolCall("t1", "collectForm", "{\"entityKeyword\":\"报价\"}")); + mem.add(ToolExecutionResultMessage.from("t1", "collectForm", + "{\"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("进行中的流程")); + + // 中途插入无关查询轮 → 卡不摘 + log.append(CONV, "user", Map.of("text", "先查下必胜客电话")); + log.append(CONV, "ai", Map.of("text", "13912345678")); + assertTrue(proj.activeCard(log.events(CONV)).contains("「报价」"), "隔轮流程卡仍在"); + + // 表单提交 + 提议 → 卡变为待确认 + log.append(CONV, "form_submit", Map.of("entity", "报价", "fields", "客户=必胜客")); + log.append(CONV, "proposal", Map.of("opId", "OP1", "summary", "新增报价")); + card = proj.activeCard(log.events(CONV)); + assertTrue(card.contains("待确认提议") && card.contains("OP1")); + + // 执行完成 → 卡摘下 + log.append(CONV, "confirm", Map.of("opId", "OP1", "status", "executed", "msg", "", "description", "新增报价")); + assertNull(proj.activeCard(log.events(CONV)), "executed 后流程卡摘下"); + } + + @Test + void skillCardPinsFullTextUntilDone() { + InMemoryLedger log = new InMemoryLedger(); + log.append(CONV, "user", Map.of("text", "报价纸盒")); + log.append(CONV, "skill_active", Map.of("name", "新建报价", "text", "【skill:新建报价】步骤:1…2…3…")); + String card = proj.activeCard(log.events(CONV)); + assertTrue(card.contains("【skill:新建报价】步骤"), "激活 skill 全文钉在流程卡"); + + log.append(CONV, "skill_done", Map.of("name", "新建报价")); + assertNull(proj.activeCard(log.events(CONV)), "skill_done 后摘下"); + } + + @Test + void legacyEventsStillRender() { + InMemoryLedger log = new InMemoryLedger(); + log.append(CONV, "user", Map.of("text", "改必胜客电话")); + log.append(CONV, "clarify", Map.of("text", "请问改成什么号码?")); + log.append(CONV, "question", Map.of("question", "哪个客户?")); + log.append(CONV, "form", Map.of("entity", "报价", "message", "请填写")); + log.append(CONV, "assistant", Map.of("text", "好的。")); + log.append(CONV, "tool", Map.of("name", "findForms", "digest", "…")); + + List> hist = proj.historyView(log.events(CONV)); + assertEquals(5, hist.size(), "tool 摘要事件不进历史,其余照旧"); + assertTrue(hist.get(3).get("content").contains("【表单】新建报价")); + + List ms = proj.project("SYS", log.events(CONV), 6000); + 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 new file mode 100644 index 0000000..fea8a27 --- /dev/null +++ b/src/test/java/com/xly/service/InMemoryLedger.java @@ -0,0 +1,76 @@ +package com.xly.service; + +import com.fasterxml.jackson.databind.ObjectMapper; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; + +/** + * 测试替身:内存版事件日志。append 与 Redis rightPush 同为「整条原子追加」, + * 用于验证投影逻辑与「任意并发写者只追加、绝无整包读改写」的设计性质。 + */ +class InMemoryLedger extends LedgerService { + + private final List raw = Collections.synchronizedList(new ArrayList<>()); + private final ObjectMapper mapper = new ObjectMapper(); + + InMemoryLedger() { + super(null, new ObjectMapper()); + } + + @Override + public void append(String convId, String type, Map data) { + Map ev = new LinkedHashMap<>(); + ev.put("t", 0L); + ev.put("type", type); + if (data != null) { + ev.putAll(data); + } + try { + raw.add(mapper.writeValueAsString(ev)); + } catch (Exception e) { + throw new RuntimeException(e); + } + } + + @Override + public List> events(String convId) { + return events(convId, 0); + } + + @Override + public List> events(String convId, int lastN) { + List snapshot; + synchronized (raw) { + snapshot = new ArrayList<>(raw); + } + int from = lastN <= 0 ? 0 : Math.max(0, snapshot.size() - lastN); + List> out = new ArrayList<>(); + for (int i = from; i < snapshot.size(); i++) { + try { + @SuppressWarnings("unchecked") + Map m = mapper.readValue(snapshot.get(i), Map.class); + out.add(m); + } catch (Exception ignore) { + } + } + return out; + } + + @Override + public boolean exists(String convId) { + return !raw.isEmpty(); + } + + @Override + public void delete(String convId) { + raw.clear(); + } + + int size() { + return raw.size(); + } +}