LedgerService.java 10.7 KB
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)—— 会话的**唯一事实源**:用户话、模型消息(含工具调用/结果)、
 * 表单/预览/保存等按钮事件,全部按发生顺序落账。前端历史与 LLM 上下文都是它的投影
 * ({@link EventProjectionService})。
 *
 * <p><b>存储分层</b>:MySQL {@code ai_chat_event} 为权威持久层(每事件一行、永不删除、可溯源,
 * 服务端赋 iSeq 会话内序号 + iTurn 轮次),Redis LIST {@code chat:ledger:{convId}}(30 天 TTL)
 * 为热缓存加速投影读;缓存失效时回源 MySQL 并回填。两边任一失败都不阻断对话(互为兜底)。
 *
 * <p>事件类型:
 * {@code user}(text, internal?)/ {@code ai}(text)/
 * {@code tool_call}(text, calls=[{id,name,args}], tools)/
 * {@code tool_result}(tcId, name, text, digest)/
 * {@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/proposal/confirm/cancel} 仍可读(兼容渲染,自然淘汰)。
 */
@Service
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<String> 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, JdbcTemplate jdbc) {
        this.redis = redis;
        this.mapper = mapper;
        this.jdbc = jdbc;
    }

    /** 追加一条事件(无身份上下文的旧签名:sMakePerson 从 convId 前缀推导)。 */
    public void append(String convId, String type, Map<String, Object> data) {
        append(convId, type, data, null);
    }

    /** 追加一条事件(绝不抛异常——账本失败不能影响对话主流程)。 */
    public void append(String convId, String type, Map<String, Object> data, AgentIdentity who) {
        if (convId == null || convId.isBlank()) {
            return;
        }
        String json;
        try {
            Map<String, Object> ev = new LinkedHashMap<>();
            ev.put("t", System.currentTimeMillis());
            ev.put("type", type);
            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;
        }
        boolean persisted = persist(convId, type, json, who);
        cache(convId, json, persisted);
    }

    /** MySQL 权威落库:服务端赋 iSeq/iTurn,唯一索引 + 重试防并发撞号。失败仅告警(Redis 兜底),返回是否成功。 */
    private boolean persist(String convId, String type, String json, AgentIdentity who) {
        if (jdbc == null) {
            return false;
        }
        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<Map<String, Object>> 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 true;
            } catch (org.springframework.dao.DuplicateKeyException dup) {
                // 并发撞号:重读末行再试
            } catch (Exception e) {
                log.warn("ledger mysql append failed (conv={}, type={}): {}", convId, type, e.getMessage());
                return false;
            }
        }
        log.warn("ledger mysql append gave up after retries (conv={}, type={})", convId, type);
        return false;
    }

    /**
     * 热缓存追加。**缓存键已失效时不 push**(persisted=true 时)——单条新事件重建的会是
     * 残缺列表且此后永不回源;留给下次读取从 MySQL 冷读回填完整历史(含本条)。
     * MySQL 落库失败时(persisted=false)无论如何都 push,宁可残缺不可丢事件。
     */
    private void cache(String convId, String json, boolean persisted) {
        try {
            String key = PREFIX + convId;
            if (persisted && !Boolean.TRUE.equals(redis.hasKey(key))) {
                return;
            }
            redis.opsForList().rightPush(key, json);
            redis.expire(key, TTL);
        } catch (Exception e) {
            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<Map<String, Object>> events(String convId) {
        return read(convId, 0);
    }

    /** 最近 lastN 条事件,供 LLM 上下文投影(超长会话不整包搬运)。 */
    public List<Map<String, Object>> events(String convId, int lastN) {
        return read(convId, lastN);
    }

    private List<Map<String, Object>> read(String convId, int lastN) {
        List<String> raw = null;
        try {
            long start = lastN <= 0 ? 0 : -lastN;
            if (Boolean.TRUE.equals(redis.hasKey(PREFIX + convId))) {
                raw = redis.opsForList().range(PREFIX + convId, start, -1);
            }
        } catch (Exception e) {
            log.warn("ledger redis read failed (conv={}): {}", convId, e.getMessage());
        }
        if (raw == null) {
            raw = coldRead(convId, lastN);
        }
        List<Map<String, Object>> out = new ArrayList<>();
        for (String s : raw) {
            try {
                @SuppressWarnings("unchecked")
                Map<String, Object> m = mapper.readValue(s, Map.class);
                out.add(m);
            } catch (Exception ignore) {
            }
        }
        return out;
    }

    /**
     * Redis 缓存失效 → 回源 MySQL 并回填缓存。**无论调用方要多小的窗口,回填一律用全量窗口**——
     * 若按调用方窗口(如去重探针的 6 条)回填,key 一旦置为存在,后续读取永不再回源,
     * 会话历史会被永久截断成那几条。回填后按 lastN 切尾返回。
     */
    private List<String> coldRead(String convId, int lastN) {
        if (jdbc == null) {
            return List.of();
        }
        try {
            List<Map<String, Object>> rows = jdbc.queryForList(
                    "SELECT sType, sPayload FROM ai_chat_event WHERE sConversationId=? ORDER BY iId DESC LIMIT " + COLD_READ_MAX,
                    convId);
            if (rows.isEmpty()) {
                return List.of();
            }
            List<String> out = new ArrayList<>(rows.size());
            for (Map<String, Object> 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 lastN > 0 && out.size() > lastN ? out.subList(out.size() - lastN, out.size()) : out;
        } catch (Exception e) {
            log.warn("ledger mysql read failed (conv={}): {}", convId, e.getMessage());
            return List.of();
        }
    }

    /**
     * 会话在热缓存或持久层是否有内容(供会话列表懒清理)。MySQL 里的事件永不删除,
     * 因此这里只看 Redis 热缓存:30 天没动过的会话从侧栏消失,但账本仍可溯源。
     */
    public boolean exists(String convId) {
        try {
            return Boolean.TRUE.equals(redis.hasKey(PREFIX + convId));
        } catch (Exception e) {
            return true; // 读失败时不误删列表项
        }
    }

    /**
     * 删除会话 = 清热缓存 + 在持久层落一条 {@code deleted} 标记(append-only,旧事件永久保留可溯源,
     * 但冷读止于标记——同名会话再启用时不会复活已删除的历史)。
     */
    public void delete(String convId) {
        try {
            redis.delete(PREFIX + convId);
        } catch (Exception ignore) {
        }
        persist(convId, "deleted", "{\"type\":\"deleted\"}", null);
    }
}