cc 1 vecka sedan
förälder
incheckning
249c4634d7

+ 26 - 0
ai-server/src/main/java/com/zsjz/ai/common/datasource/CaseDataSourceRegistry.java

@@ -6,6 +6,7 @@ import com.zsjz.ai.common.exception.ServerException;
 import lombok.extern.slf4j.Slf4j;
 import org.springframework.stereotype.Component;
 
+import java.nio.file.Path;
 import java.sql.Connection;
 import java.sql.Statement;
 import java.util.Collections;
@@ -163,6 +164,31 @@ public class CaseDataSourceRegistry {
             "ALTER TABLE file_ai_profile ADD COLUMN IF NOT EXISTS caseId BIGINT"
     );
 
+    /**
+     * 导出案件库的一份<b>独立单文件快照</b>(供 Python 子进程只读打开)。
+     *
+     * <p>为什么必须由本进程导出:DuckDB 单文件只允许一个进程读写打开,而本进程持写锁时,
+     * 子进程连只读都打不开(Windows 报「另一个程序正在使用此文件」)。
+     * 实现见 {@link DuckdbSnapshot}。
+     *
+     * @param caseId 案件 ID(未打开时返回 false)
+     * @param target 快照文件路径(覆盖已存在文件)
+     * @return 是否导出成功;失败时调用方应回退到原库路径(行为与旧版一致)
+     */
+    public boolean exportSnapshot(Long caseId, Path target) {
+        OpenedCase entry = caseId == null ? null : opened.get(caseId);
+        if (entry == null) {
+            log.debug("案件未打开,跳过快照导出: caseId={}", caseId);
+            return false;
+        }
+        try (Connection con = entry.dataSource().getConnection()) {
+            return DuckdbSnapshot.export(con, target);
+        } catch (Exception e) {
+            log.warn("导出案件快照失败: caseId={}, target={}, error={}", caseId, target, e.getMessage());
+            return false;
+        }
+    }
+
     /** 关闭并移除案件数据源;未打开时静默返回 */
     public synchronized void close(Long caseId) {
         if (caseId == null) {

+ 144 - 0
ai-server/src/main/java/com/zsjz/ai/common/datasource/DuckdbSnapshot.java

@@ -0,0 +1,144 @@
+package com.zsjz.ai.common.datasource;
+
+import lombok.extern.slf4j.Slf4j;
+
+import java.io.IOException;
+import java.nio.file.AtomicMoveNotSupportedException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.nio.file.StandardCopyOption;
+import java.sql.Connection;
+import java.sql.ResultSet;
+import java.sql.Statement;
+import java.util.UUID;
+
+/**
+ * DuckDB 库文件快照导出:把当前库(连同表与视图)复制成一个<b>独立的单文件</b>。
+ *
+ * <h3>为什么需要它</h3>
+ * DuckDB 单文件<b>只允许一个进程以读写方式打开</b>;而且只要某进程持着写锁,
+ * 其它进程连 {@code read_only=True} 都打不开(Windows 上报「另一个程序正在使用此文件」)。
+ * 案件库在本进程(Java)里是常驻打开的,于是「Python 子进程自己连案件库」必然失败 ——
+ * 这正是 {@code PythonExecutor} 之前的问题。
+ *
+ * <p>解法:由<b>持锁方的本进程</b>用 {@code ATTACH + COPY FROM DATABASE} 导出一份副本,
+ * 子进程去连副本(不同文件,不受锁约束)。为什么不用文件复制:直接 copy 一个正在被写的
+ * 库文件可能拷到撕裂页;{@code COPY FROM DATABASE} 走的是库内读取,不存在这个问题。
+ *
+ * <p>实测(DuckDB 1.5.5.1):表与视图都会随快照带过去,副本可被另一进程只读打开。
+ *
+ * <h3>并发与命名(多用户 / 多工作空间)</h3>
+ * <ul>
+ *   <li><b>导出别名每次唯一</b>:同一个库实例上并发导出时,固定别名会让第二个 ATTACH 撞名失败;</li>
+ *   <li><b>先写临时文件再替换</b>:读者(Python 子进程)不会看到半成品,
+ *       同一路径并发导出由「最后一次替换生效」兜住;</li>
+ *   <li>快照文件名由调用方按<b>工作空间 ID</b> 命名(见 {@code PythonExecutor}),
+ *       即使快照被挪到别处或只按文件名排查,也不会和别的空间/用户混淆。</li>
+ * </ul>
+ */
+@Slf4j
+public final class DuckdbSnapshot {
+
+    /** 导出期间用的临时库别名前缀(每次导出追加随机后缀;导出结束在 finally 里 DETACH) */
+    private static final String EXPORT_ALIAS_PREFIX = "snapshot_export_";
+
+    /** 临时文件后缀:导出先落临时文件,成功后再替换到目标路径 */
+    private static final String TEMP_SUFFIX = ".tmp-";
+
+    private DuckdbSnapshot() {
+    }
+
+    /**
+     * 把 {@code con} 所连的库导出为 {@code target} 指向的单文件库(覆盖已存在的同名文件)。
+     *
+     * @param con    案件库的连接(持锁方自己的连接)
+     * @param target 快照文件路径(父目录会自动创建)
+     * @return 成功返回 true;失败返回 false 并记 WARN(调用方应回退到原库路径,保持与旧行为一致)
+     */
+    public static boolean export(Connection con, Path target) {
+        if (con == null || target == null) {
+            return false;
+        }
+        String alias = EXPORT_ALIAS_PREFIX + UUID.randomUUID().toString().replace("-", "").substring(0, 8);
+        Path temp = target.resolveSibling(target.getFileName() + TEMP_SUFFIX + alias);
+        try {
+            Files.createDirectories(target.getParent());
+            Files.deleteIfExists(temp);
+
+            String mainDatabase = mainDatabaseName(con);
+            if (mainDatabase == null) {
+                log.warn("导出 DuckDB 快照失败:无法确定主库名(目标 {})", target);
+                return false;
+            }
+            String tempPath = temp.toAbsolutePath().toString().replace('\\', '/').replace("'", "''");
+            try (Statement st = con.createStatement()) {
+                st.execute("attach '" + tempPath + "' as " + alias);
+                try {
+                    st.execute("copy from database \"" + mainDatabase + "\" to " + alias);
+                } finally {
+                    st.execute("detach " + alias);
+                }
+            }
+            replace(temp, target);
+            log.debug("DuckDB 快照已导出: {} -> {} ({} bytes)", mainDatabase, target, Files.size(target));
+            return true;
+        } catch (Exception e) {
+            log.warn("导出 DuckDB 快照失败: target={}, error={}", target, e.getMessage());
+            return false;
+        } finally {
+            try {
+                Files.deleteIfExists(temp);
+            } catch (Exception e) {
+                log.debug("清理临时快照失败: {} ({})", temp, e.getMessage());
+            }
+        }
+    }
+
+    /**
+     * 把导出好的临时文件替换到目标路径。
+     *
+     * <p>优先原子替换,读者永远看不到半成品;Windows 上目标文件若正被别的进程打开
+     * (例如上一轮 Python 子进程还没退出)会失败,退化为普通替换;再失败就抛出,
+     * 由调用方回退原库。同一路径并发导出时最后一次替换生效,不会互相写坏。
+     */
+    private static void replace(Path temp, Path target) throws IOException {
+        try {
+            Files.move(temp, target, StandardCopyOption.REPLACE_EXISTING, StandardCopyOption.ATOMIC_MOVE);
+        } catch (AtomicMoveNotSupportedException e) {
+            Files.move(temp, target, StandardCopyOption.REPLACE_EXISTING);
+        }
+    }
+
+    /**
+     * 取当前连接的主库名。
+     *
+     * <p>先试 {@code current_catalog()}(DuckDB 1.x 起可用,最直接);失败再退回
+     * {@code duckdb_databases()} 里第一个非 internal 的库。
+     */
+    private static String mainDatabaseName(Connection con) {
+        try (Statement st = con.createStatement();
+             ResultSet rs = st.executeQuery("select current_catalog()")) {
+            if (rs.next()) {
+                String name = rs.getString(1);
+                if (name != null && !name.isBlank()) {
+                    return name;
+                }
+            }
+        } catch (Exception e) {
+            log.debug("current_catalog() 不可用,退回 duckdb_databases(): {}", e.getMessage());
+        }
+        try (Statement st = con.createStatement();
+             ResultSet rs = st.executeQuery(
+                     "select database_name from duckdb_databases() where internal = false")) {
+            while (rs.next()) {
+                String name = rs.getString(1);
+                if (name != null && !name.isBlank() && !name.startsWith(EXPORT_ALIAS_PREFIX)) {
+                    return name;
+                }
+            }
+        } catch (Exception e) {
+            log.warn("读取主库名失败: {}", e.getMessage());
+        }
+        return null;
+    }
+}

+ 86 - 4
ai-server/src/main/java/com/zsjz/ai/module/agent/python/PythonExecutor.java

@@ -2,6 +2,7 @@ package com.zsjz.ai.module.agent.python;
 
 import com.zsjz.ai.common.constants.PathConst;
 import com.zsjz.ai.common.context.CaseContextHolder;
+import com.zsjz.ai.common.datasource.CaseDataSourceRegistry;
 import com.zsjz.ai.common.model.plat.entity.CaseInfo;
 import com.zsjz.ai.module.agent.config.EtlProperties;
 
@@ -23,9 +24,19 @@ import java.util.stream.Stream;
  * Python 子进程沙箱执行器
  *
  * <p>每次执行使用独立临时目录({@code data/workspace/{wsId}/tmp/py-{execId}}),
- * 通过环境变量 {@code ETL_DB_PATH} 注入工作空间 DuckDB 文件路径(Python 内只读连接),
  * 超时强制终止、输出截断;执行后将沙箱内生成的图片移动到
  * {@code data/workspace/{wsId}/py-out/{execId}} 供文件服务接口输出。
+ *
+ * <h3>ETL_DB_PATH 指向的是「快照副本」而不是原库</h3>
+ * DuckDB 单文件只允许一个进程读写打开,而案件库在本进程(Java)里是常驻打开的 ——
+ * 子进程直接连原库必然失败(Windows 报「另一个程序正在使用此文件」)。
+ * 因此执行前由本进程用 {@code ATTACH + COPY FROM DATABASE} 导出一份单文件快照
+ * ({@code data/workspace/{wsId}/py-snapshot/case-{wsId}.db},见 {@link CaseDataSourceRegistry#exportSnapshot}),
+ * 把 {@code ETL_DB_PATH} 指向它。快照里表与视图都在,Python 侧用法不变(仍可 read_only 打开)。
+ * 目录按工作空间隔离、文件名再带上工作空间 ID,多用户/多空间并发时既不会重名也好辨认。
+ *
+ * <p>快照按<b>源库指纹</b>(主文件与 WAL 的大小/修改时间)缓存:源库自上次导出后没变就直接复用,
+ * 避免模型连着跑几段 Python 时反复导出整库。
  */
 @Slf4j
 @Service
@@ -37,7 +48,9 @@ public class PythonExecutor {
     private static final String CHECK_IMPORTS = "import pandas, duckdb, matplotlib";
 
     /**
-     * 环境注入的 DuckDB 文件路径变量名
+     * 环境注入的 DuckDB 文件路径变量名。
+     *
+     * <p>值是<b>快照副本</b>路径(不是案件原库):原库被本进程写锁持有,子进程连不上。
      */
     public static final String ENV_DB_PATH = "ETL_DB_PATH";
 
@@ -51,12 +64,27 @@ public class PythonExecutor {
      */
     public static final String IMAGE_URL_PREFIX = "/js/a/py/files/";
 
+    /** 快照目录名(每个工作空间一份,覆盖写) */
+    private static final String SNAPSHOT_DIR = "py-snapshot";
+
+    /**
+     * 快照文件名前缀。
+     *
+     * <p>文件名里带<b>工作空间 ID</b>({@code case-{workspaceId}.db}):本系统是多用户、多工作空间并存的,
+     * 通用文件名一旦被挪到同一目录就会互相覆盖,也不利于排障时辨认归属。
+     */
+    private static final String SNAPSHOT_PREFIX = "case-";
+
     private final EtlProperties etlProperties;
     private final CaseInfoMapper caseInfoMapper;
+    private final CaseDataSourceRegistry caseDataSourceRegistry;
 
-    public PythonExecutor(EtlProperties etlProperties, CaseInfoMapper caseInfoMapper) {
+    public PythonExecutor(EtlProperties etlProperties,
+                          CaseInfoMapper caseInfoMapper,
+                          CaseDataSourceRegistry caseDataSourceRegistry) {
         this.etlProperties = etlProperties;
         this.caseInfoMapper = caseInfoMapper;
+        this.caseDataSourceRegistry = caseDataSourceRegistry;
     }
 
     /**
@@ -119,7 +147,9 @@ public class PythonExecutor {
 
             ProcessBuilder pb = new ProcessBuilder(etlProperties.getPython().getCommand(), "main.py")
                     .directory(sandbox.toFile());
-            pb.environment().put(ENV_DB_PATH, dbFile);
+            // 指向快照副本:原库被本进程写锁持有,子进程连只读都打不开(见类注释)
+            String pythonDbFile = resolvePythonDbFile(workspaceId, dbFile);
+            pb.environment().put(ENV_DB_PATH, pythonDbFile);
             pb.environment().put("MPLBACKEND", "Agg"); // matplotlib 无头渲染
 
             Process process = pb.start();
@@ -196,6 +226,58 @@ public class PythonExecutor {
         }
     }
 
+    /**
+     * 解析注入给 Python 的库文件路径:优先「当前案件的快照副本」,失败时回退原库。
+     *
+     * <p>导出时机用源库指纹判断:主文件与 WAL 的大小/修改时间都没变 → 复用已有快照。
+     * 任何写入都会改 WAL,因此数据新鲜度与「每次重新导出」等价,代价只在真正变更时付一次。
+     */
+    private String resolvePythonDbFile(Long workspaceId, String dbFile) {
+        if (!etlProperties.getPython().isSnapshotEnabled()) {
+            return dbFile;
+        }
+        Path db = Path.of(dbFile);
+        Path snapshot = PathConst.WORKSPACE.resolve(String.valueOf(workspaceId))
+                .resolve(SNAPSHOT_DIR).resolve(SNAPSHOT_PREFIX + workspaceId + ".db");
+        Path meta = snapshot.resolveSibling(SNAPSHOT_PREFIX + workspaceId + ".meta");
+        try {
+            String fingerprint = sourceFingerprint(db);
+            if (Files.isRegularFile(snapshot) && Files.isRegularFile(meta)
+                    && fingerprint.equals(Files.readString(meta).trim())) {
+                return snapshot.toAbsolutePath().toString();
+            }
+            if (!caseDataSourceRegistry.exportSnapshot(workspaceId, snapshot)) {
+                return dbFile;
+            }
+            // 指纹在导出之后再取一次:避免把导出过程自身对源库的副作用算进来
+            Files.writeString(meta, sourceFingerprint(db));
+            return snapshot.toAbsolutePath().toString();
+        } catch (Exception e) {
+            log.warn("准备 Python 数据库快照失败,回退原库: caseId={}, error={}", workspaceId, e.getMessage());
+            return dbFile;
+        }
+    }
+
+    /**
+     * 源库指纹:主文件 + WAL 的大小与修改时间。任一写入都会动 WAL,因此指纹变化 = 数据变了。
+     */
+    private static String sourceFingerprint(Path dbFile) {
+        StringBuilder sb = new StringBuilder();
+        sb.append(statOf(dbFile)).append('|').append(statOf(Path.of(dbFile + ".wal")));
+        return sb.toString();
+    }
+
+    private static String statOf(Path file) {
+        try {
+            if (!Files.exists(file)) {
+                return "-";
+            }
+            return Files.size(file) + "@" + Files.getLastModifiedTime(file).toMillis();
+        } catch (Exception e) {
+            return "?";
+        }
+    }
+
     /**
      * 捕获沙箱内图片:移动到 py-out/{execId}/ 并生成访问 URL,LRU 清理超额执行目录
      */

+ 9 - 0
ai-server/src/main/java/com/zsjz/ai/module/agent/python/PythonProperties.java

@@ -27,4 +27,13 @@ public class PythonProperties {
 
     /** 图片输出目录保留最近执行数(LRU 清理) */
     private int maxHistoryExecs = 20;
+
+    /**
+     * 是否为每次执行导出一份案件库快照给 Python 连接。
+     *
+     * <p>默认开:DuckDB 单文件只允许一个进程读写打开,Java 持写锁时子进程连只读都打不开
+     * (见 {@code DuckdbSnapshot})。关掉的话 Python 只能用 {@code ETL_DB_PATH} 指向的原库,
+     * 而原库只要已被本进程打开就连不上。
+     */
+    private boolean snapshotEnabled = true;
 }

+ 4 - 2
ai-server/src/main/java/com/zsjz/ai/module/agent/tools/PythonAnalysisTool.java

@@ -34,9 +34,11 @@ public class PythonAnalysisTool {
      */
     @Tool(name = "execute_python",
           description = "在隔离沙箱中执行 Python 数据分析代码(pandas/numpy/duckdb/matplotlib 可用)。"
-                  + "通过环境变量 ETL_DB_PATH 获得当前工作空间 DuckDB 文件路径,"
-                  + "用 duckdb.connect(ETL_DB_PATH, read_only=True) 只读访问数据(示例:"
+                  + "通过环境变量 ETL_DB_PATH 获得当前工作空间数据的 DuckDB 路径,"
+                  + "用 duckdb.connect(ETL_DB_PATH, read_only=True) 只读访问(示例:"
                   + "import os, duckdb; con = duckdb.connect(os.environ['ETL_DB_PATH'], read_only=True))。"
+                  + "该路径指向本次执行开始时导出的一致性快照副本(表与视图都在),"
+                  + "★ 不要试图连接案件原库文件(系统已持有其写锁,会报「另一个程序正在使用此文件」)。"
                   + "用 print() 输出分析结论文本;matplotlib 图表用 plt.savefig('chart.png') 保存"
                   + "(系统自动捕获并展示给用户)。"
                   + "适合 execute_sql 无法完成的复杂分析:统计建模、多步数据变换、专业可视化。",

+ 155 - 0
ai-server/src/test/java/com/zsjz/ai/common/datasource/DuckdbSnapshotTest.java

@@ -0,0 +1,155 @@
+package com.zsjz.ai.common.datasource;
+
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.sql.Connection;
+import java.sql.DriverManager;
+import java.sql.ResultSet;
+import java.sql.Statement;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Properties;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * DuckDB 快照导出的契约测试(真实驱动,不用 mock)。
+ *
+ * <p>守两件事:
+ * <ol>
+ *   <li><b>表与视图都随快照带过去</b>:Python 侧的查询口径与直连原库一致,不用改写法;</li>
+ *   <li><b>快照能被另一个连接只读打开</b>:这正是「Java 持写锁时 Python 连原库必失败」的解药
+ *       (跨进程的锁语义已用独立 JVM 进程实测,这里用同进程的独立连接做回归)。</li>
+ * </ol>
+ */
+class DuckdbSnapshotTest {
+
+    private static final String SRC = "jdbc:duckdb:";
+
+    @TempDir
+    Path dir;
+
+    private static Connection open(Path file) throws Exception {
+        return DriverManager.getConnection(SRC + file);
+    }
+
+    private static Connection openReadOnly(Path file) throws Exception {
+        Properties props = new Properties();
+        props.setProperty("access_mode", "read_only");
+        return DriverManager.getConnection(SRC + file, props);
+    }
+
+    private static int scalar(Connection con, String sql) throws Exception {
+        try (Statement st = con.createStatement();
+             ResultSet rs = st.executeQuery(sql)) {
+            rs.next();
+            return rs.getInt(1);
+        }
+    }
+
+    @Test
+    @DisplayName("导出的快照包含表与视图,可被独立连接只读打开")
+    void exportCopiesTablesAndViews() throws Exception {
+        Path source = dir.resolve("case1.db");
+        Path snapshot = dir.resolve("py-snapshot").resolve("case-123456789.db");
+
+        try (Connection con = open(source); Statement st = con.createStatement()) {
+            st.execute("create table t as select * from range(7) tbl(i)");
+            st.execute("create view v as select i * 2 as j from t");
+
+            assertTrue(DuckdbSnapshot.export(con, snapshot), "导出应成功");
+
+            // 快照存在于另一个文件,且源库仍可继续用(导出不影响原连接)
+            assertTrue(Files.isRegularFile(snapshot));
+            assertEquals(7, scalar(con, "select count(*) from t"));
+        }
+
+        try (Connection snap = openReadOnly(snapshot)) {
+            assertEquals(7, scalar(snap, "select count(*) from t"), "表数据应完整");
+            assertEquals(7, scalar(snap, "select count(*) from v"), "视图应随快照复制");
+        }
+    }
+
+    @Test
+    @DisplayName("重复导出覆盖旧快照(数据是导出时刻的状态)")
+    void exportOverwritesExistingSnapshot() throws Exception {
+        Path source = dir.resolve("case2.db");
+        Path snapshot = dir.resolve("py-snapshot").resolve("case-123456789.db");
+
+        try (Connection con = open(source); Statement st = con.createStatement()) {
+            st.execute("create table t as select * from range(3) tbl(i)");
+            assertTrue(DuckdbSnapshot.export(con, snapshot));
+
+            st.execute("insert into t select * from range(10) tbl(i)");
+            assertTrue(DuckdbSnapshot.export(con, snapshot), "第二次导出应覆盖而不是报错");
+        }
+
+        try (Connection snap = openReadOnly(snapshot)) {
+            assertEquals(13, scalar(snap, "select count(*) from t"), "快照应是最近一次导出的内容");
+        }
+    }
+
+    @Test
+    @DisplayName("并发导出同一路径:别名不撞车、结果仍是可用快照、临时文件不残留")
+    void concurrentExportToSameTarget() throws Exception {
+        Path source = dir.resolve("case9.db");
+        Path snapshotDir = dir.resolve("py-snapshot");
+        Path snapshot = snapshotDir.resolve("case-9.db");
+        try (Connection con = open(source); Statement st = con.createStatement()) {
+            st.execute("create table t as select * from range(5) tbl(i)");
+        }
+
+        int threads = 4;
+        ExecutorService pool = Executors.newFixedThreadPool(threads);
+        try {
+            List<Future<Boolean>> futures = new ArrayList<>();
+            for (int i = 0; i < threads; i++) {
+                futures.add(pool.submit(() -> {
+                    try (Connection c = open(source)) {
+                        return DuckdbSnapshot.export(c, snapshot);
+                    }
+                }));
+            }
+            int ok = 0;
+            for (Future<Boolean> future : futures) {
+                if (Boolean.TRUE.equals(future.get())) {
+                    ok++;
+                }
+            }
+            assertTrue(ok >= 1, "并发导出一份都没成功");
+        } finally {
+            pool.shutdownNow();
+        }
+
+        // 结果仍是可用快照(最后一次替换生效),且没有 .tmp- 残留
+        try (Connection snap = openReadOnly(snapshot)) {
+            assertEquals(5, scalar(snap, "select count(*) from t"));
+        }
+        try (var files = Files.list(snapshotDir)) {
+            assertTrue(files.noneMatch(p -> p.getFileName().toString().contains(".tmp-")),
+                    "导出结束后不应留下临时文件");
+        }
+    }
+
+    @Test
+    @DisplayName("入参非法/连接不可用时返回 false,不抛异常(调用方回退原库)")
+    void exportFailsQuietly() throws Exception {
+        Path source = dir.resolve("case3.db");
+        Path snapshot = dir.resolve("py-snapshot").resolve("case-123456789.db");
+        Connection con = open(source);
+        con.close();
+
+        assertFalse(DuckdbSnapshot.export(null, snapshot), "空连接应返回 false");
+        assertFalse(DuckdbSnapshot.export(con, null), "空目标应返回 false");
+        assertFalse(DuckdbSnapshot.export(con, snapshot), "连接已关闭应返回 false 而不是抛异常");
+    }
+}