Commit 3c26a078 authored by wangteng's avatar wangteng

整体调整

parent 6eac98c4
......@@ -14,13 +14,12 @@ web/outputs/
# ── 敏感配置(API Key / 真实凭据) ──
web/configs/llm.yaml
web/configs/db_defaults.yaml
web2/configs/db_defaults.yaml
web3/backend/configs/llm.yaml
web/backend/configs/llm.yaml
workflow/config.yaml
# ── web3 运行时数据 ──
web3/data/
web3/backend/logs/
# ── web 运行时数据 ──
web/data/
web/backend/logs/
# ── Python 字节码缓存 ──
**/__pycache__/
......@@ -45,4 +44,3 @@ work-logs/
tags.lock
tags.temp
Session.vim
web2/logs/
# 说明
- 这是一个数据治理工具,现在主要开发web端
- `./web/` 和 `./web2/` 是之前开发的工具,里面的数据链接等内容可以参考
- 现在要在`./web3`里面新开发
- 这是一个数据库工具,当前代码统一位于 `./web/`
# 总体要求
- 永远不要使用数据操作语句,本工具只做查询!
......
This diff is collapsed.
This diff is collapsed.
// 抓 Vue 编译出来的 render 函数,看 group/leaf 模板的产物
const fs = require('fs');
const { JSDOM } = require('jsdom');
const dom = new JSDOM('<!DOCTYPE html>', { runScripts: 'dangerously', url: 'http://localhost:8765/' });
const win = dom.window;
const captured = [];
const origFunction = win.Function;
win.Function = function(...args) {
const body = args[args.length - 1];
if (typeof body === 'string' && body.length > 5000) {
captured.push(body);
}
return new origFunction(...args);
};
// 加载 Vue
const vue = fs.readFileSync('web/static/lib/vue.global.prod.js', 'utf8');
win.eval(vue);
const Vue = win.Vue;
// 构造一个最小 app:模板就是 tree 那段
const tpl = `
<div class="checks-tree-table">
<div v-for="row in rows" :key="row.key" class="tree-row" :class="row.is_group ? 'tree-row--group' : 'tree-row--leaf'">
<template v-if="row.is_group">
<button class="tree-caret" :aria-expanded="row.exp" @click="onClick(row)">
<span v-if="row.exp">▼</span>
<span v-else>▶</span>
</button>
<span class="checkbox-stub" @click="onGroupCheck(row, $event)">
<span class="group-title">{{ row.title }}</span>
<span class="group-count">({{ row.cc }}/{{ row.tc }})</span>
<span v-if="row.description" class="group-description">{{ row.description }}</span>
</span>
</template>
<template v-else>
<span class="tree-caret-spacer" />
<span class="checkbox-stub" @click="onLeafCheck(row, $event)">
<span>{{ row.title }}</span>
</span>
</template>
</div>
</div>
`;
const rows = [
{key:'g1', is_group:true, exp:true, cc:2, tc:3, title:'基础检查', description:'基础质量检查'},
{key:'l1', is_group:false, title:'A'},
{key:'g2', is_group:true, exp:false, cc:0, tc:2, title:'国标字段规范', description:'国标'},
];
const app = Vue.createApp({
data() { return { rows }; },
methods: { onClick: ()=>{}, onGroupCheck: ()=>{}, onLeafCheck: ()=>{} },
template: tpl,
});
const container = win.document.createElement('div');
win.document.body.appendChild(container);
try {
app.mount(container);
console.log('=== 渲染出的 DOM ===');
console.log(container.innerHTML);
} catch (e) {
console.error('mount err:', e.message);
}
console.log('\n=== 抓到的 render fn 数量 ===', captured.length);
if (captured[0]) {
console.log('render fn #0 前 2000 字符:');
console.log(captured[0].substring(0, 2000));
console.log('--- 中间 2000 字符 ---');
console.log(captured[0].substring(captured[0].length/2, captured[0].length/2 + 2000));
}
\ No newline at end of file
This diff is collapsed.
"""临时端到端测试:验证 conn_json 全部字段解析 + 密码兜底"""
import urllib.request
import json
import sqlite3
BASE = "http://localhost:8767/api"
def req(method, path, body=None):
headers = {"Content-Type": "application/json"}
data = json.dumps(body, ensure_ascii=False).encode("utf-8") if body is not None else None
r = urllib.request.Request(f"{BASE}{path}", method=method, data=data, headers=headers)
with urllib.request.urlopen(r) as resp:
raw = resp.read()
return json.loads(raw) if raw else None
# 1) 创建
payload = {
"name": "测试任务_完整连接",
"group": "身份证",
"db_type": "MySQL",
"data_source_path": None,
"source_table": "t_user_info",
"status": "启用",
"description": "测试用",
"conn_json": json.dumps(
{
"connName": "prod_db",
"host": "192.168.1.100",
"port": "3306",
"db": "customer_db",
"user": "data_check",
"password": "secret123",
"jdbcParams": "useSSL=false",
"schema": "public",
"dbType": "MySQL",
"table": "t_user_info",
},
ensure_ascii=False,
),
}
created = req("POST", "/tasks", payload)
tid = created["id"]
print(f"[1] CREATE id={tid} name={created['name']}")
# 2) GET 解析字段
got = req("GET", f"/tasks/{tid}")
print("[2] GET parsed conn fields:")
for k in ["conn_name", "host", "port", "db", "user", "jdbc_params", "schema"]:
print(f" {k} = {got.get(k)!r}")
# 校验
assert got["conn_name"] == "prod_db"
assert got["host"] == "192.168.1.100"
assert got["port"] == "3306"
assert got["db"] == "customer_db"
assert got["user"] == "data_check"
assert got["jdbc_params"] == "useSSL=false"
assert got["schema"] == "public"
print(" [OK] all conn_* fields parsed correctly")
# 3) PUT 修改名 + description,密码留空(验证密码兜底)
# 注意:status UI 已移除,但 DB 列仍在,后端接受 status 字段
upd = {
"name": "测试任务_完整连接_v2",
"description": "改名测试",
"conn_json": json.dumps(
{
"connName": "prod_db",
"host": "192.168.1.100",
"port": "3306",
"db": "customer_db",
"user": "data_check",
"password": "",
"jdbcParams": "useSSL=false",
"schema": "public",
"dbType": "MySQL",
"table": "t_user_info",
},
ensure_ascii=False,
),
}
updated = req("PUT", f"/tasks/{tid}", upd)
print(f"[3] UPDATE name={updated['name']!r} description={updated['description']!r}")
# 4) 从 DB 读 conn_json,验证密码保留
con = sqlite3.connect("web3/data/web3.db")
row = con.execute("SELECT conn_json FROM task WHERE id=?", (tid,)).fetchone()
conn = json.loads(row[0])
print(f"[4] Stored password after empty PUT: {conn.get('password')!r}")
assert conn.get("password") == "secret123", "[FAIL] password wiped!"
print(" [OK] password preserved when front sends empty")
# 5) PUT 显式带新密码
upd["conn_json"] = json.dumps({**conn, "password": "new_secret_456"}, ensure_ascii=False)
_ = req("PUT", f"/tasks/{tid}", upd)
row = con.execute("SELECT conn_json FROM task WHERE id=?", (tid,)).fetchone()
conn2 = json.loads(row[0])
print(f"[5] Stored password after explicit PUT: {conn2.get('password')!r}")
assert conn2.get("password") == "new_secret_456", "[FAIL] password not updated"
print(" [OK] password updated when explicit value sent")
# 6) 清理
req("DELETE", f"/tasks/{tid}")
print("[6] cleaned up, task deleted")
print()
print("[ALL PASSED]")
\ No newline at end of file
"""持久化验证:写入 -> 重启后端 -> 读回
证明:任务的 CRUD 真的落到 SQLite 文件(web3/data/web3.db),重启不丢。
"""
import os
import signal
import sqlite3
import subprocess
import sys
import time
import urllib.request
import json
BASE = "http://localhost:8767/api"
DB_PATH = "web3/data/web3.db"
PID_FILE = "web3/backend/logs/backend.pid"
def http(method, path, body=None):
headers = {"Content-Type": "application/json"}
data = json.dumps(body, ensure_ascii=False).encode("utf-8") if body is not None else None
r = urllib.request.Request(f"{BASE}{path}", method=method, data=data, headers=headers)
with urllib.request.urlopen(r) as resp:
raw = resp.read()
return json.loads(raw) if raw else None
def wait_backend(timeout=20):
for _ in range(timeout):
try:
http("GET", "/health")
return True
except Exception:
time.sleep(1)
return False
def start_backend():
p = subprocess.Popen(
[sys.executable, "-m", "web3.backend.start"],
stdout=open("web3/backend/logs/console.log", "ab"),
stderr=subprocess.STDOUT,
)
with open(PID_FILE, "w") as f:
f.write(str(p.pid))
print(f" start backend, pid={p.pid}")
if not wait_backend():
raise RuntimeError("backend failed to start")
print(" backend ready")
def stop_backend():
pid = int(open(PID_FILE).read().strip())
print(f" stop backend, pid={pid}")
# 用 taskkill(Windows)保险
subprocess.run(["taskkill", "/F", "/PID", str(pid)], capture_output=True)
time.sleep(2)
# ---- 阶段 1:启动后端 ----
print("[1] start backend first time ...")
start_backend()
# ---- 阶段 2:建一个任务 ----
print("[2] POST /api/tasks with full conn_json ...")
payload = {
"name": "持久化测试_身份证表",
"group": "身份证",
"db_type": "MySQL",
"data_source_path": None,
"source_table": "t_user_info",
"status": "启用",
"description": "持久化测试用",
"conn_json": json.dumps({
"connName": "persist_test_db",
"host": "10.20.30.40",
"port": "3306",
"db": "test_db",
"user": "tester",
"password": "persist_pwd_888",
"jdbcParams": "useSSL=false",
"schema": "public",
"dbType": "MySQL",
"table": "t_user_info",
}, ensure_ascii=False),
}
created = http("POST", "/tasks", payload)
tid = created["id"]
print(f" created id={tid}")
# 直接读文件确认在 SQLite 里
con = sqlite3.connect(DB_PATH)
row = con.execute("SELECT name, status, db_type, conn_json FROM task WHERE id=?", (tid,)).fetchone()
assert row is not None, "row not in DB!"
print(f" [DB] task in SQLite: {row[0]!r} | status={row[1]} | db_type={row[2]}")
# ---- 阶段 3:停后端 ----
print("[3] stop backend ...")
stop_backend()
# ---- 阶段 4:重新启动后端 ----
print("[4] start backend second time ...")
start_backend()
# ---- 阶段 5:API 再读 ----
print("[5] GET /api/tasks/<id> after restart ...")
got = http("GET", f"/tasks/{tid}")
print(f" name = {got['name']}")
print(f" status = {got['status']}")
print(f" host = {got['host']}")
print(f" db = {got['db']}")
assert got["name"] == "持久化测试_身份证表"
assert got["host"] == "10.20.30.40"
assert got["status"] == "启用"
# 直接再读一次 SQLite 文件确认磁盘上
con = sqlite3.connect(DB_PATH)
row = con.execute("SELECT name, status, db_type, conn_json FROM task WHERE id=?", (tid,)).fetchone()
assert row is not None, "row lost after restart!"
conn = json.loads(row[3])
print(f" [DB] still in SQLite: {row[0]!r} | conn.password={conn['password']!r}")
assert conn["password"] == "persist_pwd_888"
print(" [OK] data survived restart")
# ---- 阶段 6:清理 ----
print("[6] cleanup ...")
http("DELETE", f"/tasks/{tid}")
print(f" task {tid} deleted")
# 关后端
stop_backend()
try:
os.remove(PID_FILE)
except Exception:
pass
print()
print("[ALL PASSED] task CRUD is fully persisted in SQLite")
\ No newline at end of file
"""web 包入口"""
__version__ = "1.0.0"
\ No newline at end of file
"""web.api 路由模块"""
from .routes import router
__all__ = ["router"]
\ No newline at end of file
This diff is collapsed.
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""Web 工具主入口
启动方式:
python -m web.app
# 或:
cd web && python app.py
环境变量:
ANTHROPIC_API_KEY Claude API Key(启用 LLM 增强)
WEB_PORT 端口(默认 8765)
WEB_HOST 主机(默认 0.0.0.0)
"""
from __future__ import annotations
import logging
import os
import sys
from pathlib import Path
# ── 把项目根目录加入 sys.path ──
WEB_DIR = Path(__file__).resolve().parent
PROJECT_ROOT = WEB_DIR.parent
sys.path.insert(0, str(PROJECT_ROOT))
from fastapi import FastAPI
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import FileResponse, HTMLResponse
from fastapi.staticfiles import StaticFiles
# ── 日志 ──
LOG_DIR = WEB_DIR / "logs"
LOG_DIR.mkdir(parents=True, exist_ok=True)
logging.basicConfig(
level=os.environ.get("LOG_LEVEL", "INFO"),
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
handlers=[
logging.FileHandler(LOG_DIR / "app.log", encoding="utf-8"),
logging.StreamHandler(),
],
)
logger = logging.getLogger(__name__)
logger.info("=" * 60)
logger.info("数据治理 Web 工具 · 进程启动")
logger.info(f"Python: {sys.version.split()[0]}, 平台: {sys.platform}")
logger.info(f"工作目录: {PROJECT_ROOT}")
logger.info(f"日志目录: {LOG_DIR}")
# ── 创建 FastAPI 应用 ──
app = FastAPI(
title="数据治理 Web 工具",
description="基于 FastAPI + Vue 3 的数据库治理平台(支持 MySQL / 达梦 / Oracle)",
version="1.0.0",
)
logger.info("FastAPI 应用已创建")
# CORS(开发模式)
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
logger.debug("CORS 中间件已挂载(允许所有来源)")
# ── 注册路由 ──
from web.api import router as api_router
app.include_router(api_router)
logger.info(f"已注册 {len(app.routes)} 个路由")
# ── 启动时打印 LLM 配置(方便确认 Key 是否生效) ──
@app.on_event("startup")
async def show_llm_config():
logger.info("FastAPI 应用启动事件触发")
from web.core.llm import LLMConfig
cfg = LLMConfig.load()
if cfg.available:
masked_key = cfg.api_key[:6] + "..." + cfg.api_key[-4:] if len(cfg.api_key) > 12 else "***"
print(f"[LLM] provider={cfg.provider} model={cfg.effective_model}")
print(f"[LLM] base_url={cfg.effective_base_url or '(default)'}")
print(f"[LLM] api_key={masked_key}(已配置)")
logger.info(f"LLM 配置: provider={cfg.provider}, model={cfg.effective_model}")
else:
print("[LLM] 未配置 API Key → LLM 增强不可用,将使用规则推理")
logger.warning("LLM 未配置 API Key → LLM 增强不可用,将使用规则推理")
@app.on_event("shutdown")
async def shutdown_log():
logger.info("FastAPI 应用关闭")
# ── 静态资源 / 入口页 ──
STATIC_DIR = WEB_DIR / "static"
if STATIC_DIR.exists():
app.mount("/static", StaticFiles(directory=str(STATIC_DIR)), name="static")
@app.get("/", response_class=HTMLResponse)
async def root():
"""服务入口 HTML"""
index = STATIC_DIR / "index.html"
if index.exists():
return HTMLResponse(index.read_text(encoding="utf-8"))
return HTMLResponse("<h1>前端文件未找到</h1><p>请检查 web/static/index.html</p>")
# ── 启动 ──
def main():
import uvicorn
host = os.environ.get("WEB_HOST", "0.0.0.0")
port = int(os.environ.get("WEB_PORT", "8765"))
print("=" * 60)
print(f" 数据治理 Web 工具")
print(f" 访问: http://localhost:{port}")
print(f" API 文档: http://localhost:{port}/docs")
print("=" * 60)
uvicorn.run("web.app:app", host=host, port=port, reload=False, log_level="info")
if __name__ == "__main__":
main()
\ No newline at end of file
"""web3 后端统一日志配置
"""web 后端统一日志配置
从 web2/backend/_logging.py 拷贝改造(路径前缀 web2 → web3)。
项目日志配置。
- 重复调用幂等
- 所有 web3.backend.* 日志同时写 stdout + app.log
- 所有 web.backend.* 日志同时写 stdout + app.log
- 不接管 root / web.core.* logger
"""
from __future__ import annotations
......@@ -12,8 +12,8 @@ from logging.handlers import RotatingFileHandler
from pathlib import Path
# ── 路径常量 ──
WEB3_DIR = Path(__file__).resolve().parent
LOG_DIR = WEB3_DIR / "logs"
WEB_DIR = Path(__file__).resolve().parent
LOG_DIR = WEB_DIR / "logs"
LOG_DIR.mkdir(parents=True, exist_ok=True)
APP_LOG_FILE = LOG_DIR / "app.log"
......@@ -22,12 +22,12 @@ _DATE_FORMAT = None
def setup_logging(level: str = "INFO") -> None:
web3_logger = logging.getLogger("web3.backend")
if web3_logger.handlers:
web_logger = logging.getLogger("web.backend")
if web_logger.handlers:
return
web3_logger.setLevel(getattr(logging, level.upper(), logging.INFO))
web3_logger.propagate = False
web_logger.setLevel(getattr(logging, level.upper(), logging.INFO))
web_logger.propagate = False
formatter = logging.Formatter(_LOG_FORMAT, _DATE_FORMAT)
......@@ -38,11 +38,11 @@ def setup_logging(level: str = "INFO") -> None:
encoding="utf-8",
)
fh.setFormatter(formatter)
web3_logger.addHandler(fh)
web_logger.addHandler(fh)
ch = logging.StreamHandler()
ch.setFormatter(formatter)
web3_logger.addHandler(ch)
web_logger.addHandler(ch)
root = logging.getLogger()
if not root.handlers:
......@@ -50,6 +50,6 @@ def setup_logging(level: str = "INFO") -> None:
def get_logger(name: str) -> logging.Logger:
if not name.startswith("web3."):
name = f"web3.{name}"
return logging.getLogger(name)
\ No newline at end of file
if not name.startswith("web."):
name = f"web.{name}"
return logging.getLogger(name)
"""web3 后端 · FastAPI 入口"""
"""web 后端 · FastAPI 入口"""
from __future__ import annotations
import os
......@@ -6,8 +6,8 @@ import sys
from pathlib import Path
BACKEND_DIR = Path(__file__).resolve().parent
WEB3_DIR = BACKEND_DIR.parent
PROJECT_ROOT = WEB3_DIR.parent
WEB_DIR = BACKEND_DIR.parent
PROJECT_ROOT = WEB_DIR.parent
sys.path.insert(0, str(PROJECT_ROOT))
try:
......@@ -20,25 +20,25 @@ from fastapi import FastAPI
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import HTMLResponse
from web3.backend._logging import setup_logging, get_logger
from web.backend._logging import setup_logging, get_logger
setup_logging(os.environ.get("LOG_LEVEL", "INFO"))
logger = get_logger("backend")
logger.info("=" * 60)
logger.info("web3 后端 · 进程启动")
logger.info("web 后端 · 进程启动")
logger.info(f"Python: {sys.version.split()[0]}, 平台: {sys.platform}")
logger.info(f"工作目录: {PROJECT_ROOT}")
from web3.backend.db.database import init_db, DB_PATH # noqa: E402
from web3.backend import models # noqa: F401, E402 (必须 import 所有 ORM 才能被 Base.metadata 知道)
from web.backend.db.database import init_db, DB_PATH # noqa: E402
from web.backend import models # noqa: F401, E402 (必须 import 所有 ORM 才能被 Base.metadata 知道)
init_db()
logger.info(f"DB: {DB_PATH} (存在={DB_PATH.exists()})")
app = FastAPI(
title="数据治理 Web 工具 · web3 后端",
description="数据质量检测(任务 / 字段 / 规则 / 不合规)。前端在 web3/,由 Vite 5175 提供。",
title="数据治理 Web 工具 · web 后端",
description="数据质量检测(任务 / 字段 / 规则 / 不合规)。前端在 web/,由 Vite 5175 提供。",
version="0.1.0",
)
......@@ -51,30 +51,36 @@ app.add_middleware(
)
from web3.backend.routers.ai import router as ai_router # noqa: E402
from web3.backend.routers.connection_presets import router as connection_presets_router # noqa: E402
from web3.backend.routers.db import router as db_router # noqa: E402
from web3.backend.routers.queries import router as queries_router # noqa: E402
from web3.backend.routers.task_groups import router as task_groups_router # noqa: E402
from web3.backend.routers.tasks import router as tasks_router # noqa: E402
from web.backend.routers.ai import router as ai_router # noqa: E402
from web.backend.routers.connection_presets import router as connection_presets_router # noqa: E402
from web.backend.routers.db import router as db_router # noqa: E402
from web.backend.routers.queries import router as queries_router # noqa: E402
from web.backend.routers.task_groups import router as task_groups_router # noqa: E402
from web.backend.routers.tasks import router as tasks_router # noqa: E402
from web.backend.routers.data_sources import router as data_sources_router # noqa: E402
from web.backend.routers.operations import router as operations_router # noqa: E402
from web.backend.routers.rules import router as rules_router # noqa: E402
app.include_router(task_groups_router, prefix="/api")
app.include_router(tasks_router, prefix="/api")
app.include_router(db_router, prefix="/api")
app.include_router(connection_presets_router, prefix="/api")
app.include_router(ai_router, prefix="/api")
app.include_router(queries_router, prefix="/api")
logger.info("已注册路由:/api/task-groups, /api/tasks, /api/connect/*, /api/connection-presets, /api/ai/*, /api/queries/*")
app.include_router(data_sources_router, prefix="/api")
app.include_router(operations_router, prefix="/api")
app.include_router(rules_router, prefix="/api")
logger.info("已注册路由:/api/task-groups, /api/tasks, /api/data-sources, /api/connect/*, /api/connection-presets, /api/ai/*, /api/queries/*")
@app.get("/api/health")
async def health():
return {"status": "ok", "service": "web3-backend"}
return {"status": "ok", "service": "web-backend"}
@app.get("/", response_class=HTMLResponse)
async def root():
return HTMLResponse(
"<h1>web3 后端 API 服务</h1>"
"<h1>web 后端 API 服务</h1>"
"<p>开发模式请用 <code>npm run dev</code> 启动 Vite(端口 5175)。</p>"
"<p>API 文档:<a href='/docs'>/docs</a></p>"
)
"""web 后端路径与运行参数常量
- WEB_DB_PATH:主业务 SQLite 文件(任务/字段/规则/不合规)
- WEB_SOURCES_DIR:用户数据源 SQLite 所在目录(运行时只读扫描)
- WEB_PORT / WEB_HOST:FastAPI 监听端口(默认 8767)
"""
from __future__ import annotations
import os
from pathlib import Path
# backend/ 的父目录 = web/
WEB_DIR = Path(__file__).resolve().parent
# web 的父目录 = 项目根(含 data/)
PROJECT_ROOT = WEB_DIR.parent
DATA_DIR = PROJECT_ROOT / "data"
DATA_DIR.mkdir(parents=True, exist_ok=True)
SOURCES_DIR = DATA_DIR / "sources"
SOURCES_DIR.mkdir(parents=True, exist_ok=True)
DB_PATH = Path(os.environ.get("WEB_DB_PATH") or (DATA_DIR / "web.db"))
PORT = int(os.environ.get("WEB_PORT", "8767"))
HOST = os.environ.get("WEB_HOST", "0.0.0.0")
LOG_LEVEL = os.environ.get("LOG_LEVEL", "INFO")
"""AI 解释不合规数据(web3 / 2026-08-24)
"""AI 解释不合规数据(web / 2026-08-24)
入口 `explain_issue()`:把数据值 + 字段信息 + 规则 + 规则代码一起交给 LLM,让它
用自然语言解释「这条数据为什么不合规」,并审查「规则描述与实现代码是否一致」。
......@@ -15,9 +15,9 @@ from __future__ import annotations
from typing import Any, Optional
from web3.backend.core.llm import LLMUnavailable, get_llm_client
from web3.backend.models.rule import RULE_TYPES
from web3.backend._logging import get_logger
from web.backend.core.llm import LLMUnavailable, get_llm_client
from web.backend.models.rule import RULE_TYPES
from web.backend._logging import get_logger
logger = get_logger("backend.ai_explain")
......
"""AI 生成正则服务(包装 web3.backend.core.llm)
"""AI 生成正则服务(包装 web.backend.core.llm)
强制走 LLM,不做规则库兜底(2026-08-20 用户决策):
- 每次请求都过 LLM,让模型能见到真实描述(持续学习用户表述)
......@@ -14,8 +14,8 @@ from __future__ import annotations
import re
from web3.backend.core.llm import LLMUnavailable, get_llm_client
from web3.backend._logging import get_logger
from web.backend.core.llm import LLMUnavailable, get_llm_client
from web.backend._logging import get_logger
logger = get_logger("backend.ai_regex")
......
"""AI 生成规则代码(web3 / 2026-08-21)
"""AI 生成规则代码(web / 2026-08-21)
四种 rule_type 统一入口 `gen_rule()`:
......@@ -19,20 +19,20 @@ from __future__ import annotations
import re
from typing import Optional
from web3.backend.core.llm import LLMUnavailable, get_llm_client
from web3.backend.core.rule_runner import (
from web.backend.core.llm import LLMUnavailable, get_llm_client
from web.backend.core.rule_runner import (
RuleRunError,
run_rule,
validate_user_code,
)
from web3.backend.models.rule import RULE_TYPES
from web3.backend._logging import get_logger
from web.backend.models.rule import RULE_TYPES
from web.backend._logging import get_logger
logger = get_logger("backend.ai_rule")
# ── 复用老 ai_regex ─────────────────────────────────────
from web3.backend.core.ai_regex import gen_regex, test_regex # noqa: E402,F401
from web.backend.core.ai_regex import gen_regex, test_regex # noqa: E402,F401
# ── System Prompt ───────────────────────────────────────
......
"""本地备份文件存储与路径配置。"""
from __future__ import annotations
import json
from pathlib import Path
from web.backend.config import DATA_DIR
DEFAULT_BACKUP_DIR = DATA_DIR / "backups"
CONFIG_FILE = DATA_DIR / "backup-config.json"
def get_backup_dir() -> Path:
path = DEFAULT_BACKUP_DIR
try:
data = json.loads(CONFIG_FILE.read_text(encoding="utf-8"))
configured = data.get("backup_dir")
if configured:
path = Path(configured).expanduser()
except (OSError, ValueError, TypeError):
pass
path = path.resolve()
path.mkdir(parents=True, exist_ok=True)
return path
def set_backup_dir(value: str) -> Path:
if not isinstance(value, str) or not value.strip():
raise ValueError("备份目录不能为空")
path = Path(value.strip()).expanduser().resolve()
path.mkdir(parents=True, exist_ok=True)
CONFIG_FILE.write_text(json.dumps({"backup_dir": str(path)}, ensure_ascii=False, indent=2), encoding="utf-8")
return path
def safe_backup_file(filename: str) -> Path:
name = Path(filename or "").name
if not name or name != filename or name in (".", ".."):
raise ValueError("备份文件名无效")
path = (get_backup_dir() / name).resolve()
if path.parent != get_backup_dir():
raise ValueError("备份文件路径无效")
return path
"""数据库适配层(web3 自包含版)
"""数据库适配层(web 自包含版)
来源:拷贝自 web/core/db_adapter.py(2026-08-20),改为自包含:
- SQL 模板加载器从 web3.backend.core.sql_loader 导入
- SQL 模板加载器从 web.backend.core.sql_loader 导入
- 不再依赖 web.core / web.sql
提供 MySQL / 达梦 / Oracle 的统一连接与查询接口,屏蔽方言差异:
......@@ -51,16 +51,20 @@ class DBConfig:
oracle_client_dir: Optional[str] = None
def to_pymysql_kwargs(self) -> dict:
return {
kwargs = {
"host": self.host,
"port": self.port,
"user": self.user,
"password": self.password,
"database": self.database,
"charset": self.charset,
"connect_timeout": self.connect_timeout,
"cursorclass": DictCursor,
}
# MySQL supports connecting without selecting a default database;
# this is required by the schema picker which lists databases first.
if self.database:
kwargs["database"] = self.database
return kwargs
def to_dm_kwargs(self) -> dict:
"""dmPython 关键字参数映射
......@@ -212,20 +216,20 @@ def normalize_rows(rows: list[dict] | None) -> list[dict]:
# ── 信息模式查询(通过 SQL 加载器,无内联 SQL) ──────────
# 所有 SQL 已抽到 web3/backend/core/sql_templates/ 目录下的 .sql 模板
# 所有 SQL 已抽到 web/backend/core/sql_templates/ 目录下的 .sql 模板
# schema 一律走绑定参数(MySQL %s / 达梦 ?),不拼进 SQL 文本
def _render_info_columns_sql(dialect: str) -> str:
from web3.backend.core.sql_loader import get_sql_loader
from web.backend.core.sql_loader import get_sql_loader
return get_sql_loader().render("info_schema/list_columns", dialect=dialect)
def _render_info_tables_sql(dialect: str) -> str:
from web3.backend.core.sql_loader import get_sql_loader
from web.backend.core.sql_loader import get_sql_loader
return get_sql_loader().render("info_schema/list_tables", dialect=dialect)
def _render_health_sql(dialect: str) -> str:
from web3.backend.core.sql_loader import get_sql_loader
from web.backend.core.sql_loader import get_sql_loader
return get_sql_loader().render("health/check_connection", dialect=dialect)
......@@ -398,12 +402,12 @@ DATE_TYPES = {
class DBConnection:
"""统一的数据库连接封装。
用法(SQL 一律走 web3/backend/core/sql_templates/ 模板,不要在 Python 里拼 SQL):
用法(SQL 一律走 web/backend/core/sql_templates/ 模板,不要在 Python 里拼 SQL):
with DBConnection(cfg) as db:
rows = db.list_columns("mydb") # schema 走绑定参数
# 需要自定义模板时,值用 params 传,标识符才用 ${x | quote}
from web3.backend.core.sql_loader import get_sql_loader
from web.backend.core.sql_loader import get_sql_loader
sql = get_sql_loader().render("verify/count_table_rows",
dialect=cfg.db_type, table="t_project")
n = db.fetch_scalar(sql)
......@@ -627,11 +631,21 @@ class DBConnection:
finally:
cur.close()
def commit(self) -> None:
"""提交写入操作(备份还原等写库流程使用)。"""
if self._conn is None:
raise RuntimeError("DBConnection 未打开,请用 with 语句")
self._conn.commit()
def rollback(self) -> None:
if self._conn is not None:
self._conn.rollback()
# ── 信息模式便捷方法(schema 走绑定参数,不拼接) ──────
def list_columns(self, schema: str, table_name: str) -> list[dict]:
"""返回指定 schema 下指定表的所有列的元数据
SQL 通过 web3/backend/core/sql_templates/info_schema/list_columns.<dialect>.sql 加载,
SQL 通过 web/backend/core/sql_templates/info_schema/list_columns.<dialect>.sql 加载,
schema / table_name 都作为绑定参数传给 cursor,由驱动层负责转义,杜绝 SQL 注入。
"""
sql = _render_info_columns_sql(self.cfg.db_type)
......@@ -675,4 +689,4 @@ def test_connection(cfg: DBConfig) -> tuple[bool, str]:
return False, "连接成功但查询返回非预期值"
except Exception as e:
logger.error(f"[DB] 测试连接失败: {type(e).__name__}: {e}")
return False, f"连接失败: {e}"
\ No newline at end of file
return False, f"连接失败: {e}"
"""LLM 客户端封装(支持 MiniMax / Anthropic / OpenAI)
直接拷贝自 web/core/llm.py,2026-08-20 接入 web3 用作「AI 生成正则」。
改动点:default_config_path() 改成读 web3/backend/configs/llm.yaml,其他不动,
保持 web3 子项目自包含(不依赖 web.*)。
直接拷贝自 web/core/llm.py,2026-08-20 接入 web 用作「AI 生成正则」。
改动点:default_config_path() 改成读 web/backend/configs/llm.yaml,其他不动,
保持 web 子项目自包含(不依赖 web.*)。
所有方法都设计为可降级 — 如果 LLM 不可用(无 API Key / 超时 / 解析失败),返回 None 或 fallback 值,
不阻断主流程。
......@@ -14,7 +14,7 @@
配置优先级(从高到低):
1. 环境变量(最高):LLM_PROVIDER / LLM_API_KEY / LLM_BASE_URL / LLM_MODEL
2. 配置文件(中间):web3/backend/configs/llm.yaml
2. 配置文件(中间):web/backend/configs/llm.yaml
3. 内置默认值(最低)
启动时 LLM 客户端初始化失败会自动降级,主流程仍可继续。
......@@ -53,7 +53,7 @@ PROVIDER_PRESETS: dict[str, dict] = {
# ── 配置文件路径 ─────────────────────────────────────────
def default_config_path() -> Path:
"""默认 LLM 配置文件路径:web3/backend/configs/llm.yaml
"""默认 LLM 配置文件路径:web/backend/configs/llm.yaml
可通过环境变量 LLM_CONFIG_FILE 覆盖:
export LLM_CONFIG_FILE=/path/to/my-llm.yaml
......@@ -61,7 +61,7 @@ def default_config_path() -> Path:
env_path = os.environ.get("LLM_CONFIG_FILE")
if env_path:
return Path(env_path)
# web3/backend/core/llm.py → 往上两级是 web3/backend/,再拼 configs/llm.yaml
# web/backend/core/llm.py → 往上两级是 web/backend/,再拼 configs/llm.yaml
return Path(__file__).resolve().parent.parent / "configs" / "llm.yaml"
......@@ -181,7 +181,7 @@ class LLMClient:
else:
logger.info(
f"LLM 未配置 API Key(provider={self.cfg.provider}),"
f"如需启用请在 web3/backend/configs/llm.yaml 设置 api_key 或配置环境变量"
f"如需启用请在 web/backend/configs/llm.yaml 设置 api_key 或配置环境变量"
)
def _init_provider(self):
......@@ -225,9 +225,10 @@ class LLMClient:
kwargs: dict[str, Any] = {
"model": self.cfg.model,
"max_tokens": self.cfg.max_tokens,
"temperature": self.cfg.temperature,
"messages": [{"role": "user", "content": prompt}],
}
# 当前 Anthropic SDK 1.x 的 Messages API 已不接收 temperature。
# MiniMax 走同一兼容接口,传入会在本地 SDK 参数校验阶段失败。
if system:
kwargs["system"] = system
msg = self._provider.messages.create(**kwargs)
......@@ -423,4 +424,4 @@ def get_llm_client() -> LLMClient:
def reset_llm_client() -> None:
"""重置单例(配置变更时调用)"""
global _client
_client = None
\ No newline at end of file
_client = None
"""web3 数据库连接 Pydantic 模型
"""web 数据库连接 Pydantic 模型
来源:精简自 web/core/models.py(2026-08-20),只保留数据库连接所需的三个类。
"""
......
"""规则执行引擎(web3)
"""规则执行引擎(web)
四种 rule_type 走同一个入口 `run_rule()`,内部按类型 dispatch:
......
"""查询 session 管理器(web3 自包含版)
"""查询 session 管理器(web 自包含版)
为每次 ``POST /api/queries/start`` 在内存里建一个 ``QuerySession``,记录任务配置 +
DB 连接 + 分页参数的快照。前端按页号去 ``POST /api/queries/page`` 拉一页结果,可并发
......@@ -24,7 +24,7 @@ import time
from dataclasses import dataclass, field
from typing import Optional
from web3.backend.core.db_adapter import DBConfig
from web.backend.core.db_adapter import DBConfig
# ── 模块级存储 ────────────────────────────────────────────
......
"""SQL 模板加载器(web3 自包含版)
"""SQL 模板加载器(web 自包含版)
来源:精简自 web/sql/loader.py(2026-08-20)。
变化点:
- 模板根目录改为 web3/backend/core/sql_templates/
- 移除 `${var | quote}` filter(web3 用法不需要;quote 由调用方在 Python 侧处理)
- 模板根目录改为 web/backend/core/sql_templates/
- 移除 `${var | quote}` filter(web 用法不需要;quote 由调用方在 Python 侧处理)
- 移除对 web.core.db_adapter.quote_ident 的依赖
把所有 SQL 语句集中到 web3/backend/core/sql_templates/ 目录下,Python 代码不再直接拼 SQL。
把所有 SQL 语句集中到 web/backend/core/sql_templates/ 目录下,Python 代码不再直接拼 SQL。
模板使用 ${var} 占位符(避免与 SQL 自身的 {} 冲突)。
支持的占位符:
......@@ -86,7 +86,7 @@ class SQLLoader:
"""SQL 模板加载器(单例)"""
def __init__(self, base_dir: Path | None = None):
# web3/backend/core/sql_loader.py → web3/backend/core/sql_templates/ 是模板根目录
# web/backend/core/sql_loader.py → web/backend/core/sql_templates/ 是模板根目录
self.base_dir = base_dir or (Path(__file__).resolve().parent / "sql_templates")
self._cache: dict[str, str] = {}
self._lock = Lock()
......
-- ============================================================================
-- 列出指定 schema 下指定表的列元数据 (达梦方言)
-- 调用方:web3/backend/core/db_adapter.py → list_columns()
-- 调用方:web/backend/core/db_adapter.py → list_columns()
-- 参数(按位占位符):
-- [0] schema 数据库/模式名
-- [1] table_name 表名
......
-- ============================================================================
-- 列出指定 schema 下指定表的列元数据 (MySQL 方言)
-- 调用方:web3/backend/core/db_adapter.py → list_columns()
-- 调用方:web/backend/core/db_adapter.py → list_columns()
-- 参数(均为绑定参数,由调用方通过 params 传入,勿拼接):
-- [0] schema 数据库/模式名
-- [1] table_name 表名
......
-- ============================================================================
-- 列出指定 schema 下指定表的列元数据 (Oracle 方言)
-- 调用方:web3/backend/core/db_adapter.py → list_columns()
-- 调用方:web/backend/core/db_adapter.py → list_columns()
-- 参数(命名占位符 :1/:2):
-- :1 schema 用户/模式名(调用方传入的可能是小写,这里 UPPER() 兼容)
-- :2 table_name 表名
......
"""web3 后端 db 子包"""
"""web 后端 db 子包"""
from .database import engine, SessionLocal, init_db, get_session
__all__ = ["engine", "SessionLocal", "init_db", "get_session"]
\ No newline at end of file
......@@ -12,8 +12,8 @@ from pathlib import Path
from sqlalchemy import create_engine, event
from sqlalchemy.orm import sessionmaker, Session, declarative_base
from web3.backend.config import DB_PATH
from web3.backend._logging import get_logger
from web.backend.config import DB_PATH
from web.backend._logging import get_logger
logger = get_logger("backend.db")
......@@ -68,7 +68,7 @@ def init_db() -> None:
DB 已存在 → 用 SQLAlchemy 补齐缺失的表(增量迁移,不破坏已有数据),
然后 seed 各码表(每个 seed 自己判重,幂等)。
"""
from web3.backend.db import seed
from web.backend.db import seed
if not DB_PATH.exists():
# 首次建库:跑 schema.sql 全量
schema = Path(__file__).parent / "schema.sql"
......@@ -103,10 +103,13 @@ def init_db() -> None:
_migrate_task_checked_at()
# 一次性迁移:task 表加 table_comment 列(2026-08-21 接「数据表注释」UI)
_migrate_task_table_comment()
_migrate_task_data_source_id()
_migrate_legacy_rules_to_library()
logger.info(f"[init_db] DB 已存在,已执行增量建表检查:{DB_PATH}")
# 灌种子(每个 seed 函数内部判重,可重复调用)
seed.seed_task_groups()
seed.seed_connection_presets()
seed.seed_data_sources()
print(f"[init_db] 初始化完成:{DB_PATH}")
......@@ -255,4 +258,45 @@ def _migrate_task_table_comment() -> None:
}
if "table_comment" not in cols:
logger.info("[migrate] task 缺 table_comment 列,补建(默认 NULL,老任务「—」)")
conn.exec_driver_sql("ALTER TABLE task ADD COLUMN table_comment TEXT")
\ No newline at end of file
conn.exec_driver_sql("ALTER TABLE task ADD COLUMN table_comment TEXT")
def _migrate_task_data_source_id() -> None:
"""给 task 增加可复用数据源引用。"""
with engine.begin() as conn:
cols = {row[1] for row in conn.exec_driver_sql("PRAGMA table_info(task)").fetchall()}
if "data_source_id" not in cols:
conn.exec_driver_sql("ALTER TABLE task ADD COLUMN data_source_id INTEGER")
def _migrate_legacy_rules_to_library() -> None:
"""将旧 field.rule 记录登记到规则库,保证历史校验任务可继续编辑/扫描。"""
from web.backend.models.field import Field
from web.backend.models.rule import Rule
from web.backend.models.rule_library import FieldRule, RuleLibrary
with SessionLocal() as db:
legacy_rules = db.query(Rule).all()
added = 0
for old in legacy_rules:
if db.query(FieldRule).filter(FieldRule.field_id == old.field_id).first():
continue
base = (old.name or old.desc or f"历史规则 {old.id}").strip()[:110]
name = base or f"历史规则 {old.id}"
if db.query(RuleLibrary).filter(RuleLibrary.name == name).first():
name = f"{name} ({old.id})"
lib = RuleLibrary(
name=name,
description=old.desc or "",
rule_type=old.rule_type or "regex",
regex=old.regex,
code=old.code,
skip_null=old.skip_null or 0,
)
db.add(lib)
db.flush()
db.add(FieldRule(field_id=old.field_id, rule_library_id=lib.id, ord=old.ord or 0))
added += 1
if added:
db.commit()
logger.info(f"[migrate] 已将 {added} 条历史规则登记到规则库")
-- ============================================================
-- web3 主业务库 · schema.sql
-- web 主业务库 · schema.sql
-- ============================================================
-- 数据库:data/web3.db(SQLite)
-- 数据库:data/web.db(SQLite)
-- 维护方式:DB 文件首次生成时由 db/database.py::init_db() 全量执行
-- 注意:本文件是「从零建表」脚本,已存在 DB 时不会执行。
-- 表结构变更后续走迁移(plan §11 待办)。
......@@ -35,6 +35,7 @@ CREATE TABLE task (
name TEXT NOT NULL UNIQUE,
group_id INTEGER NOT NULL REFERENCES task_group(id),
data_source_path TEXT,
data_source_id INTEGER REFERENCES data_source(id),
source_table TEXT,
status TEXT NOT NULL DEFAULT '启用',
description TEXT,
......@@ -93,6 +94,27 @@ CREATE TABLE rule (
);
CREATE INDEX idx_rule_field ON rule(field_id);
-- 独立规则库;校验任务通过 field_rule 选择可复用规则。
CREATE TABLE rule_library (
id INTEGER PRIMARY KEY AUTOINCREMENT,
name TEXT NOT NULL UNIQUE,
description TEXT NOT NULL DEFAULT '',
rule_type TEXT NOT NULL DEFAULT 'regex',
regex TEXT,
code TEXT,
skip_null INTEGER NOT NULL DEFAULT 0,
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP
);
CREATE TABLE field_rule (
id INTEGER PRIMARY KEY AUTOINCREMENT,
field_id INTEGER NOT NULL REFERENCES field(id) ON DELETE CASCADE,
rule_library_id INTEGER NOT NULL REFERENCES rule_library(id) ON DELETE RESTRICT,
ord INTEGER NOT NULL DEFAULT 0,
UNIQUE(field_id, rule_library_id)
);
CREATE INDEX idx_field_rule_field ON field_rule(field_id);
-- ─────────────────────────────────────────────────────────────
-- 5) validation_run · 一次检测任务(历史持久化;运行时进度走内存)
-- ─────────────────────────────────────────────────────────────
......@@ -150,4 +172,23 @@ CREATE TABLE connection_preset (
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
UNIQUE(db_type, name)
);
CREATE INDEX idx_connection_preset_db_type ON connection_preset(db_type, ord);
\ No newline at end of file
CREATE INDEX idx_connection_preset_db_type ON connection_preset(db_type, ord);
-- 8) data_source · 可复用的数据源连接
CREATE TABLE data_source (
id INTEGER PRIMARY KEY AUTOINCREMENT,
name TEXT NOT NULL UNIQUE,
db_type TEXT NOT NULL DEFAULT 'MySQL',
host TEXT NOT NULL,
port INTEGER NOT NULL,
database TEXT NOT NULL DEFAULT '',
schema TEXT,
user TEXT NOT NULL,
password TEXT NOT NULL DEFAULT '',
oracle_client_dir TEXT,
status TEXT NOT NULL DEFAULT '启用',
last_tested_at TIMESTAMP,
last_test_ok INTEGER,
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP
);
......@@ -2,42 +2,35 @@
- seed_task_groups():6 个分组(来自 mock 需求里的左侧菜单「校验任务」)
- seed_connection_presets():13 条历史连接(来自 web2/configs/db_defaults.yaml,
硬编码避免引入 pyyaml 依赖;改完需要重启 web3 后端生效)
硬编码避免引入 pyyaml 依赖;改完需要重启 web 后端生效)
"""
from __future__ import annotations
import logging
from web3.backend.db.database import session_scope
from web3.backend.models.connection_preset import ConnectionPreset
from web3.backend.models.task_group import TaskGroup
from web.backend.db.database import session_scope
from web.backend.models.connection_preset import ConnectionPreset
from web.backend.models.task_group import TaskGroup
from web.backend.models.data_source import DataSource
logger = logging.getLogger(__name__)
# ── DML ─────────────────────────────────────────────────────
# 来自 mock web3/requirement/数据质量检测结果.html 里左侧菜单「校验任务」下的 6 个分组
# 来自 mock web/requirement/数据质量检测结果.html 里左侧菜单「校验任务」下的 6 个分组
# 注意:不要把「任务配置」加进来 —— 它是导航项不是分组
DEFAULT_TASK_GROUPS = [
("身份证", 10),
("手机号", 20),
("邮箱", 30),
("地址信息", 40),
("银行卡号", 50),
("日期格式", 60),
]
DEFAULT_TASK_GROUPS = [("校验", 10)]
def seed_task_groups() -> None:
"""首次建库时插入 6 个默认分组;已存在则跳过。"""
"""确保新版校验任务分组存在;历史分组保留给已存任务。"""
with session_scope() as db:
existed = db.query(TaskGroup).count()
if existed > 0:
logger.info(f"[seed] task_group 已有 {existed} 条,跳过")
return
added = 0
for name, ord in DEFAULT_TASK_GROUPS:
db.add(TaskGroup(name=name, ord=ord))
logger.info(f"[seed] task_group 新增 {len(DEFAULT_TASK_GROUPS)} 条")
if not db.query(TaskGroup).filter(TaskGroup.name == name).first():
db.add(TaskGroup(name=name, ord=ord)); added += 1
if added:
logger.info(f"[seed] task_group 新增 {added} 条")
# ── 历史连接预设(来自 web2/configs/db_defaults.yaml,13 条) ──
......@@ -92,3 +85,23 @@ def seed_connection_presets() -> None:
ord=cfg["ord"],
))
logger.info(f"[seed] connection_preset 新增 {len(DEFAULT_CONNECTION_PRESETS)} 条")
# 数据源页面初始示例,使用本机地址避免误连真实生产库;用户可在页面编辑。
DEFAULT_DATA_SOURCES = [
{"name": "本地开发库", "db_type": "MySQL", "host": "127.0.0.1", "port": 3306, "database": "demo", "schema": "demo", "user": "root", "password": ""},
{"name": "测试数据源", "db_type": "MySQL", "host": "127.0.0.1", "port": 3307, "database": "test", "schema": "test", "user": "test", "password": ""},
{"name": "Oracle 示例", "db_type": "Oracle", "host": "127.0.0.1", "port": 1521, "database": "ORCL", "schema": "DEMO", "user": "demo", "password": ""},
]
def seed_data_sources() -> None:
"""插入示例数据源;按名称判重,用户已有数据不受影响。"""
with session_scope() as db:
added = 0
for cfg in DEFAULT_DATA_SOURCES:
if db.query(DataSource).filter(DataSource.name == cfg["name"]).first():
continue
db.add(DataSource(**cfg)); added += 1
if added:
logger.info(f"[seed] data_source 新增 {added} 条示例数据")
"""web3 后端 ORM 模型子包
"""web 后端 ORM 模型子包
每个文件对应 schema.sql 里的一张表。
......@@ -8,7 +8,9 @@
from .connection_preset import ConnectionPreset
from .field import Field
from .rule import Rule
from .rule_library import FieldRule, RuleLibrary
from .task import Task
from .task_group import TaskGroup
from .data_source import DataSource
__all__ = ["ConnectionPreset", "Field", "Rule", "Task", "TaskGroup"]
\ No newline at end of file
__all__ = ["ConnectionPreset", "Field", "Rule", "RuleLibrary", "FieldRule", "Task", "TaskGroup", "DataSource"]
......@@ -20,7 +20,7 @@ from datetime import datetime
from sqlalchemy import Column, Integer, String, DateTime, Text, UniqueConstraint
from sqlalchemy.sql import func
from web3.backend.db.database import Base
from web.backend.db.database import Base
class ConnectionPreset(Base):
......
"""可复用的数据源连接配置。"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy import Column, Integer, String, Text, DateTime
from sqlalchemy.sql import func
from web.backend.db.database import Base
class DataSource(Base):
__tablename__ = "data_source"
id = Column(Integer, primary_key=True, autoincrement=True)
name = Column(String(100), nullable=False, unique=True)
db_type = Column(String(30), nullable=False, default="MySQL")
host = Column(String(255), nullable=False)
port = Column(Integer, nullable=False)
database = Column(String(128), nullable=False, default="")
schema = Column(String(128), nullable=True)
user = Column(String(128), nullable=False)
password = Column(Text, nullable=False, default="")
oracle_client_dir = Column(String(512), nullable=True)
status = Column(String(20), nullable=False, default="启用")
last_tested_at = Column(DateTime, nullable=True)
last_test_ok = Column(Integer, nullable=True)
created_at = Column(DateTime, nullable=False, server_default=func.current_timestamp())
updated_at = Column(DateTime, nullable=False, server_default=func.current_timestamp(), onupdate=func.current_timestamp())
def to_dict(self) -> dict:
return {
"id": self.id, "name": self.name, "db_type": self.db_type,
"host": self.host, "port": self.port, "database": self.database,
"schema": self.schema, "user": self.user,
"password": self.password, "oracle_client_dir": self.oracle_client_dir,
"status": self.status,
"last_tested_at": self.last_tested_at.isoformat() if isinstance(self.last_tested_at, datetime) else self.last_tested_at,
"last_test_ok": None if self.last_test_ok is None else bool(self.last_test_ok),
"created_at": self.created_at.isoformat() if isinstance(self.created_at, datetime) else self.created_at,
"updated_at": self.updated_at.isoformat() if isinstance(self.updated_at, datetime) else self.updated_at,
}
......@@ -16,7 +16,7 @@ from __future__ import annotations
from sqlalchemy import Column, Integer, String, ForeignKey, UniqueConstraint
from sqlalchemy.orm import relationship
from web3.backend.db.database import Base
from web.backend.db.database import Base
class Field(Base):
......
......@@ -29,7 +29,7 @@ from __future__ import annotations
from sqlalchemy import Column, Integer, Text, String, ForeignKey
from web3.backend.db.database import Base
from web.backend.db.database import Base
# rule_type 的合法值集中在一处常量,ORM / 校验 / 提示文案都引用这一份
RULE_TYPES = ("regex", "number", "date", "string")
......
"""可复用校验规则库,以及任务字段与规则的关联。"""
from __future__ import annotations
from sqlalchemy import Column, ForeignKey, Integer, String, Text, UniqueConstraint
from sqlalchemy.sql import func
from sqlalchemy import DateTime
from web.backend.db.database import Base
class RuleLibrary(Base):
__tablename__ = "rule_library"
id = Column(Integer, primary_key=True, autoincrement=True)
name = Column(String(120), nullable=False, unique=True)
description = Column(Text, nullable=False, default="")
rule_type = Column(String(20), nullable=False, default="regex")
regex = Column(Text, nullable=True)
code = Column(Text, nullable=True)
skip_null = Column(Integer, nullable=False, default=0)
created_at = Column(DateTime, nullable=False, server_default=func.current_timestamp())
updated_at = Column(DateTime, nullable=False, server_default=func.current_timestamp())
def to_dict(self) -> dict:
return {
"id": self.id, "name": self.name, "description": self.description,
"rule_type": self.rule_type, "regex": self.regex, "code": self.code,
"skip_null": bool(self.skip_null),
"created_at": self.created_at.isoformat() if self.created_at else None,
"updated_at": self.updated_at.isoformat() if self.updated_at else None,
}
class FieldRule(Base):
__tablename__ = "field_rule"
__table_args__ = (UniqueConstraint("field_id", "rule_library_id", name="uq_field_rule"),)
id = Column(Integer, primary_key=True, autoincrement=True)
field_id = Column(Integer, ForeignKey("field.id", ondelete="CASCADE"), nullable=False)
rule_library_id = Column(Integer, ForeignKey("rule_library.id", ondelete="RESTRICT"), nullable=False)
ord = Column(Integer, nullable=False, default=0)
......@@ -19,8 +19,8 @@ from sqlalchemy import Column, Integer, String, DateTime, ForeignKey, Text
from sqlalchemy.orm import relationship
from sqlalchemy.sql import func
from web3.backend.db.database import Base
from web3.backend.models.task_group import TaskGroup
from web.backend.db.database import Base
from web.backend.models.task_group import TaskGroup
class Task(Base):
......@@ -30,6 +30,7 @@ class Task(Base):
name = Column(String, nullable=False, unique=True)
group_id = Column(Integer, ForeignKey("task_group.id"), nullable=False)
data_source_path = Column(Text, nullable=True) # 例如 sources/customer_db/customer.sqlite
data_source_id = Column(Integer, ForeignKey("data_source.id"), nullable=True)
source_table = Column(String, nullable=True) # 例如 t_user_info
status = Column(String, nullable=False, default="启用") # 启用 / 停用
description = Column(Text, nullable=True)
......@@ -59,6 +60,7 @@ class Task(Base):
"group_id": self.group_id,
"group": self.group.name if self.group else None,
"data_source_path": self.data_source_path,
"data_source_id": self.data_source_id,
"source_table": self.source_table,
"status": self.status,
"description": self.description,
......@@ -73,4 +75,4 @@ class Task(Base):
}
def __repr__(self) -> str:
return f"<Task id={self.id} name={self.name!r} group_id={self.group_id}>"
\ No newline at end of file
return f"<Task id={self.id} name={self.name!r} group_id={self.group_id}>"
......@@ -12,7 +12,7 @@ from datetime import datetime
from sqlalchemy import Column, Integer, String, DateTime
from sqlalchemy.sql import func
from web3.backend.db.database import Base
from web.backend.db.database import Base
class TaskGroup(Base):
......
......@@ -2,13 +2,14 @@ fastapi>=0.110
uvicorn[standard]>=0.27
sqlalchemy>=2.0
pydantic>=2.5
# 数据库驱动(web3/backend/core/db_adapter.py 需要;dmPython 装不上不影响启动,只影响连达梦)
# 数据库驱动(web/backend/core/db_adapter.py 需要;dmPython 装不上不影响启动,只影响连达梦)
pymysql>=1.1
oracledb>=1.4
dmPython>=2.5
# LLM(web3/backend/core/ai_regex.py → core/llm.py 用)
# LLM(web/backend/core/ai_regex.py → core/llm.py 用)
# anthropic 走 Anthropic SDK(MiniMax 也走这个,兼容接口);openai 走 OpenAI SDK
# pyyaml 读 web3/backend/configs/llm.yaml
# pyyaml 读 web/backend/configs/llm.yaml
anthropic>=0.30
openai>=1.30
pyyaml>=6.0
\ No newline at end of file
pyyaml>=6.0
cryptography>=42.0
"""web 后端 routers 子包"""
\ No newline at end of file
"""AI 相关 API(web3)
"""AI 相关 API(web)
端点:
POST /api/ai/rule body: {desc, rule_type} → {ok, code, note}
......@@ -21,16 +21,18 @@
from __future__ import annotations
import time
import re
from typing import Any, Optional
from fastapi import APIRouter
from pydantic import BaseModel, ConfigDict, Field
from web3.backend.core.ai_rule import gen_rule, test_rule
from web3.backend.core.ai_explain import explain_issue as _explain_issue
from web3.backend.core.ai_regex import test_regex as _legacy_test_regex
from web3.backend.models.rule import DEFAULT_RULE_TYPE, RULE_TYPES
from web3.backend._logging import get_logger
from web.backend.core.ai_rule import gen_rule, test_rule
from web.backend.core.rule_runner import validate_user_code
from web.backend.core.ai_explain import explain_issue as _explain_issue
from web.backend.core.ai_regex import test_regex as _legacy_test_regex
from web.backend.models.rule import DEFAULT_RULE_TYPE, RULE_TYPES
from web.backend._logging import get_logger
router = APIRouter(prefix="", tags=["ai"])
logger = get_logger("backend.routers.ai")
......@@ -96,6 +98,35 @@ class TestRuleResponse(BaseModel):
error: str = ""
class ValidateRuleRequest(BaseModel):
rule_type: str = Field(DEFAULT_RULE_TYPE, description=f"规则类型:{ ' / '.join(RULE_TYPES) }")
code: str = Field(..., min_length=1, description="regex 时是正则;其他类型是 Python 校验代码")
class ValidateRuleResponse(BaseModel):
ok: bool
valid: bool = False
error: str = ""
@router.post("/rule/validate", response_model=ValidateRuleResponse, summary="验证正则表达式或校验代码")
async def ai_validate_rule(req: ValidateRuleRequest):
"""只验证表达式/代码本身,不拿业务数据执行匹配。"""
code = (req.code or "").strip()
if not code:
return ValidateRuleResponse(ok=True, valid=False, error="表达式为空")
if req.rule_type == "regex":
try:
re.compile(code)
except re.error as e:
return ValidateRuleResponse(ok=True, valid=False, error=f"正则语法错误:{e}")
return ValidateRuleResponse(ok=True, valid=True)
if req.rule_type not in RULE_TYPES:
return ValidateRuleResponse(ok=True, valid=False, error=f"不支持的规则类型:{req.rule_type}")
valid, error = validate_user_code(code)
return ValidateRuleResponse(ok=True, valid=valid, error=error)
@router.post("/rule/test", response_model=TestRuleResponse, summary="按 rule_type 校验 value")
async def ai_test_rule(req: TestRuleRequest):
"""前端「规则设置」弹窗的「测试」按钮调用。
......@@ -125,8 +156,7 @@ class GenRegexResponse(BaseModel):
async def ai_gen_regex(req: GenRegexRequest):
"""老端点,转发到 /ai/rule(rule_type='regex')。
保留原因:web3/src/api/ai.js 之前的 genRegex 调用 + 老测试(test_run_query_only_bad_rows
没用到,但 ai.js 单元测试可能用到)继续可用。
保留原因:兼容前端历史调用,避免已有客户端集成失效。
"""
logger.info(f"POST /api/ai/regex (legacy) desc={req.desc!r}")
code, note = gen_rule(req.desc, "regex")
......@@ -222,4 +252,4 @@ async def ai_explain_issue(req: ExplainIssueRequest):
f"consistency={result.get('consistency')!r} note={result.get('note')!r}"
)
logger.info("─" * 60)
return ExplainIssueResponse(**result)
\ No newline at end of file
return ExplainIssueResponse(**result)
......@@ -23,8 +23,8 @@ from pydantic import BaseModel, Field
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from web3.backend.db.database import get_session
from web3.backend.models.connection_preset import ConnectionPreset
from web.backend.db.database import get_session
from web.backend.models.connection_preset import ConnectionPreset
router = APIRouter(prefix="/connection-presets", tags=["connection-presets"])
......
"""数据源 CRUD 与连接测试。"""
from __future__ import annotations
import base64
import json
import re
from datetime import datetime
from decimal import Decimal
from typing import Literal, Optional
from fastapi import APIRouter, Depends, HTTPException
from pydantic import BaseModel, Field
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from web.backend.db.database import get_session
from web.backend.models.data_source import DataSource
from web.backend.core.db_adapter import DBConfig, DBConnection, quote_ident
from web.backend.core.backup_storage import get_backup_dir
router = APIRouter(prefix="/data-sources", tags=["data-sources"])
class DataSourceBase(BaseModel):
name: str = Field(..., min_length=1, max_length=100)
db_type: str = "MySQL"
host: str = Field(..., min_length=1)
port: int = Field(..., ge=1, le=65535)
database: str = ""
schema: Optional[str] = None
user: str = Field(..., min_length=1)
password: str = ""
oracle_client_dir: Optional[str] = None
status: str = "启用"
class DataSourceCreate(DataSourceBase):
pass
class DataSourceUpdate(BaseModel):
name: Optional[str] = None; db_type: Optional[str] = None
host: Optional[str] = None; port: Optional[int] = Field(None, ge=1, le=65535)
database: Optional[str] = None; schema: Optional[str] = None
user: Optional[str] = None; password: Optional[str] = None
oracle_client_dir: Optional[str] = None; status: Optional[str] = None
class TestResult(BaseModel):
ok: bool; message: str
class BackupRequest(BaseModel):
"""导出已保存数据源中的指定表。"""
scope: Literal["database", "tables"] = "database"
tables: list[str] = Field(default_factory=list)
database: Optional[str] = None
schema: Optional[str] = None
include_schema: bool = True
include_data: bool = True
filename: Optional[str] = None
def _cfg(s: DataSource) -> DBConfig:
labels = {"MySQL": "mysql", "Oracle": "oracle", "达梦 DM": "dameng"}
return DBConfig(db_type=labels.get(s.db_type, s.db_type.lower()), host=s.host,
port=int(s.port), user=s.user, password=s.password or "",
database=s.database, oracle_client_dir=s.oracle_client_dir)
def _json_default(value):
if isinstance(value, datetime):
return value.isoformat()
if isinstance(value, Decimal):
return str(value)
if isinstance(value, (bytes, bytearray, memoryview)):
return {"__type__": "bytes", "base64": base64.b64encode(bytes(value)).decode("ascii")}
return str(value)
def _qualified_table(cfg: DBConfig, schema: str, table: str) -> str:
"""引用表名;schema 仅用于 Oracle/达梦,MySQL 已由连接库选择。"""
table_sql = quote_ident(table, cfg.db_type)
if schema and cfg.db_type in ("oracle", "dameng"):
return f"{quote_ident(schema, cfg.db_type)}.{table_sql}"
return table_sql
def _backup_filename(value: Optional[str]) -> str:
if not value or not value.strip():
return f"db-backup-{datetime.now().strftime('%Y%m%d-%H%M%S-%f')}.json"
name = value.strip()
if name.lower().endswith(".json"):
name = name[:-5]
if not name or not re.fullmatch(r"[A-Za-z0-9_\-\u4e00-\u9fff ]{1,120}", name):
raise HTTPException(400, "备份文件名只能包含中文、字母、数字、空格、下划线和短横线")
return f"{name}.json"
@router.get("")
def list_data_sources(db: Session = Depends(get_session)):
return [s.to_dict() for s in db.query(DataSource).order_by(DataSource.id.desc()).all()]
@router.post("", status_code=201)
def create_data_source(payload: DataSourceCreate, db: Session = Depends(get_session)):
data = payload.model_dump()
data["name"] = payload.name.strip()
row = DataSource(**data)
db.add(row)
try: db.commit()
except IntegrityError:
db.rollback(); raise HTTPException(409, "数据源名称已存在")
db.refresh(row); return row.to_dict()
@router.get("/{source_id}")
def get_data_source(source_id: int, db: Session = Depends(get_session)):
row = db.get(DataSource, source_id)
if not row: raise HTTPException(404, "数据源不存在")
return row.to_dict()
@router.put("/{source_id}")
def update_data_source(source_id: int, payload: DataSourceUpdate, db: Session = Depends(get_session)):
row = db.get(DataSource, source_id)
if not row: raise HTTPException(404, "数据源不存在")
data = payload.model_dump(exclude_unset=True)
if data.get("password") == "": data.pop("password")
for k, v in data.items(): setattr(row, k, v)
try: db.commit()
except IntegrityError:
db.rollback(); raise HTTPException(409, "数据源名称已存在")
db.refresh(row); return row.to_dict()
@router.delete("/{source_id}", status_code=204)
def delete_data_source(source_id: int, db: Session = Depends(get_session)):
row = db.get(DataSource, source_id)
if not row: raise HTTPException(404, "数据源不存在")
db.delete(row); db.commit()
@router.post("/{source_id}/test", response_model=TestResult)
def test_data_source(source_id: int, db: Session = Depends(get_session)):
row = db.get(DataSource, source_id)
if not row: raise HTTPException(404, "数据源不存在")
try:
with DBConnection(_cfg(row)) as conn:
conn.fetch_scalar("SELECT 1")
row.last_test_ok = 1; row.last_tested_at = datetime.utcnow(); db.commit()
return {"ok": True, "message": "连接成功"}
except Exception as e:
row.last_test_ok = 0; row.last_tested_at = datetime.utcnow(); db.commit()
return {"ok": False, "message": f"连接失败:{type(e).__name__}: {e}"}
@router.post("/{source_id}/backup")
def backup_data_source(source_id: int, payload: BackupRequest, db: Session = Depends(get_session)):
"""从真实数据源导出 JSON 备份包。"""
row = db.get(DataSource, source_id)
if not row:
raise HTTPException(404, "数据源不存在")
cfg = _cfg(row)
if payload.database:
cfg.database = payload.database
schema = payload.schema or row.schema or cfg.database
tables = list(dict.fromkeys(t.strip() for t in payload.tables if t and t.strip()))
if payload.scope == "tables" and not tables:
raise HTTPException(400, "至少选择一张表")
result = {
"format": "db-tools-backup",
"version": 1,
"created_at": datetime.utcnow().isoformat() + "Z",
"source": {"id": row.id, "name": row.name, "db_type": row.db_type, "host": row.host,
"database": cfg.database, "schema": schema, "scope": payload.scope},
"tables": [],
}
try:
with DBConnection(cfg) as conn:
available = {str(item.get("table_name", "")): item for item in conn.list_tables(schema)}
available_upper = {key.upper(): value for key, value in available.items()}
if payload.scope == "database":
tables = list(available)
if not tables:
raise HTTPException(400, "该数据库没有可访问的数据表")
for table in tables:
meta = available.get(table) or available_upper.get(table.upper())
if not meta:
raise HTTPException(400, f"表不存在或无权访问:{table}")
actual_table = str(meta.get("table_name") or table)
columns = conn.list_columns(schema, actual_table) if payload.include_schema else []
item = {"table_name": actual_table, "table_comment": meta.get("table_comment"),
"columns": columns, "rows": []}
# 2026-09-03:MySQL 基表多存一份 SHOW CREATE TABLE 的原始 DDL ——
# 主键/自增/索引/唯一键/默认值/NOT NULL/注释/字符集/外键全保真,
# 还原到 MySQL 时优先按它建表。视图不存(table_type 过滤),走 columns 结构化路径
if (payload.include_schema and cfg.db_type == "mysql"
and str(meta.get("table_type") or "").upper() == "BASE TABLE"):
ddl_rows = conn.fetchall(
f"SHOW CREATE TABLE {quote_ident(actual_table, cfg.db_type)}")
r0 = ddl_rows[0] if ddl_rows else {}
ddl = str(r0.get("Create Table") or r0.get("create table") or "")
if ddl:
item["create_ddl"] = ddl
if payload.include_data:
sql = f"SELECT * FROM {_qualified_table(cfg, schema, actual_table)}"
item["rows"] = list(conn.iter_rows(sql))
result["tables"].append(item)
except HTTPException:
raise
except Exception as exc:
raise HTTPException(502, f"备份失败:{type(exc).__name__}: {exc}") from exc
filename = _backup_filename(payload.filename)
target = get_backup_dir() / filename
if target.exists():
raise HTTPException(409, f"备份文件已存在:{filename}")
target.write_text(json.dumps(result, ensure_ascii=False, default=_json_default), encoding="utf-8")
return {"filename": filename, "backup_dir": str(target.parent),
"tables": len(result["tables"]), "message": "备份完成"}
"""web3 数据库连接 API
"""web 数据库连接 API
端点:
POST /api/connect/test 测试连接(MySQL / Oracle / 达梦)
......@@ -7,7 +7,7 @@
POST /api/connect/columns 列出指定表的全部字段(数据字典)
实现:
- 复用 web3/backend/core/db_adapter(自包含,不依赖 web.core)
- 复用 web/backend/core/db_adapter(自包含,不依赖 web.core)
- db_type 兼容前端中文 label("MySQL"/"Oracle"/"达梦 DM"),内部统一小写
- 错误用 logger.exception 拿全栈,响应 message 给前端弹窗
- 错误分类提示(WinError 10061 / timeout / Access denied)写到 logs/app.log
......@@ -21,10 +21,10 @@ from typing import Any
from fastapi import APIRouter
from pydantic import Field
from web3.backend.core.db_adapter import DBConfig, DBConnection
from web3.backend.core.sql_loader import get_sql_loader
from web3.backend.core.models import TestConnectionRequest, TestConnectionResponse
from web3.backend._logging import get_logger
from web.backend.core.db_adapter import DBConfig, DBConnection
from web.backend.core.sql_loader import get_sql_loader
from web.backend.core.models import TestConnectionRequest, TestConnectionResponse
from web.backend._logging import get_logger
router = APIRouter(prefix="/connect", tags=["database"])
logger = get_logger("backend.routers.db")
......@@ -97,7 +97,7 @@ def _log_failure_hint(msg: str) -> None:
# ── 端点 ─────────────────────────────────────────────────
@router.post("/test", response_model=TestConnectionResponse, summary="测试数据库连接")
async def connect_test(req: TestConnectionRequest):
"""复用 web3.backend.core.db_adapter,验证 db_type/host/port/user/password/database
"""复用 web.backend.core.db_adapter,验证 db_type/host/port/user/password/database
返回 TestConnectionResponse(ok, message, db_type, database_name)
"""
......
This diff is collapsed.
"""web3 任务查询 + 校验 API(分页版)
"""web 任务查询 + 校验 API(分页版)
端点:
POST /api/queries/start 创建查询 session:COUNT → 编译规则 → 入 session_manager
......@@ -6,7 +6,7 @@
POST /api/queries/cancel 摘掉 session(best-effort 中断)
设计:
- 复用 web3.backend.core.db_adapter 建连(每页新连接,简单)
- 复用 web.backend.core.db_adapter 建连(每页新连接,简单)
- 复用 routers/db._normalize_db_type 做中文 label → 内部小写
- 列名按 task.field_list 显式 SELECT(不 SELECT * 防止暴露未配置列)
- 校验引擎:rule_runner.run_rule() 按 rule_type dispatch
......@@ -40,22 +40,23 @@ from fastapi import APIRouter, Depends, HTTPException
from pydantic import BaseModel, Field as PydField
from sqlalchemy.orm import Session
from web3.backend._logging import get_logger
from web3.backend.core.db_adapter import DBConfig, DBConnection, paginate_sql, quote_ident
from web3.backend.core.rule_runner import RuleRunError, run_rule
from web3.backend.core.session_manager import (
from web.backend._logging import get_logger
from web.backend.core.db_adapter import DBConfig, DBConnection, paginate_sql, quote_ident
from web.backend.core.rule_runner import RuleRunError, run_rule
from web.backend.core.session_manager import (
CompiledFieldRules,
QuerySession,
get as sm_get,
pop as sm_pop,
put as sm_put,
)
from web3.backend.db.database import get_session
from web3.backend.models.field import Field
from web3.backend.models.rule import Rule
from web3.backend.models.task import Task
from web3.backend.routers.db import _normalize_db_type
from web3.backend.routers.tasks import _parse_conn_json
from web.backend.db.database import get_session
from web.backend.models.field import Field
from web.backend.models.rule import Rule
from web.backend.models.rule_library import FieldRule, RuleLibrary
from web.backend.models.task import Task
from web.backend.routers.db import _normalize_db_type
from web.backend.routers.tasks import _parse_conn_json
router = APIRouter(prefix="/queries", tags=["queries"])
logger = get_logger("backend.routers.queries")
......@@ -127,10 +128,14 @@ def _load_task_compiled(db: Session, task_id: int):
return None, None, None, None, None
fields = db.query(Field).filter(Field.task_id == task.id).order_by(Field.ord).all()
rules_by_field: dict[int, list[Rule]] = {
f.id: db.query(Rule).filter(Rule.field_id == f.id).order_by(Rule.ord).all()
for f in fields
}
# 新任务优先使用独立规则库关联;没有关联的历史任务继续读取旧 rule 表。
rules_by_field = {}
for f in fields:
links = db.query(FieldRule).filter(FieldRule.field_id == f.id).order_by(FieldRule.ord).all()
library_rules = [db.get(RuleLibrary, link.rule_library_id) for link in links]
rules_by_field[f.id] = [r for r in library_rules if r] or (
db.query(Rule).filter(Rule.field_id == f.id).order_by(Rule.ord).all()
)
fields_with_rules: list[tuple[Field, list[Rule]]] = [
(f, rules_by_field[f.id]) for f in fields if rules_by_field[f.id]
]
......@@ -146,7 +151,7 @@ def _load_task_compiled(db: Session, task_id: int):
"rule_list": [
{
"id": r.id,
"desc": r.desc,
"desc": getattr(r, "description", "") or getattr(r, "desc", ""),
# 2026-08-24:规则名称,老规则空串;前端用 desc 兜底显示
"name": r.name or "",
"rule_type": r.rule_type or "regex",
......@@ -170,7 +175,7 @@ def _load_task_compiled(db: Session, task_id: int):
snapshots: list[tuple[str, Optional[str], Optional[str], str, bool, str, int]] = []
for r in rules:
rt = str(r.rule_type or "regex")
desc = str(r.desc or "")
desc = str(getattr(r, "description", "") or getattr(r, "desc", "") or "")
skip_null = bool(r.skip_null)
name = str(r.name or "")
rule_id = int(r.id or 0)
......@@ -451,4 +456,4 @@ async def cancel_query(req: CancelQueryRequest):
cancelled=True,
pages_scanned=sess.pages_scanned,
bad_rows_so_far=sess.bad_rows_total,
)
\ No newline at end of file
)
"""独立校验规则库的 CRUD 接口。"""
from __future__ import annotations
from datetime import datetime, timezone
from typing import Optional
from fastapi import APIRouter, Depends, HTTPException
from pydantic import BaseModel, Field
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from web.backend.db.database import get_session
from web.backend.models.rule_library import FieldRule, RuleLibrary
router = APIRouter(prefix="/rules", tags=["rules"])
_TYPES = {"regex", "number", "date", "string"}
class RulePayload(BaseModel):
name: str = Field(..., min_length=1, max_length=120)
description: str = ""
rule_type: str = "regex"
regex: Optional[str] = None
code: Optional[str] = None
skip_null: bool = False
def _apply(rule: RuleLibrary, payload: RulePayload) -> None:
if payload.rule_type not in _TYPES:
raise HTTPException(status_code=400, detail="规则类型仅支持 regex、number、date、string")
if payload.rule_type == "regex" and not (payload.regex or "").strip():
raise HTTPException(status_code=400, detail="正则规则必须填写正则表达式")
if payload.rule_type != "regex" and not (payload.code or "").strip():
raise HTTPException(status_code=400, detail="自定义规则必须填写 check(value) 代码")
rule.name = payload.name.strip()
rule.description = payload.description.strip()
rule.rule_type = payload.rule_type
rule.regex = (payload.regex or "").strip() or None
rule.code = (payload.code or "").strip() or None
rule.skip_null = 1 if payload.skip_null else 0
@router.get("")
def list_rules(db: Session = Depends(get_session)):
return [r.to_dict() for r in db.query(RuleLibrary).order_by(RuleLibrary.id.desc()).all()]
@router.post("", status_code=201)
def create_rule(payload: RulePayload, db: Session = Depends(get_session)):
rule = RuleLibrary()
_apply(rule, payload)
db.add(rule)
try:
db.commit()
except IntegrityError:
db.rollback()
raise HTTPException(status_code=409, detail=f"规则名称已存在:{payload.name}")
db.refresh(rule)
return rule.to_dict()
@router.put("/{rule_id}")
def update_rule(rule_id: int, payload: RulePayload, db: Session = Depends(get_session)):
rule = db.get(RuleLibrary, rule_id)
if not rule:
raise HTTPException(status_code=404, detail="规则不存在")
_apply(rule, payload)
rule.updated_at = datetime.now(timezone.utc).replace(tzinfo=None)
try:
db.commit()
except IntegrityError:
db.rollback()
raise HTTPException(status_code=409, detail=f"规则名称已存在:{payload.name}")
db.refresh(rule)
return rule.to_dict()
@router.delete("/{rule_id}", status_code=204)
def delete_rule(rule_id: int, db: Session = Depends(get_session)):
rule = db.get(RuleLibrary, rule_id)
if not rule:
raise HTTPException(status_code=404, detail="规则不存在")
if db.query(FieldRule).filter(FieldRule.rule_library_id == rule_id).first():
raise HTTPException(status_code=409, detail="该规则正在被校验任务使用,请先在任务中移除")
db.delete(rule)
db.commit()
......@@ -16,8 +16,8 @@ from pydantic import BaseModel, Field
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from web3.backend.db.database import get_session
from web3.backend.models.task_group import TaskGroup
from web.backend.db.database import get_session
from web.backend.models.task_group import TaskGroup
router = APIRouter(prefix="/task-groups", tags=["task-groups"])
......
......@@ -26,13 +26,14 @@ from pydantic import BaseModel, Field as PydField
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from web3.backend.db.database import get_session
from web3.backend.core.db_adapter import DBConfig, DBConnection
from web3.backend.models.field import Field
from web3.backend.models.rule import Rule
from web3.backend.models.task import Task
from web3.backend.models.task_group import TaskGroup
from web3.backend._logging import get_logger
from web.backend.db.database import get_session
from web.backend.core.db_adapter import DBConfig, DBConnection
from web.backend.models.field import Field
from web.backend.models.rule import Rule
from web.backend.models.rule_library import FieldRule, RuleLibrary
from web.backend.models.task import Task
from web.backend.models.task_group import TaskGroup
from web.backend._logging import get_logger
router = APIRouter(prefix="/tasks", tags=["tasks"])
......@@ -60,6 +61,8 @@ class FieldPayload(BaseModel):
key: str = PydField(..., min_length=1, description="英文字段名(与数据源对齐)")
show_default: bool = PydField(True, description="结果表是否默认展示")
rules: list[RulePayload] = PydField(default_factory=list, description="字段下的规则列表")
# 新版校验任务只引用规则库;rules 保留给已有任务兼容读取。
rule_ids: list[int] = PydField(default_factory=list, description="规则库 ID,可多选")
class TaskBase(BaseModel):
......@@ -67,6 +70,7 @@ class TaskBase(BaseModel):
group: str = PydField(..., description="分组名(前端传名,服务端查 id)")
db_type: str = PydField("MySQL", description="UI 冗余,不参与检测")
data_source_path: Optional[str] = PydField(None, description="如 sources/customer_db/customer.sqlite")
data_source_id: Optional[int] = None
source_table: Optional[str] = PydField(None, description="如 t_user_info")
status: str = PydField("启用", description="启用 / 停用")
description: Optional[str] = None
......@@ -85,6 +89,7 @@ class TaskUpdate(BaseModel):
group: Optional[str] = None
db_type: Optional[str] = None
data_source_path: Optional[str] = None
data_source_id: Optional[int] = None
source_table: Optional[str] = None
status: Optional[str] = None
description: Optional[str] = None
......@@ -134,6 +139,7 @@ class TaskOut(BaseModel):
schema: Optional[str] = None # Schema
oracle_client_dir: Optional[str] = None # Oracle Instant Client 路径(仅 Oracle 有效)
data_source_path: Optional[str]
data_source_id: Optional[int] = None
source_table: Optional[str]
status: str
description: Optional[str]
......@@ -150,6 +156,10 @@ class TaskOut(BaseModel):
def _resolve_group_id(db: Session, group_name: str) -> int:
g = db.query(TaskGroup).filter(TaskGroup.name == group_name).first()
if not g and group_name == "校验":
g = TaskGroup(name="校验", ord=10)
db.add(g)
db.flush()
if not g:
raise HTTPException(status_code=400, detail=f"分组不存在:{group_name}")
return g.id
......@@ -196,6 +206,12 @@ def _fetch_field_comments(t: Task, conn: dict) -> dict[str, str]:
user = conn.get("user")
password = conn.get("password")
database = conn.get("db")
# 2026-09-03:MySQL 的 schema 就是 database。老任务 conn_json 可能只存了 schema(如 db-test)
# 而没存 db —— 互为回退,否则下面的连接信息校验会提前放弃,详情接口带不回字段注释
# (「扫描后返回再编辑」时注释变「—」的后端一半原因)
if db_type_internal == "mysql":
database = database or schema
schema = schema or database
if not (host and port and user and database):
# 连接信息不全(典型:密码缺失 / 任务用预设建但预设不存在)—— 不报错,直接降级
return {}
......@@ -251,6 +267,16 @@ def _replace_fields(db: Session, task: Task, fields_payload: list[FieldPayload])
)
db.add(field)
db.flush() # 拿到 field.id
# 校验任务选择的是独立规则库,不复制规则定义。
rule_ids = list(dict.fromkeys(f.rule_ids))
if rule_ids:
found = db.query(RuleLibrary).filter(RuleLibrary.id.in_(rule_ids)).count()
if found != len(rule_ids):
raise ValueError("所选规则不存在或已删除")
for r_idx, rule_id in enumerate(rule_ids):
db.add(FieldRule(field_id=field.id, rule_library_id=rule_id, ord=r_idx))
# 兼容旧请求中的内嵌规则。新版页面不会再写入这个表。
for r_idx, r in enumerate(f.rules):
rt = (r.rule_type or "regex").strip() or "regex"
regex_val = (r.regex or "").strip() or None if rt == "regex" else None
......@@ -280,6 +306,8 @@ def _copy_fields(db: Session, src_task: Task, dst_task: Task) -> None:
)
db.add(new_f)
db.flush()
for link in db.query(FieldRule).filter(FieldRule.field_id == src_f.id).order_by(FieldRule.ord).all():
db.add(FieldRule(field_id=new_f.id, rule_library_id=link.rule_library_id, ord=link.ord))
for r_idx, src_r in enumerate(
db.query(Rule).filter(Rule.field_id == src_f.id).order_by(Rule.ord).all()
):
......@@ -307,17 +335,26 @@ def _row_to_out(db: Session, t: Task, *, include_field_list: bool = False) -> Ta
fields_count = (
db.query(Field).filter(Field.task_id == t.id).count() if t.id else 0
)
task_field_ids = [f.id for f in db.query(Field).filter(Field.task_id == t.id).all()]
rules_count = (
db.query(Rule).join(Field, Rule.field_id == Field.id)
.filter(Field.task_id == t.id).count() if t.id else 0
)
db.query(FieldRule).filter(FieldRule.field_id.in_(task_field_ids)).count()
+ db.query(Rule).join(Field, Rule.field_id == Field.id).filter(Field.task_id == t.id).count()
) if task_field_ids else 0
field_list: list[FieldOut] = []
if include_field_list:
# 2026-08-24:详情接口顺带从数据源拉一次列注释(连不上 / 没 source_table 时降级为空 dict)
comment_map = _fetch_field_comments(t, conn)
for f in db.query(Field).filter(Field.task_id == t.id).order_by(Field.ord).all():
links = db.query(FieldRule).filter(FieldRule.field_id == f.id).order_by(FieldRule.ord).all()
library_rules = [db.get(RuleLibrary, link.rule_library_id) for link in links]
rules = [
RuleOut(
id=r.id, desc=r.description, name=r.name, regex=r.regex,
rule_type=r.rule_type, code=r.code, skip_null=bool(r.skip_null), ord=idx,
)
for idx, r in enumerate(library_rules) if r
] or [
RuleOut(
id=r.id, desc=r.desc,
name=r.name or "", # 2026-08-24:规则名称(老规则空串,UI desc 兜底)
......@@ -327,8 +364,7 @@ def _row_to_out(db: Session, t: Task, *, include_field_list: bool = False) -> Ta
skip_null=bool(r.skip_null),
ord=r.ord,
)
for r in db.query(Rule).filter(Rule.field_id == f.id).order_by(Rule.ord).all()
]
for r in db.query(Rule).filter(Rule.field_id == f.id).order_by(Rule.ord).all()]
field_list.append(FieldOut(
id=f.id,
field_key=f.field_key,
......@@ -358,6 +394,7 @@ def _row_to_out(db: Session, t: Task, *, include_field_list: bool = False) -> Ta
schema=conn.get("schema"),
oracle_client_dir=conn.get("oracleClientDir"),
data_source_path=d["data_source_path"],
data_source_id=d.get("data_source_id"),
source_table=d["source_table"],
status=d["status"],
description=d["description"],
......@@ -399,6 +436,7 @@ def create_task(payload: TaskCreate, db: Session = Depends(get_session)) -> Task
name=payload.name.strip(),
group_id=group_id,
data_source_path=payload.data_source_path,
data_source_id=payload.data_source_id,
source_table=payload.source_table,
status=payload.status or "启用",
description=payload.description,
......@@ -505,6 +543,7 @@ def copy_task(task_id: int, db: Session = Depends(get_session)) -> TaskOut:
name=new_name,
group_id=t.group_id,
data_source_path=t.data_source_path,
data_source_id=t.data_source_id,
source_table=t.source_table,
status=t.status or "启用",
description=t.description,
......@@ -519,4 +558,4 @@ def copy_task(task_id: int, db: Session = Depends(get_session)) -> TaskOut:
db.commit()
db.refresh(clone)
return _row_to_out(db, clone, include_field_list=True)
\ No newline at end of file
return _row_to_out(db, clone, include_field_list=True)
"""web3 后端启动脚本(开发模式)
"""web 后端启动脚本(开发模式)
用法:
cd 数据治理
python -m web3.backend.start
python -m web.backend.start
# 或
python web3/backend/start.py
python web/backend/start.py
"""
from __future__ import annotations
......@@ -25,27 +25,27 @@ except Exception:
import uvicorn
from web3.backend.config import HOST, PORT
from web3.backend._logging import setup_logging, get_logger, LOG_DIR, APP_LOG_FILE
from web.backend.config import HOST, PORT
from web.backend._logging import setup_logging, get_logger, LOG_DIR, APP_LOG_FILE
setup_logging(os.environ.get("LOG_LEVEL", "INFO"))
logger = get_logger("backend.start")
logger.info("=" * 60)
logger.info("web3 后端 · 进程启动")
logger.info("web 后端 · 进程启动")
logger.info(f"Python: {sys.version.split()[0]}, 平台: {sys.platform}")
logger.info(f"工作目录: {PROJECT_ROOT}")
logger.info(f"应用日志: {APP_LOG_FILE}")
if __name__ == "__main__":
print("=" * 60)
print(f" web3 后端 · FastAPI")
print(f" web 后端 · FastAPI")
print(f" API: http://localhost:{PORT}/api")
print(f" API 文档: http://localhost:{PORT}/docs")
print(f" 前端 dev: http://localhost:5175 (Vite 单独启动)")
print(f" 应用日志: {APP_LOG_FILE}")
print("=" * 60)
uvicorn.run(
"web3.backend.app:app",
"web.backend.app:app",
host=HOST,
port=PORT,
reload=False,
......
{
"_comment": "检查项分组配置。叶子节点用 step_id 引用 orchestrator 已注册的 step。改分组只动这个文件;新增 step 只动 orchestrator。两个源头互相解耦。2026-08-14 增加子组(subgroups)层级:字段长度检查拆为 7 个独立 step,放进 subgroups[].children。",
"groups": [
{
"key": "basic",
"title": "基础检查",
"description": "不需要调用大模型的基础数据质量检查",
"default_expand": true,
"children": [
{ "step_id": "empty_fields" },
{ "step_id": "missing_comments" },
{
"key": "length_check_group",
"title": "字段长度检查",
"description": "按字段类型拆分(身份证 / USCC / 手机 / 区划 / 邮编 / 邮箱 / 银行卡),单独勾选 + 单独配置(names / comments)。",
"default_expand": true,
"children": [
{ "step_id": "len_id_card" },
{ "step_id": "len_uscc" },
{ "step_id": "len_mobile" },
{ "step_id": "len_xzqh" },
{ "step_id": "len_postal" },
{ "step_id": "len_email" },
{ "step_id": "len_bank_card" }
]
}
]
},
{
"key": "standards",
"title": "国标字段规范",
"description": "IND-001 ~ IND-019 系列:按 GB / GA / 央行 / 国统字 / ISO 标准做字段值级别合规校验(格式 / 校验位 / 出生日期 / 号段 / 编码存在性 / 邮编 / 银行账号 Luhn / 单位类型 / 经济类型 / 行业代码 等)",
"default_expand": true,
"children": [
{ "step_id": "std_ind_001" },
{ "step_id": "std_ind_002" },
{ "step_id": "std_ind_003" },
{ "step_id": "std_ind_004" },
{ "step_id": "std_ind_005" },
{ "step_id": "std_ind_006" },
{ "step_id": "std_ind_016" },
{ "step_id": "std_ind_017" },
{ "step_id": "std_ind_018" },
{ "step_id": "std_ind_019" }
]
},
{
"key": "business_main",
"title": "业务字段规范",
"description": "IND-401 / IND-007-a:核心业务字段(姓名 / 个人公积金账号)",
"default_expand": true,
"children": [
{ "step_id": "std_ind_401" },
{ "step_id": "std_ind_007_a" }
]
}
]
}
\ No newline at end of file
# ============================================================
# Web 工具默认配置
# 用户通过前端表单提交的内容会覆盖这里的连接信息
# ============================================================
# 数据库类型选项
db_types:
- mysql
- dameng
- oracle
# 采样与阈值(与 CLI 工作流一致)
thresholds:
empty_field:
high: 0.80
mid: 0.50
min_rows_to_check: 5
similarity:
merge_threshold: 0.80
# 标准库
standards:
enabled: true
auto_apply: true
strict_mode: false
# LLM(Claude API)
llm:
model: claude-sonnet-5 # 默认模型
max_tokens: 1024
temperature: 0.2
# API Key 从环境变量 ANTHROPIC_API_KEY 读取,不在此保存
# Web 服务
web:
host: 0.0.0.0
port: 8765
title: "数据分析工具"
\ No newline at end of file
This diff is collapsed.
"""web.core 核心模块
- db_adapter: MySQL / 达梦 / Oracle 适配
- llm: Claude / OpenAI 客户端
- job_manager: 后台任务 + SSE
- orchestrator: 调度 8 步流程
- models: Pydantic 数据模型
"""
from .db_adapter import DBConfig, DBConnection, open_db, test_connection, quote_ident
from .llm import LLMClient, LLMConfig, get_llm_client, LLMUnavailable
__all__ = [
"DBConfig",
"DBConnection",
"open_db",
"test_connection",
"quote_ident",
"LLMClient",
"LLMConfig",
"get_llm_client",
"LLMUnavailable",
]
\ No newline at end of file
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
"""LIKE 模式转义工具(跨 DB 共享)
escape char 用 `!`(单字符),原因:
- MySQL / 达梦 / Oracle 都接受 ESCAPE '<任意单字符>'
- 不需要 SQL 反斜杠转义(避免达梦 [CODE:-6106] 无效的转义字符长度)
- 用户输入里 `!` 极少;万一有,先转义成 `!!`
供 sql_builder / step9 共享,避免再在 step9 内嵌一份。
"""
from __future__ import annotations
ESCAPE_CHAR = "!"
def escape_like(s: str) -> str:
"""对 LIKE 模式里的 ! \\ % _ 做反转义,保证用户输入的字面量不被当通配符。
转义顺序:! 必须最先,否则后面产生的 \\! 会被二次转义成 \\!!。
"""
return (
s.replace("!", "!!") # ! → !! (escape char itself, must first)
.replace("\\", "!\\") # \ → !\
.replace("%", "!%") # % → !%
.replace("_", "!_") # _ → !_
)
\ No newline at end of file
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
"""Step 实现包 — Web 版(兼容 MySQL + 达梦 + Oracle,接 LLM)
每个 step_* 函数接收通用 cfg + 上游 step_outputs,返回结构化 dict。
所有 Step 通过 web.core.orchestrator.run_governance_workflow() 调度。
"""
\ No newline at end of file
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
{
"name": "data-governance-web2",
"name": "data-governance-web",
"private": true,
"version": "0.1.0",
"type": "module",
......
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
Markdown is supported
0%
or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment