"""旧源只读加载:仅合成SQLite和显式隔离PG;不访问真实旧库或导入业务目标。""" from __future__ import annotations import base64 import hashlib import json import sqlite3 import subprocess import sys from pathlib import Path from uuid import uuid4 import psycopg import pytest from psycopg import sql from 建立映射 import 迁移错误 from 读取旧PG import _打开只读连接 as 打开只读PG from 读取旧PG import 读取旧库 as 读取PG from 读取旧SQLite import _打开只读连接 as 打开只读SQLite from 读取旧SQLite import 读取旧库 as 读取SQLite @pytest.fixture def sqlite副本(tmp_path): 路径 = tmp_path / "合成旧库.db" with sqlite3.connect(路径) as 连: 连.executescript(""" CREATE TABLE runs(id TEXT PRIMARY KEY, output_text TEXT NOT NULL); CREATE TABLE events(run_id TEXT NOT NULL,seq INTEGER NOT NULL, payload_json TEXT,PRIMARY KEY(run_id,seq)); CREATE TABLE cards(id TEXT PRIMARY KEY,kind TEXT,payload_json TEXT,embedding BLOB); INSERT INTO runs VALUES('r1','仅供测试的原运行输出'); INSERT INTO events VALUES('r1',1,'{"片段":"甲"}'),('r1',2,'{"片段":"乙"}'); """) 连.execute("INSERT INTO cards VALUES(?,?,?,?)", ("c1", "embedding", "{}", b"\x00\xff")) return 路径 def _sqlite(路径, **参数): return 读取SQLite( 路径, 数据集="合成旧源", 表=("runs", "events", "cards"), 快照版本="export-v1", **参数 ) def test_SQLite原行复合主键字节与只读性__21b003(sqlite副本): 原哈希 = hashlib.sha256(sqlite副本.read_bytes()).hexdigest() 结果 = _sqlite(sqlite副本) assert len(结果) == 4 事件 = [r for r in 结果 if r.source.table == "events"] assert [json.loads(r.source.id) for r in 事件] == [ {"run_id": "r1", "seq": 1}, {"run_id": "r1", "seq": 2}, ] assert all(r.source.revision == "export-v1" for r in 结果) assert all(str(sqlite副本) not in r.source.database for r in 结果) 卡 = next(r for r in 结果 if r.source.table == "cards") assert base64.b64decode(卡.原行["embedding"]["__bytea_base64__"]) == b"\x00\xff" with 打开只读SQLite(sqlite副本) as 连: assert 连.execute("PRAGMA query_only").fetchone()[0] == 1 with pytest.raises(sqlite3.OperationalError, match="readonly"): 连.execute("DELETE FROM runs") assert hashlib.sha256(sqlite副本.read_bytes()).hexdigest() == 原哈希 def test_SQLite缺文件未知表和视图不能暗读__21b004(sqlite副本, tmp_path): 缺失 = tmp_path / "不存在.db" with pytest.raises(FileNotFoundError): _sqlite(缺失) assert not 缺失.exists() with pytest.raises(迁移错误, match="source_table_denied"): 读取SQLite(sqlite副本, 数据集="合成旧源", 表=("private_keys",), 快照版本="x") with sqlite3.connect(sqlite副本) as 连: 连.execute("DROP TABLE cards") 连.execute("CREATE VIEW cards AS SELECT id FROM runs") with pytest.raises(迁移错误, match="source_table_denied"): _sqlite(sqlite副本) def test_SQLite范围行数与字节上限不静默截断__21b005(sqlite副本): 结果 = 读取SQLite( sqlite副本, 数据集="合成旧源", 表=("events",), 快照版本="x", 范围={"events": {"start": ["r1", 2], "end": ["r1", 3]}}, ) assert len(结果) == 1 and 结果[0].原行["seq"] == 2 with pytest.raises(迁移错误, match="source_range_invalid"): _sqlite(sqlite副本, 范围={"events": {"start": [1]}}) with pytest.raises(迁移错误, match="snapshot_limit"): _sqlite(sqlite副本, 最大行数=1) with sqlite3.connect(sqlite副本) as 连: 连.execute("UPDATE runs SET output_text=?", ("x" * (16 * 1024 * 1024),)) with pytest.raises(迁移错误, match="snapshot_limit"): _sqlite(sqlite副本) def test_读取命令接纯分流且不覆盖源__21b006(sqlite副本, tmp_path): 脚本 = Path(__file__).resolve().parents[2] / "数据库/旧库迁移/入口.py" 输出, 映射, 计划 = [tmp_path / n for n in ("快照.json", "映射.json", "计划.json")] 命令 = [ sys.executable, str(脚本), "读取SQLite", "--副本", str(sqlite副本), "--数据集", "合成旧源", "--表", "runs", "--表", "events", "--快照版本", "export-v1", "--输出", str(输出), ] for _ in range(2): 运行 = subprocess.run(命令, capture_output=True, text=True, timeout=15) assert 运行.returncode == 0, 运行.stderr assert "仅供测试的原运行输出" not in 运行.stdout assert len(json.loads(输出.read_text())) == 3 and 输出.stat().st_mode & 0o777 == 0o600 映射.write_text('{"works": []}', encoding="utf-8") 运行 = subprocess.run( [ sys.executable, str(脚本), "分流", "--快照", str(输出), "--映射", str(映射), "--输出", str(计划), ], capture_output=True, text=True, timeout=15, ) assert 运行.returncode == 0, 运行.stderr assert json.loads(计划.read_text())["dispositions"] == {"历史保留": 3} 哈希 = hashlib.sha256(sqlite副本.read_bytes()).hexdigest() assert subprocess.run([*命令[:-1], str(sqlite副本)], capture_output=True).returncode == 1 assert hashlib.sha256(sqlite副本.read_bytes()).hexdigest() == 哈希 @pytest.fixture def pg旧源(隔离数据库URL, tmp_path): schema = "w21_source_" + uuid4().hex with psycopg.connect(隔离数据库URL) as 连: 连.execute(sql.SQL("CREATE SCHEMA {}").format(sql.Identifier(schema))) 连.execute( sql.SQL( "CREATE TABLE {}(id bigint PRIMARY KEY, tenant_id bigint NOT NULL, " "revision int NOT NULL, description text)" ).format(sql.Identifier(schema, "muse_knowledge_entity")) ) 连.execute( sql.SQL( "INSERT INTO {} VALUES (1,7,3,'合成甲'),(2,8,4,'合成乙'),(3,7,5,'合成丙')" ).format(sql.Identifier(schema, "muse_knowledge_entity")) ) 连.execute( sql.SQL( "CREATE TABLE {}(run_id text, seq int, data text, PRIMARY KEY(run_id,seq))" ).format(sql.Identifier(schema, "example_run_receipt")) ) 连.execute( sql.SQL("INSERT INTO {} VALUES ('r1',1,'合成事件')").format( sql.Identifier(schema, "example_run_receipt") ) ) 文件 = tmp_path / "连接说明.txt" 文件.write_text(隔离数据库URL, encoding="utf-8") 文件.chmod(0o600) try: yield 文件, schema finally: with psycopg.connect(隔离数据库URL) as 连: 连.execute(sql.SQL("DROP SCHEMA {} CASCADE").format(sql.Identifier(schema))) @pytest.mark.数据库 def test_PG按租户范围读取且保留真实修订__21b007(pg旧源): 文件, schema = pg旧源 参数 = { "数据集": "合成PG", "schema": schema, "表": ("muse_knowledge_entity",), "快照版本": "export-v1", } 一 = 读取PG(文件, 租户=7, **参数) 二 = 读取PG(文件, 租户=8, **参数) assert [r.原行["id"] for r in 一] == [1, 3] and [r.原行["id"] for r in 二] == [2] assert [r.source.revision for r in 一] == ["3", "5"] assert 一[0].source.database != 二[0].source.database 限定 = 读取PG(文件, 租户=7, 范围={"muse_knowledge_entity": {"start": [3], "end": [4]}}, **参数) assert len(限定) == 1 and 限定[0].原行["id"] == 3 with pytest.raises(迁移错误, match="source_scope_invalid"): 读取PG(文件, 租户=None, **参数) with pytest.raises(迁移错误, match="snapshot_limit"): 读取PG(文件, 租户=7, 最大行数=1, **参数) 全局 = 读取PG( 文件, 数据集="合成PG", schema=schema, 租户=None, 表=("example_run_receipt",), 快照版本="export-v1", ) assert json.loads(全局[0].source.id) == {"run_id": "r1", "seq": 1} assert 全局[0].source.revision == "export-v1" @pytest.mark.数据库 def test_PG源事务真正只读且声明错误不读取__21b008(pg旧源): 文件, schema = pg旧源 with pytest.raises(迁移错误, match="25006"): with 打开只读PG(文件) as 连: assert ( 连.execute("SHOW transaction_read_only").fetchone()["transaction_read_only"] == "on" ) assert ( 连.execute("SHOW transaction_isolation").fetchone()["transaction_isolation"] == "repeatable read" ) 连.execute( sql.SQL("DELETE FROM {}").format(sql.Identifier(schema, "muse_knowledge_entity")) ) with psycopg.connect(文件.read_text()) as 连: assert ( 连.execute( sql.SQL("SELECT count(*) FROM {}").format( sql.Identifier(schema, "muse_knowledge_entity") ) ).fetchone()[0] == 3 ) 连.execute( sql.SQL("CREATE VIEW {} AS SELECT 1 AS id").format( sql.Identifier(schema, "example_run") ) ) with pytest.raises(迁移错误, match="source_table_denied"): 读取PG(文件, 数据集="合成PG", schema=schema, 租户=None, 表=("example_run",), 快照版本="v1") def test_连接说明不回落环境且错误不泄露凭据__21b009(tmp_path, monkeypatch): 文件 = tmp_path / "连接说明.txt" monkeypatch.setattr( psycopg.Connection, "connect", lambda *a, **k: pytest.fail("配置非法时不得连接") ) for 内容 in ("", "dbname=somewhere", "host=localhost user=someone"): 文件.write_text(内容, encoding="utf-8") with pytest.raises(迁移错误, match="source_connection_invalid"): with 打开只读PG(文件): pytest.fail("不得取得连接") 文件.write_text("password='never-print-this", encoding="utf-8") with pytest.raises(迁移错误) as 错: with 打开只读PG(文件): pytest.fail("不得取得连接") assert "never-print-this" not in str(错.value)