From 945b4ace860c5edfecd56e42fd6e8342d027f558 Mon Sep 17 00:00:00 2001 From: wangjiming Date: Thu, 30 Jul 2026 14:38:56 +0800 Subject: [PATCH] update --- UI测试手册.md | 88 +- backend/app/api/v1/endpoints/data_process.py | 2261 ++++++++++++++ backend/app/api/v1/router.py | 10 + backend/app/core/config.py | 2 +- backend/app/db/session.py | 2 +- backend/app/db/sql/002_data_process.sql | 343 ++ backend/app/db/sql/002_governance.sql | 2 + .../app/modules/data_process/algorithms.py | 2763 +++++++++++++++++ backend/app/modules/data_process/constants.py | 7 + .../modules/data_process/dataset_format.py | 148 + .../modules/data_process/document_chunking.py | 444 +++ .../app/modules/data_process/generation.py | 562 ++++ .../modules/data_process/office_preview.py | 308 ++ .../app/modules/data_process/schema_cli.py | 76 + backend/app/modules/data_process/storage.py | 499 +++ backend/app/modules/data_process/store.py | 2526 +++++++++++++++ backend/app/schemas/data_process.py | 373 +++ backend/pyproject.toml | 8 + backend/requirements.txt | 8 + .../dpsf_73355fb704644768b557/v1/111.json | 16 + compute/requirements.txt | 1 + devserver.err | 47 - frontend/package-lock.json | 263 ++ frontend/package.json | 1 + frontend/src/api/modules/dataProcess.ts | 306 ++ frontend/src/api/modules/dataset.ts | 4 +- frontend/src/api/modules/model.ts | 2 +- frontend/src/router/index.ts | 14 +- frontend/src/types/dataProcess.ts | 398 +++ frontend/src/utils/dataProcessStatus.d.ts | 13 + frontend/src/utils/dataProcessStatus.js | 24 + .../src/views/dashboard/DashboardView.vue | 12 +- .../data-process/DataProcessCreateView.vue | 1004 ++++-- .../data-process/DataProcessDetailView.vue | 973 ++++-- .../data-process/DataProcessListView.vue | 174 +- .../create/GenerationOptionsPanel.vue | 246 +- .../data-process/create/GenerationStep.vue | 6 +- .../create/OfficeSourceViewer.vue | 622 ++++ .../data-process/create/PdfSourceViewer.vue | 667 ++++ .../create/PreviewCompareStep.vue | 73 +- .../data-process/create/ResultEditorStep.vue | 243 +- .../data-process/create/SourceUploadStep.vue | 232 +- .../create/StructuredOptionsPanel.vue | 45 +- .../data-process/create/TaskSetupStep.vue | 23 +- .../create/UnstructuredOptionsPanel.vue | 68 +- .../create/data-process-create.scss | 6 + .../create/dataProcessCreateState.ts | 292 +- .../views/data-process/create/previewModel.ts | 525 +--- .../src/views/data-process/create/types.ts | 67 +- .../create/useDataProcessGeneration.ts | 425 ++- .../create/useDataProcessPreviewBuild.ts | 139 + .../create/useDataProcessRegeneration.ts | 254 ++ .../create/useDataProcessSourceUpload.ts | 198 ++ frontend/vite.config.ts | 4 +- 前端功能失效问题排查.md | 134 + 55 files changed, 16633 insertions(+), 1318 deletions(-) create mode 100644 backend/app/api/v1/endpoints/data_process.py create mode 100644 backend/app/db/sql/002_data_process.sql create mode 100644 backend/app/modules/data_process/algorithms.py create mode 100644 backend/app/modules/data_process/constants.py create mode 100644 backend/app/modules/data_process/dataset_format.py create mode 100644 backend/app/modules/data_process/document_chunking.py create mode 100644 backend/app/modules/data_process/generation.py create mode 100644 backend/app/modules/data_process/office_preview.py create mode 100644 backend/app/modules/data_process/schema_cli.py create mode 100644 backend/app/modules/data_process/storage.py create mode 100644 backend/app/modules/data_process/store.py create mode 100644 backend/app/schemas/data_process.py create mode 100644 backend/storage/data-process/dpt_3ed8684864a34688b8ee/dpsf_73355fb704644768b557/v1/111.json delete mode 100644 devserver.err create mode 100644 frontend/src/api/modules/dataProcess.ts create mode 100644 frontend/src/types/dataProcess.ts create mode 100644 frontend/src/utils/dataProcessStatus.d.ts create mode 100644 frontend/src/utils/dataProcessStatus.js create mode 100644 frontend/src/views/data-process/create/OfficeSourceViewer.vue create mode 100644 frontend/src/views/data-process/create/PdfSourceViewer.vue create mode 100644 frontend/src/views/data-process/create/useDataProcessPreviewBuild.ts create mode 100644 frontend/src/views/data-process/create/useDataProcessRegeneration.ts create mode 100644 frontend/src/views/data-process/create/useDataProcessSourceUpload.ts create mode 100644 前端功能失效问题排查.md diff --git a/UI测试手册.md b/UI测试手册.md index a481e28..41b80c6 100644 --- a/UI测试手册.md +++ b/UI测试手册.md @@ -1,6 +1,6 @@ # 第 1~4 周功能 · 界面人工测试手册 -> 用途:你按这份手册在浏览器里点一遍,验证第 1~4 周的功能(登录、租户/项目、资源 ACL、审计日志、审批中心、写操作自动审计与审批拦截)。 +> 用途:你按这份手册在浏览器里点一遍,验证第 1~4 周的功能(登录、租户/项目、资源 ACL、审计日志、审批中心、数据处理 C 模块、写操作自动审计与审批拦截)。 > 功能代码层均已联调通过(含此前修复的 `audit.ts`/`approval.ts` 双重解包、审计接口 `/system` 前缀、审批模块 `include_router` 启动崩溃)。下面是给你的人工回归步骤。 --- @@ -276,7 +276,77 @@ npm run dev --- -## 6. 验收清单(打勾) +## 6. 数据处理(C 模块 · 第 2 周) + +> 入口:左侧菜单「数据治理 → 数据处理」(`/data-process`),对应后端 `/modelTF/data-process/*`。该模块为本次合并从 `yg_ft1` 并入的可运行子系统(后端 `app/modules/data_process/` + 前端 `views/data-process/`)。 +> 创建任务为**六步向导**:`创建任务 → 大模型选择 → 上传文件 → 数据预览 → 开始生成 → 结果编辑与保存`。后端启动时会自动建表(`002_data_process.sql`),无需手动迁移。 + +### 6.1 列表与入口 +1. 左侧菜单「数据治理 → 数据处理」 +2. 预期:进入列表页 `/data-process`,展示已有数据处理任务(空状态显示「暂无数据」属正常,不要当成 bug) +3. 点右上「新建数据处理」进入创建向导 + +### 6.2 创建任务(六步向导) + +**步骤 1 · 创建任务** +1. 填任务名称(如 `test-dp-manual`)、描述(可选) +2. 选处理类型: + - `结构化数据`(structured):按行生成 QA 对 + - `非结构化数据`(unstructured):文档切片(chunk 方法/大小/重叠/保留表格代码块等) + - `外部数据源`(external):接 PostgreSQL 等外部库 +3. 配置处理选项(预处理、语义增强、数据集切分比例、输出类型、温度等,保持默认即可) +4. 点「继续:选择大模型」 + +**步骤 2 · 大模型选择** +1. 从模型下拉选一个生成模型(需「模型管理」里已登记可用模型;若下拉为空,先到「模型管理」登记一个基座/API 模型) +2. 设输出要求(生成模型、提示词、输出类型、reasoning、温度、max_tokens、质量过滤开关等) +3. 点「继续:上传文件」 + +**步骤 3 · 上传文件** +- 结构化/非结构化:点「上传」选本地文件(jsonl/csv/pdf/docx/xlsx/txt 等,按处理类型校验扩展名);或点「使用样例文件」快速载入示例 `finance_qa.jsonl` +- 外部数据源:填连接(postgresql URL、认证模式、只读 `SELECT` 查询),点「测试连接」→ 联通后「拉取数据」 +1. 选/拉取至少一个源文件,等待上传完成(状态变 `ready`) +2. 点「继续:数据预览」,系统自动按配置**切分**(进度条,可能耗时;失败会提示原因,可重试) + +**步骤 4 · 数据预览** +1. 切分完成后进入预览,左侧为源文件 / 分片列表,右侧为预览内容 +2. 可编辑预览条目内容、新增 / 删除条目、还原为原文(不影响源文件) +3. 切换不同源文件核对切分结果 +4. 点「确认预览并继续」 + +**步骤 5 · 开始生成** +1. 进入生成页,显示任务摘要(任务名、处理类型、文件、预览条目数、修改条数) +2. 点「开始生成」启动处理(调用所选大模型) +3. 等待生成完成:状态 `running → completed`(失败显示原因,可「重新生成」) +4. 生成成功后点「查看生成结果」 + +**步骤 6 · 结果编辑与保存** +1. 进入结果编辑页,逐条检查结果(问题 / 答案等),可编辑字段、单条或批量「重新生成」 +2. 校验通过后点「保存任务」 +3. 预期:提示「生成结果已确认」,跳转到任务详情页 `/data-process/<任务ID>` + +### 6.3 任务详情 +1. 列表点目标任务「详情」(或保存后自动进入) +2. 预期:详情页展示任务信息、处理配置、源文件、预览条目、生成结果 +3. 支持「重新生成」(`/data-process/:id/regenerate`)与「任务进度」(`/data-process/:id/workflow`)两个子页,均复用创建向导 + +### 6.4 重新生成(可选) +1. 详情页或列表进入「重新生成」子页 +2. 调整配置(保持原处理类型),保存后按新配置重新切分 / 生成;原生成结果与已发布数据在点击「开始生成」前保持不变 +3. 验证:提交重新生成后,进度可在「任务进度」页查看 + +### 6.5 发布数据集(第 2 周验收点) +1. 结果确认后,在详情页点「发布数据集」(对应 `POST /modelTF/data-process/{id}/publish`) +2. 填目标数据集元信息(名称 / 项目归属等),确认发布 +3. 预期:生成数据集记录,可在「数据集管理」看到该发布数据集,来源链路保留 + +### 6.6 删除任务(可选) +1. 列表行点「删除」,二次确认后删除 +2. 预期:列表不再显示该任务 + +--- + +## 7. 验收清单(打勾) **第 1~2 周** - [ ] 1. 登录成功,进仪表盘,显示 `admin` @@ -287,6 +357,14 @@ npm run dev - [ ] 6. 归档:点后状态变 `archived` - [ ] 7. 用户设置:创建用户(账号/显示名/初始密码/角色/状态/页面权限均可设,列表出现新账号)、重置密码(自定义生效 / 留空回退 `platform123` / `admin` 被拒)、启用禁用账号(禁用后登录失败、启用后恢复)、页面权限精细控制(取消权限后登录即生效)、删除账号(二次确认后列表移除、登录失败) +**第 2 周(数据处理 · C 模块)** +- [ ] 16. 数据处理入口可见(数据治理 → 数据处理),列表 / 空状态正常 +- [ ] 17. 创建向导六步可走通:创建任务 → 选模型 → 上传/拉取源文件 → 切分预览(可编辑/增删/还原)→ 生成 → 结果保存并跳详情 +- [ ] 18. 结构 / 非结构化 / 外部数据源三类处理类型均可配置并产出预览条目 +- [ ] 19. 生成结果可编辑、可单条 / 批量重新生成,保存后跳详情 +- [ ] 20. 发布数据集成功,「数据集管理」可见且保留来源链路 +- [ ] 21. 全程浏览器控制台(F12 → Console)在数据处理流程中无红色报错 + **第 3 周(审计)** - [ ] 8. 审计日志:按 动作/操作人/项目 过滤均能返回正确结果 - [ ] 9. 审计日志:导出 CSV 成功,列与内容正确 @@ -301,19 +379,19 @@ npm run dev --- -## 7. 我自测已覆盖(你不用重复,除非想验证) +## 8. 我自测已覆盖(你不用重复,除非想验证) - 后端真实导入:`import app.main` → `IMPORT_OK`,启动崩溃已修复(`approval/__init__.py` 补 re-export `router`)。 - 前端 `type-check` 全绿:`audit.ts`/`approval.ts` 双重解包已改 `get/post`;审计接口已加 `/system` 前缀(`/system/audit-logs`、`/system/audit-logs/export`)。 - 接口链路已用真实代码核对:审计查询/导出(`system`)、审批模板/实例/逐步决策(`approvals`)、项目写操作自动 `record_audit` 与 `_require_no_pending_approval` 拦截均按上述行为实现。 -## 8. 已知非 bug / 注意事项(仅供参考) +## 9. 已知非 bug / 注意事项(仅供参考) 1. 前端由你自己在 **Windows 终端** 跑 `npm run dev` 启动(默认 **16801**)。不要从 WSL 终端启动(会慢 8 倍)。 2. 「归档项目」当前是**直接执行无确认弹窗**——功能正确,建议后续补个二次确认,避免误操作。 4. 控制台偶见的 `ERR_ABORTED` 是导航时浏览器正常中止旧 CSS 请求,无害;Google Fonts 外网字体加载失败不影响功能。 5. 审计「操作人」列 = 登录 token(当前登录用户标识),由前端 `Authorization: Bearer ` 透传,非真实姓名。 -## 9. 清理测试数据(可选) +## 10. 清理测试数据(可选) 手动建的 `test-tenant-manual` / `test-project-manual`、审批实例/模板可在对应列表里删除,或告诉我帮你清库(后端连 PostgreSQL,`PlatformStore` 启动时自动建表)。 diff --git a/backend/app/api/v1/endpoints/data_process.py b/backend/app/api/v1/endpoints/data_process.py new file mode 100644 index 0000000..1500a64 --- /dev/null +++ b/backend/app/api/v1/endpoints/data_process.py @@ -0,0 +1,2261 @@ +from __future__ import annotations + +import hashlib +import ipaddress +import json +import logging +import os +import re +import socket +import time +from concurrent.futures import ThreadPoolExecutor, as_completed +from contextlib import contextmanager +from dataclasses import asdict +from pathlib import Path +from threading import BoundedSemaphore, Lock +from typing import Any, Iterator, Literal +from urllib.parse import quote, urlsplit + +import httpx +import psycopg +from fastapi import ( + APIRouter, + BackgroundTasks, + Body, + Depends, + File, + Header, + HTTPException, + Query, + UploadFile, +) +from fastapi.responses import StreamingResponse +from psycopg.rows import dict_row + +from app.modules.data_process.algorithms import ( + ParsedText, + canonical_record_json, + content_quality_flags, + desensitize_pii, + desensitize_structured_record, + detect_pdf_document_noise, + estimate_token_count, + extract_pdf_page_texts, + generate_standard_records, + is_near_duplicate, + near_duplicate_fingerprint, + parse_text_content, + preprocess_structured_records, + remove_document_noise, + score_quality, +) +from app.modules.data_process.document_chunking import ( + DocumentChunk, + chunk_fixed_text, + chunk_layout_document, + chunk_semantic_text, + merge_short_chunks, +) +from app.modules.data_process.generation import generate_model_records +from app.modules.data_process.office_preview import ( + MAX_XLSX_PREVIEW_ROWS, + build_docx_preview, + build_xlsx_preview, +) +from app.modules.data_process.storage import ( + LocalDataProcessStorage, + StagedSourceObject, + get_data_process_storage, +) +from app.modules.data_process.store import ( + ConflictError, + DataProcessStore, + DataProcessStoreError, + InvalidStateError, + NotFoundError, + get_data_process_store, + new_id, +) +from app.schemas.data_process import ( + DataProcessRegenerateRequest, + DataProcessStatus, + DataProcessTaskCreate, + DataProcessTaskUpdate, + DataProcessWorkflowStepUpdate, + ExternalPullRequest, + ExternalSourceRequest, + GenerateRequest, + PreviewBuildRequest, + PreviewItemCreate, + PreviewItemUpdate, + ProcessType, + PublishRequest, + ResultBatchRegenerateRequest, + ResultRegenerateRequest, + ResultUpdate, +) + +router = APIRouter(prefix="/data-process") +logger = logging.getLogger(__name__) +MAX_SOURCE_FILE_BYTES = 200 * 1024 * 1024 +MAX_SOURCE_FILE_COUNT = 20 +MAX_SOURCE_BATCH_BYTES = 500 * 1024 * 1024 +MAX_EXTERNAL_PULL_BYTES = 50 * 1024 * 1024 +STRUCTURED_SOURCE_SUFFIXES = frozenset( + {".json", ".jsonl", ".ndjson", ".csv", ".tsv", ".xlsx"} +) +UNSTRUCTURED_SOURCE_SUFFIXES = frozenset( + { + ".txt", + ".md", + ".markdown", + ".pdf", + ".docx", + ".pptx", + ".json", + ".jsonl", + ".ndjson", + } +) +SUPPORTED_SOURCE_SUFFIXES = STRUCTURED_SOURCE_SUFFIXES | UNSTRUCTURED_SOURCE_SUFFIXES +LEGACY_OFFICE_CONVERSIONS = { + ".doc": ".docx", + ".xls": ".xlsx", + ".ppt": ".pptx", +} +RAW_INLINE_PREVIEW_MEDIA_TYPES = { + "pdf": "application/pdf", + "docx": "application/vnd.openxmlformats-officedocument.wordprocessingml.document", + "xlsx": "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet", + "pptx": "application/vnd.openxmlformats-officedocument.presentationml.presentation", +} +RESULT_REGENERATION_CONCURRENCY = 4 +RESULT_REGENERATION_RETRIES = 0 +RESULT_REGENERATION_TIMEOUT_SECONDS = 60.0 +_result_regeneration_slots = BoundedSemaphore(RESULT_REGENERATION_CONCURRENCY) +_result_regeneration_claims_lock = Lock() +_active_result_regenerations: set[tuple[str, str]] = set() + + +def ok(data: Any = None, message: str = "ok") -> dict[str, Any]: + return {"code": 0, "message": message, "data": data} + + +def fail(status_code: int, message: str) -> HTTPException: + return HTTPException( + status_code=status_code, + detail={"code": status_code, "message": message, "data": None}, + ) + + +@contextmanager +def api_errors() -> Iterator[None]: + try: + yield + except NotFoundError as exc: + raise fail(404, str(exc)) from exc + except ConflictError as exc: + raise fail(409, str(exc)) from exc + except InvalidStateError as exc: + raise fail(409, str(exc)) from exc + except (DataProcessStoreError, ValueError) as exc: + raise fail(400, str(exc)) from exc + except (psycopg.errors.UndefinedTable, psycopg.errors.UndefinedColumn) as exc: + raise fail( + 503, + "data process schema is missing or out of date; run schema_cli --check", + ) from exc + except psycopg.OperationalError as exc: + raise fail(503, "data process database is unavailable") from exc + + +def _safe_file_name(value: str | None, fallback: str) -> str: + name = Path((value or "").replace("\\", "/")).name.replace("\x00", "").strip() + return name if name not in {"", ".", ".."} else fallback + + +def _range_not_satisfiable(size: int) -> HTTPException: + return HTTPException( + status_code=416, + detail={"code": 416, "message": "invalid source byte range", "data": None}, + headers={"Content-Range": f"bytes */{size}"}, + ) + + +def _source_byte_range(value: str | None, size: int) -> tuple[int, int] | None: + if value is None: + return None + match = re.fullmatch(r"bytes=(\d*)-(\d*)", value.strip()) + if match is None or size <= 0: + raise _range_not_satisfiable(size) + start_text, end_text = match.groups() + if not start_text and not end_text: + raise _range_not_satisfiable(size) + if start_text: + start = int(start_text) + end = int(end_text) if end_text else size - 1 + if start >= size or end < start: + raise _range_not_satisfiable(size) + else: + suffix_length = int(end_text) + if suffix_length <= 0: + raise _range_not_satisfiable(size) + start = max(0, size - suffix_length) + end = size - 1 + return start, min(end, size - 1) + + +def _commit_source_batch( + store: DataProcessStore, + storage: LocalDataProcessStorage, + task_id: str, + prepared: list[dict[str, Any]], + staged: list[StagedSourceObject], +) -> list[dict[str, Any]]: + storage.publish(staged) + try: + return store.add_source_files(task_id, prepared) + except Exception: + for item in staged: + try: + storage.delete(item.reference) + except Exception: + # 文件系统回滚失败不能覆盖数据库抛出的根因,并继续清理其余对象。 + logger.exception( + "failed to roll back data process source object task_id=%s", + task_id, + ) + raise + + +def _value(config: dict[str, Any], snake_name: str, camel_name: str, default: Any) -> Any: + if snake_name in config: + return config[snake_name] + return config.get(camel_name, default) + + +def _preprocess_options(config: dict[str, Any]) -> set[str]: + values = _value(config, "preprocess_options", "preprocessOptions", []) + return {str(item) for item in values} if isinstance(values, list) else set() + + +def _preview_quality(content: str, config: dict[str, Any]) -> dict[str, Any]: + records = generate_standard_records( + [{"id": "quality-preview", "edited_content": content}], + split={"train": 100, "validation": 0, "test": 0}, + ) + record = records[0] if records else {"instruction": "", "input": "", "output": ""} + minimum = int(_value(config, "min_output_length", "minOutputLength", 20) or 20) + return asdict( + score_quality( + record, + min_output_length=max(1, minimum), + source_content=content, + ) + ) + + +def _parse_stored_source(source: dict[str, Any]) -> ParsedText: + """重新读取已入库的规范化正文,避免把二进制格式当二进制重复解析。""" + + content = str(source.get("content") or "") + file_format = str(source.get("file_format") or "").lower() + if file_format == "xlsx": + # XLSX 上传阶段已安全解析为 JSONL 后入库。 + return parse_text_content(content, file_format="jsonl") + if file_format in {"pdf", "docx", "pptx"}: + # 文档上传阶段已抽取文本,预览阶段只需要对正文切片。 + return ParsedText(format=file_format, text=content, records=()) + return parse_text_content( + content, + filename=source.get("name"), + file_format=file_format or None, + ) + + +def _chunk_source_text( + source: dict[str, Any], + config: dict[str, Any], + preprocess_options: set[str], +) -> list[DocumentChunk]: + """根据任务配置调用真实的 Docling/LlamaIndex 切分器。""" + + method = str(_value(config, "chunk_method", "chunkMethod", "layout_hybrid")) + preserve_context = "preserve_context" in preprocess_options + chunk_size = int(_value(config, "chunk_size", "chunkSize", 800)) + overlap = ( + int(_value(config, "chunk_overlap", "chunkOverlap", 100)) + if preserve_context + else 0 + ) + text = str(source.get("content") or "") + if method == "layout_hybrid": + raw = source.get("raw_content") + if not isinstance(raw, bytes): + raise InvalidStateError("版面结构混合切分需要原始文件,请重新上传后再处理") + chunks = chunk_layout_document( + raw, + filename=str(source.get("name") or "document.pdf"), + source_text=text, + chunk_size=chunk_size, + ) + elif method == "semantic": + chunks = chunk_semantic_text( + text, + chunk_size=chunk_size, + chunk_overlap=overlap, + breakpoint_percentile_threshold=int( + _value( + config, + "semantic_breakpoint_percentile", + "semanticBreakpointPercentile", + 95, + ) + ), + ) + elif method == "fixed": + chunks = chunk_fixed_text(text, chunk_size=chunk_size, chunk_overlap=overlap) + else: + raise ValueError(f"unsupported chunk method: {method}") + if "merge_short_content" in preprocess_options: + chunks = merge_short_chunks( + chunks, + source_text=text, + min_token_count=int(_value(config, "min_chunk_size", "minChunkSize", 100)), + max_token_count=chunk_size, + ) + return chunks + + +_NEGATION_MARKERS = frozenset({"不", "无", "未", "否", "没有", "并非", "not", "no", "never"}) + + +def _safe_near_duplicate(left: str, right: str) -> bool: + """保守判断近重复,数字或否定含义变化时始终保留。""" + + if min(estimate_token_count(left), estimate_token_count(right)) < 20: + return False + if re.findall(r"\d+(?:\.\d+)?", left) != re.findall(r"\d+(?:\.\d+)?", right): + return False + left_lower = left.casefold() + right_lower = right.casefold() + left_negations = {marker for marker in _NEGATION_MARKERS if marker in left_lower} + right_negations = {marker for marker in _NEGATION_MARKERS if marker in right_lower} + if left_negations != right_negations: + return False + return is_near_duplicate( + left, + right, + similarity_threshold=0.92, + max_hamming_distance=2, + ) + + +def _near_duplicate_band_keys(content: str) -> tuple[tuple[int, int], ...]: + """将 64 位 SimHash 分为三段,汉明距离不超过 2 时至少命中一段。""" + + fingerprint = int(near_duplicate_fingerprint(content), 16) + widths = (22, 21, 21) + shift = 0 + keys: list[tuple[int, int]] = [] + for index, width in enumerate(widths): + keys.append((index, (fingerprint >> shift) & ((1 << width) - 1))) + shift += width + return tuple(keys) + + +def _build_preview_items( + task: dict[str, Any], source_files: list[dict[str, Any]] +) -> list[dict[str, Any]]: + config = task.get("config") or {} + process_type = task["process_type"] + preprocess_options = _preprocess_options(config) + should_desensitize = "desensitize" in preprocess_options + should_clean_invalid = bool( + preprocess_options & {"clean_invalid", "clean_invalid_content"} + ) + should_deduplicate = bool( + preprocess_options & {"deduplicate", "deduplicate_content"} + ) + seen_content_hashes: set[str] = set() + seen_near_duplicate_bands: dict[tuple[int, int], list[str]] = {} + items: list[dict[str, Any]] = [] + + def append_item(item: dict[str, Any]) -> None: + content = str(item.get("edited_content") or "").strip() + if should_clean_invalid and not content: + return + content_hash = hashlib.sha256(content.encode("utf-8")).hexdigest() + if should_deduplicate and content_hash in seen_content_hashes: + return + seen_content_hashes.add(content_hash) + if not content: + item["status"] = "invalid" + items.append(item) + + for source in source_files: + parsed = _parse_stored_source(source) + if process_type == "unstructured": + document_noise_spans = tuple(source.get("document_noise_spans") or ()) + chunks = _chunk_source_text(source, config, preprocess_options) + for chunk in chunks: + content = ( + remove_document_noise( + chunk.contextualized_content, + document_noise_spans, + source_offset=chunk.source_start or 0, + ) + if ( + should_clean_invalid + and document_noise_spans + and chunk.source_start is not None + ) + else chunk.contextualized_content + ) + preprocess_flags = content_quality_flags( + content, + min_chars=0, + min_tokens=0, + ) + if content != chunk.contextualized_content: + preprocess_flags = (*preprocess_flags, "document_noise_removed") + flag_set = set(preprocess_flags) + if "clean_invalid_content" in preprocess_options and flag_set & { + "empty_content", + "low_printable_ratio", + "repetitive_content", + }: + continue + if "filter_low_quality" in preprocess_options and flag_set & { + "content_too_long", + "mojibake", + "low_printable_ratio", + "repetitive_content", + }: + continue + if "deduplicate_content" in preprocess_options: + band_keys = _near_duplicate_band_keys(content) + candidates = { + previous + for key in band_keys + for previous in seen_near_duplicate_bands.get(key, ()) + } + if any( + _safe_near_duplicate(content, previous) + for previous in candidates + ): + continue + for key in band_keys: + seen_near_duplicate_bands.setdefault(key, []).append(content) + pii_counts: dict[str, int] = {} + if should_desensitize: + content, pii_counts = desensitize_pii(content) + quality = _preview_quality(content, config) + quality["pii_replacements"] = pii_counts + quality["preprocess_flags"] = list(preprocess_flags) + quality["chunk_method"] = str( + _value(config, "chunk_method", "chunkMethod", "layout_hybrid") + ) + quality["heading_path"] = list(chunk.heading_path) + quality["source_pages"] = list(chunk.source_pages) + quality["doc_item_refs"] = list(chunk.doc_item_refs) + quality["source_bboxes"] = list(chunk.source_bboxes) + append_item( + { + "source_file_id": source["id"], + "original_content": chunk.original_content, + "edited_content": content, + "source_start": chunk.source_start, + "source_end": chunk.source_end, + "source_start_line": chunk.source_start_line, + "source_end_line": chunk.source_end_line, + "token_count": estimate_token_count(content), + "status": ( + "modified" + if content != chunk.original_content + else "original" + ), + "quality_score": quality, + } + ) + continue + + structured_options = preprocess_options & { + "clean_invalid", + "detect_structure", + "deduplicate", + "normalize_format", + "filter_anomaly", + } + source_records = list(parsed.records) + processed_records = preprocess_structured_records( + source_records, + structured_options, + ) + if not processed_records and parsed.text and not source_records: + processed_records = [{"value": parsed.text}] + same_cardinality = len(processed_records) == len(source_records) + for index, record in enumerate(processed_records): + original_record = source_records[index] if same_cardinality else record + original_content = json.dumps( + original_record, + ensure_ascii=False, + separators=(",", ":"), + ) + pii_counts: dict[str, int] = {} + edited_record = record + if should_desensitize: + edited_record, pii_counts = desensitize_structured_record(record) + content = ( + canonical_record_json(edited_record) + if "normalize_format" in preprocess_options + else json.dumps( + edited_record, + ensure_ascii=False, + separators=(",", ":"), + ) + ) + quality = _preview_quality(content, config) + quality["pii_replacements"] = pii_counts + append_item( + { + "source_file_id": source["id"], + "original_content": original_content, + "edited_content": content, + "source_start": None, + "source_end": None, + "source_start_line": None, + "source_end_line": None, + "token_count": estimate_token_count(content), + "status": "modified" if content != original_content else "original", + "quality_score": quality, + } + ) + return items + + +def _all_preview_items(store: DataProcessStore, task_id: str) -> list[dict[str, Any]]: + """分页读取全部预览项,避免固定上限静默截断任务。""" + + items: list[dict[str, Any]] = [] + page = 1 + page_size = 5_000 + while True: + result = store.list_preview_items(task_id, page=page, page_size=page_size) + batch = result["items"] + items.extend(batch) + if len(items) >= int(result["total"]) or not batch: + return items + page += 1 + + +def _run_generation( + store: DataProcessStore, task_id: str, generation_run_id: str +) -> None: + started_at = time.perf_counter() + logger.info( + "data process generation worker started task_id=%s generation_run_id=%s", + task_id, + generation_run_id, + ) + try: + task = store.get_task(task_id) + if not store.generation_is_running(task_id, generation_run_id): + logger.info( + "data process generation worker skipped inactive run task_id=%s " + "generation_run_id=%s", + task_id, + generation_run_id, + ) + return + all_preview_items = _all_preview_items(store, task_id) + preview_items = [ + item + for item in all_preview_items + if item.get("status") != "invalid" + and str(item.get("edited_content") or item.get("original_content") or "").strip() + ] + pre_filtered_count = len(all_preview_items) - len(preview_items) + config = task.get("config") or {} + model_id = _value(config, "generation_model_id", "generationModelId", None) + generation_model: dict[str, Any] | None = None + if model_id: + generation_model = store.get_generation_model(str(model_id)) + task = store.save_generation_model_snapshot( + task_id, + generation_model, + generation_run_id=generation_run_id, + ) + config = task.get("config") or config + split = _value( + config, + "dataset_split", + "datasetSplit", + {"train": 80, "validation": 10, "test": 10}, + ) + pairs = ( + _value(config, "qa_pairs_per_chunk", "qaPairsPerChunk", 1) + if task["process_type"] == "unstructured" + else _value(config, "qa_pairs_per_row", "qaPairsPerRow", 1) + ) + output_type = str( + _value(config, "output_type", "outputType", "standard") + ).strip().lower() + if generation_model: + runtime_config = { + **config, + "generation_prompt": _value( + config, "generation_prompt", "generationPrompt", "" + ), + "output_type": output_type, + "reasoning_detail": _value( + config, "reasoning_detail", "reasoningDetail", "normal" + ), + "max_tokens": _value(config, "max_tokens", "maxTokens", 1024), + "json_mode": _value(config, "json_mode", "jsonMode", False), + } + def report_progress(processed_count: int, total_count: int) -> None: + if not store.update_generation_progress( + task_id, + generation_run_id, + processed_count, + total_count, + ): + raise InvalidStateError("generation run is no longer active") + + generated = generate_model_records( + preview_items, + model=generation_model, + config=runtime_config, + task_id=task_id, + split=split, + qa_pairs_per_item=int(pairs or 1), + on_progress=report_progress, + ) + elif output_type == "reasoning": + raise InvalidStateError("思维链输出必须配置可用的数据生成模型") + else: + generated = generate_standard_records( + preview_items, + qa_pairs_per_item=int(pairs or 1), + semantic_enrichment=bool( + _value(config, "semantic_enrichment", "semanticEnrichment", False) + ), + split=split, + split_seed=task_id, + ) + if not store.update_generation_progress( + task_id, + generation_run_id, + len(preview_items), + len(preview_items), + ): + logger.info( + "data process generation stopped before completion task_id=%s " + "generation_run_id=%s", + task_id, + generation_run_id, + ) + return + + known_fingerprints: set[str] = set() + accepted: list[dict[str, Any]] = [] + filtered_count = pre_filtered_count + duplicate_count = 0 + error_count = 0 + quality_filter = bool( + _value(config, "quality_filter_enabled", "qualityFilterEnabled", False) + ) + filter_low = bool(_value(config, "filter_low_quality", "filterLowQuality", True)) + filter_short = bool( + _value(config, "filter_short_content", "filterShortContent", True) + ) + deduplicate = bool( + _preprocess_options(config) + & {"deduplicate", "deduplicate_content"} + ) + minimum = max(1, int(_value(config, "min_output_length", "minOutputLength", 20) or 20)) + preview_sources = { + str(item["id"]): str( + item.get("edited_content") or item.get("original_content") or "" + ) + for item in preview_items + } + + for record in generated: + quality = score_quality( + record, + min_output_length=minimum, + source_content=preview_sources.get(str(record.get("preview_item_id") or ""), ""), + known_fingerprints=known_fingerprints, + ) + if "duplicate_record" in quality.flags: + duplicate_count += 1 + else: + # 即使首条随后因短文本/低质量被过滤,也要阻止同批后续重复结果。 + known_fingerprints.add(quality.fingerprint) + should_filter = ( + (deduplicate and "duplicate_record" in quality.flags) + or ( + quality_filter + and filter_short + and "output_too_short" in quality.flags + ) + or (quality_filter and filter_low and not quality.is_valid) + ) + if not quality.is_valid: + error_count += 1 + record["status"] = "invalid" + record["error"] = ", ".join(quality.flags) or "quality validation failed" + if should_filter: + filtered_count += 1 + continue + record["quality_score"] = asdict(quality) + accepted.append(record) + + # stop 请求可能在纯函数计算期间到达,最终写入前再次检查状态。 + if store.generation_is_running(task_id, generation_run_id): + completed = store.complete_generation( + task_id, + accepted, + generation_run_id=generation_run_id, + filtered_count=filtered_count, + duplicate_count=duplicate_count, + error_count=error_count, + ) + logger.info( + "data process generation completed task_id=%s generation_run_id=%s " + "output_count=%s filtered_count=%s duplicate_count=%s error_count=%s " + "duration_ms=%.2f", + task_id, + generation_run_id, + completed.get("output_count", len(accepted)), + completed.get("filtered_count", filtered_count), + completed.get("duplicate_count", duplicate_count), + completed.get("error_count", error_count), + (time.perf_counter() - started_at) * 1000, + ) + else: + logger.info( + "data process generation stopped before result persistence task_id=%s " + "generation_run_id=%s", + task_id, + generation_run_id, + ) + except Exception as exc: + logger.exception( + "data process generation failed task_id=%s generation_run_id=%s duration_ms=%.2f", + task_id, + generation_run_id, + (time.perf_counter() - started_at) * 1000, + ) + try: + if store.generation_is_running(task_id, generation_run_id): + store.mark_failed( + task_id, + str(exc), + generation_run_id=generation_run_id, + ) + except Exception: + logger.exception( + "failed to persist data process generation failure task_id=%s " + "generation_run_id=%s", + task_id, + generation_run_id, + ) + + +@router.get("") +def list_tasks( + page: int = Query(default=1, ge=1), + page_size: int = Query(default=20, ge=1, le=200), + keyword: str | None = Query(default=None), + status: DataProcessStatus | None = Query(default=None), + process_type: ProcessType | None = Query(default=None), + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + return ok( + store.list_tasks( + page=page, + page_size=page_size, + keyword=keyword, + status=status, + process_type=process_type, + ) + ) + + +@router.post("") +def create_task( + payload: DataProcessTaskCreate, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + task = store.create_task(payload.model_dump(mode="json")) + return ok(task, "data process task created") + + +@router.get("/{task_id}") +def task_detail( + task_id: str, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + # 兼容旧版曾在第五步前清空结果的异常任务;严格特征匹配且幂等。 + store.recover_legacy_aborted_regeneration(task_id) + task = store.get_task(task_id) + source_files = store.list_source_files(task_id) + task["source_files"] = source_files + task["source_file_count"] = len(source_files) + return ok(task) + + +@router.put("/{task_id}") +def update_task( + task_id: str, + payload: DataProcessTaskUpdate, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + return ok( + store.update_task(task_id, payload.model_dump(exclude_unset=True, mode="json")), + "data process task updated", + ) + + +@router.put("/{task_id}/workflow-step") +def update_workflow_step( + task_id: str, + payload: DataProcessWorkflowStepUpdate, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + return ok( + store.update_workflow_step(task_id, payload.workflow_step.value), + "data process workflow step updated", + ) + + +@router.post("/{task_id}/regenerate") +def prepare_regeneration( + task_id: str, + payload: DataProcessRegenerateRequest, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + return ok( + store.prepare_regeneration(task_id, payload.model_dump(mode="json")), + "data process task prepared for regeneration", + ) + + +@router.delete("/{task_id}") +def delete_task( + task_id: str, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + store.delete_task(task_id) + return ok({"deleted": task_id}, "data process task deleted") + + +@router.get("/{task_id}/source-files") +def source_files( + task_id: str, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + return ok({"files": store.list_source_files(task_id)}) + + +@router.post("/{task_id}/source-files") +async def upload_source_files( + task_id: str, + files: list[UploadFile] = File(...), + store: DataProcessStore = Depends(get_data_process_store), + storage: LocalDataProcessStorage = Depends(get_data_process_storage), +) -> dict[str, Any]: + if not files: + raise fail(400, "at least one source file is required") + if len(files) > MAX_SOURCE_FILE_COUNT: + raise fail(413, f"a source batch may contain at most {MAX_SOURCE_FILE_COUNT} files") + prepared: list[dict[str, Any]] = [] + staged: list[StagedSourceObject] = [] + batch_id = storage.new_batch_id() + batch_size = 0 + commit_attempted = False + try: + with api_errors(): + task = store.get_task(task_id) + process_type = str(task["process_type"]) + if process_type == "external": + raise InvalidStateError( + "external tasks must import data through the external source endpoint" + ) + allowed_suffixes = ( + UNSTRUCTURED_SOURCE_SUFFIXES + if process_type == "unstructured" + else STRUCTURED_SOURCE_SUFFIXES + ) + for upload in files: + raw = await upload.read(MAX_SOURCE_FILE_BYTES + 1) + if len(raw) > MAX_SOURCE_FILE_BYTES: + raise fail(413, f"source file exceeds {MAX_SOURCE_FILE_BYTES} bytes") + name = _safe_file_name(upload.filename, "source.txt") + suffix = Path(name).suffix.lower() + if suffix in LEGACY_OFFICE_CONVERSIONS: + replacement = LEGACY_OFFICE_CONVERSIONS[suffix] + raise fail( + 415, + f"legacy {suffix} format is not supported; " + f"convert the file to {replacement} and upload again", + ) + if suffix not in SUPPORTED_SOURCE_SUFFIXES: + raise fail(415, f"unsupported source file format: {suffix or 'none'}") + if suffix not in allowed_suffixes: + raise fail( + 415, + f"{suffix} is not supported for {process_type} data processing", + ) + parsed = parse_text_content(raw, filename=name) + if not parsed.text: + raise fail(400, f"source file is empty: {name}") + batch_size += len(raw) + if batch_size > MAX_SOURCE_BATCH_BYTES: + raise fail(413, f"source batch exceeds {MAX_SOURCE_BATCH_BYTES} bytes") + source_file_id = new_id("dpsf") + staged_object = storage.stage_bytes( + batch_id=batch_id, + task_id=task_id, + source_file_id=source_file_id, + version=1, + name=name, + content=raw, + ) + staged.append(staged_object) + record_count = len(parsed.records) or (1 if parsed.text else 0) + prepared.append( + { + "id": source_file_id, + "storage_object_id": staged_object.reference, + "name": name, + "content": parsed.text, + "raw_size": len(raw), + "checksum_sha256": hashlib.sha256(raw).hexdigest(), + "file_format": parsed.format, + "record_count": record_count, + "metadata": { + "content_type": upload.content_type or "text/plain", + "original_size_bytes": len(raw), + "original_checksum_sha256": hashlib.sha256(raw).hexdigest(), + }, + "created_by": None, + } + ) + # publish 无论成功或失败都会消费并清理暂存对象,外层不能再次 discard。 + commit_attempted = True + created = _commit_source_batch(store, storage, task_id, prepared, staged) + finally: + if not commit_attempted: + storage.discard(staged) + return ok({"files": created}, "source files uploaded") + + +@router.get("/{task_id}/source-files/{file_id}/content") +def source_file_content( + task_id: str, + file_id: str, + start_line: int | None = Query(default=None, ge=1), + line_count: int = Query(default=200, ge=1, le=10_000), + offset: int = Query(default=0, ge=0), + limit: int = Query(default=100_000, ge=1, le=1_000_000), + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + if start_line is not None: + return ok(store.source_content_lines(task_id, file_id, start_line, line_count)) + return ok(store.source_content_window(task_id, file_id, offset, limit)) + + +@router.get("/{task_id}/source-files/{file_id}/raw") +def source_file_raw( + task_id: str, + file_id: str, + range_header: str | None = Header(default=None, alias="Range"), + store: DataProcessStore = Depends(get_data_process_store), + storage: LocalDataProcessStorage = Depends(get_data_process_storage), +) -> StreamingResponse: + with api_errors(): + source = store.get_source_file(task_id, file_id, include_content=False) + file_format = str(source.get("file_format") or "").lower() + media_type = RAW_INLINE_PREVIEW_MEDIA_TYPES.get(file_format) + if media_type is None: + raise fail( + 415, + "raw inline preview is only available for PDF and modern Office source files", + ) + storage_object_id = str(source.get("storage_object_id") or "") + actual_size = storage.file_size( + storage_object_id, + expected_task_id=task_id, + expected_source_file_id=file_id, + ) + if actual_size is None: + raise fail(410, "the original file is unavailable for this legacy source file") + expected_size = int(source.get("size_bytes") or 0) + if actual_size != expected_size: + raise ValueError("source object size does not match metadata") + selected_range = _source_byte_range(range_header, actual_size) + start, end = selected_range or (0, actual_size - 1) + length = end - start + 1 + default_name = f"source.{file_format}" + name = _safe_file_name(str(source.get("name") or default_name), default_name) + headers = { + "Accept-Ranges": "bytes", + "Cache-Control": "private, no-store", + "Content-Disposition": f"inline; filename*=UTF-8''{quote(name, safe='')}", + "Content-Length": str(length), + "X-Accel-Buffering": "no", + "X-Content-Type-Options": "nosniff", + } + checksum = str(source.get("checksum_sha256") or "") + if checksum: + headers["ETag"] = f'"{checksum}"' + if selected_range is not None: + headers["Content-Range"] = f"bytes {start}-{end}/{actual_size}" + body = storage.iter_bytes( + storage_object_id, + expected_task_id=task_id, + expected_source_file_id=file_id, + expected_size=actual_size, + start=start, + length=length, + ) + return StreamingResponse( + body, + status_code=206 if selected_range is not None else 200, + media_type=media_type, + headers=headers, + ) + + +@router.get("/{task_id}/source-files/{file_id}/office-preview") +def source_file_office_preview( + task_id: str, + file_id: str, + sheet_index: int = Query(default=0, ge=0), + offset: int = Query(default=0, ge=0), + limit: int = Query(default=100, ge=1, le=MAX_XLSX_PREVIEW_ROWS), + store: DataProcessStore = Depends(get_data_process_store), + storage: LocalDataProcessStorage = Depends(get_data_process_storage), +) -> dict[str, Any]: + """返回 Word 版式块或 Excel 工作表网格,不把二进制内容下发给组件解析。""" + + with api_errors(): + source = store.get_source_file(task_id, file_id, include_content=False) + file_format = str(source.get("file_format") or "").lower() + if file_format not in {"docx", "xlsx"}: + raise fail(415, "Office preview is only available for DOCX and XLSX source files") + storage_object_id = str(source.get("storage_object_id") or "") + actual_size = storage.file_size( + storage_object_id, + expected_task_id=task_id, + expected_source_file_id=file_id, + ) + if actual_size is None: + raise fail(410, "the original Office file is unavailable for this legacy source file") + expected_size = int(source.get("size_bytes") or 0) + if actual_size != expected_size: + raise ValueError("source object size does not match metadata") + raw = b"".join( + storage.iter_bytes( + storage_object_id, + expected_task_id=task_id, + expected_source_file_id=file_id, + expected_size=actual_size, + ) + ) + preview = ( + build_docx_preview(raw) + if file_format == "docx" + else build_xlsx_preview( + raw, + sheet_index=sheet_index, + offset=offset, + limit=limit, + ) + ) + preview["file_name"] = str(source.get("name") or f"source.{file_format}") + return ok(preview) + + +@router.get("/{task_id}/source-files/{file_id}/pdf-pages") +def source_file_pdf_pages( + task_id: str, + file_id: str, + store: DataProcessStore = Depends(get_data_process_store), + storage: LocalDataProcessStorage = Depends(get_data_process_storage), +) -> dict[str, Any]: + """返回切片全文偏移对应的 PDF 物理页范围。""" + + with api_errors(): + source = store.get_source_file(task_id, file_id, include_content=False) + if str(source.get("file_format") or "").lower() != "pdf": + raise fail(415, "PDF page mapping is only available for PDF source files") + storage_object_id = str(source.get("storage_object_id") or "") + actual_size = storage.file_size( + storage_object_id, + expected_task_id=task_id, + expected_source_file_id=file_id, + ) + if actual_size is None: + raise fail(410, "the original PDF is unavailable for this legacy source file") + expected_size = int(source.get("size_bytes") or 0) + if actual_size != expected_size: + raise ValueError("source object size does not match metadata") + raw = b"".join( + storage.iter_bytes( + storage_object_id, + expected_task_id=task_id, + expected_source_file_id=file_id, + expected_size=actual_size, + ) + ) + pages = extract_pdf_page_texts(raw) + return ok( + { + "page_count": len(pages), + "pages": [ + { + "page_number": page.page_number, + "source_start": page.source_start, + "source_end": page.source_end, + } + for page in pages + ], + } + ) + + +@router.delete("/{task_id}/source-files/{file_id}") +def delete_source_file( + task_id: str, + file_id: str, + store: DataProcessStore = Depends(get_data_process_store), + storage: LocalDataProcessStorage = Depends(get_data_process_storage), +) -> dict[str, Any]: + with api_errors(): + source = store.get_source_file(task_id, file_id, include_content=False) + storage_object_id = str(source.get("storage_object_id") or "") + storage.validate_owner( + storage_object_id, + expected_task_id=task_id, + expected_source_file_id=file_id, + ) + store.delete_source_file(task_id, file_id) + cleanup_pending = False + try: + storage.delete( + storage_object_id, + expected_task_id=task_id, + expected_source_file_id=file_id, + ) + except Exception: + # 数据库软删除已经提交,不能再向客户端返回可重试的失败;保留逻辑引用, + # 由后续存储清理任务重试物理删除。 + cleanup_pending = True + logger.exception( + "failed to remove data process source object after soft deletion", + extra={"task_id": task_id, "source_file_id": file_id}, + ) + return ok( + {"deleted": file_id, "storage_cleanup_pending": cleanup_pending}, + "source file removed", + ) + + +def _external_postgres_connection(payload: ExternalSourceRequest) -> psycopg.Connection[Any]: + kind = payload.type.strip().lower() + parsed_url = urlsplit(payload.url) + scheme = parsed_url.scheme.lower() + if kind not in {"postgres", "postgresql"} or scheme not in {"postgres", "postgresql"}: + raise fail(501, f"external data source type is not supported: {payload.type}") + if parsed_url.username or parsed_url.password: + raise fail(400, "database credentials must use the account and password fields") + if payload.auth_mode not in {"none", "basic"}: + raise fail(400, "PostgreSQL supports only none or basic authentication") + if payload.auth_mode == "basic" and not payload.username: + raise fail(400, "database username is required for basic authentication") + hostname = parsed_url.hostname + if not hostname: + raise fail(400, "external PostgreSQL URL must include a hostname") + allow_private = os.getenv("DATA_PROCESS_ALLOW_PRIVATE_EXTERNAL_DB", "").lower() in { + "1", + "true", + "yes", + } + if not allow_private: + try: + addresses = { + item[4][0] + for item in socket.getaddrinfo( + hostname, + parsed_url.port or 5432, + type=socket.SOCK_STREAM, + ) + } + except socket.gaierror as exc: + raise fail(400, "external PostgreSQL hostname cannot be resolved") from exc + if any( + (address := ipaddress.ip_address(value)).is_private + or address.is_loopback + or address.is_link_local + or address.is_reserved + or address.is_unspecified + for value in addresses + ): + raise fail( + 403, + "private or local database addresses are disabled; " + "set DATA_PROCESS_ALLOW_PRIVATE_EXTERNAL_DB=true only in a trusted deployment", + ) + kwargs: dict[str, Any] = { + "connect_timeout": 5, + "row_factory": dict_row, + "application_name": "yg-ft-data-process-readonly", + "options": "-c default_transaction_read_only=on -c statement_timeout=30000", + } + if payload.auth_mode == "basic" and payload.username: + kwargs["user"] = payload.username + if payload.auth_mode == "basic" and payload.password: + kwargs["password"] = payload.password + return psycopg.connect(payload.url, **kwargs) + + +@router.post("/{task_id}/external/test") +def test_external_source( + task_id: str, + payload: ExternalSourceRequest, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + task = store.get_task(task_id) + if str(task.get("process_type")) != "external": + raise InvalidStateError( + "external source access requires an external data processing task" + ) + try: + with _external_postgres_connection(payload) as conn: + conn.execute("SELECT 1 AS ok").fetchone() + except psycopg.Error as exc: + raise fail(502, "external PostgreSQL connection test failed") from exc + return ok({"connected": True, "type": payload.type}) + + +@router.post("/{task_id}/external/pull") +def pull_external_source( + task_id: str, + payload: ExternalPullRequest, + store: DataProcessStore = Depends(get_data_process_store), + storage: LocalDataProcessStorage = Depends(get_data_process_storage), +) -> dict[str, Any]: + query = (payload.query or "").strip() + if query.endswith(";"): + query = query[:-1].rstrip() + if ";" in query: + raise fail(400, "external pull accepts exactly one read-only query") + first_token = query.split(maxsplit=1)[0].lower() if query else "" + if first_token not in {"select", "with"}: + raise fail(400, "a read-only SELECT or WITH query is required for external pull") + with api_errors(): + task = store.get_task(task_id) + if str(task.get("process_type")) != "external": + raise InvalidStateError("external pull requires an external data processing task") + try: + with _external_postgres_connection(payload) as conn: + conn.execute("SET TRANSACTION READ ONLY") + conn.execute("SET LOCAL statement_timeout = '30s'") + cursor = conn.execute(query) + rows: list[dict[str, Any]] = [] + content_parts: list[str] = [] + content_size = 0 + while len(rows) < payload.limit: + batch = cursor.fetchmany(min(1_000, payload.limit - len(rows))) + if not batch: + break + for row in batch: + line = json.dumps(row, ensure_ascii=False, default=str) + "\n" + content_size += len(line.encode("utf-8")) + if content_size > MAX_EXTERNAL_PULL_BYTES: + raise fail(413, "external pull result exceeds the 50 MiB safety limit") + rows.append(row) + content_parts.append(line) + conn.rollback() + except psycopg.Error as exc: + raise fail(502, "external PostgreSQL query failed") from exc + if not rows: + raise fail(400, "external query returned no rows") + content = "".join(content_parts) + raw = content.encode("utf-8") + name = _safe_file_name(payload.file_name, "external-data.jsonl") + source_file_id = new_id("dpsf") + staged = storage.stage_bytes( + batch_id=storage.new_batch_id(), + task_id=task_id, + source_file_id=source_file_id, + version=1, + name=name, + content=raw, + ) + sources = _commit_source_batch( + store, + storage, + task_id, + [ + { + "id": source_file_id, + "storage_object_id": staged.reference, + "name": name, + "content": content, + "raw_size": len(raw), + "checksum_sha256": hashlib.sha256(raw).hexdigest(), + "file_format": "jsonl", + "record_count": len(rows), + "metadata": { + "external_type": payload.type, + "external_host": urlsplit(payload.url).hostname, + "external_limit": payload.limit, + }, + } + ], + [staged], + ) + source = sources[0] + return ok({"files": [source]}, "external source pulled") + + +def _prepare_preview_items( + task_id: str, + store: DataProcessStore, + storage: LocalDataProcessStorage, + source_file_ids: list[str] | None = None, +) -> list[dict[str, Any]]: + task = store.get_task(task_id) + source_summaries = store.list_source_files(task_id) + if source_file_ids is not None: + requested = set(source_file_ids) + source_summaries = [item for item in source_summaries if item["id"] in requested] + found = {item["id"] for item in source_summaries} + missing = requested - found + if missing: + raise NotFoundError(f"source files not found: {', '.join(sorted(missing))}") + sources = [ + store.get_source_file(task_id, item["id"], include_content=True) + for item in source_summaries + ] + if not sources: + raise InvalidStateError("at least one source file is required") + config = task.get("config") or {} + preprocess_options = _preprocess_options(config) + chunk_method = str( + _value(config, "chunk_method", "chunkMethod", "layout_hybrid") + ) + is_unstructured = task.get("process_type") == "unstructured" + if is_unstructured and ( + chunk_method == "layout_hybrid" + or preprocess_options & {"clean_invalid", "clean_invalid_content"} + ): + for index, source in enumerate(sources): + if ( + chunk_method != "layout_hybrid" + and str(source.get("file_format") or "").lower() != "pdf" + ): + continue + storage_object_id = str(source.get("storage_object_id") or "") + actual_size = storage.file_size( + storage_object_id, + expected_task_id=task_id, + expected_source_file_id=str(source["id"]), + ) + if actual_size is None: + if chunk_method == "layout_hybrid": + raise InvalidStateError( + "版面结构混合切分无法读取原始文件,请重新上传后再处理" + ) + continue + expected_size = int(source.get("size_bytes") or 0) + if expected_size and actual_size != expected_size: + raise ValueError("source object size does not match metadata") + raw = b"".join( + storage.iter_bytes( + storage_object_id, + expected_task_id=task_id, + expected_source_file_id=str(source["id"]), + expected_size=actual_size, + ) + ) + enriched = dict(source) + if chunk_method == "layout_hybrid": + enriched["raw_content"] = raw + sources[index] = enriched + continue + if str(source.get("file_format") or "").lower() != "pdf": + continue + pages = extract_pdf_page_texts(raw) + extracted_text = "\n\n".join(page.text for page in pages if page.text) + if extracted_text != str(source.get("content") or ""): + logger.warning( + "skip PDF document noise detection because stored offsets differ for %s", + source["id"], + ) + continue + enriched["document_noise_spans"] = detect_pdf_document_noise(pages) + sources[index] = enriched + items = _build_preview_items(task, sources) + if not items and source_file_ids is None: + raise InvalidStateError("source files did not produce preview items") + return items + + +def _run_preview( + store: DataProcessStore, + storage: LocalDataProcessStorage, + task_id: str, + preview_run_id: str, + source_file_ids: list[str], +) -> None: + """后台逐文件切分;所有写入均由 preview_run_id 保护。""" + + started_at = time.perf_counter() + logger.info( + "data process preview started task_id=%s preview_run_id=%s total_files=%s", + task_id, + preview_run_id, + len(source_file_ids), + ) + try: + if not store.mark_preview_running(task_id, preview_run_id): + logger.info( + "data process preview skipped inactive run task_id=%s preview_run_id=%s", + task_id, + preview_run_id, + ) + return + total_files = len(source_file_ids) + total_items = 0 + for completed_files, source_file_id in enumerate(source_file_ids, start=1): + if not store.preview_is_running(task_id, preview_run_id): + logger.info( + "data process preview cancelled task_id=%s preview_run_id=%s " + "completed_files=%s total_files=%s", + task_id, + preview_run_id, + completed_files - 1, + total_files, + ) + return + items = _prepare_preview_items( + task_id, + store, + storage, + [source_file_id], + ) + if not items: + raise InvalidStateError( + f"source file did not produce preview items: {source_file_id}" + ) + created = store.replace_preview_items( + task_id, + items, + source_file_ids=[source_file_id], + preview_run_id=preview_run_id, + ) + total_items += len(created) + if not store.update_preview_progress( + task_id, + preview_run_id, + completed_files, + total_files, + ): + logger.info( + "data process preview stopped before progress update task_id=%s " + "preview_run_id=%s completed_files=%s total_files=%s", + task_id, + preview_run_id, + completed_files, + total_files, + ) + return + if store.complete_preview(task_id, preview_run_id): + logger.info( + "data process preview completed task_id=%s preview_run_id=%s " + "total_files=%s total_items=%s duration_ms=%.2f", + task_id, + preview_run_id, + total_files, + total_items, + (time.perf_counter() - started_at) * 1000, + ) + else: + logger.info( + "data process preview completion ignored for inactive run task_id=%s " + "preview_run_id=%s", + task_id, + preview_run_id, + ) + except Exception as exc: + logger.exception( + "data process preview failed task_id=%s preview_run_id=%s duration_ms=%.2f", + task_id, + preview_run_id, + (time.perf_counter() - started_at) * 1000, + ) + try: + if store.preview_is_running(task_id, preview_run_id): + store.mark_preview_failed( + task_id, + str(exc), + preview_run_id=preview_run_id, + ) + except Exception: + logger.exception( + "failed to persist data process preview failure task_id=%s " + "preview_run_id=%s", + task_id, + preview_run_id, + ) + + +@router.post("/{task_id}/preview/start", status_code=202) +def start_preview( + task_id: str, + background_tasks: BackgroundTasks, + payload: PreviewBuildRequest = Body(default_factory=PreviewBuildRequest), + store: DataProcessStore = Depends(get_data_process_store), + storage: LocalDataProcessStorage = Depends(get_data_process_storage), +) -> dict[str, Any]: + with api_errors(): + task, selected_ids = store.start_preview( + task_id, + source_file_ids=payload.source_file_ids, + ) + preview_run_id = str(task["preview_run_id"]) + progress = store.preview_progress(task_id) + background_tasks.add_task( + _run_preview, + store, + storage, + task_id, + preview_run_id, + selected_ids, + ) + return ok(progress, "data process preview started") + + +@router.get("/{task_id}/preview/progress") +def preview_progress( + task_id: str, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + return ok(store.preview_progress(task_id)) + + +@router.post("/{task_id}/preview/build") +def build_preview( + task_id: str, + payload: PreviewBuildRequest = Body(default_factory=PreviewBuildRequest), + store: DataProcessStore = Depends(get_data_process_store), + storage: LocalDataProcessStorage = Depends(get_data_process_storage), +) -> dict[str, Any]: + with api_errors(): + selected_ids = payload.source_file_ids + items = _prepare_preview_items(task_id, store, storage, selected_ids) + created = store.replace_preview_items( + task_id, + items, + source_file_ids=selected_ids, + ) + target_ids = selected_ids or [ + str(source["id"]) for source in store.list_source_files(task_id) + ] + file_counts = dict.fromkeys(target_ids, 0) + for item in created: + source_file_id = str(item.get("source_file_id") or "") + if source_file_id in file_counts: + file_counts[source_file_id] += 1 + files = [ + { + "source_file_id": source_file_id, + "preview_count": count, + "status": "completed" if count else "empty", + } + for source_file_id, count in file_counts.items() + ] + return ok( + { + "items": created, + "total": len(created), + "page": 1, + "page_size": len(created), + "file_counts": file_counts, + "files": files, + }, + "preview built", + ) + + +@router.get("/{task_id}/preview") +def preview_items( + task_id: str, + source_file_id: str | None = Query(default=None), + page: int = Query(default=1, ge=1), + page_size: int = Query(default=200, ge=1, le=1000), + keyword: str | None = Query(default=None), + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + return ok( + store.list_preview_items( + task_id, + source_file_id=source_file_id, + page=page, + page_size=page_size, + keyword=keyword, + ) + ) + + +@router.post("/{task_id}/preview") +def create_preview_item( + task_id: str, + payload: PreviewItemCreate, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + task = store.get_task(task_id) + item_payload = payload.model_dump(mode="json") + content = payload.edited_content + item_payload["token_count"] = estimate_token_count(content) + item_payload["quality_score"] = _preview_quality( + content, + task.get("config") or {}, + ) + if not content.strip(): + item_payload["status"] = "invalid" + item = store.create_preview_item(task_id, item_payload) + return ok(item, "preview item created") + + +@router.put("/{task_id}/preview/{preview_id}") +def update_preview_item( + task_id: str, + preview_id: str, + payload: PreviewItemUpdate, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + task = store.get_task(task_id) + update = payload.model_dump(exclude_unset=True, mode="json") + update["quality_score"] = _preview_quality(payload.edited_content, task.get("config") or {}) + item = store.update_preview_item( + task_id, + preview_id, + update, + ) + return ok(item, "preview item updated") + + +@router.delete("/{task_id}/preview/{preview_id}") +def delete_preview_item( + task_id: str, + preview_id: str, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + store.delete_preview_item(task_id, preview_id) + return ok({"deleted": preview_id}, "preview item deleted") + + +def _start_generation( + task_id: str, + payload: GenerateRequest, + background_tasks: BackgroundTasks, + store: DataProcessStore, +) -> dict[str, Any]: + with api_errors(): + task = store.start_generation(task_id, replace_existing=payload.replace_existing) + background_tasks.add_task( + _run_generation, + store, + task_id, + str(task["generation_run_id"]), + ) + return ok(store.progress(task_id), "data process generation started") + + +@router.post("/{task_id}/generate") +def generate( + task_id: str, + background_tasks: BackgroundTasks, + payload: GenerateRequest = Body(default_factory=GenerateRequest), + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + return _start_generation(task_id, payload, background_tasks, store) + + +@router.post("/{task_id}/start") +def start( + task_id: str, + background_tasks: BackgroundTasks, + payload: GenerateRequest = Body(default_factory=GenerateRequest), + store: DataProcessStore = Depends(get_data_process_store), + storage: LocalDataProcessStorage = Depends(get_data_process_storage), +) -> dict[str, Any]: + with api_errors(): + items = _prepare_preview_items(task_id, store, storage) + store.replace_preview_items(task_id, items) + return _start_generation(task_id, payload, background_tasks, store) + + +@router.post("/{task_id}/stop") +def stop( + task_id: str, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + store.stop_task(task_id) + return ok(store.progress(task_id), "data process task stopped") + + +@router.get("/{task_id}/progress") +def progress( + task_id: str, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + return ok(store.progress(task_id)) + + +@router.get("/{task_id}/results") +def results( + task_id: str, + page: int = Query(default=1, ge=1), + page_size: int = Query(default=100, ge=1, le=1000), + status: Literal["valid", "modified", "invalid"] | None = Query(default=None), + split: Literal["train", "validation", "test"] | None = Query(default=None), + keyword: str | None = Query(default=None), + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + return ok( + store.list_results( + task_id, + page=page, + page_size=page_size, + status=status, + split=split, + keyword=keyword, + ) + ) + + +@router.post("/{task_id}/confirm-results") +def confirm_results( + task_id: str, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + return ok(store.confirm_results(task_id), "data process results confirmed") + + +@router.put("/{task_id}/results/{result_id}") +def update_result( + task_id: str, + result_id: str, + payload: ResultUpdate, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + task = store.get_task(task_id) + current = store.get_result(task_id, result_id) + update = payload.model_dump(exclude_unset=True, mode="json") + merged = {**current, **update} + preview_id = current.get("preview_item_id") + source_content = "" + if preview_id: + preview = store.get_preview_item(task_id, str(preview_id)) + source_content = str( + preview.get("edited_content") or preview.get("original_content") or "" + ) + minimum = max( + 1, + int( + _value( + task.get("config") or {}, + "min_output_length", + "minOutputLength", + 20, + ) + or 20 + ), + ) + quality = score_quality( + merged, + min_output_length=minimum, + source_content=source_content, + ) + update["quality_score"] = asdict(quality) + result = store.update_result( + task_id, + result_id, + update, + ) + return ok(result, "data process result updated") + + +@router.post("/{task_id}/results/{result_id}/restore") +def restore_result( + task_id: str, + result_id: str, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + task = store.get_task(task_id) + current = store.get_result(task_id, result_id) + restored = { + **current, + "instruction": current.get("original_instruction") or current.get("instruction") or "", + "input": current.get("original_input") or current.get("input") or "", + "output": current.get("original_output") or current.get("output") or "", + } + preview_id = current.get("preview_item_id") + source_content = "" + if preview_id: + preview = store.get_preview_item(task_id, str(preview_id)) + source_content = str( + preview.get("edited_content") or preview.get("original_content") or "" + ) + minimum = max( + 1, + int( + _value( + task.get("config") or {}, + "min_output_length", + "minOutputLength", + 20, + ) + or 20 + ), + ) + quality = score_quality( + restored, + min_output_length=minimum, + source_content=source_content, + ) + restored = store.update_result( + task_id, + result_id, + { + "instruction": restored["instruction"], + "input": restored["input"], + "output": restored["output"], + "quality_score": asdict(quality), + "expected_updated_at": current.get("updated_at"), + }, + ) + return ok(restored, "data process result restored") + + +class _ResultRegenerationFailed(InvalidStateError): + """模型返回或质量校验失败,原失败结果必须保持不变。""" + + +def _assert_result_regeneration_allowed(task: dict[str, Any]) -> None: + if task.get("status") != "completed" or task.get("workflow_step") != "results": + raise InvalidStateError("task is not editing generation results") + if task.get("results_confirmed"): + raise InvalidStateError("confirmed results cannot be regenerated") + if task.get("output_dataset_id"): + raise InvalidStateError("published results cannot be regenerated") + + +def _result_regeneration_model( + task: dict[str, Any], + store: DataProcessStore, +) -> tuple[dict[str, Any], dict[str, Any]]: + config = task.get("config") or {} + model_id = _value(config, "generation_model_id", "generationModelId", None) + if not model_id: + raise InvalidStateError("task does not have a generation model") + return config, store.get_generation_model(str(model_id)) + + +def _result_regeneration_timeout(config: dict[str, Any]) -> float: + configured = float( + _value(config, "request_timeout_seconds", "requestTimeoutSeconds", 60) + ) + return max(1.0, min(RESULT_REGENERATION_TIMEOUT_SECONDS, configured)) + + +@contextmanager +def _claim_result_regeneration(task_id: str, result_id: str) -> Iterator[None]: + key = (task_id, result_id) + with _result_regeneration_claims_lock: + if key in _active_result_regenerations: + raise ConflictError("data process result regeneration is already running") + _active_result_regenerations.add(key) + try: + yield + finally: + with _result_regeneration_claims_lock: + _active_result_regenerations.discard(key) + + +def _generate_result_replacement( + task_id: str, + current: dict[str, Any], + preview: dict[str, Any], + config: dict[str, Any], + generation_model: dict[str, Any], + model_client: httpx.Client | None = None, +) -> dict[str, Any]: + source_content = str( + preview.get("edited_content") or preview.get("original_content") or "" + ).strip() + if not source_content or preview.get("status") == "invalid": + raise InvalidStateError("result source preview item is invalid or empty") + + output_type = str( + _value(config, "output_type", "outputType", "standard") + ).strip().lower() + previous_instruction = str(current.get("instruction") or "")[:1000] + previous_output = str(current.get("output") or "")[:1000] + base_prompt = str( + _value(config, "generation_prompt", "generationPrompt", "") or "" + ) + regeneration_instruction = ( + "这是一次失败结果的重新生成。请使用新的提问角度和表达," + "不要复述旧结果。旧问题:" + f"{previous_instruction or '无'};旧答案:{previous_output or '无'}。" + ) + runtime_config = { + **config, + "generation_prompt": f"{base_prompt}\n{regeneration_instruction}".strip(), + "output_type": output_type, + "reasoning_detail": _value( + config, "reasoning_detail", "reasoningDetail", "normal" + ), + "max_tokens": _value(config, "max_tokens", "maxTokens", 1024), + "json_mode": _value(config, "json_mode", "jsonMode", False), + # 交互式重新生成只做一次新尝试,避免继承整任务的重试配置后长时间等待。 + "generation_retries": RESULT_REGENERATION_RETRIES, + "request_timeout_seconds": _result_regeneration_timeout(config), + } + with _result_regeneration_slots: + generated = generate_model_records( + [preview], + model=generation_model, + config=runtime_config, + task_id=task_id, + split={"train": 100, "validation": 0, "test": 0}, + qa_pairs_per_item=1, + client=model_client, + ) + if not generated or generated[0].get("status") == "invalid": + reason = str( + (generated[0] if generated else {}).get("error") + or "model generation failed" + ) + raise _ResultRegenerationFailed(f"重新生成结果仍无效:{reason}") + + minimum = max( + 1, + int(_value(config, "min_output_length", "minOutputLength", 20) or 20), + ) + replacement = generated[0] + quality = score_quality( + replacement, + min_output_length=minimum, + source_content=source_content, + ) + if not quality.is_valid: + reason = ", ".join(quality.flags) or "quality validation failed" + raise _ResultRegenerationFailed(f"重新生成结果未通过质量校验:{reason}") + replacement["quality_score"] = asdict(quality) + replacement["status"] = "valid" + replacement["error"] = None + return replacement + + +def _regenerate_result_in_place( + task_id: str, + current: dict[str, Any], + preview: dict[str, Any], + config: dict[str, Any], + generation_model: dict[str, Any], + store: DataProcessStore, + *, + expected_updated_at: str, + model_client: httpx.Client | None = None, +) -> dict[str, Any]: + result_id = str(current["id"]) + with _claim_result_regeneration(task_id, result_id): + replacement = _generate_result_replacement( + task_id, + current, + preview, + config, + generation_model, + model_client, + ) + return store.replace_generated_result( + task_id, + result_id, + replacement, + expected_updated_at=expected_updated_at, + ) + + +def _safe_regeneration_error(exc: Exception) -> str: + return re.sub(r"\s+", " ", str(exc)).strip()[:500] or "result regeneration failed" + + +@router.post("/{task_id}/results/regenerate-batch") +def regenerate_results_batch( + task_id: str, + payload: ResultBatchRegenerateRequest, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + """并发重新生成一批失败结果;每条独立提交并允许部分成功。""" + + started_at = time.perf_counter() + batch_id = new_id("dprb") + with api_errors(): + task = store.get_task(task_id) + _assert_result_regeneration_allowed(task) + config, generation_model = _result_regeneration_model(task, store) + prepared: list[tuple[int, dict[str, Any], dict[str, Any], str]] = [] + failures: list[tuple[int, dict[str, str]]] = [] + + for index, requested in enumerate(payload.items): + try: + current = store.get_result(task_id, requested.result_id) + if current.get("status") != "invalid": + raise InvalidStateError("only an invalid result can be regenerated") + if requested.expected_updated_at != str(current.get("updated_at") or ""): + raise ConflictError("data process result was modified by another request") + preview_id = current.get("preview_item_id") + if not preview_id: + raise InvalidStateError( + "result is not associated with a source preview item" + ) + preview = store.get_preview_item(task_id, str(preview_id)) + prepared.append( + (index, current, preview, requested.expected_updated_at) + ) + except ConflictError as exc: + failures.append((index, { + "result_id": requested.result_id, + "code": "conflict", + "message": _safe_regeneration_error(exc), + })) + except (NotFoundError, InvalidStateError) as exc: + failures.append((index, { + "result_id": requested.result_id, + "code": "skipped", + "message": _safe_regeneration_error(exc), + })) + + logger.info( + "data process result batch regeneration started batch_id=%s task_id=%s " + "requested=%s prepared=%s concurrency=%s", + batch_id, + task_id, + len(payload.items), + len(prepared), + min(RESULT_REGENERATION_CONCURRENCY, len(prepared)), + ) + successes: list[tuple[int, dict[str, Any]]] = [] + if prepared: + request_timeout = _result_regeneration_timeout(config) + model_timeout = httpx.Timeout( + request_timeout, + connect=min(10.0, request_timeout), + ) + model_limits = httpx.Limits( + max_connections=RESULT_REGENERATION_CONCURRENCY, + max_keepalive_connections=RESULT_REGENERATION_CONCURRENCY, + ) + # httpx.Client 支持跨线程复用,批次内共享连接池可减少重复建连开销。 + with ( + httpx.Client(timeout=model_timeout, limits=model_limits) as model_client, + ThreadPoolExecutor( + max_workers=min(RESULT_REGENERATION_CONCURRENCY, len(prepared)), + thread_name_prefix="data-result-regeneration", + ) as executor, + ): + futures = { + executor.submit( + _regenerate_result_in_place, + task_id, + current, + preview, + config, + generation_model, + store, + expected_updated_at=expected_updated_at, + model_client=model_client, + ): (index, str(current["id"]), time.perf_counter()) + for index, current, preview, expected_updated_at in prepared + } + for future in as_completed(futures): + index, result_id, item_started_at = futures[future] + try: + regenerated = future.result() + successes.append((index, regenerated)) + outcome = "succeeded" + except ConflictError as exc: + outcome = "conflict" + failures.append((index, { + "result_id": result_id, + "code": outcome, + "message": _safe_regeneration_error(exc), + })) + except _ResultRegenerationFailed as exc: + outcome = "generation_failed" + failures.append((index, { + "result_id": result_id, + "code": outcome, + "message": _safe_regeneration_error(exc), + })) + except (NotFoundError, InvalidStateError) as exc: + outcome = "skipped" + failures.append((index, { + "result_id": result_id, + "code": outcome, + "message": _safe_regeneration_error(exc), + })) + except Exception as exc: # pragma: no cover - defensive boundary + outcome = "internal_error" + logger.exception( + "data process result batch regeneration crashed " + "batch_id=%s task_id=%s result_id=%s", + batch_id, + task_id, + result_id, + ) + failures.append((index, { + "result_id": result_id, + "code": outcome, + "message": _safe_regeneration_error(exc), + })) + logger.info( + "data process result batch item finished batch_id=%s task_id=%s " + "result_id=%s outcome=%s duration_ms=%.2f", + batch_id, + task_id, + result_id, + outcome, + (time.perf_counter() - item_started_at) * 1000, + ) + + success_items = [item for _, item in sorted(successes, key=lambda pair: pair[0])] + failure_items = [item for _, item in sorted(failures, key=lambda pair: pair[0])] + remaining_invalid_count = int( + store.list_results( + task_id, + page=1, + page_size=1, + status="invalid", + )["total"] + ) + duration_ms = (time.perf_counter() - started_at) * 1000 + logger.info( + "data process result batch regeneration completed batch_id=%s task_id=%s " + "succeeded=%s failed=%s remaining_invalid=%s duration_ms=%.2f", + batch_id, + task_id, + len(success_items), + len(failure_items), + remaining_invalid_count, + duration_ms, + ) + return ok( + { + "batch_id": batch_id, + "total": len(payload.items), + "succeeded": len(success_items), + "failed": len(failure_items), + "remaining_invalid_count": remaining_invalid_count, + "duration_ms": round(duration_ms, 2), + "items": success_items, + "failures": failure_items, + }, + "data process results regenerated", + ) + + +@router.post("/{task_id}/results/{result_id}/regenerate") +def regenerate_result( + task_id: str, + result_id: str, + payload: ResultRegenerateRequest, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + """只重新生成一个失败结果,成功后原位替换且不影响其他结果。""" + + started_at = time.perf_counter() + with api_errors(): + task = store.get_task(task_id) + _assert_result_regeneration_allowed(task) + current = store.get_result(task_id, result_id) + if current.get("status") != "invalid": + raise InvalidStateError("only an invalid result can be regenerated") + if payload.expected_updated_at != str(current.get("updated_at") or ""): + raise ConflictError("data process result was modified by another request") + preview_id = current.get("preview_item_id") + if not preview_id: + raise InvalidStateError("result is not associated with a source preview item") + preview = store.get_preview_item(task_id, str(preview_id)) + config, generation_model = _result_regeneration_model(task, store) + try: + result = _regenerate_result_in_place( + task_id, + current, + preview, + config, + generation_model, + store, + expected_updated_at=payload.expected_updated_at, + ) + except _ResultRegenerationFailed as exc: + logger.warning( + "data process result regeneration failed task_id=%s result_id=%s " + "duration_ms=%.2f reason=%s", + task_id, + result_id, + (time.perf_counter() - started_at) * 1000, + _safe_regeneration_error(exc), + ) + raise + logger.info( + "data process result regenerated task_id=%s result_id=%s duration_ms=%.2f", + task_id, + result_id, + (time.perf_counter() - started_at) * 1000, + ) + return ok(result, "data process result regenerated") + + +@router.post("/{task_id}/publish") +def publish( + task_id: str, + payload: PublishRequest, + store: DataProcessStore = Depends(get_data_process_store), +) -> dict[str, Any]: + with api_errors(): + result = store.publish(task_id, payload.model_dump(mode="json")) + message = "dataset published" if result["created"] else "dataset already published" + return ok(result, message) diff --git a/backend/app/api/v1/router.py b/backend/app/api/v1/router.py index 87185f7..0169442 100644 --- a/backend/app/api/v1/router.py +++ b/backend/app/api/v1/router.py @@ -8,6 +8,9 @@ from app.modules.tenant import router as tenant_router from app.modules.project import router as project_router from app.modules.resource import router as resource_router from app.modules.approval import router as approval_router +from app.api.v1.endpoints.data_process import router as data_process_router + +import logging api_router = APIRouter() api_router.include_router(health_router, tags=["health"]) @@ -18,4 +21,11 @@ api_router.include_router(tenant_router, tags=["tenant"]) api_router.include_router(project_router, tags=["project"]) api_router.include_router(resource_router, tags=["resource"]) api_router.include_router(approval_router, tags=["approval"]) +api_router.include_router(data_process_router, tags=["data_process"]) + +try: + from app.modules.data_process.store import DataProcessStore + DataProcessStore().ensure_schema() +except Exception as exc: # best-effort at startup; feature degrades if DB unavailable + logging.getLogger(__name__).warning("data_process schema ensure failed: %s", exc) diff --git a/backend/app/core/config.py b/backend/app/core/config.py index e027dad..048df89 100644 --- a/backend/app/core/config.py +++ b/backend/app/core/config.py @@ -23,7 +23,7 @@ class Settings: app_env: str = os.getenv("APP_ENV", "local") route_prefix: str = os.getenv("MODELTF_ROUTE_PREFIX", "/modelTF") app_mode: str = os.getenv("APP_MODE", "local") - database_url: str = os.getenv("DATABASE_URL", "postgresql+psycopg://yg_ft:change_me@localhost:15432/yg_ft") + database_url: str = os.getenv("DATABASE_URL", "postgresql+psycopg://root:8811614287327Leo@www.caoxiaozhu.com:5432/yg_ft") cors_allow_origins: list[str] = None # type: ignore[assignment] compute_mode: str = os.getenv("COMPUTE_MODE", "real") compute_status_sync_mode: str = os.getenv("COMPUTE_STATUS_SYNC_MODE", "polling") diff --git a/backend/app/db/session.py b/backend/app/db/session.py index befe080..89bf3cc 100644 --- a/backend/app/db/session.py +++ b/backend/app/db/session.py @@ -8,7 +8,7 @@ from sqlalchemy import create_engine from sqlalchemy.orm import Session, sessionmaker -DATABASE_URL = os.getenv("DATABASE_URL", "postgresql+psycopg://yg_ft:change_me@localhost:15432/yg_ft") +DATABASE_URL = os.getenv("DATABASE_URL", "postgresql+psycopg://root:8811614287327Leo@www.caoxiaozhu.com:5432/yg_ft") engine = create_engine( DATABASE_URL, diff --git a/backend/app/db/sql/002_data_process.sql b/backend/app/db/sql/002_data_process.sql new file mode 100644 index 0000000..b05cb81 --- /dev/null +++ b/backend/app/db/sql/002_data_process.sql @@ -0,0 +1,343 @@ +-- Data processing migration. +-- +-- IMPORTANT: This file is intentionally NOT wired into application startup. +-- Apply it explicitly in a controlled deployment, or call +-- DataProcessStore.ensure_schema() from an administrative command. + +BEGIN; + +-- This migration targets the current runtime schema created by +-- 001_platform_runtime.sql. Refuse the UUID/JSONB target-design schema instead +-- of partially altering it with incompatible TEXT foreign keys. +DO $$ +DECLARE + datasets_id_type TEXT; +BEGIN + SELECT format_type(a.atttypid, a.atttypmod) + INTO datasets_id_type + FROM pg_attribute a + JOIN pg_class c ON c.oid = a.attrelid + JOIN pg_namespace n ON n.oid = c.relnamespace + WHERE n.nspname = current_schema() + AND c.relname = 'datasets' + AND a.attname = 'id' + AND a.attnum > 0 + AND NOT a.attisdropped; + IF datasets_id_type IS NULL THEN + RAISE EXCEPTION '002_data_process.sql requires 001_platform_runtime.sql first'; + END IF; + IF datasets_id_type <> 'text' THEN + RAISE EXCEPTION + '002_data_process.sql supports only the current TEXT runtime schema; found datasets.id type %', + datasets_id_type; + END IF; +END $$; + +ALTER TABLE datasets ADD COLUMN IF NOT EXISTS source_task_id TEXT; +ALTER TABLE datasets ADD COLUMN IF NOT EXISTS size_bytes BIGINT NOT NULL DEFAULT 0; +ALTER TABLE datasets ADD COLUMN IF NOT EXISTS record_count BIGINT NOT NULL DEFAULT 0; +ALTER TABLE datasets ADD COLUMN IF NOT EXISTS metadata TEXT NOT NULL DEFAULT '{}'; +ALTER TABLE datasets ADD COLUMN IF NOT EXISTS tenant_id TEXT; +ALTER TABLE datasets ADD COLUMN IF NOT EXISTS project_id TEXT; +ALTER TABLE datasets ADD COLUMN IF NOT EXISTS owner_id TEXT; +ALTER TABLE datasets ADD COLUMN IF NOT EXISTS created_by TEXT; +ALTER TABLE datasets ADD COLUMN IF NOT EXISTS created_at TIMESTAMPTZ NOT NULL DEFAULT now(); +ALTER TABLE datasets ADD COLUMN IF NOT EXISTS updated_at TIMESTAMPTZ NOT NULL DEFAULT now(); +ALTER TABLE datasets ADD COLUMN IF NOT EXISTS deleted_at TIMESTAMPTZ; + +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS storage_object_id TEXT; +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS current_version_id TEXT; +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS size_bytes BIGINT NOT NULL DEFAULT 0; +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS record_count BIGINT NOT NULL DEFAULT 0; +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS file_format VARCHAR(40); +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS checksum_sha256 CHAR(64); +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS version_no INTEGER NOT NULL DEFAULT 1; +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS source_task_id TEXT; +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS tenant_id TEXT; +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS project_id TEXT; +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS created_by TEXT; +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS metadata TEXT NOT NULL DEFAULT '{}'; +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS created_at TIMESTAMPTZ NOT NULL DEFAULT now(); +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS updated_at TIMESTAMPTZ NOT NULL DEFAULT now(); +ALTER TABLE dataset_files ADD COLUMN IF NOT EXISTS deleted_at TIMESTAMPTZ; + +CREATE TABLE IF NOT EXISTS data_process_tasks ( + id TEXT PRIMARY KEY, + name VARCHAR(150) NOT NULL, + description TEXT, + status VARCHAR(20) NOT NULL DEFAULT 'pending' + CHECK (status IN ('pending', 'running', 'completed', 'failed', 'stopped')), + process_type VARCHAR(20) NOT NULL + CHECK (process_type IN ('structured', 'unstructured', 'external')), + source_dataset_id TEXT REFERENCES datasets(id) ON DELETE SET NULL, + output_dataset_id TEXT REFERENCES datasets(id) ON DELETE SET NULL, + config TEXT NOT NULL DEFAULT '{}', + progress NUMERIC(5,2) NOT NULL DEFAULT 0 CHECK (progress >= 0 AND progress <= 100), + input_count BIGINT NOT NULL DEFAULT 0 CHECK (input_count >= 0), + output_count BIGINT NOT NULL DEFAULT 0 CHECK (output_count >= 0), + filtered_count BIGINT NOT NULL DEFAULT 0 CHECK (filtered_count >= 0), + duplicate_count BIGINT NOT NULL DEFAULT 0 CHECK (duplicate_count >= 0), + error_count BIGINT NOT NULL DEFAULT 0 CHECK (error_count >= 0), + failure_reason TEXT, + generation_run_id TEXT, + results_confirmed BOOLEAN NOT NULL DEFAULT TRUE, + workflow_step VARCHAR(20) NOT NULL DEFAULT 'create' + CHECK (workflow_step IN ('create', 'model', 'upload', 'preview', 'generate', 'results')), + preview_status VARCHAR(20) NOT NULL DEFAULT 'idle' + CHECK (preview_status IN ('idle', 'queued', 'running', 'completed', 'failed', 'cancelled')), + preview_progress NUMERIC(5,2) NOT NULL DEFAULT 0 + CHECK (preview_progress >= 0 AND preview_progress <= 100), + preview_run_id TEXT, + preview_failure_reason TEXT, + preview_total_files INTEGER NOT NULL DEFAULT 0 CHECK (preview_total_files >= 0), + preview_completed_files INTEGER NOT NULL DEFAULT 0 CHECK (preview_completed_files >= 0), + tenant_id TEXT, + project_id TEXT, + owner_id TEXT, + approval_status VARCHAR(30) NOT NULL DEFAULT 'not_required', + created_by TEXT, + updated_by TEXT, + deleted_by TEXT, + started_at TIMESTAMPTZ, + completed_at TIMESTAMPTZ, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), + deleted_at TIMESTAMPTZ +); + +ALTER TABLE data_process_tasks ADD COLUMN IF NOT EXISTS generation_run_id TEXT; +-- 历史任务在引入六步确认流程前已经完成审核,默认保留为已确认; +-- 新任务由创建接口显式写入 FALSE,并在第六步确认后转为 TRUE。 +ALTER TABLE data_process_tasks + ADD COLUMN IF NOT EXISTS results_confirmed BOOLEAN NOT NULL DEFAULT TRUE; +UPDATE data_process_tasks +SET results_confirmed=FALSE +WHERE status <> 'completed' AND results_confirmed=TRUE; + +-- 先以可空列接入旧库,才能只回填历史行;随后再收紧默认值与约束。 +ALTER TABLE data_process_tasks ADD COLUMN IF NOT EXISTS workflow_step VARCHAR(20); +ALTER TABLE data_process_tasks ADD COLUMN IF NOT EXISTS preview_status VARCHAR(20); +ALTER TABLE data_process_tasks ADD COLUMN IF NOT EXISTS preview_progress NUMERIC(5,2); +ALTER TABLE data_process_tasks ADD COLUMN IF NOT EXISTS preview_run_id TEXT; +ALTER TABLE data_process_tasks ADD COLUMN IF NOT EXISTS preview_failure_reason TEXT; +ALTER TABLE data_process_tasks ADD COLUMN IF NOT EXISTS preview_total_files INTEGER; +ALTER TABLE data_process_tasks ADD COLUMN IF NOT EXISTS preview_completed_files INTEGER; + +CREATE TEMP TABLE data_process_workflow_backfill_ids ON COMMIT DROP AS +SELECT id FROM data_process_tasks WHERE workflow_step IS NULL; + +UPDATE data_process_tasks task +SET workflow_step = CASE + WHEN task.status IN ('running', 'failed', 'stopped') THEN 'generate' + WHEN task.status = 'completed' AND task.results_confirmed=FALSE THEN 'generate' + WHEN task.status = 'completed' THEN 'results' + ELSE 'create' +END +WHERE task.workflow_step IS NULL; +UPDATE data_process_tasks +SET preview_status='idle', preview_progress=0, + preview_total_files=0, preview_completed_files=0 +WHERE preview_status IS NULL OR preview_progress IS NULL + OR preview_total_files IS NULL OR preview_completed_files IS NULL; + +ALTER TABLE data_process_tasks ALTER COLUMN workflow_step SET DEFAULT 'create'; +ALTER TABLE data_process_tasks ALTER COLUMN workflow_step SET NOT NULL; +ALTER TABLE data_process_tasks ALTER COLUMN preview_status SET DEFAULT 'idle'; +ALTER TABLE data_process_tasks ALTER COLUMN preview_status SET NOT NULL; +ALTER TABLE data_process_tasks ALTER COLUMN preview_progress SET DEFAULT 0; +ALTER TABLE data_process_tasks ALTER COLUMN preview_progress SET NOT NULL; +ALTER TABLE data_process_tasks ALTER COLUMN preview_total_files SET DEFAULT 0; +ALTER TABLE data_process_tasks ALTER COLUMN preview_total_files SET NOT NULL; +ALTER TABLE data_process_tasks ALTER COLUMN preview_completed_files SET DEFAULT 0; +ALTER TABLE data_process_tasks ALTER COLUMN preview_completed_files SET NOT NULL; + +DO $$ +BEGIN + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conrelid='data_process_tasks'::regclass + AND conname='ck_data_process_tasks_workflow_step' + ) THEN + ALTER TABLE data_process_tasks ADD CONSTRAINT ck_data_process_tasks_workflow_step + CHECK (workflow_step IN ('create', 'model', 'upload', 'preview', 'generate', 'results')); + END IF; + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conrelid='data_process_tasks'::regclass + AND conname='ck_data_process_tasks_preview_status' + ) THEN + ALTER TABLE data_process_tasks ADD CONSTRAINT ck_data_process_tasks_preview_status + CHECK (preview_status IN ('idle', 'queued', 'running', 'completed', 'failed', 'cancelled')); + END IF; + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conrelid='data_process_tasks'::regclass + AND conname='ck_data_process_tasks_preview_progress' + ) THEN + ALTER TABLE data_process_tasks ADD CONSTRAINT ck_data_process_tasks_preview_progress + CHECK (preview_progress >= 0 AND preview_progress <= 100); + END IF; + IF NOT EXISTS ( + SELECT 1 FROM pg_constraint + WHERE conrelid='data_process_tasks'::regclass + AND conname='ck_data_process_tasks_preview_file_counts' + ) THEN + ALTER TABLE data_process_tasks ADD CONSTRAINT ck_data_process_tasks_preview_file_counts + CHECK (preview_total_files >= 0 AND preview_completed_files >= 0 + AND preview_completed_files <= preview_total_files); + END IF; +END $$; + +CREATE UNIQUE INDEX IF NOT EXISTS uq_data_process_tasks_name_alive + ON data_process_tasks(name) WHERE deleted_at IS NULL; +CREATE INDEX IF NOT EXISTS idx_data_process_tasks_scope_status + ON data_process_tasks(tenant_id, project_id, status, created_at DESC) + WHERE deleted_at IS NULL; +CREATE INDEX IF NOT EXISTS idx_data_process_tasks_creator_created + ON data_process_tasks(created_by, created_at DESC) WHERE deleted_at IS NULL; + +CREATE TABLE IF NOT EXISTS data_process_source_files ( + id TEXT PRIMARY KEY, + task_id TEXT NOT NULL REFERENCES data_process_tasks(id) ON DELETE CASCADE, + storage_object_id TEXT, + name TEXT NOT NULL, + size_bytes BIGINT NOT NULL DEFAULT 0 CHECK (size_bytes >= 0), + record_count BIGINT NOT NULL DEFAULT 0 CHECK (record_count >= 0), + file_format VARCHAR(40), + checksum_sha256 CHAR(64) NOT NULL, + version_no INTEGER NOT NULL DEFAULT 1 CHECK (version_no > 0), + content TEXT NOT NULL, + content_preview TEXT, + metadata TEXT NOT NULL DEFAULT '{}', + tenant_id TEXT, + project_id TEXT, + created_by TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), + deleted_at TIMESTAMPTZ +); + +CREATE INDEX IF NOT EXISTS idx_data_process_source_files_task + ON data_process_source_files(task_id, created_at) WHERE deleted_at IS NULL; +CREATE UNIQUE INDEX IF NOT EXISTS uq_data_process_source_checksum_alive + ON data_process_source_files(task_id, checksum_sha256) WHERE deleted_at IS NULL; + +CREATE TABLE IF NOT EXISTS data_process_preview_items ( + id TEXT PRIMARY KEY, + task_id TEXT NOT NULL REFERENCES data_process_tasks(id) ON DELETE CASCADE, + source_file_id TEXT REFERENCES data_process_source_files(id) ON DELETE CASCADE, + original_content TEXT NOT NULL DEFAULT '', + edited_content TEXT NOT NULL DEFAULT '', + source_start INTEGER CHECK (source_start IS NULL OR source_start >= 0), + source_end INTEGER CHECK (source_end IS NULL OR source_end >= 0), + source_start_line INTEGER CHECK (source_start_line IS NULL OR source_start_line > 0), + source_end_line INTEGER CHECK (source_end_line IS NULL OR source_end_line > 0), + token_count INTEGER NOT NULL DEFAULT 0 CHECK (token_count >= 0), + status VARCHAR(20) NOT NULL DEFAULT 'original' + CHECK (status IN ('original', 'modified', 'manual', 'invalid')), + quality_score TEXT NOT NULL DEFAULT '{}', + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), + CHECK (source_start IS NULL OR source_end IS NULL OR source_end >= source_start), + CHECK (source_start_line IS NULL OR source_end_line IS NULL OR source_end_line >= source_start_line) +); + +CREATE INDEX IF NOT EXISTS idx_data_process_preview_task_file + ON data_process_preview_items(task_id, source_file_id, created_at); + +-- 子表在新库中到这里才存在;只修复本次新增 workflow_step 前的历史任务。 +UPDATE data_process_tasks task +SET workflow_step = CASE + WHEN EXISTS ( + SELECT 1 FROM data_process_preview_items preview + WHERE preview.task_id=task.id + ) THEN 'preview' + WHEN EXISTS ( + SELECT 1 FROM data_process_source_files source_file + WHERE source_file.task_id=task.id AND source_file.deleted_at IS NULL + ) THEN 'upload' + ELSE task.workflow_step +END +WHERE task.id IN (SELECT id FROM data_process_workflow_backfill_ids) + AND task.status='pending'; + +CREATE TABLE IF NOT EXISTS data_process_results ( + id TEXT PRIMARY KEY, + task_id TEXT NOT NULL REFERENCES data_process_tasks(id) ON DELETE CASCADE, + preview_item_id TEXT REFERENCES data_process_preview_items(id) ON DELETE SET NULL, + instruction TEXT NOT NULL, + input TEXT NOT NULL DEFAULT '', + output TEXT NOT NULL, + original_instruction TEXT, + original_input TEXT, + original_output TEXT, + status VARCHAR(20) NOT NULL DEFAULT 'valid' + CHECK (status IN ('valid', 'modified', 'invalid')), + error TEXT, + split VARCHAR(20) CHECK (split IS NULL OR split IN ('train', 'validation', 'test')), + quality_score TEXT NOT NULL DEFAULT '{}', + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +CREATE INDEX IF NOT EXISTS idx_data_process_results_task_status + ON data_process_results(task_id, status, id); +CREATE INDEX IF NOT EXISTS idx_data_process_results_task_split + ON data_process_results(task_id, split); + +CREATE TABLE IF NOT EXISTS dataset_file_versions ( + id TEXT PRIMARY KEY, + dataset_file_id TEXT NOT NULL REFERENCES dataset_files(id) ON DELETE CASCADE, + version_no INTEGER NOT NULL CHECK (version_no > 0), + storage_object_id TEXT NOT NULL, + content_preview TEXT, + description TEXT, + base_version_id TEXT REFERENCES dataset_file_versions(id) ON DELETE SET NULL, + size_bytes BIGINT NOT NULL DEFAULT 0 CHECK (size_bytes >= 0), + record_count BIGINT NOT NULL DEFAULT 0 CHECK (record_count >= 0), + checksum_sha256 CHAR(64) NOT NULL, + source_task_id TEXT REFERENCES data_process_tasks(id) ON DELETE SET NULL, + metadata TEXT NOT NULL DEFAULT '{}', + created_by TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +ALTER TABLE dataset_file_versions ADD COLUMN IF NOT EXISTS source_task_id TEXT; +ALTER TABLE dataset_file_versions ADD COLUMN IF NOT EXISTS metadata TEXT NOT NULL DEFAULT '{}'; +CREATE UNIQUE INDEX IF NOT EXISTS uq_dataset_file_versions_no_002 + ON dataset_file_versions(dataset_file_id, version_no); +CREATE INDEX IF NOT EXISTS idx_dataset_file_versions_source_task_002 + ON dataset_file_versions(source_task_id) WHERE source_task_id IS NOT NULL; + +CREATE TABLE IF NOT EXISTS dataset_records ( + id TEXT PRIMARY KEY, + dataset_id TEXT NOT NULL REFERENCES datasets(id) ON DELETE CASCADE, + dataset_file_id TEXT REFERENCES dataset_files(id) ON DELETE CASCADE, + version_id TEXT REFERENCES dataset_file_versions(id) ON DELETE CASCADE, + line_no INTEGER, + split VARCHAR(20) CHECK (split IS NULL OR split IN ('train', 'validation', 'test')), + instruction TEXT, + input TEXT, + output TEXT, + raw TEXT NOT NULL DEFAULT '{}', + status VARCHAR(20) NOT NULL DEFAULT 'valid' + CHECK (status IN ('valid', 'modified', 'invalid')), + source_task_id TEXT REFERENCES data_process_tasks(id) ON DELETE SET NULL, + source_result_id TEXT REFERENCES data_process_results(id) ON DELETE SET NULL, + preview_item_id TEXT REFERENCES data_process_preview_items(id) ON DELETE SET NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +ALTER TABLE dataset_records ADD COLUMN IF NOT EXISTS source_task_id TEXT; +ALTER TABLE dataset_records ADD COLUMN IF NOT EXISTS source_result_id TEXT; +ALTER TABLE dataset_records ADD COLUMN IF NOT EXISTS preview_item_id TEXT; +CREATE INDEX IF NOT EXISTS idx_dataset_records_dataset_002 + ON dataset_records(dataset_id, id); +CREATE INDEX IF NOT EXISTS idx_dataset_records_source_task_002 + ON dataset_records(source_task_id, source_result_id); +CREATE INDEX IF NOT EXISTS idx_datasets_source_task_002 + ON datasets(source_task_id) WHERE source_task_id IS NOT NULL; +CREATE INDEX IF NOT EXISTS idx_dataset_files_source_task_002 + ON dataset_files(source_task_id) WHERE source_task_id IS NOT NULL; + +COMMIT; diff --git a/backend/app/db/sql/002_governance.sql b/backend/app/db/sql/002_governance.sql index 3c17207..c335c9e 100644 --- a/backend/app/db/sql/002_governance.sql +++ b/backend/app/db/sql/002_governance.sql @@ -35,6 +35,8 @@ CREATE TABLE IF NOT EXISTS tenants ( code TEXT NOT NULL UNIQUE, status TEXT NOT NULL, owner_user_id TEXT, + quota TEXT, + retention_policy_id TEXT, create_time TEXT NOT NULL ); diff --git a/backend/app/modules/data_process/algorithms.py b/backend/app/modules/data_process/algorithms.py new file mode 100644 index 0000000..07d4b38 --- /dev/null +++ b/backend/app/modules/data_process/algorithms.py @@ -0,0 +1,2763 @@ +"""数据处理模块使用的无副作用算法。 + +本模块不访问数据库、文件系统或网络,便于 API、后台任务和测试共同复用。 +所有偏移量均为 Python 字符串偏移量。 +""" + +from __future__ import annotations + +import csv +import hashlib +import io +import json +import math +import re +import unicodedata +import xml.etree.ElementTree as ET +import zipfile +from collections import Counter +from collections.abc import Iterable, Mapping, Sequence +from copy import deepcopy +from dataclasses import dataclass +from datetime import date, datetime, time +from pathlib import Path, PurePosixPath +from typing import Any, Literal +from urllib.parse import unquote, urlsplit + +from docx import Document +from docx.oxml.table import CT_Tbl +from docx.oxml.text.paragraph import CT_P +from docx.table import Table +from docx.text.paragraph import Paragraph +from openpyxl import load_workbook +from openpyxl.utils.cell import range_boundaries +from pptx import Presentation +from pypdf import PdfReader + +from app.modules.data_process.constants import MAX_QA_PAIRS_PER_ITEM + +TextFormat = Literal[ + "json", + "jsonl", + "csv", + "markdown", + "txt", + "pdf", + "docx", + "xlsx", + "pptx", +] +DatasetSplit = Literal["train", "validation", "test"] +StructuredPreprocessOption = Literal[ + "clean_invalid", + "detect_structure", + "deduplicate", + "normalize_format", + "filter_anomaly", + "desensitize", +] + +SUPPORTED_TEXT_FORMATS: tuple[TextFormat, ...] = ( + "json", + "jsonl", + "csv", + "markdown", + "txt", + "pdf", + "docx", + "xlsx", + "pptx", +) + +_FORMAT_ALIASES: dict[str, TextFormat] = { + "json": "json", + "jsonl": "jsonl", + "ndjson": "jsonl", + "csv": "csv", + "tsv": "csv", + "md": "markdown", + "markdown": "markdown", + "txt": "txt", + "text": "txt", + "pdf": "pdf", + "docx": "docx", + "xlsx": "xlsx", + "pptx": "pptx", +} +_LEGACY_OFFICE_FORMATS: dict[str, str] = { + "doc": "docx", + "xls": "xlsx", + "ppt": "pptx", +} +_OFFICE_OPEN_XML_FORMATS = {"docx", "xlsx", "pptx"} +_MAX_ARCHIVE_ENTRIES = 10_000 +_MAX_ARCHIVE_UNCOMPRESSED_BYTES = 512 * 1024 * 1024 +_MAX_ARCHIVE_ENTRY_BYTES = 128 * 1024 * 1024 +_MAX_ARCHIVE_COMPRESSION_RATIO = 200 +_MAX_EXTRACTED_TEXT_CHARS = 20_000_000 +_MAX_PDF_PAGES = 2_000 +_MAX_PRESENTATION_SLIDES = 2_000 +_MAX_WORKBOOK_SHEETS = 100 +_MAX_WORKBOOK_ROWS = 100_000 +_MAX_WORKBOOK_SCANNED_ROWS = 200_000 +_MAX_WORKBOOK_COLUMNS = 256 +_MAX_WORKBOOK_CELLS = 2_000_000 +_MAX_WORKBOOK_HEADER_ROWS = 8 +_MAX_WORKBOOK_HEADER_SCAN_ROWS = 64 +_MAX_WORKBOOK_MERGED_RANGES = 100_000 +_MAX_STRUCTURED_FIELDS = 1_024 +_MAX_STRUCTURED_DEPTH = 16 +_MAX_ANOMALY_TEXT_CHARS = 1_000_000 +_STRUCTURED_OPTIONS = { + "clean_invalid", + "detect_structure", + "deduplicate", + "normalize_format", + "filter_anomaly", + "desensitize", +} +_IDENTITY_FIELD_PATTERN = re.compile(r"(?:^|[._])(?:id|uuid|key|code)$|(?:^|[._]).+_id$") +_MOJIBAKE_MARKERS = ("\ufffd", "锟斤拷", "烫烫烫", "屯屯屯", "Ã", "Â", "â€") +_NAME_FIELD_NAMES = { + "name", + "full_name", + "fullname", + "real_name", + "contact_name", + "customer_name", + "recipient_name", + "姓名", + "真实姓名", + "联系人", + "联系人姓名", + "收件人", +} + +_EMAIL_PATTERN = re.compile( + r"(?姓名|真实姓名|联系人(?:姓名)?|收件人)" + r"(?P\s*(?:[::=]|为)\s*|\s+)" + r"(?P[\u3400-\u4dbf\u4e00-\u9fff·]{2,8})" +) +_ENGLISH_NAME_CONTEXT_PATTERN = re.compile( + r"(?im)(?P