313 lines
16 KiB
Python
313 lines
16 KiB
Python
#!/usr/bin/env python3
|
||
"""parse-book skill:升格全量重抽准备——清一本书的升格产出与串行状态(2026-07-18 批9c 收尾)。
|
||
|
||
为什么需要它(锚点「全量重抽的管线坑」第 1 条):升格是**串行有状态**管线,直接重跑会被
|
||
旧状态卡死——窗全是 done 不再处理、预扫/判重底册装着旧代卡把新观察并进废卡、旧审计行会被
|
||
撤销机制错还原到已软删的旧卡上、别名表全行唯一索引让新别名被"冲突即跳过"静默吞掉。
|
||
|
||
清理边界(逐表核实过约束后定的,2026-07-18):
|
||
· muse_knowledge_draft 升格卡 → **软删**(deleted=TRUE,红线:知识数据绝不物删,可回滚)
|
||
· example_upgrade_window → 保留窗行(窗参数未变、from_chapter 锚幂等),仅重置 status=pending
|
||
· example_upgrade_alias → **硬删**该书行——uk(tenant_id,work_id,alias) 是全行唯一索引(非部分索引),
|
||
软删旧行会挡住重抽新 INSERT(ON CONFLICT DO NOTHING 静默吞)→ 别名判重空转
|
||
· example_upgrade_presence → 硬删该书行(运行状态留档,重抽全量重生成;undo 机制既有做法即硬删)
|
||
· example_upgrade_card_state → 硬删该书行(卡水位挂旧卡 draft_id,新卡新水位)
|
||
· example_upgrade_audit → 硬删该书行——undo_window 的还原 JOIN **不过滤 d.deleted**,
|
||
留着旧审计行,重抽期间任何同窗撤销都会把旧代软删卡错还原(污染)
|
||
· example_knowledge_embedding → **软删**该书升格卡的全部活向量(不物删,保留恢复能力);同时覆盖
|
||
上次旧版 reset 遗留在已软删卡上的活向量,避免 HNSW 候选持续膨胀
|
||
|
||
用法:
|
||
reset_upgrade_work.py --work-id N 预览(只读统计,不写库)
|
||
reset_upgrade_work.py --work-id N --execute --backup-dir DIR --backup-id ID \
|
||
--confirmation-sha SHA 真清(锁内验备份,单事务全清或全不清)
|
||
"""
|
||
import sys
|
||
import pathlib
|
||
|
||
import click
|
||
import psycopg
|
||
from psycopg.rows import dict_row
|
||
|
||
# 复用管线的连接与租户常量(与 parse_upgrade 同源)
|
||
sys.path.insert(0, str(pathlib.Path(__file__).resolve().parent))
|
||
import backup_upgrade_work as backup # noqa: E402
|
||
from parse_llm import DSN, TENANT # noqa: E402
|
||
from upgrade_work_lock import UpgradeWorkLockUnavailable, upgrade_work_lock # noqa: E402
|
||
|
||
SOURCE_TYPE = "upgrade_book"
|
||
|
||
# capture_input_snapshot 实际读取 chapter/block 与 schema/schema_version;它们必须和七域一起先锁定,
|
||
# 否则完整输入核对后仍可能被正文或 active 合同写入穿透。固定顺序采用“输入源→派生域”:潜在写者
|
||
# 可能先持输入源表再写派生域,若 reset 反向先持派生表再等待输入源会形成互等。派生域内部沿用既有
|
||
# 顺序,尤其保持 draft→embedding,与 embed/parse 写事务的固定表锁顺序一致。
|
||
RESET_TABLE_LOCK_SQL = """LOCK TABLE
|
||
muse_content_chapter,
|
||
muse_content_block,
|
||
muse_meta_schema,
|
||
muse_meta_schema_version,
|
||
muse_knowledge_draft,
|
||
example_upgrade_window,
|
||
example_upgrade_alias,
|
||
example_upgrade_presence,
|
||
example_upgrade_card_state,
|
||
example_upgrade_audit,
|
||
example_knowledge_embedding
|
||
IN SHARE ROW EXCLUSIVE MODE"""
|
||
|
||
|
||
def _assert_snapshot_matches_manifest(manifest, current_domains):
|
||
"""复用备份模块摘要规则,逐域核对当前事务快照与已验证 manifest。"""
|
||
|
||
for name, spec in backup.DOMAIN_SPECS.items():
|
||
rows = backup.sort_rows(name, current_domains[name])
|
||
actual = {
|
||
"rowCount": len(rows),
|
||
"primaryKeySha": backup._primary_key_sha(spec, rows),
|
||
"contentSha": backup._content_sha(rows),
|
||
}
|
||
expected = manifest["domains"][name]
|
||
for field, value in actual.items():
|
||
if expected.get(field) != value:
|
||
raise click.ClickException(
|
||
f"备份七域快照不一致:{name}.{field} 已漂移,拒绝 reset"
|
||
)
|
||
|
||
|
||
def _require_execute_backup(backup_dir, backup_id, confirmation_sha):
|
||
"""execute 的三项备份确认缺一不可;preview 不受此约束。"""
|
||
|
||
supplied = {
|
||
"--backup-dir": backup_dir,
|
||
"--backup-id": backup_id,
|
||
"--confirmation-sha": confirmation_sha,
|
||
}
|
||
missing = [option for option, value in supplied.items() if not value]
|
||
if missing:
|
||
raise click.ClickException(
|
||
f"--execute 必须提供备份确认参数:{', '.join(missing)}"
|
||
)
|
||
|
||
|
||
def _assert_code_identity_matches_manifest(manifest, current_identity):
|
||
"""要求 execute 使用的完整代码 fileSha 与 confirmation 绑定值逐项一致。"""
|
||
|
||
manifest_files = manifest.get("input", {}).get("codeFiles")
|
||
if not isinstance(manifest_files, dict) or "reset_upgrade_work.py" not in manifest_files:
|
||
raise click.ClickException(
|
||
"备份 manifest.input.codeFiles 缺少 reset_upgrade_work.py,拒绝 reset"
|
||
)
|
||
current_files = current_identity.get("codeFiles")
|
||
if not isinstance(current_files, dict):
|
||
raise click.ClickException("无法读取当前完整 codeFiles 身份,拒绝 reset")
|
||
if manifest_files != current_files:
|
||
changed = sorted(
|
||
name for name in set(manifest_files) | set(current_files)
|
||
if manifest_files.get(name) != current_files.get(name)
|
||
)
|
||
raise click.ClickException(
|
||
f"当前代码 fileSha 与备份不一致:{', '.join(changed)}"
|
||
)
|
||
|
||
|
||
def _reset(work_id, execute, manifest=None, code_identity=None):
|
||
"""在调用方已持有同书锁时预览或执行重抽清理。"""
|
||
|
||
with psycopg.connect(DSN) as conn:
|
||
if execute:
|
||
# 输入源与七域表锁、七域核对及完整输入绑定均在同一事务内,先于 destructive SQL。
|
||
conn.execute(RESET_TABLE_LOCK_SQL)
|
||
with conn.cursor(row_factory=dict_row) as cursor:
|
||
title, current_domains = backup._read_snapshot(cursor, work_id, TENANT)
|
||
_assert_snapshot_matches_manifest(manifest, current_domains)
|
||
current_input = backup.capture_input_snapshot(
|
||
cursor,
|
||
work_id,
|
||
TENANT,
|
||
code_identity,
|
||
manifest["input"]["expectedChapters"],
|
||
manifest["input"]["expectedWindows"],
|
||
windows=current_domains["windows"],
|
||
)
|
||
backup.assert_restore_input(
|
||
manifest, current_input, work_id=work_id, tenant=TENANT
|
||
)
|
||
else:
|
||
title = conn.execute("SELECT title FROM muse_content_work WHERE id=%s",
|
||
(work_id,)).fetchone()[0]
|
||
# 预览统计:各表将被处理的行数
|
||
cards = conn.execute(
|
||
"""SELECT count(*) FROM muse_knowledge_draft
|
||
WHERE tenant_id=%s AND work_id=%s AND source_type=%s AND deleted=FALSE""",
|
||
(TENANT, work_id, SOURCE_TYPE)).fetchone()[0]
|
||
wins = conn.execute(
|
||
"""SELECT count(*), count(*) FILTER (WHERE status='done') FROM example_upgrade_window
|
||
WHERE tenant_id=%s AND work_id=%s AND deleted=FALSE""",
|
||
(TENANT, work_id)).fetchone()
|
||
ali = conn.execute("SELECT count(*) FROM example_upgrade_alias WHERE tenant_id=%s AND work_id=%s",
|
||
(TENANT, work_id)).fetchone()[0]
|
||
pres = conn.execute("SELECT count(*) FROM example_upgrade_presence WHERE tenant_id=%s AND work_id=%s",
|
||
(TENANT, work_id)).fetchone()[0]
|
||
# card_state 无 work 冗余错?——有 work_id 列(建表即有);audit 无 work_id,靠 JOIN 卡定位
|
||
cs = conn.execute("SELECT count(*) FROM example_upgrade_card_state WHERE tenant_id=%s AND work_id=%s",
|
||
(TENANT, work_id)).fetchone()[0]
|
||
aud = conn.execute(
|
||
"""SELECT count(*) FROM example_upgrade_audit a
|
||
WHERE a.tenant_id=%s AND a.draft_id IN
|
||
(SELECT id FROM muse_knowledge_draft
|
||
WHERE tenant_id=%s AND work_id=%s AND source_type=%s)""",
|
||
(TENANT, TENANT, work_id, SOURCE_TYPE)).fetchone()[0]
|
||
# 不限制 d.deleted:除本轮活卡外,也要收口旧版 reset 已软删卡遗留的活向量。
|
||
embeddings = conn.execute(
|
||
"""SELECT count(*) FROM example_knowledge_embedding e
|
||
WHERE e.tenant_id=%s AND e.deleted=FALSE AND e.entity_id IS NULL AND EXISTS (
|
||
SELECT 1 FROM muse_knowledge_draft d
|
||
WHERE d.id=e.draft_id AND d.tenant_id=%s
|
||
AND d.work_id=%s AND d.source_type=%s
|
||
)""",
|
||
(TENANT, TENANT, work_id, SOURCE_TYPE)).fetchone()[0]
|
||
click.echo(f"《{title}》work={work_id} 重抽准备{'(执行)' if execute else '(预览)'}:")
|
||
click.echo(f" 软删活升格卡 {cards} | 软删活向量 {embeddings} | "
|
||
f"重置窗 {wins[0]}(其中done {wins[1]}) | "
|
||
f"硬删 别名{ali} 留档{pres} 卡水位{cs} 审计{aud}")
|
||
if not execute:
|
||
click.echo(" (预览模式未写库;加 --execute 真清)")
|
||
return
|
||
# 目标 draft 只要曾绑定确认实体就失败关闭,避免把已确认实体向量降格为 reset 产物处理。
|
||
confirmed_vectors = conn.execute(
|
||
"""SELECT count(*) FROM example_knowledge_embedding e
|
||
WHERE e.tenant_id=%s AND e.entity_id IS NOT NULL AND EXISTS (
|
||
SELECT 1 FROM muse_knowledge_draft d
|
||
WHERE d.id=e.draft_id AND d.tenant_id=%s
|
||
AND d.work_id=%s AND d.source_type=%s
|
||
)""",
|
||
(TENANT, TENANT, work_id, SOURCE_TYPE)).fetchone()[0]
|
||
if confirmed_vectors:
|
||
raise click.ClickException(
|
||
f"目标升格 draft 存在 {confirmed_vectors} 条 entity_id 非空向量,拒绝 reset"
|
||
)
|
||
# 单事务执行:全清或全不清(中途失败自动回滚,不留半清状态)
|
||
conn.execute(
|
||
"""UPDATE muse_knowledge_draft SET deleted=TRUE, updater='upgrade-reset'
|
||
WHERE tenant_id=%s AND work_id=%s AND source_type=%s AND deleted=FALSE""",
|
||
(TENANT, work_id, SOURCE_TYPE))
|
||
# 只软删目标作品升格卡的活向量;双租户边界避免软引用串租户时误伤。
|
||
conn.execute(
|
||
"""UPDATE example_knowledge_embedding e SET deleted=TRUE, updater='upgrade-reset'
|
||
WHERE e.tenant_id=%s AND e.deleted=FALSE AND e.entity_id IS NULL AND EXISTS (
|
||
SELECT 1 FROM muse_knowledge_draft d
|
||
WHERE d.id=e.draft_id AND d.tenant_id=%s
|
||
AND d.work_id=%s AND d.source_type=%s
|
||
)""",
|
||
(TENANT, TENANT, work_id, SOURCE_TYPE))
|
||
conn.execute(
|
||
"""UPDATE example_upgrade_window SET status='pending', error_message=NULL,
|
||
updater='upgrade-reset' WHERE tenant_id=%s AND work_id=%s AND deleted=FALSE""",
|
||
(TENANT, work_id))
|
||
# 审计先删(靠卡 JOIN 定位,卡还查得到——虽然刚软删,JOIN 不看 deleted)
|
||
conn.execute(
|
||
"""DELETE FROM example_upgrade_audit
|
||
WHERE tenant_id=%s AND draft_id IN
|
||
(SELECT id FROM muse_knowledge_draft
|
||
WHERE tenant_id=%s AND work_id=%s AND source_type=%s)""",
|
||
(TENANT, TENANT, work_id, SOURCE_TYPE))
|
||
conn.execute("DELETE FROM example_upgrade_alias WHERE tenant_id=%s AND work_id=%s",
|
||
(TENANT, work_id))
|
||
conn.execute("DELETE FROM example_upgrade_presence WHERE tenant_id=%s AND work_id=%s",
|
||
(TENANT, work_id))
|
||
conn.execute("DELETE FROM example_upgrade_card_state WHERE tenant_id=%s AND work_id=%s",
|
||
(TENANT, work_id))
|
||
# 提交前复核:任一断言失败都由同一业务事务回滚,不留下半清状态。
|
||
left = conn.execute(
|
||
"""SELECT count(*) FROM muse_knowledge_draft
|
||
WHERE tenant_id=%s AND work_id=%s AND source_type=%s AND deleted=FALSE""",
|
||
(TENANT, work_id, SOURCE_TYPE)).fetchone()[0]
|
||
pend = conn.execute(
|
||
"""SELECT count(*) FILTER (WHERE status='pending'), count(*) FROM example_upgrade_window
|
||
WHERE tenant_id=%s AND work_id=%s AND deleted=FALSE""", (TENANT, work_id)).fetchone()
|
||
aliases_left = conn.execute(
|
||
"SELECT count(*) FROM example_upgrade_alias WHERE tenant_id=%s AND work_id=%s",
|
||
(TENANT, work_id)).fetchone()[0]
|
||
presence_left = conn.execute(
|
||
"SELECT count(*) FROM example_upgrade_presence WHERE tenant_id=%s AND work_id=%s",
|
||
(TENANT, work_id)).fetchone()[0]
|
||
state_left = conn.execute(
|
||
"SELECT count(*) FROM example_upgrade_card_state WHERE tenant_id=%s AND work_id=%s",
|
||
(TENANT, work_id)).fetchone()[0]
|
||
audit_left = conn.execute(
|
||
"""SELECT count(*) FROM example_upgrade_audit a
|
||
WHERE a.tenant_id=%s AND a.draft_id IN
|
||
(SELECT id FROM muse_knowledge_draft
|
||
WHERE tenant_id=%s AND work_id=%s AND source_type=%s)""",
|
||
(TENANT, TENANT, work_id, SOURCE_TYPE)).fetchone()[0]
|
||
vectors_left = conn.execute(
|
||
"""SELECT count(*) FROM example_knowledge_embedding e
|
||
WHERE e.tenant_id=%s AND e.deleted=FALSE AND EXISTS (
|
||
SELECT 1 FROM muse_knowledge_draft d
|
||
WHERE d.id=e.draft_id AND d.tenant_id=%s
|
||
AND d.work_id=%s AND d.source_type=%s
|
||
)""",
|
||
(TENANT, TENANT, work_id, SOURCE_TYPE)).fetchone()[0]
|
||
failures = []
|
||
if left:
|
||
failures.append(f"active升格卡={left}")
|
||
if pend[0] != pend[1]:
|
||
failures.append(f"windows pending={pend[0]}/{pend[1]}")
|
||
for label, value in (
|
||
("alias", aliases_left), ("presence", presence_left),
|
||
("state", state_left), ("audit", audit_left),
|
||
("目标draft活向量", vectors_left)):
|
||
if value:
|
||
failures.append(f"{label}={value}")
|
||
if failures:
|
||
raise click.ClickException(
|
||
"清后复核失败:" + ";".join(failures)
|
||
)
|
||
conn.commit()
|
||
click.echo(f" ✅ 七域清理断言全部通过:活升格卡0 | "
|
||
f"窗pending {pend[0]}/{pend[1]} | alias/presence/state/audit 0 | "
|
||
f"目标draft活向量0")
|
||
|
||
|
||
@click.command()
|
||
@click.option("--work-id", type=int, required=True)
|
||
@click.option("--execute", is_flag=True, help="真执行(默认只预览统计)")
|
||
@click.option("--backup-dir", type=click.Path(path_type=pathlib.Path, exists=True,
|
||
file_okay=False, resolve_path=False))
|
||
@click.option("--backup-id")
|
||
@click.option("--confirmation-sha")
|
||
def main(work_id, execute, backup_dir, backup_id, confirmation_sha):
|
||
"""先获取同书会话锁,再预览或执行升格全量重抽准备。"""
|
||
|
||
if execute:
|
||
_require_execute_backup(backup_dir, backup_id, confirmation_sha)
|
||
try:
|
||
with upgrade_work_lock(DSN, TENANT, work_id):
|
||
manifest = None
|
||
code_identity = None
|
||
if execute:
|
||
manifest = backup.verify_backup(backup_dir)
|
||
backup.validate_execute_confirmation(
|
||
True,
|
||
backup_id,
|
||
manifest["backup_id"],
|
||
confirmation_sha,
|
||
manifest["confirmationSha"],
|
||
)
|
||
if manifest.get("tenant") != TENANT or manifest.get("work") != work_id:
|
||
raise click.ClickException("备份 manifest 的 tenant/work 与 reset 目标不一致")
|
||
# 仍在同书 advisory lock 内,且尚未创建业务连接或执行任何 reset SQL。
|
||
code_identity = backup.capture_code_identity()
|
||
_assert_code_identity_matches_manifest(manifest, code_identity)
|
||
return _reset(work_id, execute, manifest, code_identity)
|
||
except UpgradeWorkLockUnavailable as exc:
|
||
raise click.ClickException(str(exc)) from exc
|
||
except click.ClickException:
|
||
raise
|
||
except Exception as exc:
|
||
raise click.ClickException(str(exc)) from exc
|
||
|
||
|
||
if __name__ == "__main__":
|
||
main()
|