db.py 99 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041
  1. # -*- coding: utf-8 -*-
  2. """mode_workflow · MySQL 持久化(DB 为唯一事实源)
  3. ================================================================================
  4. 读 .env 的 MYSQL_* 连接 MySQL。四张表:
  5. search_process —— 每行一个 (query, 帖子):工序方向的搜索 + llm 评估结果
  6. search_tools —— 同结构,工具方向的搜索结果(方向由表区分,不再用 mode_type 列)
  7. mode_process —— 每行一个解构出的工序(steps 等嵌套结构存 JSON 列)
  8. mode_tools —— 每行一个解构出的工具
  9. 与旧 fixed_query_eval/db.py 的关键差异:本系统 DB 是主存储,写入失败直接 raise,
  10. 不做"失败不阻断"。读侧保留防御(返回空/None)。
  11. 用法:
  12. python db.py init # 建表(幂等)
  13. python db.py check # 打印四表行数
  14. python db.py clear # 清空四表数据(TRUNCATE)
  15. """
  16. import json
  17. import os
  18. import sys
  19. from datetime import datetime
  20. from pathlib import Path
  21. PROJECT_ROOT = Path(__file__).resolve().parents[2]
  22. sys.path.insert(0, str(PROJECT_ROOT))
  23. from dotenv import load_dotenv
  24. load_dotenv()
  25. import pymysql
  26. from pymysql.cursors import DictCursor
  27. from dbutils.pooled_db import PooledDB
  28. # ── 连接池 ──────────────────────────────────────────────────────────────────
  29. # MySQL 是远程 RDS,每次 pymysql.connect() 的 TCP+鉴权握手 ~0.5s。旧实现每个
  30. # 请求新建一条连接,一次"点开帖子"要 2~3 个请求 = 2~3 次握手 ≈ 1s。改用连接池
  31. # 复用长连接后,握手只在池初始化时各发生一次,后续取连接近乎零开销。
  32. # server.py 是 ThreadingHTTPServer(每请求一线程),PooledDB 线程安全,正好匹配。
  33. # 注意:fetch_* 里的 conn.close() 在池连接上语义是"归还池中"而非真正断开。
  34. _POOL = None
  35. def _pool():
  36. global _POOL
  37. if _POOL is None:
  38. if not os.getenv("MYSQL_HOST"):
  39. raise RuntimeError("缺 MYSQL_HOST:检查 .env 的 MYSQL_* 配置")
  40. _POOL = PooledDB(
  41. creator=pymysql,
  42. mincached=2, # 启动即预热 2 条,首点不再吃冷握手
  43. maxcached=5, # 空闲保留上限
  44. maxconnections=20, # 并发上限(ThreadingHTTPServer 线程数)
  45. blocking=True, # 连接耗尽时等待而非报错
  46. ping=1, # 取用前 ping,自动剔除被 RDS 掐断的死连接
  47. host=os.getenv("MYSQL_HOST"),
  48. port=int(os.getenv("MYSQL_PORT", 3306)),
  49. user=os.getenv("MYSQL_USER"),
  50. password=os.getenv("MYSQL_PASSWORD"),
  51. database=os.getenv("MYSQL_DATABASE"),
  52. charset="utf8mb4", cursorclass=DictCursor,
  53. autocommit=True, connect_timeout=10,
  54. )
  55. return _POOL
  56. def _conn():
  57. """从池取一条连接;用法不变(with cursor / conn.close() 归还池)。"""
  58. return _pool().connection()
  59. # ── DDL ──────────────────────────────────────────────────────────────────────
  60. SEARCH_TABLES = {"process": "search_process", "tools": "search_tools"}
  61. MODE_TABLES = {"process": "mode_process", "tools": "mode_tools"}
  62. def _search_table(mode_or_table):
  63. """mode(process/tools)或表名 → 合法搜索表名(白名单,防 SQL 注入)。"""
  64. t = SEARCH_TABLES.get(mode_or_table, mode_or_table)
  65. if t not in SEARCH_TABLES.values():
  66. raise ValueError(f"未知搜索表/模式: {mode_or_table!r}")
  67. return t
  68. def _mode_table(mode_or_table):
  69. """mode(process/tools)或表名 → 合法解构表名(白名单,防 SQL 注入)。"""
  70. t = MODE_TABLES.get(mode_or_table, mode_or_table)
  71. if t not in MODE_TABLES.values():
  72. raise ValueError(f"未知解构表/模式: {mode_or_table!r}")
  73. return t
  74. def _ddl_search(table, direction):
  75. return f"""
  76. CREATE TABLE IF NOT EXISTS {table} (
  77. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  78. query_id VARCHAR(32) NOT NULL COMMENT 'q0000',
  79. query_text VARCHAR(512) NULL,
  80. case_id VARCHAR(128) NOT NULL COMMENT 'platform_channelContentId',
  81. platform VARCHAR(32) NULL,
  82. channel_content_id VARCHAR(128) NULL,
  83. title VARCHAR(512) NULL,
  84. url VARCHAR(1024) NULL,
  85. content_type VARCHAR(32) NULL,
  86. body LONGTEXT NULL,
  87. images JSON NULL,
  88. videos JSON NULL,
  89. like_count INT NULL,
  90. publish_time VARCHAR(64) NULL,
  91. quality_score FLOAT NULL COMMENT 'post._quality_score',
  92. quality_grade VARCHAR(8) NULL,
  93. found_by JSON NULL COMMENT '命中的措辞数组',
  94. knowledge_type JSON NULL COMMENT '["能力","工序","工具"] 子集',
  95. overall_score FLOAT NULL COMMENT '(相关均值+质量均值)/2',
  96. llm_evaluation JSON NULL COMMENT '评估全量 blob',
  97. created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
  98. updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  99. UNIQUE KEY uk_qid_case (query_id, case_id),
  100. KEY idx_platform (platform)
  101. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='搜索+评估结果({direction})';
  102. """
  103. DDL_PROCESS = """
  104. CREATE TABLE IF NOT EXISTS mode_process (
  105. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  106. query_id VARCHAR(32) NOT NULL,
  107. case_id VARCHAR(128) NOT NULL,
  108. platform VARCHAR(32) NULL,
  109. post_title VARCHAR(512) NULL,
  110. source JSON NULL COMMENT '解构返回的 source 块',
  111. procedure_id VARCHAR(16) NULL COMMENT 'p1,p2…',
  112. name VARCHAR(255) NULL,
  113. purpose TEXT NULL,
  114. category VARCHAR(32) NULL COMMENT '产物创造/资产建设/自动化/分析/学习',
  115. declarations JSON NULL,
  116. type_registry JSON NULL,
  117. steps JSON NULL COMMENT '步骤数组全量',
  118. step_count INT NULL,
  119. tools_used JSON NULL COMMENT '从 steps[].via 去重提取',
  120. model VARCHAR(64) NULL,
  121. version VARCHAR(32) NULL COMMENT 'v_MMDDHHMM,保留历史;link_* 为跨 query 复制(cost=0)',
  122. cost_usd DECIMAL(10,6) NULL COMMENT '本次解构调用成本(同版本各行相同,聚合需按 case+version 去重)',
  123. duration_s FLOAT NULL,
  124. seq SMALLINT NULL COMMENT '帖内序号(0-based);与 (query_id,case_id,version) 组唯一键防并发/重复写',
  125. created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
  126. UNIQUE KEY uk_q_case_ver_seq (query_id, case_id, version, seq),
  127. KEY idx_case_ver (case_id, version),
  128. KEY idx_qid (query_id)
  129. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='工序解构结果(每行一个工序)';
  130. """
  131. DDL_TOOLS = """
  132. CREATE TABLE IF NOT EXISTS mode_tools (
  133. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  134. query_id VARCHAR(32) NOT NULL,
  135. case_id VARCHAR(128) NOT NULL,
  136. platform VARCHAR(32) NULL,
  137. post_title VARCHAR(512) NULL,
  138. source JSON NULL COMMENT '解构时帖子来源块(tool_extract._row_to_source 产出)',
  139. tool_name VARCHAR(255) NULL,
  140. substance_scope JSON NULL COMMENT '实质作用域(数组)',
  141. form_scope JSON NULL COMMENT '形式作用域(数组或null)',
  142. creation_layer VARCHAR(32) NULL COMMENT '制作层/创作层',
  143. source_link VARCHAR(1024) NULL,
  144. input_desc TEXT NULL,
  145. output_desc TEXT NULL,
  146. usage_json JSON NULL,
  147. cases_json JSON NULL,
  148. defects_json JSON NULL,
  149. updated_time VARCHAR(64) NULL COMMENT '工具最新更新时间',
  150. model VARCHAR(64) NULL,
  151. version VARCHAR(32) NULL COMMENT 'v_MMDDHHMM;link_* 为跨 query 复制(cost=0)',
  152. cost_usd DECIMAL(10,6) NULL COMMENT '同 mode_process,聚合按 case+version 去重',
  153. duration_s FLOAT NULL,
  154. seq SMALLINT NULL COMMENT '帖内序号(0-based);与 (query_id,case_id,version) 组唯一键防并发/重复写',
  155. created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
  156. UNIQUE KEY uk_q_case_ver_seq (query_id, case_id, version, seq),
  157. KEY idx_case_ver (case_id, version),
  158. KEY idx_qid (query_id),
  159. KEY idx_tool_name (tool_name)
  160. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='工具解构结果(每行一个工具)';
  161. """
  162. # 工序知识「已导入知识库」台账:防重复上传(stages/import_process_knowledge.py 用)。
  163. # 每条知识 = 某 case 的某个工序(proc_index 1-based)。记录导入时的 mode_process 版本:
  164. # 版本变了(重解构)说明内容已变,应重导;版本不变即视为「已传过」,跳过。
  165. # 选 DB 台账而非本地文件,是为了换机器/换链接后也不会重复写知识库。
  166. # 注:工具知识用独立的 tools_ingest_log,不与本表混用(case_id 是帖子物理身份,
  167. # 同帖可能既被工序解构又被工具解构,共表会在 (case_id, index) 上撞键)。
  168. DDL_INGEST_LOG = """
  169. CREATE TABLE IF NOT EXISTS knowledge_ingest_log (
  170. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  171. case_id VARCHAR(128) NOT NULL,
  172. proc_index INT NOT NULL COMMENT '工序序号(1-based),对齐导入脚本枚举',
  173. version VARCHAR(32) NULL COMMENT '导入时 mode_process 版本;变了应重导',
  174. knowledge_id VARCHAR(128) NULL COMMENT '接口返回的 knowledge_id',
  175. api_url VARCHAR(255) NULL,
  176. ingested_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
  177. updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  178. UNIQUE KEY uk_case_proc (case_id, proc_index)
  179. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='工序知识已导入台账(防重复上传)';
  180. """
  181. # 工具知识「已导入知识库」台账:语义同 knowledge_ingest_log,但针对工具方向独立成表
  182. # (stages/import_tools_knowledge.py 用)。每条知识 = 某 case 的某个工具(tool_index 1-based),
  183. # 版本记录导入时的 mode_tools 版本;变了(重解构)应重导,不变即「已传过」跳过。
  184. DDL_TOOLS_INGEST_LOG = """
  185. CREATE TABLE IF NOT EXISTS tools_ingest_log (
  186. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  187. case_id VARCHAR(128) NOT NULL,
  188. tool_index INT NOT NULL COMMENT '工具序号(1-based),对齐导入脚本枚举',
  189. version VARCHAR(32) NULL COMMENT '导入时 mode_tools 版本;变了应重导',
  190. knowledge_id VARCHAR(128) NULL COMMENT '接口返回的 knowledge_id',
  191. api_url VARCHAR(255) NULL,
  192. ingested_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
  193. updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  194. UNIQUE KEY uk_case_tool (case_id, tool_index)
  195. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='工具知识已导入台账(防重复上传)';
  196. """
  197. # 「解构知识(工具·仅实质/形式作用域 → what-制作还原)」已导入台账。
  198. # 与 tools_ingest_log 分表:两者同源 mode_tools、同键 (case_id, tool_index),但上传到不同维度
  199. # (本表 = what-制作还原,tools_ingest_log = how工具),共表会在 (case_id, tool_index) 撞键,
  200. # 使一方上传后另一方被误判「已传过」而跳过。故独立成表(stages/import_destruction_knowledge.py 用)。
  201. DDL_DESTRUCTION_INGEST_LOG = """
  202. CREATE TABLE IF NOT EXISTS destruction_ingest_log (
  203. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  204. case_id VARCHAR(128) NOT NULL,
  205. tool_index INT NOT NULL COMMENT '工具序号(1-based),对齐导入脚本枚举',
  206. version VARCHAR(32) NULL COMMENT '导入时 mode_tools 版本;变了应重导',
  207. knowledge_id VARCHAR(128) NULL COMMENT '接口返回的 knowledge_id',
  208. api_url VARCHAR(255) NULL,
  209. ingested_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
  210. updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  211. UNIQUE KEY uk_case_tool (case_id, tool_index)
  212. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='解构知识(what-制作还原)已导入台账(防重复上传)';
  213. """
  214. # step 维度归类结果独立表(3-I,见 docs/step_classification_3-I_设计.md):一行 = 某 case 某版本
  215. # 某工序某 step 某维度某子项的命中分类。与 mode_process 解耦,归类重跑只覆盖本表、不重写 steps blob;
  216. # 存 matched_path 全路径,使「某节点及其子树的所有帖」一句 SQL(LIKE 前缀)即可,且与显示同源。
  217. DDL_STEP_CLASSIFICATION = """
  218. CREATE TABLE IF NOT EXISTS step_classification (
  219. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  220. case_id VARCHAR(128) NOT NULL,
  221. query_id VARCHAR(32) NULL COMMENT 'record 时的 post_id,便于溯源',
  222. version VARCHAR(32) NULL COMMENT '对齐 mode_process 最新真实版,重跑覆盖',
  223. procedure_id VARCHAR(16) NULL COMMENT 'p1,p2…',
  224. step_id VARCHAR(16) NULL COMMENT 's1,s2…',
  225. dimension VARCHAR(8) NOT NULL COMMENT '实质/形式/类型/作用/动作',
  226. sub_index SMALLINT NOT NULL COMMENT '原值「、」拆分后的子项下标',
  227. raw_term VARCHAR(255) NULL COMMENT '原始词',
  228. matched_name VARCHAR(255) NULL COMMENT '命中分类名(单一最优;无命中不入库)',
  229. matched_path VARCHAR(512) NULL COMMENT '命中分类全路径,如 /表象/视觉/空间/空间环境',
  230. matched_id INT NULL COMMENT 'cat-api stable_id',
  231. score FLOAT NULL,
  232. match_run_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  233. UNIQUE KEY uk_step_dim_sub (case_id, version, procedure_id, step_id, dimension, sub_index),
  234. KEY idx_case (case_id),
  235. KEY idx_dim_name (dimension, matched_name),
  236. KEY idx_dim_path (dimension, matched_path)
  237. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='step 维度归类结果(与 mode_process 解耦)';
  238. """
  239. # 工序解构拆表(3-II,见 docs/拆表_3-II_设计.md):process_run → process_procedure → process_step 三层,
  240. # 与旧 mode_process 并行。step_json 逐字保留原 step → ProcessPayload 还原零漂移;抽取列供 DB 层查询。
  241. DDL_PROCESS_RUN = """
  242. CREATE TABLE IF NOT EXISTS process_run (
  243. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  244. query_id VARCHAR(32) NOT NULL,
  245. case_id VARCHAR(128) NOT NULL,
  246. version VARCHAR(32) NOT NULL,
  247. is_latest TINYINT NOT NULL DEFAULT 0 COMMENT '该 case 最新真实版=1(link_ 恒 0)',
  248. platform VARCHAR(32) NULL,
  249. post_title VARCHAR(512) NULL,
  250. source JSON NULL,
  251. model VARCHAR(64) NULL,
  252. cost_usd DECIMAL(10,6) NULL,
  253. duration_s FLOAT NULL,
  254. created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
  255. UNIQUE KEY uk_q_case_ver (query_id, case_id, version),
  256. KEY idx_case_latest (case_id, is_latest),
  257. KEY idx_qid (query_id)
  258. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='工序解构运行级(一行=一帖一版本)';
  259. """
  260. DDL_PROCESS_PROCEDURE = """
  261. CREATE TABLE IF NOT EXISTS process_procedure (
  262. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  263. run_id BIGINT NOT NULL,
  264. case_id VARCHAR(128) NOT NULL,
  265. version VARCHAR(32) NOT NULL,
  266. procedure_id VARCHAR(16) NULL,
  267. seq SMALLINT NOT NULL,
  268. name VARCHAR(255) NULL,
  269. purpose TEXT NULL,
  270. category VARCHAR(32) NULL,
  271. declarations JSON NULL,
  272. type_registry JSON NULL,
  273. tools_used JSON NULL,
  274. step_count INT NULL,
  275. UNIQUE KEY uk_run_seq (run_id, seq),
  276. KEY idx_run (run_id),
  277. KEY idx_case_ver (case_id, version)
  278. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='工序解构工序级(一行=一个 procedure)';
  279. """
  280. DDL_PROCESS_STEP = """
  281. CREATE TABLE IF NOT EXISTS process_step (
  282. id BIGINT AUTO_INCREMENT PRIMARY KEY,
  283. procedure_pk BIGINT NOT NULL,
  284. case_id VARCHAR(128) NOT NULL,
  285. version VARCHAR(32) NOT NULL,
  286. step_id VARCHAR(16) NULL,
  287. seq SMALLINT NOT NULL,
  288. kind VARCHAR(16) NULL,
  289. grp VARCHAR(16) NULL COMMENT '原 step.group(group 为保留字)',
  290. via VARCHAR(255) NULL,
  291. effect VARCHAR(255) NULL,
  292. action VARCHAR(255) NULL,
  293. substance VARCHAR(512) NULL,
  294. form VARCHAR(512) NULL,
  295. step_json JSON NOT NULL COMMENT '原 step 逐字 JSON,payload 还原即用它',
  296. UNIQUE KEY uk_proc_seq (procedure_pk, seq),
  297. KEY idx_procedure (procedure_pk),
  298. KEY idx_case_ver (case_id, version),
  299. KEY idx_action (action), KEY idx_substance (substance), KEY idx_form (form)
  300. ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='工序解构步骤级(一行=一个 step)';
  301. """
  302. def _ensure_column(cur, table, column, column_ddl):
  303. """给已存在的表幂等补列:列已存在则跳过(MySQL ADD COLUMN 无 IF NOT EXISTS)。
  304. column_ddl 为 ADD COLUMN 后的完整定义,如 \"source JSON NULL ... AFTER post_title\"。"""
  305. cur.execute("""SELECT COUNT(*) AS n FROM information_schema.columns
  306. WHERE table_schema=DATABASE() AND table_name=%s AND column_name=%s""",
  307. (table, column))
  308. if cur.fetchone()["n"] == 0:
  309. cur.execute(f"ALTER TABLE {table} ADD COLUMN {column_ddl}")
  310. def _ensure_unique_index(cur, table, index_name, cols):
  311. """幂等加唯一索引:已存在则跳过(MySQL ADD INDEX 无 IF NOT EXISTS)。
  312. cols 为列表达式,如 "query_id, case_id, version, seq"。加之前需保证无冲突数据。"""
  313. cur.execute("""SELECT COUNT(*) AS n FROM information_schema.statistics
  314. WHERE table_schema=DATABASE() AND table_name=%s AND index_name=%s""",
  315. (table, index_name))
  316. if cur.fetchone()["n"] == 0:
  317. cur.execute(f"ALTER TABLE {table} ADD UNIQUE KEY {index_name} ({cols})")
  318. def init_tables():
  319. conn = _conn()
  320. try:
  321. with conn.cursor() as cur:
  322. cur.execute(_ddl_search("search_process", "工序方向"))
  323. cur.execute(_ddl_search("search_tools", "工具方向"))
  324. cur.execute(DDL_PROCESS)
  325. cur.execute(DDL_TOOLS)
  326. cur.execute(DDL_INGEST_LOG)
  327. cur.execute(DDL_TOOLS_INGEST_LOG)
  328. cur.execute(DDL_DESTRUCTION_INGEST_LOG)
  329. cur.execute(DDL_STEP_CLASSIFICATION)
  330. cur.execute(DDL_PROCESS_RUN)
  331. cur.execute(DDL_PROCESS_PROCEDURE)
  332. cur.execute(DDL_PROCESS_STEP)
  333. # 历史库迁移:version 由 VARCHAR(16) 放宽到 32,容纳 link_v_mopN_* 复制版本。
  334. # MODIFY 幂等(已是 32 则 MySQL 元数据无操作),建表后表必存在,可安全执行。
  335. for t in ("mode_process", "mode_tools"):
  336. cur.execute(f"ALTER TABLE {t} MODIFY COLUMN version VARCHAR(32) NULL")
  337. # 历史库迁移:给老 mode_tools 补 source 列(MySQL 的 ADD COLUMN 无 IF NOT EXISTS,
  338. # 故先查 information_schema 判存在,缺了才 ADD,幂等)。
  339. _ensure_column(cur, "mode_tools", "source",
  340. "source JSON NULL COMMENT '解构时帖子来源块' AFTER post_title")
  341. # 历史库迁移:加 seq(帖内序号)+ (query_id,case_id,version,seq) 唯一键,防并发/重复
  342. # 写入产生重复行。顺序必须是 加列 → 回填 → 加唯一键。MySQL 5.7 无窗口函数,seq 在
  343. # 应用层按 (query_id,case_id,version) 内 id 升序回填(现有数据该粒度已无重复)。
  344. for t in ("mode_process", "mode_tools"):
  345. _ensure_column(cur, t, "seq",
  346. "seq SMALLINT NULL COMMENT '帖内序号(0-based)' AFTER duration_s")
  347. for t in ("mode_process", "mode_tools"):
  348. cur.execute(f"""SELECT id, query_id, case_id, version FROM {t}
  349. WHERE seq IS NULL ORDER BY query_id, case_id, version, id""")
  350. key, n, ups = None, 0, []
  351. for r in cur.fetchall():
  352. k = (r["query_id"], r["case_id"], r["version"])
  353. if k != key:
  354. key, n = k, 0
  355. ups.append((n, r["id"])); n += 1
  356. if ups:
  357. cur.executemany(f"UPDATE {t} SET seq=%s WHERE id=%s", ups)
  358. print(f" ↳ {t}: 回填 seq {len(ups)} 行")
  359. for t in ("mode_process", "mode_tools"):
  360. _ensure_unique_index(cur, t, "uk_q_case_ver_seq",
  361. "query_id, case_id, version, seq")
  362. print("✅ 建表完成:search_process, search_tools, mode_process, mode_tools, "
  363. "knowledge_ingest_log, tools_ingest_log, step_classification")
  364. finally:
  365. conn.close()
  366. def clear_tables():
  367. """清空四张表的数据(TRUNCATE,表结构保留)。"""
  368. conn = _conn()
  369. try:
  370. with conn.cursor() as cur:
  371. for t in ("search_process", "search_tools", "mode_process", "mode_tools",
  372. "process_run", "process_procedure", "process_step", "step_classification"):
  373. cur.execute(f"TRUNCATE TABLE {t}")
  374. print(f"🧹 已清空 {t}")
  375. finally:
  376. conn.close()
  377. # ── 工具函数 ──────────────────────────────────────────────────────────────────
  378. def _loads(v, default=None):
  379. """pymysql 的 JSON 列可能返回字符串,统一解析。"""
  380. if v is None:
  381. return default
  382. if isinstance(v, (list, dict)):
  383. return v
  384. try:
  385. return json.loads(v)
  386. except Exception:
  387. return default
  388. def _j(v):
  389. """写入 JSON 列:None 保持 NULL,其余 dumps。"""
  390. return None if v is None else json.dumps(v, ensure_ascii=False)
  391. def _collect_scores(node):
  392. """递归收集嵌套评估里所有「得分」。LLM 直出的得分多为字符串("1"/"4"),
  393. 个别为数字(如 时效性 10),统一按 float 解析;非数值(如 "N/A")跳过不计入。"""
  394. out = []
  395. if isinstance(node, dict):
  396. for k, v in node.items():
  397. if k == "得分":
  398. try:
  399. out.append(float(v))
  400. except (TypeError, ValueError):
  401. pass
  402. else:
  403. out.extend(_collect_scores(v))
  404. elif isinstance(node, list):
  405. for v in node:
  406. out.extend(_collect_scores(v))
  407. return out
  408. def overall_score(e):
  409. """综合分 = (相关性各项均值 + 质量各项均值) / 可得部分数。算不出返回 None。"""
  410. parts = []
  411. for key in ("相关性", "质量"):
  412. scores = _collect_scores((e or {}).get(key))
  413. if scores:
  414. parts.append(sum(scores) / len(scores))
  415. return round(sum(parts) / len(parts), 2) if parts else None
  416. def _recency_hard(date_str):
  417. """硬时效(同 mode_procedure/server.py:_recency_hard):半年内=3 / 两年内=2 / 更早=1。
  418. publish_time 头 10 字符按 YYYY-MM-DD 解析,失败返回 None(不参与判定)。"""
  419. try:
  420. d = datetime.strptime(str(date_str or "")[:10], "%Y-%m-%d")
  421. except (ValueError, TypeError):
  422. return None
  423. days = (datetime.now() - d).days
  424. if days <= 180:
  425. return 3
  426. if days <= 730:
  427. return 2
  428. return 1
  429. def _fixed_dim_score(evaluation, name):
  430. """取 质量.固定维度.<name>.得分 标量,缺失/非数值返回 None(不参与判定)。"""
  431. v = (((evaluation or {}).get("质量") or {}).get("固定维度") or {}).get(name)
  432. if isinstance(v, dict):
  433. v = v.get("得分")
  434. try:
  435. return float(v) if v is not None else None
  436. except (TypeError, ValueError):
  437. return None
  438. def _impl_score(evaluation):
  439. """取 质量.动态维度.工序.字段完整性.实现完整性.得分 标量,缺失/非数值返回 None。
  440. 新版 prompt 把旧「可复现性」的硬封顶规则并入了「实现完整性」,故采纳门槛改读此处。"""
  441. v = ((((((evaluation or {}).get("质量") or {}).get("动态维度") or {})
  442. .get("工序") or {}).get("字段完整性") or {}).get("实现完整性"))
  443. if isinstance(v, dict):
  444. v = v.get("得分")
  445. try:
  446. return float(v) if v is not None else None
  447. except (TypeError, ValueError):
  448. return None
  449. def _repro_score(evaluation):
  450. """采纳门槛用的「可复现/可实现」得分:优先旧版「可复现性」(固定维度),
  451. 缺失则回退新版「实现完整性」(动态维度.工序)。这样新旧两套评估 blob 都能正确判定。"""
  452. v = _fixed_dim_score(evaluation, "可复现性")
  453. return v if v is not None else _impl_score(evaluation)
  454. def is_adopted(overall, evaluation, publish_time):
  455. """采纳/命中判定,口径对齐 mode_procedure 的 decision=="report":
  456. 制作相关性<4、可复现/实现完整性<4、发布超两年、综合分<6 —— 任一命中即不采纳;指标缺失不参与判定。
  457. (意图可控性暂只采分不设门槛,留待阈值标定后再开。)
  458. 可复现/实现门槛兼容新旧 schema:旧版读「可复现性」,新版读「实现完整性」(见 _repro_score)。
  459. fail-closed:评估失败(_error)、blob 缺失/为空、或综合分算不出(None)→ 直接判不采纳。
  460. 评不出的帖子不该混进命中集(此前 fail-open 会因各指标取不到值而误判采纳)。"""
  461. if not isinstance(evaluation, dict) or not evaluation or evaluation.get("_error"):
  462. return False
  463. if overall is None:
  464. return False
  465. rel = None
  466. v = ((evaluation or {}).get("相关性") or {}).get("和内容制作知识相关")
  467. if isinstance(v, dict):
  468. v = v.get("得分")
  469. try:
  470. rel = float(v) if v is not None else None
  471. except (TypeError, ValueError):
  472. rel = None
  473. if rel is not None and rel < 4:
  474. return False
  475. repro = _repro_score(evaluation)
  476. if repro is not None and repro < 4:
  477. return False
  478. rh = _recency_hard(publish_time)
  479. if rh is not None and rh < 2:
  480. return False
  481. if overall is not None and float(overall) < 6:
  482. return False
  483. return True
  484. def is_adopted_rel(overall, rel, publish_time, repro=None):
  485. """is_adopted 的轻量版:相关性得分(rel)、可复现/实现门槛(repro)已由 SQL JSON_EXTRACT
  486. 直接取出(repro 由 _REPRO_SQL 兼容新旧 schema 取值),无需传输/解析整块 llm_evaluation。
  487. 判定口径与 is_adopted 完全一致(含 fail-closed:综合分算不出→不采纳;失败帖的 overall_score 列为 NULL)。"""
  488. if overall is None:
  489. return False
  490. try:
  491. rel = float(rel) if rel is not None else None
  492. except (TypeError, ValueError):
  493. rel = None
  494. if rel is not None and rel < 4:
  495. return False
  496. try:
  497. repro = float(repro) if repro is not None else None
  498. except (TypeError, ValueError):
  499. repro = None
  500. if repro is not None and repro < 4:
  501. return False
  502. rh = _recency_hard(publish_time)
  503. if rh is not None and rh < 2:
  504. return False
  505. if overall is not None and float(overall) < 6:
  506. return False
  507. return True
  508. # ── search_process / search_tools ────────────────────────────────────────────
  509. def upsert_search_posts(query_id, query_text, results, table="search_process"):
  510. """一组搜索结果写入指定搜索表(按 (query_id, case_id) upsert)。返回写入条数。
  511. table:search_process(工序方向) / search_tools(工具方向)。"""
  512. table = _search_table(table)
  513. if not results:
  514. return 0
  515. rows = []
  516. for r in results:
  517. post = r.get("post") or {}
  518. e = r.get("llm_evaluation") or {}
  519. rows.append((
  520. query_id, query_text, r.get("case_id"), r.get("platform"),
  521. r.get("channel_content_id"),
  522. (post.get("title") or post.get("desc") or "")[:500],
  523. r.get("source_url"), post.get("content_type"),
  524. post.get("body_text") or post.get("desc") or "",
  525. _j(post.get("images") or []), _j(post.get("videos") or []),
  526. post.get("like_count"),
  527. str(post.get("publish_time") or post.get("publish_timestamp") or "")[:64],
  528. post.get("_quality_score"), post.get("_quality_grade"),
  529. _j(r.get("found_by_queries") or []),
  530. _j(e.get("知识类型") or []),
  531. overall_score(e),
  532. _j(e),
  533. ))
  534. sql = f"""
  535. INSERT INTO {table}
  536. (query_id, query_text, case_id, platform, channel_content_id, title, url,
  537. content_type, body, images, videos, like_count, publish_time,
  538. quality_score, quality_grade, found_by, knowledge_type,
  539. overall_score, llm_evaluation)
  540. VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)
  541. ON DUPLICATE KEY UPDATE
  542. query_text=VALUES(query_text), platform=VALUES(platform),
  543. channel_content_id=VALUES(channel_content_id), title=VALUES(title), url=VALUES(url),
  544. content_type=VALUES(content_type), body=VALUES(body), images=VALUES(images),
  545. videos=VALUES(videos), like_count=VALUES(like_count), publish_time=VALUES(publish_time),
  546. quality_score=VALUES(quality_score), quality_grade=VALUES(quality_grade),
  547. found_by=VALUES(found_by), knowledge_type=VALUES(knowledge_type),
  548. overall_score=VALUES(overall_score), llm_evaluation=VALUES(llm_evaluation);
  549. """
  550. conn = _conn()
  551. try:
  552. with conn.cursor() as cur:
  553. cur.executemany(sql, rows)
  554. return len(rows)
  555. finally:
  556. conn.close()
  557. # 占位帖 case_id:query 列表由 search_process 按 query_id 聚合得出(无独立 query 主表),
  558. # 一个 query 要进列表必须至少有一行。为支持「只登记 query、不触发搜索」,给这类 query 写
  559. # 一行哨兵帖,只承载 query_id+query_text。该哨兵行不属于任何真实帖子,故所有「帖子视图 /
  560. # 统计」读取点都用 _REAL_POST 过滤掉它(fetch_queries 的 post_count、fetch_posts、
  561. # fetch_all_posts、count_executed_queries、fetch_dashboard_rows)。真搜不会用到此 case_id。
  562. PENDING_CASE_ID = "__pending__"
  563. _REAL_POST = f"case_id <> '{PENDING_CASE_ID}'"
  564. def add_pending_process_queries(texts):
  565. """把一批 query 词作为「占位 query」加入工序 query 列表(search_process),不触发搜索/解构。
  566. 每条新增写一行哨兵帖(case_id=PENDING_CASE_ID,只填 query_id/query_text)。
  567. 去重:① 文件内重复保序去重;② query_text 已存在于 search_process(含此前占位)则跳过。
  568. query_id 跨 process/tools 统一续号,避免与工具方向撞号。返回 (added, skipped)。"""
  569. seen, cleaned = set(), []
  570. for t in texts:
  571. t = (t or "").strip()
  572. if t and t not in seen:
  573. seen.add(t)
  574. cleaned.append(t)
  575. conn = _conn()
  576. try:
  577. with conn.cursor() as cur:
  578. cur.execute("SELECT DISTINCT query_text FROM search_process WHERE query_text IS NOT NULL")
  579. existing = {r["query_text"] for r in cur.fetchall()}
  580. cur.execute("SELECT query_id FROM search_process "
  581. "UNION SELECT query_id FROM search_tools")
  582. nums = [int(r["query_id"][1:]) for r in cur.fetchall()
  583. if r["query_id"] and r["query_id"].startswith("q") and r["query_id"][1:].isdigit()]
  584. nxt = (max(nums) + 1) if nums else 0
  585. rows = []
  586. for t in cleaned:
  587. if t in existing:
  588. continue
  589. rows.append((f"q{nxt:04d}", t, PENDING_CASE_ID))
  590. nxt += 1
  591. if rows:
  592. cur.executemany(
  593. "INSERT INTO search_process (query_id, query_text, case_id) "
  594. "VALUES (%s,%s,%s)", rows)
  595. return len(rows), len(cleaned) - len(rows)
  596. finally:
  597. conn.close()
  598. def fetch_queries(mode="process"):
  599. """某方向搜索表的 query 列表 + 帖子数 + 采纳/命中数 + 解构进度。"""
  600. table = _search_table(mode)
  601. conn = _conn()
  602. try:
  603. with conn.cursor() as cur:
  604. # post_count 只数真实帖,占位哨兵行不计(占位 query 显示为 0 帖);
  605. # GROUP BY 仍含占位 query_id,故无搜索的 query 也会出现在列表里。
  606. cur.execute(f"""SELECT query_id, MAX(query_text) AS query_text,
  607. COUNT(CASE WHEN {_REAL_POST} THEN 1 END) AS post_count
  608. FROM {table} GROUP BY query_id ORDER BY query_id""")
  609. queries = cur.fetchall()
  610. # 采纳数:SQL 直取 rel/repro 标量算,**不拉整表 llm_evaluation**(旧版全表 blob,切 tab 巨慢)
  611. cur.execute(f"""SELECT query_id, overall_score, publish_time,
  612. {_REL_SQL} AS rel, {_REPRO_SQL} AS repro FROM {table}""")
  613. hits = {}
  614. for r in cur.fetchall():
  615. if is_adopted_rel(r["overall_score"], r["rel"], r["publish_time"], r["repro"]):
  616. hits[r["query_id"]] = hits.get(r["query_id"], 0) + 1
  617. cur.execute("SELECT query_id, COUNT(DISTINCT case_id) AS n FROM process_run GROUP BY query_id") # 3-II
  618. np = {r["query_id"]: r["n"] for r in cur.fetchall()}
  619. cur.execute("SELECT query_id, COUNT(DISTINCT case_id) AS n FROM mode_tools GROUP BY query_id")
  620. nt = {r["query_id"]: r["n"] for r in cur.fetchall()}
  621. finally:
  622. conn.close()
  623. for q in queries:
  624. q["hit_count"] = hits.get(q["query_id"], 0)
  625. q["process_done"] = np.get(q["query_id"], 0)
  626. q["tools_done"] = nt.get(q["query_id"], 0)
  627. return queries
  628. def fetch_posts(query_id, mode="process"):
  629. """列表用:只取列表所需列 + SQL 直取 adopted 标量,**不拉 body/videos/llm_evaluation 大字段**
  630. (llm_evaluation ~1.5MB/帖,旧版 SELECT * 切 tab/选 query 要几十 MB 过远程 RDS,故慢)。
  631. 正文/评分等详情按需走 fetch_post。带 adopted/has_process/has_tools;adopted 口径用
  632. is_adopted_rel(与 is_adopted 完全一致,rel/repro 由 _REL_SQL/_REPRO_SQL 直取标量)。"""
  633. table = _search_table(mode)
  634. conn = _conn()
  635. try:
  636. with conn.cursor() as cur:
  637. cur.execute(f"""SELECT id, query_id, query_text, case_id, platform, channel_content_id,
  638. title, url, content_type, images, like_count, publish_time,
  639. quality_score, quality_grade, found_by, knowledge_type, overall_score,
  640. {_REL_SQL} AS rel, {_REPRO_SQL} AS repro
  641. FROM {table} WHERE query_id=%s AND {_REAL_POST}
  642. ORDER BY overall_score DESC, id""", (query_id,))
  643. rows = cur.fetchall()
  644. cur.execute("SELECT DISTINCT case_id FROM process_run WHERE query_id=%s", (query_id,)) # 3-II
  645. hp = {r["case_id"] for r in cur.fetchall()}
  646. cur.execute("SELECT DISTINCT case_id FROM mode_tools WHERE query_id=%s", (query_id,))
  647. ht = {r["case_id"] for r in cur.fetchall()}
  648. # 已归类(工序):hp 中各 case 最新真实版的工序里有任一 step 含 substanceMatch(同归类回写口径)。
  649. # 3-II:改查三表——run 级聚合「该 run 是否有任一 step_json 含 substanceMatch」,逻辑同 _categorized_from_rows。
  650. hc = set()
  651. if hp:
  652. ph = ",".join(["%s"] * len(hp))
  653. cur.execute(f"""SELECT r.case_id, r.version, r.id,
  654. (LEFT(r.version,5)='link_') AS islink,
  655. EXISTS(SELECT 1 FROM process_step s
  656. JOIN process_procedure p ON s.procedure_pk=p.id
  657. WHERE p.run_id=r.id AND s.step_json LIKE %s) AS cat
  658. FROM process_run r WHERE r.case_id IN ({ph})""",
  659. ['%substanceMatch%'] + list(hp))
  660. hc = _categorized_from_rows(cur.fetchall())
  661. finally:
  662. conn.close()
  663. for r in rows:
  664. for col in ("images", "found_by", "knowledge_type"):
  665. r[col] = _loads(r[col])
  666. r["adopted"] = is_adopted_rel(r["overall_score"], r.pop("rel", None),
  667. r["publish_time"], r.pop("repro", None))
  668. r["has_process"] = r["case_id"] in hp
  669. r["has_tools"] = r["case_id"] in ht
  670. r["has_category"] = r["case_id"] in hc
  671. return rows
  672. def fetch_post(query_id, case_id, table="search_process"):
  673. """指定搜索表的单帖完整行(给 pipeline 脚本重建 source 用)。无则 None。"""
  674. table = _search_table(table)
  675. conn = _conn()
  676. try:
  677. with conn.cursor() as cur:
  678. cur.execute(f"SELECT * FROM {table} WHERE query_id=%s AND case_id=%s",
  679. (query_id, case_id))
  680. row = cur.fetchone()
  681. finally:
  682. conn.close()
  683. if not row:
  684. return None
  685. for col in ("images", "videos", "found_by", "knowledge_type", "llm_evaluation"):
  686. row[col] = _loads(row[col])
  687. return row
  688. def fetch_post_by_case(case_id):
  689. """按 case_id 取单帖原文(标题/正文/配图/链接),**不限 query、不限方向**。
  690. 供知识库(search.html)按 source_id=case_id 展示原文用——那里只有 source_id,没有
  691. query_id/mode。先查 search_tools 再查 search_process(同一帖两表内容一致),取任一命中行。
  692. 只回列表展示所需列(title/url/body/images/platform/publish_time),无则 None。
  693. """
  694. conn = _conn()
  695. try:
  696. with conn.cursor() as cur:
  697. for table in ("search_tools", "search_process"):
  698. cur.execute(f"""SELECT case_id, title, url, body, images, platform, publish_time
  699. FROM {table} WHERE case_id=%s LIMIT 1""", (case_id,))
  700. row = cur.fetchone()
  701. if row:
  702. row["images"] = _loads(row["images"])
  703. return row
  704. finally:
  705. conn.close()
  706. return None
  707. def fetch_all_posts(mode="process", *, query_ids=None, case_ids=None, adopted_only=False, distinct=False,
  708. limit=None, offset=0):
  709. """某方向「全部帖子」:跨所有 query 的列表(瘦身列,口径同 fetch_posts,不拉
  710. body/videos/llm_evaluation 大字段)。fetch_posts 限定单 query,本函数默认取全表。
  711. - query_ids:选填 query_id 列表,传了就 WHERE query_id IN(...) 只取这些 query
  712. 的帖子(SQL 层过滤,不拉全表);None=全部,[]=空结果。
  713. - adopted_only=True:只返回采纳帖(is_adopted_rel 口径,rel/repro 由
  714. _REL_SQL/_REPRO_SQL 直取标量算,不拉整表 blob)。
  715. - distinct=True:按 case_id 去重(同一帖被多个 query 搜到时,只保留
  716. overall_score 最高的一行——已按 score 降序,取首次出现即最高分)。
  717. - limit/offset:分页(limit=None 不分页)。
  718. 返回 (total, rows):total 为过滤(+去重)后的总条数,rows 为本页切片。"""
  719. table = _search_table(mode)
  720. # 始终排除占位哨兵行(无搜索的 query 不在帖子视图里出现)
  721. where, params = f" WHERE {_REAL_POST}", []
  722. if query_ids is not None:
  723. if not query_ids:
  724. return 0, [] # 显式空列表:直接空结果,不必查库
  725. where += " AND query_id IN (" + ",".join(["%s"] * len(query_ids)) + ")"
  726. params = list(query_ids)
  727. if case_ids is not None:
  728. # case_ids:选填 case_id 列表(知识归类联动用——按「分类树节点的 knowledge_ids」
  729. # 取被归类进该节点的帖子,而非按 query 搜索结果)。None=不过滤,[]=空结果。
  730. if not case_ids:
  731. return 0, []
  732. where += " AND case_id IN (" + ",".join(["%s"] * len(case_ids)) + ")"
  733. params += list(case_ids)
  734. conn = _conn()
  735. try:
  736. with conn.cursor() as cur:
  737. cur.execute(f"""SELECT id, query_id, query_text, case_id, platform, channel_content_id,
  738. title, url, content_type, images, like_count, publish_time,
  739. quality_score, quality_grade, found_by, knowledge_type, overall_score,
  740. {_REL_SQL} AS rel, {_REPRO_SQL} AS repro
  741. FROM {table}{where}
  742. ORDER BY overall_score DESC, id""", params)
  743. rows = cur.fetchall()
  744. # has_process/has_tools 全局判定:跨 query 的「该帖是否已解构」,两张解构表各取一次
  745. cur.execute("SELECT DISTINCT case_id FROM process_run") # 3-II
  746. hp = {r["case_id"] for r in cur.fetchall()}
  747. cur.execute("SELECT DISTINCT case_id FROM mode_tools")
  748. ht = {r["case_id"] for r in cur.fetchall()}
  749. finally:
  750. conn.close()
  751. out, seen = [], set()
  752. for r in rows:
  753. for col in ("images", "found_by", "knowledge_type"):
  754. r[col] = _loads(r[col])
  755. r["adopted"] = is_adopted_rel(r["overall_score"], r.pop("rel", None),
  756. r["publish_time"], r.pop("repro", None))
  757. if adopted_only and not r["adopted"]:
  758. continue
  759. if distinct:
  760. if r["case_id"] in seen:
  761. continue
  762. seen.add(r["case_id"])
  763. r["has_process"] = r["case_id"] in hp
  764. r["has_tools"] = r["case_id"] in ht
  765. out.append(r)
  766. total = len(out)
  767. if limit is not None:
  768. out = out[offset:offset + limit]
  769. elif offset:
  770. out = out[offset:]
  771. return total, out
  772. def count_executed_queries(mode="process"):
  773. """该方向「已执行」的 query 数 = 搜索表里出现过的 distinct query_id 个数。
  774. 注:一次搜索若 0 命中则不写任何行,故不计入(口径为「已产出结果的 query」)。"""
  775. table = _search_table(mode)
  776. conn = _conn()
  777. try:
  778. with conn.cursor() as cur:
  779. cur.execute(f"SELECT COUNT(DISTINCT query_id) AS n FROM {table} WHERE {_REAL_POST}")
  780. return cur.fetchone()["n"]
  781. finally:
  782. conn.close()
  783. # ── mode_process ─────────────────────────────────────────────────────────────
  784. def replace_process(query_id, case_id, platform, post_title, payload,
  785. model, version, cost_usd, duration_s):
  786. """写入一帖某版本的工序解构结果(payload = {source, procedures})。
  787. 删 (case_id, version) 旧行再插,同版本重跑幂等、跨版本保留历史。返回工序条数。"""
  788. source = payload.get("source")
  789. procedures = payload.get("procedures") or []
  790. conn = _conn()
  791. try:
  792. conn.begin() # DELETE+INSERT 原子化:配合 uk_q_case_ver_seq,并发/重复写入不会留下重复行
  793. with conn.cursor() as cur:
  794. cur.execute("DELETE FROM mode_process WHERE case_id=%s AND version=%s",
  795. (case_id, version))
  796. if procedures:
  797. rows = []
  798. for i, p in enumerate(procedures):
  799. steps = p.get("steps") or []
  800. vias = []
  801. for s in steps:
  802. v = s.get("via")
  803. if v and v not in vias:
  804. vias.append(v)
  805. rows.append((
  806. query_id, case_id, platform, (post_title or "")[:500],
  807. _j(source), p.get("id"), (p.get("name") or "")[:250],
  808. p.get("purpose"), p.get("category"),
  809. _j(p.get("declarations")), _j(p.get("type_registry")),
  810. _j(steps), len(steps), _j(vias),
  811. model, version, cost_usd, duration_s, i,
  812. ))
  813. cur.executemany("""
  814. INSERT INTO mode_process
  815. (query_id, case_id, platform, post_title, source, procedure_id, name,
  816. purpose, category, declarations, type_registry, steps, step_count,
  817. tools_used, model, version, cost_usd, duration_s, seq)
  818. VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)
  819. """, rows)
  820. # 3-II 双写:同一事务内覆盖写 process_run/procedure/step(新表为主,mode_process 兜底)
  821. _split_write_run(cur, query_id, case_id, platform, post_title, source,
  822. procedures, model, version, cost_usd, duration_s)
  823. conn.commit()
  824. return len(procedures)
  825. except Exception:
  826. conn.rollback()
  827. raise
  828. finally:
  829. conn.close()
  830. def fetch_process_versions(case_id):
  831. # 3-II:版本列表读自三表(n=该版本工序数,口径同旧 mode_process 行数)
  832. conn = _conn()
  833. try:
  834. with conn.cursor() as cur:
  835. cur.execute("""SELECT r.version, COUNT(p.id) AS n, MAX(r.model) AS model
  836. FROM process_run r
  837. LEFT JOIN process_procedure p ON p.run_id = r.id
  838. WHERE r.case_id=%s
  839. GROUP BY r.version
  840. ORDER BY (LEFT(r.version,5)='link_') ASC, MAX(r.id) DESC""", (case_id,))
  841. return cur.fetchall()
  842. finally:
  843. conn.close()
  844. def fetch_process(case_id, version=None):
  845. """重建 {case_id, version, model, source, procedures:[...]}。version=None 取最新。
  846. 3-II:读自 process_run/procedure/step(返回形状与旧 mode_process 口径完全一致)。"""
  847. return rebuild_process_from_split(case_id, version)
  848. def fetch_all_process(case_ids=None, lite=False):
  849. """批量取多帖工序解构(每帖取最新真实版,link_ 排后),一次查询拍平。
  850. - case_ids:选填 case_id 列表;None=全表所有有解构的帖,[]=空(直接返回空)。
  851. - lite:True 走精简投影(丢大字段 + 截断 value),供工序库平铺表快速首屏。
  852. 返回 {case_id: _proc_payload(...)};无解构记录的 case_id 不出现在结果里。"""
  853. if case_ids is not None and not case_ids:
  854. return {}
  855. # 3-II:读自三表,3 条批量 SQL(run → procedure → step),避免按 case 逐次往返(保首屏性能)。
  856. where, params = "", []
  857. if case_ids is not None:
  858. where = " WHERE case_id IN (" + ",".join(["%s"] * len(case_ids)) + ")"
  859. params = list(case_ids)
  860. conn = _conn()
  861. try:
  862. with conn.cursor() as cur:
  863. cur.execute(f"SELECT * FROM process_run{where} ORDER BY case_id, id", params)
  864. runs = cur.fetchall()
  865. best = {} # case_id -> (key, run);口径同 fetch_process(is_latest→真实版→id 最大)
  866. for r in runs:
  867. c = r["case_id"]
  868. k = ((r.get("is_latest") or 0), not str(r["version"]).startswith("link_"), r["id"])
  869. if c not in best or k > best[c][0]:
  870. best[c] = (k, r)
  871. run_ids = [run["id"] for _k, run in best.values()]
  872. procs, steps_by_proc = [], {}
  873. if run_ids:
  874. ph = ",".join(["%s"] * len(run_ids))
  875. cur.execute(f"SELECT * FROM process_procedure WHERE run_id IN ({ph}) "
  876. "ORDER BY run_id, seq, id", run_ids)
  877. procs = cur.fetchall()
  878. pid = [p["id"] for p in procs]
  879. if pid:
  880. ph2 = ",".join(["%s"] * len(pid))
  881. cur.execute(f"SELECT procedure_pk, step_json FROM process_step "
  882. f"WHERE procedure_pk IN ({ph2}) ORDER BY procedure_pk, seq, id", pid)
  883. for s in cur.fetchall():
  884. steps_by_proc.setdefault(s["procedure_pk"], []).append(
  885. _loads(s["step_json"], {}))
  886. finally:
  887. conn.close()
  888. procs_by_run = {}
  889. for p in procs:
  890. procs_by_run.setdefault(p["run_id"], []).append(p)
  891. out = {}
  892. for cid, (_k, run) in best.items():
  893. rprocs = procs_by_run.get(run["id"], [])
  894. if lite:
  895. procedures = [{"id": p["procedure_id"], "name": p["name"],
  896. "steps": _lite_steps(steps_by_proc.get(p["id"], []))} for p in rprocs]
  897. out[cid] = {"case_id": cid, "version": run["version"],
  898. "title": run["post_title"], "procedures": procedures}
  899. else:
  900. procedures = [{
  901. "id": p["procedure_id"], "name": p["name"], "purpose": p["purpose"],
  902. "category": p["category"], "declarations": _loads(p["declarations"]),
  903. "type_registry": _loads(p["type_registry"]),
  904. "steps": steps_by_proc.get(p["id"], []),
  905. "tools_used": _loads(p["tools_used"], [])} for p in rprocs]
  906. out[cid] = {"case_id": cid, "version": run["version"], "platform": run["platform"],
  907. "title": run["post_title"], "model": run["model"],
  908. "cost_usd": float(run["cost_usd"]) if run["cost_usd"] is not None else None,
  909. "duration_s": run["duration_s"], "source": _loads(run["source"]),
  910. "procedures": procedures}
  911. return out
  912. def fetch_all_process_pairs():
  913. """全量重跑归类用:每个有解构记录的 case 取其「最新真实版」(link_ 排后、id 最大)
  914. 所在行的 (query_id, case_id)。选版口径与 fetch_process / fetch_all_process 完全一致,
  915. 保证回写的版本即前端展示的版本。query_id 作 post_id 供 category-match record。
  916. 返回 [(query_id, case_id)],按 case_id 排序、按 case 去重。"""
  917. conn = _conn()
  918. try:
  919. with conn.cursor() as cur:
  920. cur.execute("SELECT id, query_id, case_id, version FROM process_run ORDER BY case_id, id")
  921. rows = cur.fetchall()
  922. finally:
  923. conn.close()
  924. best = {} # case_id -> ((is_real, id), query_id);取 max,与 fetch_all_process 同口径
  925. for r in rows:
  926. c = r["case_id"]
  927. key = (not str(r["version"]).startswith("link_"), r["id"])
  928. if c not in best or key > best[c][0]:
  929. best[c] = (key, r["query_id"])
  930. return [(qid, c) for c, (_k, qid) in sorted(best.items())]
  931. # ── 拆表(3-II):mode_process → process_run/procedure/step ────────────────────────
  932. def _jn(v):
  933. """JSON 列搬运规整:先 _loads 再 _j,无论 pymysql 返回 str 还是已解析对象都得到合法 JSON 文本/None。"""
  934. return _j(_loads(v))
  935. def _t(s, n):
  936. """抽取列按字符数截断(只用于可查询列;完整值在 step_json,不丢)。None 原样。"""
  937. return s if (s is None or len(str(s)) <= n) else str(s)[:n]
  938. def migrate_to_split_tables(truncate=True):
  939. """把 mode_process 灌入 process_run/procedure/step(加法,不动 mode_process)。
  940. 幂等:默认先 TRUNCATE 三表再灌。step_json 逐字保留原 step。返回 {runs, procedures, steps}。"""
  941. conn = _conn()
  942. try:
  943. with conn.cursor() as cur:
  944. cur.execute("SELECT * FROM mode_process ORDER BY query_id, case_id, version, seq, id")
  945. rows = cur.fetchall()
  946. finally:
  947. conn.close()
  948. # 每 case 最新真实版(口径同 fetch_all_process),用于置 is_latest
  949. best = {}
  950. for r in rows:
  951. c = r["case_id"]; k = (not str(r["version"]).startswith("link_"), r["id"])
  952. if c not in best or k > best[c][0]:
  953. best[c] = (k, r["version"])
  954. latest_ver = {c: v for c, (_k, v) in best.items()}
  955. # 按 (query_id, case_id, version) 分组 = 一个 run(保序)
  956. runs = {}
  957. order = []
  958. for r in rows:
  959. key = (r["query_id"], r["case_id"], r["version"])
  960. if key not in runs:
  961. runs[key] = []; order.append(key)
  962. runs[key].append(r)
  963. nrun = nproc = nstep = 0
  964. conn = _conn()
  965. try:
  966. conn.begin()
  967. with conn.cursor() as cur:
  968. if truncate:
  969. for t in ("process_step", "process_procedure", "process_run"):
  970. cur.execute(f"TRUNCATE TABLE {t}")
  971. for key in order:
  972. qid, cid, ver = key
  973. prows = sorted(runs[key],
  974. key=lambda r: (r["seq"] if r["seq"] is not None else 0, r["id"]))
  975. first = prows[0]
  976. is_latest = 1 if (not str(ver).startswith("link_") and latest_ver.get(cid) == ver) else 0
  977. cur.execute("""INSERT INTO process_run
  978. (query_id, case_id, version, is_latest, platform, post_title, source,
  979. model, cost_usd, duration_s, created_at)
  980. VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)""",
  981. (qid, cid, ver, is_latest, first["platform"], first["post_title"],
  982. _jn(first["source"]), first["model"], first["cost_usd"], first["duration_s"],
  983. first["created_at"])) # 保留原解构时间(Dashboard 成本趋势按此)
  984. run_id = cur.lastrowid; nrun += 1
  985. for pseq, pr in enumerate(prows):
  986. steps = _loads(pr["steps"], [])
  987. cur.execute("""INSERT INTO process_procedure
  988. (run_id, case_id, version, procedure_id, seq, name, purpose, category,
  989. declarations, type_registry, tools_used, step_count)
  990. VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)""",
  991. (run_id, cid, ver, pr["procedure_id"], pseq, pr["name"], pr["purpose"],
  992. pr["category"], _jn(pr["declarations"]), _jn(pr["type_registry"]),
  993. _jn(pr["tools_used"]), pr["step_count"]))
  994. proc_pk = cur.lastrowid; nproc += 1
  995. for sseq, st in enumerate(steps if isinstance(steps, list) else []):
  996. if not isinstance(st, dict):
  997. st = {"value": st}
  998. cur.execute("""INSERT INTO process_step
  999. (procedure_pk, case_id, version, step_id, seq, kind, grp, via,
  1000. effect, action, substance, form, step_json)
  1001. VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)""",
  1002. (proc_pk, cid, ver, _t(st.get("id"), 16), sseq,
  1003. _t(st.get("kind"), 16), _t(st.get("group"), 16), _t(st.get("via"), 255),
  1004. _t(st.get("effect"), 255), _t(st.get("action"), 255),
  1005. _t(st.get("substance"), 512), _t(st.get("form"), 512),
  1006. _j(st)))
  1007. nstep += 1
  1008. conn.commit()
  1009. except Exception:
  1010. conn.rollback(); raise
  1011. finally:
  1012. conn.close()
  1013. return {"runs": nrun, "procedures": nproc, "steps": nstep}
  1014. def _steps_by_proc(cur, run_id):
  1015. """取某 run 下所有 step 的 step_json,按 procedure_pk 分组(组内按 seq 序)。"""
  1016. cur.execute("""SELECT procedure_pk, step_json FROM process_step
  1017. WHERE procedure_pk IN (SELECT id FROM process_procedure WHERE run_id=%s)
  1018. ORDER BY procedure_pk, seq, id""", (run_id,))
  1019. out = {}
  1020. for s in cur.fetchall():
  1021. out.setdefault(s["procedure_pk"], []).append(_loads(s["step_json"], {}))
  1022. return out
  1023. def _split_payload(cur, run, lite=False):
  1024. """run 行 → ProcessPayload(形状与 _proc_payload 完全一致)。需开放 cursor。无 run 返回 None。"""
  1025. if not run:
  1026. return None
  1027. cur.execute("SELECT * FROM process_procedure WHERE run_id=%s ORDER BY seq, id", (run["id"],))
  1028. procs = cur.fetchall()
  1029. sbp = _steps_by_proc(cur, run["id"])
  1030. case_id, version = run["case_id"], run["version"]
  1031. if lite:
  1032. procedures = [{
  1033. "id": p["procedure_id"], "name": p["name"],
  1034. "steps": _lite_steps(sbp.get(p["id"], [])),
  1035. } for p in procs]
  1036. return {"case_id": case_id, "version": version,
  1037. "title": run["post_title"], "procedures": procedures}
  1038. procedures = [{
  1039. "id": p["procedure_id"], "name": p["name"], "purpose": p["purpose"],
  1040. "category": p["category"], "declarations": _loads(p["declarations"]),
  1041. "type_registry": _loads(p["type_registry"]), "steps": sbp.get(p["id"], []),
  1042. "tools_used": _loads(p["tools_used"], []),
  1043. } for p in procs]
  1044. return {"case_id": case_id, "version": version, "platform": run["platform"],
  1045. "title": run["post_title"], "model": run["model"],
  1046. "cost_usd": float(run["cost_usd"]) if run["cost_usd"] is not None else None,
  1047. "duration_s": run["duration_s"],
  1048. "source": _loads(run["source"]), "procedures": procedures}
  1049. def _select_run(cur, case_id, version=None, query_id=None):
  1050. """选 run:version 指定则精确取;否则取最新(is_latest 优先 → 非 link_ → id 最大,
  1051. 口径同旧 fetch_process)。query_id 给定则限定该 query。无则 None。"""
  1052. where = "case_id=%s"; params = [case_id]
  1053. if query_id is not None:
  1054. where += " AND query_id=%s"; params.append(query_id)
  1055. if version is not None:
  1056. cur.execute(f"SELECT * FROM process_run WHERE {where} AND version=%s LIMIT 1",
  1057. params + [version])
  1058. else:
  1059. cur.execute(f"""SELECT * FROM process_run WHERE {where}
  1060. ORDER BY is_latest DESC, (LEFT(version,5)='link_') ASC, id DESC
  1061. LIMIT 1""", params)
  1062. return cur.fetchone()
  1063. def rebuild_process_from_split(case_id, version=None):
  1064. """从三表重建 ProcessPayload(迁移校验/读侧共用)。version=None 取最新。"""
  1065. conn = _conn()
  1066. try:
  1067. with conn.cursor() as cur:
  1068. return _split_payload(cur, _select_run(cur, case_id, version))
  1069. finally:
  1070. conn.close()
  1071. def _refresh_is_latest(cur, case_id):
  1072. """重算某 case 各 run 的 is_latest:最新真实版(非 link_、id 最大)置 1,其余 0。
  1073. 口径同 fetch_process 选版。只有 link_ 时无 1(读侧 rebuild 仍能按排序兜底)。需在事务内调用。"""
  1074. cur.execute("SELECT id, version FROM process_run WHERE case_id=%s", (case_id,))
  1075. rows = cur.fetchall()
  1076. if not rows:
  1077. return
  1078. cur.execute("UPDATE process_run SET is_latest=0 WHERE case_id=%s", (case_id,))
  1079. best = max(rows, key=lambda r: (not str(r["version"]).startswith("link_"), r["id"]))
  1080. if not str(best["version"]).startswith("link_"):
  1081. cur.execute("UPDATE process_run SET is_latest=1 WHERE id=%s", (best["id"],))
  1082. def _split_write_run(cur, query_id, case_id, platform, post_title, source,
  1083. procedures, model, version, cost_usd, duration_s):
  1084. """在三表覆盖写某 (query_id, case_id, version) 的 run/procedure/step(需在调用方事务内)。
  1085. 与 replace_process 写 mode_process 同口径(tools_used 从 steps[].via 去重)。维护 is_latest。"""
  1086. cur.execute("SELECT id FROM process_run WHERE query_id=%s AND case_id=%s AND version=%s",
  1087. (query_id, case_id, version))
  1088. old = cur.fetchone()
  1089. if old:
  1090. cur.execute("""DELETE FROM process_step WHERE procedure_pk IN
  1091. (SELECT id FROM process_procedure WHERE run_id=%s)""", (old["id"],))
  1092. cur.execute("DELETE FROM process_procedure WHERE run_id=%s", (old["id"],))
  1093. cur.execute("DELETE FROM process_run WHERE id=%s", (old["id"],))
  1094. cur.execute("""INSERT INTO process_run
  1095. (query_id, case_id, version, is_latest, platform, post_title, source,
  1096. model, cost_usd, duration_s)
  1097. VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)""",
  1098. (query_id, case_id, version, 0, platform, (post_title or "")[:500],
  1099. _j(source), model, cost_usd, duration_s))
  1100. run_id = cur.lastrowid
  1101. for pseq, p in enumerate(procedures or []):
  1102. steps = p.get("steps") or []
  1103. vias = []
  1104. for s in steps:
  1105. v = s.get("via") if isinstance(s, dict) else None
  1106. if v and v not in vias:
  1107. vias.append(v)
  1108. cur.execute("""INSERT INTO process_procedure
  1109. (run_id, case_id, version, procedure_id, seq, name, purpose, category,
  1110. declarations, type_registry, tools_used, step_count)
  1111. VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)""",
  1112. (run_id, case_id, version, p.get("id"), pseq, (p.get("name") or "")[:250],
  1113. p.get("purpose"), p.get("category"), _j(p.get("declarations")),
  1114. _j(p.get("type_registry")), _j(vias), len(steps)))
  1115. proc_pk = cur.lastrowid
  1116. for sseq, st in enumerate(steps):
  1117. if not isinstance(st, dict):
  1118. st = {"value": st}
  1119. cur.execute("""INSERT INTO process_step
  1120. (procedure_pk, case_id, version, step_id, seq, kind, grp, via,
  1121. effect, action, substance, form, step_json)
  1122. VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)""",
  1123. (proc_pk, case_id, version, _t(st.get("id"), 16), sseq,
  1124. _t(st.get("kind"), 16), _t(st.get("group"), 16), _t(st.get("via"), 255),
  1125. _t(st.get("effect"), 255), _t(st.get("action"), 255),
  1126. _t(st.get("substance"), 512), _t(st.get("form"), 512), _j(st)))
  1127. _refresh_is_latest(cur, case_id)
  1128. def _split_update_steps(cur, case_id, version, steps_in_order):
  1129. """归类回写:按工序顺序覆盖某 (case_id, version) 各 procedure 的 step_json(及抽取列)。
  1130. steps_in_order 与 process_procedure(按 seq) 一一对应。需在事务内。返回更新的 step 行数。"""
  1131. cur.execute("""SELECT id FROM process_procedure
  1132. WHERE case_id=%s AND version=%s ORDER BY seq, id""", (case_id, version))
  1133. proc_ids = [r["id"] for r in cur.fetchall()]
  1134. if len(proc_ids) != len(steps_in_order):
  1135. raise ValueError(f"process_procedure 行数({len(proc_ids)})与工序数({len(steps_in_order)})不一致")
  1136. n = 0
  1137. for proc_pk, steps in zip(proc_ids, steps_in_order):
  1138. cur.execute("DELETE FROM process_step WHERE procedure_pk=%s", (proc_pk,))
  1139. for sseq, st in enumerate(steps or []):
  1140. if not isinstance(st, dict):
  1141. st = {"value": st}
  1142. cur.execute("""INSERT INTO process_step
  1143. (procedure_pk, case_id, version, step_id, seq, kind, grp, via,
  1144. effect, action, substance, form, step_json)
  1145. VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)""",
  1146. (proc_pk, case_id, version, _t(st.get("id"), 16), sseq,
  1147. _t(st.get("kind"), 16), _t(st.get("group"), 16), _t(st.get("via"), 255),
  1148. _t(st.get("effect"), 255), _t(st.get("action"), 255),
  1149. _t(st.get("substance"), 512), _t(st.get("form"), 512), _j(st)))
  1150. n += 1
  1151. return n
  1152. # ── step_classification(归类独立表,3-I)────────────────────────────────────────
  1153. def replace_step_classification(case_id, version, records):
  1154. """覆盖写某 (case_id, version) 的 step 维度归类结果(DELETE+INSERT 原子,语义同 replace_process)。
  1155. records:[{query_id, procedure_id, step_id, dimension, sub_index, raw_term,
  1156. matched_name, matched_path, matched_id, score}](只含有命中的子项)。
  1157. 同版本重跑幂等、洗净旧归属。返回写入条数。"""
  1158. conn = _conn()
  1159. try:
  1160. conn.begin()
  1161. with conn.cursor() as cur:
  1162. cur.execute("DELETE FROM step_classification WHERE case_id=%s AND version=%s",
  1163. (case_id, version))
  1164. if records:
  1165. rows = [(
  1166. case_id, r.get("query_id"), version,
  1167. r.get("procedure_id"), r.get("step_id"),
  1168. r.get("dimension"), r.get("sub_index"), r.get("raw_term"),
  1169. r.get("matched_name"), r.get("matched_path"),
  1170. r.get("matched_id"), r.get("score"),
  1171. ) for r in records]
  1172. cur.executemany("""
  1173. INSERT INTO step_classification
  1174. (case_id, query_id, version, procedure_id, step_id, dimension,
  1175. sub_index, raw_term, matched_name, matched_path, matched_id, score)
  1176. VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)
  1177. """, rows)
  1178. conn.commit()
  1179. return len(records)
  1180. except Exception:
  1181. conn.rollback()
  1182. raise
  1183. finally:
  1184. conn.close()
  1185. def fetch_classified_cases(dimension, path=None, name=None, subtree=False):
  1186. """维度联动取帖:返回某维度下命中分类的 distinct case_id 集合。
  1187. - path:给定分类全路径。**默认精确匹配本节点**(matched_path = path),与 cat-api 树徽标
  1188. (节点 knowledge_count)同口径,保证取到的帖工序解构 substanceMatch 就是该节点名;
  1189. subtree=True 时改取 **该节点及整棵子树**(matched_path = path 或 LIKE 'path/%')。
  1190. - name:给定分类名精确匹配(matched_name = name)。
  1191. path 与 name 至少给一个;都给则取并集。
  1192. 与 mode_process.steps[].substanceMatch 同源(同一次归类双写),不依赖 cat-api 树。"""
  1193. conds, params = [], [dimension]
  1194. if path:
  1195. if subtree:
  1196. conds.append("(matched_path = %s OR matched_path LIKE %s)")
  1197. params += [path, path.rstrip("/") + "/%"]
  1198. else:
  1199. conds.append("matched_path = %s")
  1200. params.append(path)
  1201. if name:
  1202. conds.append("matched_name = %s")
  1203. params.append(name)
  1204. if not conds:
  1205. return set()
  1206. where = "dimension=%s AND (" + " OR ".join(conds) + ")"
  1207. conn = _conn()
  1208. try:
  1209. with conn.cursor() as cur:
  1210. cur.execute(f"SELECT DISTINCT case_id FROM step_classification WHERE {where}", params)
  1211. return {r["case_id"] for r in cur.fetchall()}
  1212. finally:
  1213. conn.close()
  1214. def fetch_process_by_query(query_id, case_id, version=None):
  1215. """同 fetch_process,但用 (query_id, case_id) 精确定位某 query 下该帖的工序
  1216. (category-match 用:post_id=query_id / knowledge_id=case_id)。
  1217. version=None 取该 (query_id, case_id) 下最新真实版(link_ 排后)。无行返回 None。"""
  1218. conn = _conn()
  1219. try:
  1220. with conn.cursor() as cur:
  1221. return _split_payload(cur, _select_run(cur, case_id, version, query_id=query_id))
  1222. finally:
  1223. conn.close()
  1224. def update_process_steps_by_query(query_id, case_id, version, steps_in_order):
  1225. """按工序顺序覆盖某 (query_id, case_id, version) 各行的 steps JSON 列。
  1226. steps_in_order 必须与 fetch_process_by_query 返回的 procedures 同序(均按 seq, id 升序);
  1227. 按行 id 一一对应更新,稳健于 seq 不连续。行数与工序数不符则报错回滚。返回更新行数。"""
  1228. conn = _conn()
  1229. try:
  1230. conn.begin()
  1231. with conn.cursor() as cur:
  1232. cur.execute("""SELECT id FROM mode_process
  1233. WHERE query_id=%s AND case_id=%s AND version=%s
  1234. ORDER BY seq, id""", (query_id, case_id, version))
  1235. ids = [r["id"] for r in cur.fetchall()]
  1236. if len(ids) != len(steps_in_order):
  1237. raise ValueError(f"行数({len(ids)})与工序数({len(steps_in_order)})不一致")
  1238. n = 0
  1239. for row_id, steps in zip(ids, steps_in_order):
  1240. cur.execute("UPDATE mode_process SET steps=%s WHERE id=%s", (_j(steps), row_id))
  1241. n += cur.rowcount
  1242. # 3-II 双写:同步更新 process_step.step_json(同事务)
  1243. _split_update_steps(cur, case_id, version, steps_in_order)
  1244. conn.commit()
  1245. return n
  1246. except Exception:
  1247. conn.rollback()
  1248. raise
  1249. finally:
  1250. conn.close()
  1251. def update_process_steps(case_id, version, steps_in_order):
  1252. """按工序顺序覆盖某 (case_id, version) 各行的 steps JSON 列(不限 query_id)。
  1253. 与 fetch_process / fetch_extract 同口径(按 case 的某版本),保证归类回写的版本
  1254. 与前端 /api/extract 展示的版本一致(否则 link_ 复制帖会写错版本、前端看不到)。
  1255. steps_in_order 须与 fetch_process(case_id, version).procedures 同序(按 id 升序)。
  1256. 行数与工序数不符则报错回滚。返回更新行数。"""
  1257. conn = _conn()
  1258. try:
  1259. conn.begin()
  1260. with conn.cursor() as cur:
  1261. cur.execute("""SELECT id FROM mode_process WHERE case_id=%s AND version=%s
  1262. ORDER BY id""", (case_id, version))
  1263. ids = [r["id"] for r in cur.fetchall()]
  1264. if len(ids) != len(steps_in_order):
  1265. raise ValueError(f"行数({len(ids)})与工序数({len(steps_in_order)})不一致")
  1266. n = 0
  1267. for row_id, steps in zip(ids, steps_in_order):
  1268. cur.execute("UPDATE mode_process SET steps=%s WHERE id=%s", (_j(steps), row_id))
  1269. n += cur.rowcount
  1270. # 3-II 双写:同步更新 process_step.step_json(同事务)
  1271. _split_update_steps(cur, case_id, version, steps_in_order)
  1272. conn.commit()
  1273. return n
  1274. except Exception:
  1275. conn.rollback()
  1276. raise
  1277. finally:
  1278. conn.close()
  1279. def _categorized_from_rows(rows):
  1280. """rows:[{case_id, version, id, islink(0/1), cat(0/1)}]。返回已归类 case 集合。
  1281. 口径:每 case 取最新真实版(真实版优先、id 最大),该版本**任一行** cat=1 即已归类。
  1282. 关键——不能只看「id 最大的那一行」:工序里可能有 steps 为空的 procedure(step_count=0),
  1283. 其行永远不含 substanceMatch,若恰好 id 最大会误判整条 case 未归类(见该函数修复缘由)。"""
  1284. best, has = {}, {}
  1285. for r in rows:
  1286. c, v = r["case_id"], r["version"]
  1287. sk = (1 if r["islink"] else 0, -r["id"]) # 真实版(islink=0)优先,其次 id 大;取 min
  1288. if c not in best or sk < best[c][0]:
  1289. best[c] = (sk, v)
  1290. k = (c, v)
  1291. has[k] = has.get(k, False) or bool(r["cat"])
  1292. return {c for c, (sk, v) in best.items() if has.get((c, v))}
  1293. def fetch_categorized_cases(case_ids, mode="process"):
  1294. """返回 case_ids 中「已归类」的子集:该 case 最新真实版(link_ 排后)的工序里有任一步骤
  1295. 含 substanceMatch(归类跑过的非空 step 一定带此 key)。与归类回写/前端展示同口径。
  1296. 供前端判断「是否已全部归类 → 提示重新归类」。仅工序方向有意义(mode_process)。"""
  1297. if not case_ids:
  1298. return set()
  1299. ph = ",".join(["%s"] * len(case_ids))
  1300. conn = _conn()
  1301. try:
  1302. with conn.cursor() as cur:
  1303. if mode == "process":
  1304. # 3-II:run 级聚合「该 run 是否有任一 step_json 含 substanceMatch」
  1305. cur.execute(f"""SELECT r.case_id, r.version, r.id,
  1306. (LEFT(r.version,5)='link_') AS islink,
  1307. EXISTS(SELECT 1 FROM process_step s
  1308. JOIN process_procedure p ON s.procedure_pk=p.id
  1309. WHERE p.run_id=r.id AND s.step_json LIKE %s) AS cat
  1310. FROM process_run r WHERE r.case_id IN ({ph})""",
  1311. ['%substanceMatch%'] + list(case_ids))
  1312. else:
  1313. table = _mode_table(mode)
  1314. cur.execute(f"""SELECT case_id, version, id,
  1315. (LEFT(version,5)='link_') AS islink, (steps LIKE %s) AS cat
  1316. FROM {table} WHERE case_id IN ({ph})""",
  1317. ['%substanceMatch%'] + list(case_ids))
  1318. rows = cur.fetchall()
  1319. finally:
  1320. conn.close()
  1321. return _categorized_from_rows(rows)
  1322. _LITE_VALUE_LIMIT = 300 # lite 模式 输入/输出 value 截断字节上限
  1323. def _trunc(s, limit=_LITE_VALUE_LIMIT):
  1324. """按 UTF-8 字节截断长文本(不切坏多字节字符),超限追加省略号。非字符串原样返回。"""
  1325. if not isinstance(s, str):
  1326. return s
  1327. b = s.encode("utf-8")
  1328. if len(b) <= limit:
  1329. return s
  1330. return b[:limit].decode("utf-8", "ignore") + "…"
  1331. def _lite_steps(steps):
  1332. """lite 模式:仅截断每步 inputs/outputs 的 value(其余字段供平铺表/分组用,保留)。"""
  1333. if not isinstance(steps, list):
  1334. return steps
  1335. for st in steps:
  1336. if not isinstance(st, dict):
  1337. continue
  1338. for io_key in ("inputs", "outputs"):
  1339. ios = st.get(io_key)
  1340. if isinstance(ios, list):
  1341. for io in ios:
  1342. if isinstance(io, dict) and "value" in io:
  1343. io["value"] = _trunc(io.get("value"))
  1344. return steps
  1345. def _proc_payload(case_id, version, rows, lite=False):
  1346. """mode_process 行集 → {case_id, version, …, procedures:[...]}。无行返回 None。
  1347. lite=True:工序库平铺表只需 procedures[].{id,name,steps},丢弃大字段
  1348. (source/declarations/type_registry/tools_used)并截断输入/输出 value;
  1349. 完整值由前端展开时按 case 调 /api/process 懒加载。"""
  1350. if not rows:
  1351. return None
  1352. if lite:
  1353. procedures = [{
  1354. "id": r["procedure_id"], "name": r["name"],
  1355. "steps": _lite_steps(_loads(r["steps"], [])),
  1356. } for r in rows]
  1357. return {"case_id": case_id, "version": version,
  1358. "title": rows[0]["post_title"], "procedures": procedures}
  1359. procedures = [{
  1360. "id": r["procedure_id"], "name": r["name"], "purpose": r["purpose"],
  1361. "category": r["category"], "declarations": _loads(r["declarations"]),
  1362. "type_registry": _loads(r["type_registry"]), "steps": _loads(r["steps"], []),
  1363. "tools_used": _loads(r["tools_used"], []),
  1364. } for r in rows]
  1365. return {"case_id": case_id, "version": version, "platform": rows[0]["platform"],
  1366. "title": rows[0]["post_title"], "model": rows[0]["model"],
  1367. "cost_usd": float(rows[0]["cost_usd"]) if rows[0]["cost_usd"] is not None else None,
  1368. "duration_s": rows[0]["duration_s"],
  1369. "source": _loads(rows[0]["source"]), "procedures": procedures}
  1370. # ── mode_tools ───────────────────────────────────────────────────────────────
  1371. def replace_tools(query_id, case_id, platform, post_title, tools,
  1372. model, version, cost_usd, duration_s, source=None):
  1373. """写入一帖某版本的工具解构结果。语义同 replace_process。返回工具条数。
  1374. source:帖子来源块(同 mode_process,每行重复存),供知识上传脚本重建 source 用。"""
  1375. src = _j(source)
  1376. conn = _conn()
  1377. try:
  1378. conn.begin() # DELETE+INSERT 原子化:配合 uk_q_case_ver_seq,并发/重复写入不会留下重复行
  1379. with conn.cursor() as cur:
  1380. cur.execute("DELETE FROM mode_tools WHERE case_id=%s AND version=%s",
  1381. (case_id, version))
  1382. if tools:
  1383. rows = [(
  1384. query_id, case_id, platform, (post_title or "")[:500], src,
  1385. (t.get("工具名称") or "")[:250],
  1386. _j(t.get("实质作用域")), _j(t.get("形式作用域")),
  1387. t.get("创作层级"), t.get("来源链接"), t.get("输入"), t.get("输出"),
  1388. _j(t.get("用法")), _j(t.get("案例")), _j(t.get("缺点")),
  1389. t.get("最新更新时间"), model, version, cost_usd, duration_s, i,
  1390. ) for i, t in enumerate(tools)]
  1391. cur.executemany("""
  1392. INSERT INTO mode_tools
  1393. (query_id, case_id, platform, post_title, source, tool_name, substance_scope,
  1394. form_scope, creation_layer, source_link, input_desc, output_desc,
  1395. usage_json, cases_json, defects_json, updated_time, model, version,
  1396. cost_usd, duration_s, seq)
  1397. VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)
  1398. """, rows)
  1399. conn.commit()
  1400. return len(tools)
  1401. except Exception:
  1402. conn.rollback()
  1403. raise
  1404. finally:
  1405. conn.close()
  1406. def fetch_tools_versions(case_id):
  1407. conn = _conn()
  1408. try:
  1409. with conn.cursor() as cur:
  1410. cur.execute("""SELECT version, COUNT(*) AS n, MAX(model) AS model
  1411. FROM mode_tools WHERE case_id=%s
  1412. GROUP BY version
  1413. ORDER BY (LEFT(version,5)='link_') ASC, MAX(id) DESC""", (case_id,))
  1414. return cur.fetchall()
  1415. finally:
  1416. conn.close()
  1417. def fetch_tools(case_id, version=None):
  1418. """重建 {case_id, version, model, tool_count, tools:[...]}。version=None 取最新。"""
  1419. conn = _conn()
  1420. try:
  1421. with conn.cursor() as cur:
  1422. if version is None:
  1423. cur.execute("""SELECT version FROM mode_tools WHERE case_id=%s
  1424. ORDER BY (LEFT(version,5)='link_') ASC, id DESC LIMIT 1""", (case_id,))
  1425. row = cur.fetchone()
  1426. if not row:
  1427. return None
  1428. version = row["version"]
  1429. cur.execute("""SELECT * FROM mode_tools WHERE case_id=%s AND version=%s
  1430. ORDER BY id""", (case_id, version))
  1431. rows = cur.fetchall()
  1432. finally:
  1433. conn.close()
  1434. return _tools_payload(case_id, version, rows)
  1435. def _tools_payload(case_id, version, rows):
  1436. """mode_tools 行集 → {case_id, version, …, tools:[...]}。无行返回 None。"""
  1437. if not rows:
  1438. return None
  1439. tools = [{
  1440. "工具名称": r["tool_name"], "实质作用域": _loads(r["substance_scope"]),
  1441. "形式作用域": _loads(r["form_scope"]), "创作层级": r["creation_layer"],
  1442. "来源链接": r["source_link"], "输入": r["input_desc"], "输出": r["output_desc"],
  1443. "用法": _loads(r["usage_json"]), "案例": _loads(r["cases_json"]),
  1444. "缺点": _loads(r["defects_json"]), "最新更新时间": r["updated_time"],
  1445. } for r in rows]
  1446. return {"case_id": case_id, "version": version, "platform": rows[0]["platform"],
  1447. "title": rows[0]["post_title"], "model": rows[0]["model"],
  1448. "cost_usd": float(rows[0]["cost_usd"]) if rows[0]["cost_usd"] is not None else None,
  1449. "duration_s": rows[0]["duration_s"],
  1450. "source": _loads(rows[0].get("source")),
  1451. "tool_count": len(tools), "tools": tools}
  1452. # ── 点击帖子合一查询(单连接,最少往返;远程 RDS 每次往返 ~80ms,故按次数优化)──
  1453. def fetch_extract(mode, case_id, version=None):
  1454. """一次取版本列表 + 解构详情,复用同一条池连接、最少往返。
  1455. 返回 {versions, data, missing}。mode: process / tools。"""
  1456. is_proc = mode != "tools"
  1457. if is_proc:
  1458. # 3-II:工序方向版本列表 + 详情都读自三表(形状不变)
  1459. conn = _conn()
  1460. try:
  1461. with conn.cursor() as cur:
  1462. cur.execute("""SELECT r.version, COUNT(p.id) AS n, MAX(r.model) AS model
  1463. FROM process_run r
  1464. LEFT JOIN process_procedure p ON p.run_id = r.id
  1465. WHERE r.case_id=%s GROUP BY r.version
  1466. ORDER BY (LEFT(r.version,5)='link_') ASC, MAX(r.id) DESC""", (case_id,))
  1467. versions = cur.fetchall()
  1468. target = version or (versions[0]["version"] if versions else None)
  1469. payload = _split_payload(cur, _select_run(cur, case_id, target)) if target else None
  1470. finally:
  1471. conn.close()
  1472. return {"versions": versions, "data": payload, "missing": payload is None}
  1473. # 工具方向:仍读 mode_tools(本期不拆)
  1474. mtable = _mode_table("tools")
  1475. conn = _conn()
  1476. try:
  1477. with conn.cursor() as cur:
  1478. cur.execute(f"""SELECT version, COUNT(*) AS n, MAX(model) AS model
  1479. FROM {mtable} WHERE case_id=%s
  1480. GROUP BY version
  1481. ORDER BY (LEFT(version,5)='link_') ASC, MAX(id) DESC""", (case_id,))
  1482. versions = cur.fetchall()
  1483. target = version or (versions[0]["version"] if versions else None)
  1484. rows = []
  1485. if target is not None:
  1486. cur.execute(f"SELECT * FROM {mtable} WHERE case_id=%s AND version=%s ORDER BY id",
  1487. (case_id, target))
  1488. rows = cur.fetchall()
  1489. finally:
  1490. conn.close()
  1491. payload = _tools_payload(case_id, target, rows)
  1492. return {"versions": versions, "data": payload, "missing": payload is None}
  1493. # ── 跨 query 去重 / link 复制(方案A:解构前先去重,避免重复花钱)──────────────
  1494. # case_id 是帖子物理身份(platform_channelContentId),与 query 无关。同一帖被多个
  1495. # query 搜到时只需真实解构一次;其余 query 用 link_* 复制行补齐关联(cost=0)。
  1496. def latest_real_version(case_id, mode="process"):
  1497. """该 case 是否已有「真实」解构(任意 query;link_* 是复制品,不算源)。
  1498. 返回最新一行 {"version","query_id"} 或 None。给解构前去重判定用。"""
  1499. table = "process_run" if mode == "process" else _mode_table(mode) # 3-II:工序读三表的 run
  1500. conn = _conn()
  1501. try:
  1502. with conn.cursor() as cur:
  1503. cur.execute(f"""SELECT version, query_id FROM {table}
  1504. WHERE case_id=%s AND LEFT(version,5) <> 'link_'
  1505. ORDER BY id DESC LIMIT 1""", (case_id,))
  1506. return cur.fetchone()
  1507. finally:
  1508. conn.close()
  1509. def link_process(query_id, case_id, mode="process"):
  1510. """把 case 在别处最新「真实」版本的解构行复制到目标 query
  1511. (version='link_'+源版本, cost_usd=0)。幂等(先删目标同版本)。
  1512. 返回复制行数;该 case 从未真实解构过则返回 0(无源可复制)。"""
  1513. table = _mode_table(mode)
  1514. conn = _conn()
  1515. try:
  1516. with conn.cursor() as cur:
  1517. cur.execute(f"""SELECT version FROM {table}
  1518. WHERE case_id=%s AND LEFT(version,5) <> 'link_'
  1519. ORDER BY id DESC LIMIT 1""", (case_id,))
  1520. r = cur.fetchone()
  1521. if not r:
  1522. return 0
  1523. srcver = r["version"]
  1524. newver = ("link_" + srcver)[:32] # version 列 VARCHAR(32)
  1525. # 复制除自增 id / 时间戳外的全部列,改写 query_id / version / cost。
  1526. cur.execute(f"SHOW COLUMNS FROM {table}")
  1527. cols = [c["Field"] for c in cur.fetchall()
  1528. if c["Field"] not in ("id", "created_at", "updated_at")]
  1529. cur.execute(f"SELECT {','.join(cols)} FROM {table} WHERE case_id=%s AND version=%s",
  1530. (case_id, srcver))
  1531. rows = cur.fetchall()
  1532. cur.execute(f"DELETE FROM {table} WHERE query_id=%s AND case_id=%s AND version=%s",
  1533. (query_id, case_id, newver))
  1534. for row in rows:
  1535. row = dict(row)
  1536. row["query_id"] = query_id
  1537. row["version"] = newver
  1538. row["cost_usd"] = 0
  1539. cur.execute(
  1540. f"INSERT INTO {table} ({','.join(cols)}) VALUES ({','.join(['%s']*len(cols))})",
  1541. [row[k] for k in cols])
  1542. # 3-II 双写(仅工序方向):把 link_ 复制版同步写进三表
  1543. if mode == "process" and rows:
  1544. srows = sorted(rows, key=lambda r: (r.get("seq") if r.get("seq") is not None else 0))
  1545. f = srows[0]
  1546. procedures = [{
  1547. "id": r.get("procedure_id"), "name": r.get("name"), "purpose": r.get("purpose"),
  1548. "category": r.get("category"), "declarations": _loads(r.get("declarations")),
  1549. "type_registry": _loads(r.get("type_registry")), "steps": _loads(r.get("steps"), []),
  1550. } for r in srows]
  1551. _split_write_run(cur, query_id, case_id, f.get("platform"), f.get("post_title"),
  1552. _loads(f.get("source")), procedures, f.get("model"), newver,
  1553. 0, f.get("duration_s"))
  1554. return len(rows)
  1555. finally:
  1556. conn.close()
  1557. # ── Dashboard 原始行(指标计算在 server.py)─────────────────────────────────────
  1558. # 采纳判定只需「和内容制作知识相关」的得分,用 SQL JSON_EXTRACT 直取这一个标量,
  1559. # 避免把整块 llm_evaluation(本库 ~1.5MB)拉到 Python 再解析。得分可能直接是数字,
  1560. # 也可能裹在 {"得分": x} 里,COALESCE 两条路径覆盖两种存法,口径同 is_adopted。
  1561. _REL_SQL = ("JSON_UNQUOTE(COALESCE("
  1562. "JSON_EXTRACT(llm_evaluation,'$.\"相关性\".\"和内容制作知识相关\".\"得分\"'),"
  1563. "JSON_EXTRACT(llm_evaluation,'$.\"相关性\".\"和内容制作知识相关\"')))")
  1564. # 可复现/实现门槛标量直取(口径同 is_adopted 的 _repro_score):兼容新旧 schema——
  1565. # 旧版「质量.固定维度.可复现性」,新版「质量.动态维度.工序.字段完整性.实现完整性」,COALESCE 依次回退。
  1566. _REPRO_SQL = ("JSON_UNQUOTE(COALESCE("
  1567. "JSON_EXTRACT(llm_evaluation,'$.\"质量\".\"固定维度\".\"可复现性\".\"得分\"'),"
  1568. "JSON_EXTRACT(llm_evaluation,'$.\"质量\".\"固定维度\".\"可复现性\"'),"
  1569. "JSON_EXTRACT(llm_evaluation,'$.\"质量\".\"动态维度\".\"工序\".\"字段完整性\".\"实现完整性\".\"得分\"'),"
  1570. "JSON_EXTRACT(llm_evaluation,'$.\"质量\".\"动态维度\".\"工序\".\"字段完整性\".\"实现完整性\"')))")
  1571. def fetch_adopted_process_cases(query_id=None):
  1572. """返回「已采纳且有工序解构」的 case_id 列表(供知识上传脚本用)。
  1573. 采纳是帖子级属性(评估存在 search_process),工序解构存在 mode_process,故二者 JOIN:
  1574. 只取两边都有的 case,再用 is_adopted_rel(口径同 Dashboard)在 Python 侧过滤。
  1575. relevance 得分由 _REL_SQL 直取标量,不传整块 llm_evaluation。
  1576. query_id 给定时只看该搜索任务下的 case。返回去重、按 case_id 排序的列表。
  1577. """
  1578. sql = (f"SELECT DISTINCT s.case_id, s.overall_score, s.publish_time, "
  1579. f"{_REL_SQL} AS rel, {_REPRO_SQL} AS repro "
  1580. "FROM search_process s "
  1581. "JOIN (SELECT DISTINCT case_id FROM process_run) m ON s.case_id = m.case_id") # 3-II
  1582. params = ()
  1583. if query_id:
  1584. sql += " WHERE s.query_id=%s"
  1585. params = (query_id,)
  1586. conn = _conn()
  1587. try:
  1588. with conn.cursor() as cur:
  1589. cur.execute(sql, params)
  1590. rows = cur.fetchall()
  1591. finally:
  1592. conn.close()
  1593. cases = [r["case_id"] for r in rows
  1594. if is_adopted_rel(r["overall_score"], r["rel"], r["publish_time"], r["repro"])]
  1595. return sorted(set(cases))
  1596. def fetch_destructed_tools_cases(query_ids=None):
  1597. """返回「有工具解构」的 case_id 列表(**不看采纳**),按 mode_tools.query_id 过滤。
  1598. 与 fetch_adopted_tools_cases 的区别:不 JOIN search_tools、不做 is_adopted_rel 过滤,
  1599. 直接取 mode_tools 里存在解构行的帖子。query_ids=None 取全部;给列表则 IN 过滤
  1600. (可跨多个 query;同一 case 在多 query 下解构时按 case_id 去重)。去重、排序返回。
  1601. 供「所有已解构(含未采纳)」上传口径用。
  1602. """
  1603. sql = "SELECT DISTINCT case_id FROM mode_tools"
  1604. params = ()
  1605. if query_ids is not None:
  1606. if not query_ids:
  1607. return [] # 显式空列表:直接空结果
  1608. sql += " WHERE query_id IN (" + ",".join(["%s"] * len(query_ids)) + ")"
  1609. params = tuple(query_ids)
  1610. conn = _conn()
  1611. try:
  1612. with conn.cursor() as cur:
  1613. cur.execute(sql, params)
  1614. return sorted({r["case_id"] for r in cur.fetchall()})
  1615. finally:
  1616. conn.close()
  1617. def fetch_adopted_tools_cases(query_id=None):
  1618. """返回「已采纳且有工具解构」的 case_id 列表(供工具知识上传脚本用)。
  1619. 与 fetch_adopted_process_cases 完全同构,只把搜索/解构表换成工具方向:
  1620. 采纳是帖子级属性(评估存在 search_tools),工具解构存在 mode_tools,故二者 JOIN,
  1621. 只取两边都有的 case,再用 is_adopted_rel(口径同 Dashboard)在 Python 侧过滤。
  1622. query_id 给定时只看该搜索任务下的 case。返回去重、按 case_id 排序的列表。
  1623. """
  1624. sql = (f"SELECT DISTINCT s.case_id, s.overall_score, s.publish_time, "
  1625. f"{_REL_SQL} AS rel, {_REPRO_SQL} AS repro "
  1626. "FROM search_tools s "
  1627. "JOIN (SELECT DISTINCT case_id FROM mode_tools) m ON s.case_id = m.case_id")
  1628. params = ()
  1629. if query_id:
  1630. sql += " WHERE s.query_id=%s"
  1631. params = (query_id,)
  1632. conn = _conn()
  1633. try:
  1634. with conn.cursor() as cur:
  1635. cur.execute(sql, params)
  1636. rows = cur.fetchall()
  1637. finally:
  1638. conn.close()
  1639. cases = [r["case_id"] for r in rows
  1640. if is_adopted_rel(r["overall_score"], r["rel"], r["publish_time"], r["repro"])]
  1641. return sorted(set(cases))
  1642. def route_tables(knowledge_types):
  1643. """知识类型标签 → 落表列表(有序去重)。
  1644. 工序/能力 → search_process;工具 → search_tools;两者都含写两表;空/None 兜底 search_process。
  1645. 评估是统一一套(同一 llm_evaluation blob),故同帖落多表不重复打分,只是多写一行。"""
  1646. kt = set(knowledge_types or [])
  1647. tables = []
  1648. if kt & {"工具"}:
  1649. tables.append("search_tools")
  1650. if (kt & {"工序", "能力"}) or not tables: # 工序/能力,或没命中任何已知标签 → 兜底 process
  1651. tables.insert(0, "search_process")
  1652. return tables
  1653. # ── 评估去重:复用 query 无关分,只重算 query 相关分(search_eval.py 用)──────────
  1654. def fetch_existing_eval(case_id, table="search_process"):
  1655. """返回该 case 在搜索表里最近一条「有效」评估 blob(任意 query)。
  1656. 评估去重用:同帖在别的相似 query 下评过时,复用其 query 无关分(质量/通用相关/时效),
  1657. 只重算「和 query 相关」。无有效评估(全是 _error 或没评过)返回 None。
  1658. 取最近若干条逐一挑出首个非 error、结构完整的 blob。"""
  1659. table = _search_table(table)
  1660. conn = _conn()
  1661. try:
  1662. with conn.cursor() as cur:
  1663. cur.execute(f"""SELECT llm_evaluation FROM {table}
  1664. WHERE case_id=%s AND llm_evaluation IS NOT NULL
  1665. ORDER BY updated_at DESC, id DESC LIMIT 5""", (case_id,))
  1666. rows = cur.fetchall()
  1667. finally:
  1668. conn.close()
  1669. for r in rows:
  1670. e = _loads(r["llm_evaluation"])
  1671. if isinstance(e, dict) and not e.get("_error") and isinstance(e.get("相关性"), dict):
  1672. return e
  1673. return None
  1674. def fetch_existing_eval_any(case_id):
  1675. """跨两张搜索表找该 case 最近一条有效评估 blob。
  1676. 评估与表无关(统一一套),任一表评过即可复用,避免同帖在两表各评一次。无则 None。"""
  1677. for table in ("search_process", "search_tools"):
  1678. e = fetch_existing_eval(case_id, table)
  1679. if e:
  1680. return e
  1681. return None
  1682. def update_post_eval(query_id, case_id, evaluation, table="search_process"):
  1683. """用新的评估 blob 覆盖某 (query, case) 行的 llm_evaluation,并同步重算派生列
  1684. overall_score、knowledge_type(口径同 upsert_search_posts)。返回受影响行数。"""
  1685. table = _search_table(table)
  1686. overall = overall_score(evaluation)
  1687. ktype = evaluation.get("知识类型") if isinstance(evaluation, dict) else None
  1688. conn = _conn()
  1689. try:
  1690. with conn.cursor() as cur:
  1691. n = cur.execute(
  1692. f"UPDATE {table} SET llm_evaluation=%s, overall_score=%s, knowledge_type=%s "
  1693. "WHERE query_id=%s AND case_id=%s",
  1694. (_j(evaluation), overall, _j(ktype), query_id, case_id))
  1695. return n
  1696. finally:
  1697. conn.close()
  1698. # ── 上传去重:知识库已导入台账(stages/import_process_knowledge.py 用)────────────────
  1699. def fetch_ingested_map(case_id):
  1700. """返回 {proc_index: version} —— 该 case 各工序已导入知识库的版本。空表示没传过。"""
  1701. conn = _conn()
  1702. try:
  1703. with conn.cursor() as cur:
  1704. cur.execute("SELECT proc_index, version FROM knowledge_ingest_log WHERE case_id=%s",
  1705. (case_id,))
  1706. return {r["proc_index"]: r["version"] for r in cur.fetchall()}
  1707. finally:
  1708. conn.close()
  1709. def mark_ingested(case_id, proc_index, version, knowledge_id=None, api_url=None):
  1710. """记一条「已导入」台账(case_id+proc_index 唯一,重导同序号则更新版本/knowledge_id)。"""
  1711. conn = _conn()
  1712. try:
  1713. with conn.cursor() as cur:
  1714. cur.execute("""INSERT INTO knowledge_ingest_log
  1715. (case_id, proc_index, version, knowledge_id, api_url)
  1716. VALUES (%s,%s,%s,%s,%s)
  1717. ON DUPLICATE KEY UPDATE version=VALUES(version),
  1718. knowledge_id=VALUES(knowledge_id), api_url=VALUES(api_url)""",
  1719. (case_id, proc_index, version, knowledge_id, api_url))
  1720. finally:
  1721. conn.close()
  1722. def fetch_tools_ingested_map(case_id):
  1723. """返回 {tool_index: version} —— 该 case 各工具已导入知识库的版本。空表示没传过。
  1724. 工具方向独立台账(tools_ingest_log),与工序的 knowledge_ingest_log 互不干扰。"""
  1725. conn = _conn()
  1726. try:
  1727. with conn.cursor() as cur:
  1728. cur.execute("SELECT tool_index, version FROM tools_ingest_log WHERE case_id=%s",
  1729. (case_id,))
  1730. return {r["tool_index"]: r["version"] for r in cur.fetchall()}
  1731. finally:
  1732. conn.close()
  1733. def mark_tools_ingested(case_id, tool_index, version, knowledge_id=None, api_url=None):
  1734. """记一条工具「已导入」台账(case_id+tool_index 唯一,重导同序号则更新版本/knowledge_id)。"""
  1735. conn = _conn()
  1736. try:
  1737. with conn.cursor() as cur:
  1738. cur.execute("""INSERT INTO tools_ingest_log
  1739. (case_id, tool_index, version, knowledge_id, api_url)
  1740. VALUES (%s,%s,%s,%s,%s)
  1741. ON DUPLICATE KEY UPDATE version=VALUES(version),
  1742. knowledge_id=VALUES(knowledge_id), api_url=VALUES(api_url)""",
  1743. (case_id, tool_index, version, knowledge_id, api_url))
  1744. finally:
  1745. conn.close()
  1746. def fetch_destruction_ingested_map(case_id):
  1747. """返回 {tool_index: version} —— 该 case 各工具的解构知识(what-制作还原)已导入的版本。
  1748. 独立台账(destruction_ingest_log),与工具方向 tools_ingest_log 互不干扰(见 DDL 注释)。"""
  1749. conn = _conn()
  1750. try:
  1751. with conn.cursor() as cur:
  1752. cur.execute("SELECT tool_index, version FROM destruction_ingest_log WHERE case_id=%s",
  1753. (case_id,))
  1754. return {r["tool_index"]: r["version"] for r in cur.fetchall()}
  1755. finally:
  1756. conn.close()
  1757. def mark_destruction_ingested(case_id, tool_index, version, knowledge_id=None, api_url=None):
  1758. """记一条解构知识「已导入」台账(case_id+tool_index 唯一,重导同序号则更新版本/knowledge_id)。"""
  1759. conn = _conn()
  1760. try:
  1761. with conn.cursor() as cur:
  1762. cur.execute("""INSERT INTO destruction_ingest_log
  1763. (case_id, tool_index, version, knowledge_id, api_url)
  1764. VALUES (%s,%s,%s,%s,%s)
  1765. ON DUPLICATE KEY UPDATE version=VALUES(version),
  1766. knowledge_id=VALUES(knowledge_id), api_url=VALUES(api_url)""",
  1767. (case_id, tool_index, version, knowledge_id, api_url))
  1768. finally:
  1769. conn.close()
  1770. def fetch_dashboard_rows():
  1771. """拉 Dashboard 计算所需的轻量行。数据量级:百~千行,Python 聚合足够。
  1772. 优化:① 不传 llm_evaluation 整块,SQL 只取采纳判定要的相关性得分;
  1773. ② steps 只取每个 case 的最新版本(覆盖度只看最新版),历史/link_ 版本不传 steps。"""
  1774. conn = _conn()
  1775. try:
  1776. with conn.cursor() as cur:
  1777. # 进度分母走「采纳」口径;mode 标方向(工序帖来自 search_process)。
  1778. cols = (f"query_id, case_id, platform, overall_score, publish_time, "
  1779. f"{_REL_SQL} AS rel, {_REPRO_SQL} AS repro")
  1780. cur.execute(f"SELECT {cols} FROM search_process WHERE {_REAL_POST}")
  1781. posts = cur.fetchall()
  1782. for p in posts:
  1783. p["mode"] = "process"
  1784. cur.execute(f"SELECT {cols} FROM search_tools")
  1785. st = cur.fetchall()
  1786. for p in st:
  1787. p["mode"] = "tools"
  1788. posts += st
  1789. # 3-II:成本/耗时按 run(=case+version)取自 process_run;steps 仅最新版(is_latest)需要,
  1790. # 其余 run steps=[] 省传输。返回形状与旧 mode_process 段一致(id/case_id/version/cost/dur/created_at/steps)。
  1791. cur.execute("""SELECT id, case_id, version, cost_usd, duration_s, created_at, is_latest
  1792. FROM process_run ORDER BY id""")
  1793. run_rows = cur.fetchall()
  1794. latest_ids = [r["id"] for r in run_rows if r["is_latest"]]
  1795. steps_by_run = {}
  1796. if latest_ids:
  1797. ph = ",".join(["%s"] * len(latest_ids))
  1798. cur.execute(f"""SELECT p.run_id, s.step_json FROM process_step s
  1799. JOIN process_procedure p ON s.procedure_pk = p.id
  1800. WHERE p.run_id IN ({ph})
  1801. ORDER BY p.run_id, p.seq, s.seq, s.id""", latest_ids)
  1802. for row in cur.fetchall():
  1803. steps_by_run.setdefault(row["run_id"], []).append(_loads(row["step_json"], {}))
  1804. procs = [{"id": r["id"], "case_id": r["case_id"], "version": r["version"],
  1805. "cost_usd": r["cost_usd"], "duration_s": r["duration_s"],
  1806. "created_at": r["created_at"],
  1807. "steps": steps_by_run.get(r["id"], [])} for r in run_rows]
  1808. cur.execute("""SELECT id, case_id, version, tool_name, substance_scope,
  1809. form_scope, cost_usd, duration_s, created_at
  1810. FROM mode_tools""")
  1811. tools = cur.fetchall()
  1812. finally:
  1813. conn.close()
  1814. for p in posts:
  1815. # 采纳判定:口径同帖子列表(is_adopted),作为「需解构」分母依据
  1816. p["adopted"] = is_adopted_rel(p["overall_score"], p["rel"], p["publish_time"], p["repro"])
  1817. for r in procs:
  1818. r["steps"] = _loads(r["steps"], [])
  1819. r["cost_usd"] = float(r["cost_usd"]) if r["cost_usd"] is not None else None
  1820. r["created_at"] = str(r["created_at"]) if r["created_at"] else None
  1821. for r in tools:
  1822. r["substance_scope"] = _loads(r["substance_scope"], [])
  1823. r["form_scope"] = _loads(r["form_scope"], [])
  1824. r["cost_usd"] = float(r["cost_usd"]) if r["cost_usd"] is not None else None
  1825. r["created_at"] = str(r["created_at"]) if r["created_at"] else None
  1826. return posts, procs, tools
  1827. def check():
  1828. conn = _conn()
  1829. try:
  1830. with conn.cursor() as cur:
  1831. for t in ("search_process", "search_tools", "mode_process", "mode_tools",
  1832. "process_run", "process_procedure", "process_step", "step_classification"):
  1833. cur.execute(f"SELECT COUNT(*) AS n FROM {t}")
  1834. print(f"{t}: {cur.fetchone()['n']} 行")
  1835. finally:
  1836. conn.close()
  1837. if __name__ == "__main__":
  1838. cmd = sys.argv[1] if len(sys.argv) > 1 else ""
  1839. if cmd == "init":
  1840. init_tables()
  1841. elif cmd == "check":
  1842. check()
  1843. elif cmd == "clear":
  1844. clear_tables()
  1845. else:
  1846. print("用法:\n python db.py init # 建表\n python db.py check # 四表行数\n python db.py clear # 清空四表数据")