From 4537faf2c4a8a1dd72cdfc72b88bd4626ec63f37 Mon Sep 17 00:00:00 2001 From: chenjw28 <792430652@qq.com> Date: Wed, 16 Sep 2026 17:55:04 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BC=98=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .env.example | 5 + AGENTS.md | 2 +- README.md | 2 +- alembic.ini | 2 +- deploy/Dockerfile.base | 6 +- deploy/Dockerfile.moldinsight | 37 ++- docker-compose.yml | 33 ++- docs/API_CONTRACT.md | 4 +- docs/ARCHITECTURE.md | 2 +- docs/OPERATIONS.md | 3 +- docs/ROADMAP.md | 28 ++- docs/STATUS.md | 6 + docs/TECH_DEBT.md | 130 +++++----- {alembic => migrations}/README | 0 {alembic => migrations}/env.py | 0 {alembic => migrations}/script.py.mako | 0 ...06c18c51b0d_add_product_id_to_stp_files.py | 0 .../versions/9928d7f8c1ef_initial_schema.py | 0 ...d91e47_add_batch_id_to_processing_tasks.py | 38 +++ src/celery_tasks.py | 7 +- src/moldinsight/api/advanced_router.py | 24 +- src/moldinsight/api/batch_router.py | 123 ++++----- src/moldinsight/api/task_router.py | 13 +- src/moldinsight/api/upload_router.py | 20 +- .../services/processing_service.py | 146 ++++++++--- .../services/storage_integration_rustfs.py | 45 ++-- src/moldinsight/services/task_dispatcher.py | 17 +- .../services/task_query_service.py | 43 +++- src/moldinsight/storage/rustfs_storage.py | 14 ++ src/shared/config/settings.py | 12 +- src/shared/database/init_db.py | 15 +- src/shared/models/database.py | 4 + src/shared/services/auth_routes.py | 12 +- src/shared/services/auth_service.py | 23 +- src/shared/services/redis_task_manager.py | 234 +++++++----------- tests/test_batch_status_pg.py | 126 ++++++++++ tests/test_deployment_config.py | 37 +++ tests/test_redis_no_fallback.py | 57 +++++ tests/test_status_endpoint_auth.py | 139 +++++++++++ 39 files changed, 986 insertions(+), 423 deletions(-) rename {alembic => migrations}/README (100%) rename {alembic => migrations}/env.py (100%) rename {alembic => migrations}/script.py.mako (100%) rename {alembic => migrations}/versions/006c18c51b0d_add_product_id_to_stp_files.py (100%) rename {alembic => migrations}/versions/9928d7f8c1ef_initial_schema.py (100%) create mode 100644 migrations/versions/a3f8c2d91e47_add_batch_id_to_processing_tasks.py create mode 100644 tests/test_batch_status_pg.py create mode 100644 tests/test_deployment_config.py create mode 100644 tests/test_redis_no_fallback.py create mode 100644 tests/test_status_endpoint_auth.py diff --git a/.env.example b/.env.example index c37db1a..3d73076 100644 --- a/.env.example +++ b/.env.example @@ -38,6 +38,11 @@ PARALLEL_PROCESSING=true # 数据库配置(服务器已部署,请填写真实地址) DB_HOST=localhost DB_PORT=5432 + +# 是否在应用启动时自动执行 alembic 迁移(默认 true,保持单机开发体验)。 +# 多副本/容器编排部署建议设为 false:多个实例同时启动会并发迁移, +# 改由部署流程单点执行 alembic CLI 或 python -m shared.database.init_db +AUTO_MIGRATE=true DB_NAME=moldinsight DB_USER=moldinsight_user DB_PASSWORD=moldinsight_password diff --git a/AGENTS.md b/AGENTS.md index b81818c..2e8dcbc 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -104,7 +104,7 @@ src/ celery_app.py # Celery app(Redis broker,task_acks_late) celery_tasks.py # moldinsight 异步分析任务 frontend/ # Vue 3 独立工程:src/modules 按域组织(moldinsight/inventory/users/login/home);src/types/api.ts 为 openapi 生成物,勿手改 -alembic/ # 数据库迁移 +migrations/ # 数据库迁移 scripts/ # 一次性迁移与工具脚本(migrations/ 数据迁移、db/ 索引与审计 SQL、tools/ 检查工具),非运行时代码 tests/ # pytest:sqlite+aiosqlite 临时库;pythonocc 缺失时 OCC 契约测试自动 skip deploy/ # Dockerfile.* / nginx / build 脚本 diff --git a/README.md b/README.md index a2cd5e0..0f0c278 100644 --- a/README.md +++ b/README.md @@ -92,7 +92,7 @@ geMoldInsight/ │ ├── celery_app.py # Celery app │ └── celery_tasks.py # moldinsight 异步任务 ├── frontend/ # 独立前端工程 -├── alembic/ # 数据库迁移 +├── migrations/ # 数据库迁移 ├── deploy/ # 镜像、Nginx、部署辅助文件 ├── docs/ ├── tests/ diff --git a/alembic.ini b/alembic.ini index 807ded2..08f1d9f 100644 --- a/alembic.ini +++ b/alembic.ini @@ -5,7 +5,7 @@ # this is typically a path given in POSIX (e.g. forward slashes) # format, relative to the token %(here)s which refers to the location of this # ini file -script_location = %(here)s/alembic +script_location = %(here)s/migrations # template used to generate migration file names; The default value is %%(rev)s_%%(slug)s # Uncomment the line below if you want the files to be prepended with date and time diff --git a/deploy/Dockerfile.base b/deploy/Dockerfile.base index a3ecce0..c116bba 100644 --- a/deploy/Dockerfile.base +++ b/deploy/Dockerfile.base @@ -1,4 +1,4 @@ -FROM python:3.12-slim +FROM python:3.12-slim-bookworm WORKDIR /app @@ -13,5 +13,9 @@ RUN pip install --no-cache-dir -r requirements-base.txt && \ COPY src/shared/ /app/src/shared/ +# init_db 启动期自动迁移(AUTO_MIGRATE)需要迁移脚本随镜像分发 +COPY migrations/ /app/migrations/ +COPY alembic.ini /app/alembic.ini + ENV PYTHONPATH=/app/src ENV PYTHONUNBUFFERED=1 diff --git a/deploy/Dockerfile.moldinsight b/deploy/Dockerfile.moldinsight index 581eb16..a06c205 100644 --- a/deploy/Dockerfile.moldinsight +++ b/deploy/Dockerfile.moldinsight @@ -1,24 +1,43 @@ -FROM continuumio/miniconda3:latest AS pythonocc +# PythonOCC 仅经 conda-forge 提供,且其动态库与 conda Python 的 ABI 绑定。 +# 旧方式(conda 环境装好后把 site-packages 拷入 python:slim 系统 python)依赖 +# 两侧 Python ABI 恰好兼容,属脆弱做法(TECH_DEBT D13);现改为直接以同一 +# conda 运行时作为最终镜像的执行环境,自带全部动态库。 +FROM continuumio/miniconda3:24.7.1-0 -RUN conda update -n base -c defaults conda -y && \ - conda create -n moldinsight python=3.12 pythonocc-core=7.9.0 -c conda-forge -y +# 锁定几何栈核心版本;pip 侧全量版本锁待首次镜像构建成功后由 +# `pip freeze > deploy/requirements-moldinsight.lock.txt` 生成(D13 遗留项) +RUN conda create -n moldinsight -c conda-forge -y \ + python=3.12 \ + pythonocc-core=7.9.0 \ + && conda clean -afy -FROM gemold-base:latest +ENV PATH=/opt/conda/envs/moldinsight/bin:$PATH \ + PYTHONUNBUFFERED=1 -COPY --from=pythonocc /opt/conda/envs/moldinsight/lib/python3.12/site-packages/ /usr/local/lib/python3.12/site-packages/ +WORKDIR /app -COPY deploy/requirements-moldinsight.txt . - -RUN pip install --no-cache-dir -r requirements-moldinsight.txt && \ - rm requirements-moldinsight.txt +# 自包含构建:不再基于 gemold-base,基础依赖与模块依赖一并安装 +COPY deploy/requirements-base.txt deploy/requirements-moldinsight.txt ./ +RUN pip install --no-cache-dir \ + -r requirements-base.txt \ + -r requirements-moldinsight.txt \ + && rm requirements-base.txt requirements-moldinsight.txt +COPY src/shared/ /app/src/shared/ COPY src/moldinsight/ /app/src/moldinsight/ COPY src/inventory/ /app/src/inventory/ COPY src/entrypoints/ /app/src/entrypoints/ COPY src/celery_app.py src/celery_tasks.py /app/src/ + +# init_db 启动期自动迁移(AUTO_MIGRATE)需要迁移脚本随镜像分发 +COPY migrations/ /app/migrations/ +COPY alembic.ini /app/alembic.ini + COPY uploads/ /app/uploads/ COPY html_output/ /app/html_output/ +ENV PYTHONPATH=/app/src + RUN mkdir -p /app/logs EXPOSE 8000 diff --git a/docker-compose.yml b/docker-compose.yml index e7b63f0..7fab610 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -40,9 +40,9 @@ services: RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} RUSTFS_TIMEOUT: ${RUSTFS_TIMEOUT:-30} - SECRET_KEY: ${SECRET_KEY:-change-me-in-production} + SECRET_KEY: ${SECRET_KEY:?SECRET_KEY 未配置:请在 .env 中设置} ADMIN_USERNAME: ${ADMIN_USERNAME:-admin} - ADMIN_PASSWORD: ${ADMIN_PASSWORD:-admin123} + ADMIN_PASSWORD: ${ADMIN_PASSWORD:?ADMIN_PASSWORD 未配置:请在 .env 中设置} ADMIN_EMAIL: ${ADMIN_EMAIL:-admin@gemold.com} ADMIN_FULL_NAME: ${ADMIN_FULL_NAME:-系统管理员} ALGORITHM: ${ALGORITHM:-HS256} @@ -64,6 +64,12 @@ services: LLM_MODEL: ${LLM_MODEL:-gpt-4o-mini} LLM_TIMEOUT: ${LLM_TIMEOUT:-60} LLM_MAX_TOKENS: ${LLM_MAX_TOKENS:-2000} + AUTO_MIGRATE: ${AUTO_MIGRATE:-true} + # 共享卷过渡兜底(D6/D11):主链路已改走 RustFS,本地卷仅为 + # RustFS 异常时的本地路径回退与 HTML 产物互通保留,后续批次移除 + volumes: + - uploads_data:/app/uploads + - html_data:/app/html_output restart: unless-stopped profiles: - full @@ -89,7 +95,7 @@ services: RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} RUSTFS_TIMEOUT: ${RUSTFS_TIMEOUT:-30} - SECRET_KEY: ${SECRET_KEY:-change-me-in-production} + SECRET_KEY: ${SECRET_KEY:?SECRET_KEY 未配置:请在 .env 中设置} ALGORITHM: ${ALGORITHM:-HS256} ACCESS_TOKEN_EXPIRE_MINUTES: ${ACCESS_TOKEN_EXPIRE_MINUTES:-1440} DEBUG: ${DEBUG:-false} @@ -109,6 +115,10 @@ services: LLM_MODEL: ${LLM_MODEL:-gpt-4o-mini} LLM_TIMEOUT: ${LLM_TIMEOUT:-60} LLM_MAX_TOKENS: ${LLM_MAX_TOKENS:-2000} + # 与 backend 共享本地卷(过渡兜底,见 D6/D11):worker 下载回退与 HTML 产物写读 + volumes: + - uploads_data:/app/uploads + - html_data:/app/html_output depends_on: - backend restart: unless-stopped @@ -142,9 +152,9 @@ services: RUSTFS_ACCESS_KEY: ${RUSTFS_ACCESS_KEY} RUSTFS_SECRET_KEY: ${RUSTFS_SECRET_KEY} RUSTFS_TIMEOUT: ${RUSTFS_TIMEOUT:-30} - SECRET_KEY: ${SECRET_KEY:-change-me-in-production} + SECRET_KEY: ${SECRET_KEY:?SECRET_KEY 未配置:请在 .env 中设置} ADMIN_USERNAME: ${ADMIN_USERNAME:-admin} - ADMIN_PASSWORD: ${ADMIN_PASSWORD:-admin123} + ADMIN_PASSWORD: ${ADMIN_PASSWORD:?ADMIN_PASSWORD 未配置:请在 .env 中设置} ADMIN_EMAIL: ${ADMIN_EMAIL:-admin@gemold.com} ADMIN_FULL_NAME: ${ADMIN_FULL_NAME:-系统管理员} ALGORITHM: ${ALGORITHM:-HS256} @@ -166,6 +176,10 @@ services: LLM_MODEL: ${LLM_MODEL:-gpt-4o-mini} LLM_TIMEOUT: ${LLM_TIMEOUT:-60} LLM_MAX_TOKENS: ${LLM_MAX_TOKENS:-2000} + AUTO_MIGRATE: ${AUTO_MIGRATE:-true} + volumes: + - uploads_data:/app/uploads + - html_data:/app/html_output restart: unless-stopped profiles: - moldinsight @@ -191,15 +205,16 @@ services: REDIS_HOST: ${REDIS_HOST} REDIS_PORT: ${REDIS_PORT:-6379} REDIS_PASSWORD: ${REDIS_PASSWORD:-} - SECRET_KEY: ${SECRET_KEY:-change-me-in-production} + SECRET_KEY: ${SECRET_KEY:?SECRET_KEY 未配置:请在 .env 中设置} ADMIN_USERNAME: ${ADMIN_USERNAME:-admin} - ADMIN_PASSWORD: ${ADMIN_PASSWORD:-admin123} + ADMIN_PASSWORD: ${ADMIN_PASSWORD:?ADMIN_PASSWORD 未配置:请在 .env 中设置} ADMIN_EMAIL: ${ADMIN_EMAIL:-admin@gemold.com} ADMIN_FULL_NAME: ${ADMIN_FULL_NAME:-系统管理员} ALGORITHM: ${ALGORITHM:-HS256} ACCESS_TOKEN_EXPIRE_MINUTES: ${ACCESS_TOKEN_EXPIRE_MINUTES:-1440} DEBUG: ${DEBUG:-false} SERVE_FRONTEND_STATIC: ${SERVE_FRONTEND_STATIC:-false} + AUTO_MIGRATE: ${AUTO_MIGRATE:-true} restart: unless-stopped profiles: - inventory @@ -209,3 +224,7 @@ services: networks: gemold_network: driver: bridge + +volumes: + uploads_data: + html_data: diff --git a/docs/API_CONTRACT.md b/docs/API_CONTRACT.md index ae00ad8..ca6ab65 100644 --- a/docs/API_CONTRACT.md +++ b/docs/API_CONTRACT.md @@ -44,8 +44,8 @@ | 域 | 端点 | 文件 | |---|---|---| | 上传 | `/api/upload` | upload_router.py | -| 批量分析 | `/api/batch-upload`、`/api/batch/{batch_id}` | batch_router.py | -| 任务状态 | `/api/status/{task_id}` | task_router.py | +| 批量分析 | `/api/batch-upload`、`/api/batch/{batch_id}`(聚合状态以 PG 为准;响应含 `current_step`;他人批次 403、不存在 404) | batch_router.py | +| 任务状态 | `/api/status/{task_id}`(需登录;仅任务所有者可访问,他人/无主任务 403,不存在 404) | task_router.py | | 历史结果 | `/api/history`、`/api/history/{filename}` | history_router.py | | CAM | `/api/cam/plan` | cam_router.py | | 铝价(模拟数据) | `/api/aluminum-price/current`、`/api/aluminum-price/history` | aluminum_price_routes.py | diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 64e8db7..11bc54c 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -96,7 +96,7 @@ geMoldInsight/ │ ├── celery_app.py │ └── celery_tasks.py ├── frontend/ -├── alembic/ +├── migrations/ ├── deploy/ ├── docs/ └── tests/ diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index 151a93f..374b131 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -12,6 +12,7 @@ - **Compose 运行**:compose 文件用 `${VAR}` 从同目录 `.env` 注入容器环境变量(见 [docker-compose.yml](../docker-compose.yml))。 - **键值约定**: - `DB_HOST / DB_PORT / DB_NAME / DB_USER / DB_PASSWORD`:**惰性校验、无代码默认**——缺失时 import 不报错(便于测试/静态分析),真正连库时才失败。生产必须显式配置。 + - `AUTO_MIGRATE`:应用启动时是否自动执行 alembic 迁移,默认 `true`(单机开发语义);**多副本 / 容器编排部署应设 `false`**,改由部署流程单点执行 `alembic upgrade head` 或 `python -m shared.database.init_db`(迁移脚本已随镜像分发于 `/app/migrations/`)。 - `SECRET_KEY`:JWT 签名密钥,**无默认**;生产必须 ≥32 字符强随机。 - `ADMIN_PASSWORD`:初始管理员密码,**无默认**;首次建库前必须设置。 - `RUSTFS_*`:对象存储(兼容 `MINIO_*` 别名写法);本地开发缺省值仅为占位,连不上会在用到存储的链路报错。 @@ -26,7 +27,7 @@ - 后端依赖:`pip install -r requirements.txt`。 - **OCC 几何能力**:PythonOCC 不走 pip 主路径,通过 conda 环境提供(本项目实践环境名 `gemold`)。无 OCC 环境时项目可启动,但几何分析契约测试自动 skip。 - 前端:`cd frontend && npm install`。 -- 数据库迁移:`alembic/`(`alembic.ini` 在仓库根);数据修复类一次性脚本在 `scripts/migrations/` 与 `scripts/db/`,**不是运行时代码**,勿在服务内引用。 +- 数据库迁移:`migrations/`(`alembic.ini` 在仓库根;2026-09-16 由 `alembic/` 改名——原目录名与 alembic 包重名,应用内 import 会被遮蔽导致启动期迁移静默失败);数据修复类一次性脚本在 `scripts/migrations/` 与 `scripts/db/`,**不是运行时代码**,勿在服务内引用。 ## 3. 本地启动 diff --git a/docs/ROADMAP.md b/docs/ROADMAP.md index 8894f92..94489bd 100644 --- a/docs/ROADMAP.md +++ b/docs/ROADMAP.md @@ -24,11 +24,13 @@ geMoldInsight 已从历史单体逐步演进为“双业务模块 + 共享平台 ### 2.1 主线一:模块化架构收口 目标: + - 继续巩固 `moldinsight / inventory / frontend / shared` 的边界 - 减少历史单体遗留语义 - 让 README、架构文档、部署文档与代码结构一致 重点方向: + - 继续收敛 `shared` 的职责 - 逐步明确 identity / platform 的边界语义 - 收敛历史文档与旧部署叙事 @@ -36,11 +38,13 @@ geMoldInsight 已从历史单体逐步演进为“双业务模块 + 共享平台 ### 2.2 主线二:moldinsight 工程化增强 目标: + - 让 STEP/STP 分析链路更稳定 - 让导出、批量分析、成本估算、任务状态等链路更可靠 - 继续提高 OCC 相关处理的可维护性与可测试性 重点方向: + - `advanced_router` 拆分与请求模型规范化 - 模具分析链路的结构继续收口 - OCC 依赖场景下的契约测试/集成测试继续补齐 @@ -48,11 +52,13 @@ geMoldInsight 已从历史单体逐步演进为“双业务模块 + 共享平台 ### 2.3 主线三:inventory 业务层继续沉淀 目标: + - 让 inventory 从“可用”继续走向“可扩展” - 继续将路由中的业务逻辑下沉为 service 层 - 保持与 moldinsight 的桥接模型清晰 重点方向: + - 业务 service 复用强化 - 数据模型归属进一步清晰化 - 前后端契约持续减少手写漂移 @@ -60,10 +66,12 @@ geMoldInsight 已从历史单体逐步演进为“双业务模块 + 共享平台 ### 2.4 主线四:部署与运维一致性 目标: + - 让推荐部署模式、Compose 入口、运维文档、Nginx/端口说明不再冲突 - 让前端、后端、异步任务链路在部署说明上形成单一叙事 重点方向: + - 继续收口部署文档 - 把历史部署迁移方案移入归档 - 保持同域前端 + unified backend 的默认认知清晰 @@ -103,16 +111,18 @@ geMoldInsight 已从历史单体逐步演进为“双业务模块 + 共享平台 > 2026-09-15 完成 moldinsight 后端设计审查,产出的具体治理批次是当前下一阶段最具体的执行计划。 > 债务明细与逐项现状见 [TECH_DEBT.md](TECH_DEBT.md) §3(D5–D14);本小节只描述批次、顺序与每批归属。 -| 批次 | 主题 | 内容 | 对应债务 | -|------|------|------|---------| -| 批次 0 | 安全与诚实(0.5–1 天) | `/api/status/{task_id}` 补鉴权 + 任务归属校验;`pythonocc_available` 真实检测;bcrypt 超长密码拒绝;SECRET_KEY / RUSTFS_* 惰性校验补齐 | D5 | -| 批次 1 | 部署正确性(1–2 天) | 主链路改走 RustFS(分派入参 `file_path` → `stp_file_id`,worker 按 object_key 下载解析);compose 共享卷兜底(过渡);alembic 移出 startup(`AUTO_MIGRATE` 开关);OCC 镜像引入方式修正 + 依赖锁文件 | D6、D12、D13 | -| 批次 2 | 任务一致性模型(2–4 天) | PG 为单一事实源、Redis 仅热缓存;去掉多进程内存回退;批量元数据入库;型腔失败标 failed;持久化事务边界收口 | D7、D8、D9、D11 | -| 批次 3 | API 与代码结构(3–5 天) | `_safe_include` 失败显式化(/health 暴露缺失路由);advanced_router 拆分 + Pydantic 请求模型;async 重计算统一 executor;StorageIntegrationService 拆分;配置治理 | D1、D14 | -| 批次 4 | 架构演进(5 天+) | 共享 ORM 按模块拆分;OCC 吞吐方案设计先行;文档 / 契约同步 | D3、D10 | +| 批次 | 主题 | 内容 | 对应债务 | +| ------ | ------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | --------------- | +| 批次 0 | 安全与诚实(0.5–1 天) | `/api/status/{task_id}` 补鉴权 + 任务归属校验;`pythonocc_available` 真实检测;bcrypt 超长密码拒绝;SECRET_KEY / RUSTFS_* 惰性校验补齐 | D5 | +| 批次 1 | 部署正确性(1–2 天) | 主链路改走 RustFS(分派入参`file_path` → `stp_file_id`,worker 按 object_key 下载解析);compose 共享卷兜底(过渡);alembic 移出 startup(`AUTO_MIGRATE` 开关);OCC 镜像引入方式修正 + 依赖锁文件 | D6、D12、D13 | +| 批次 2 | 任务一致性模型(2–4 天) | PG 为单一事实源、Redis 仅热缓存;去掉多进程内存回退;批量元数据入库;型腔失败标 failed;持久化事务边界收口 | D7、D8、D9、D11 | +| 批次 3 | API 与代码结构(3–5 天) | `_safe_include` 失败显式化(/health 暴露缺失路由);advanced_router 拆分 + Pydantic 请求模型;async 重计算统一 executor;StorageIntegrationService 拆分;配置治理 | D1、D14 | +| 批次 4 | 架构演进(5 天+) | 共享 ORM 按模块拆分;OCC 吞吐方案设计先行;文档 / 契约同步 | D3、D10 | **执行顺序建议**:批次 0 与批次 1 的 D6(RustFS 主链路)先行——前者是确认的安全漏洞,后者是部署根本性缺陷,两者互不依赖、改动可控。其余按批次顺序推进,每批完成同步 STATUS / TECH_DEBT / API_CONTRACT。 +> 进度:批次 0 / 1 / 2 已于 2026-09-16 完成(D13 的 pip 全量锁文件为批次 1 遗留项,随下次镜像构建补齐;D11 留待后续批次,正确性已由批次 1 共享卷兜底);完成明细见 [STATUS.md](STATUS.md) 与 [TECH_DEBT.md](TECH_DEBT.md) §2.5–2.6。 + --- ## 4. 中长期方向 @@ -120,6 +130,7 @@ geMoldInsight 已从历史单体逐步演进为“双业务模块 + 共享平台 ### 4.1 平台层语义收敛 长期仍建议将 `shared` 逐步收敛为更清晰的平台层语义,但这应建立在: + - 当前模块边界稳定 - 共享职责分层足够清晰 - 文档与部署已经同步收口 @@ -127,6 +138,7 @@ geMoldInsight 已从历史单体逐步演进为“双业务模块 + 共享平台 ### 4.2 文档体系持续治理 后续文档治理原则: + - README 只做入口 - 当前状态只在 [STATUS.md](STATUS.md) - 历史材料统一入 `docs/archive/` @@ -135,6 +147,7 @@ geMoldInsight 已从历史单体逐步演进为“双业务模块 + 共享平台 ### 4.3 测试能力继续增强 重点继续放在: + - OCC 相关集成验证 - 跨模块关键链路回归测试 - 关键契约的自动化保护 @@ -144,6 +157,7 @@ geMoldInsight 已从历史单体逐步演进为“双业务模块 + 共享平台 ### 4.4 专题文档持续分级 后续还会继续把专题文档区分为三类: + - 当前仍有参考价值的专题文档(保留并补定位) - 纯阶段性任务/检查单/迁移计划(迁入 archive) - 可被主骨架吸收的重复说明(逐步收口) diff --git a/docs/STATUS.md b/docs/STATUS.md index 4132f82..64cfb42 100644 --- a/docs/STATUS.md +++ b/docs/STATUS.md @@ -3,6 +3,12 @@ > 文档定位:**唯一的「现在到哪了」**。README / AGENTS / 各主文档只链接到这里,不复制状态内容。 > 维护规则:每完整完成一个需求,**倒序在本文顶部加一条**(日期 + 主题 + 关键事实);其余主文档(架构 / 规划 / 技术债 / 部署)维护各自的"当前有效说法",本文只记录"什么时候做到了哪一步"。维护规则出处见根目录 [AGENTS.md](../AGENTS.md)。 +> 2026-09-16(**批次 2(任务一致性模型)完成**:① D7 清偿——Redis 进程内存回退**彻底删除**(写 no-op / 读 None,查询路径自然落 PG),PG 为任务状态单一事实源;批量元数据入库:`processing_tasks` 新增 `batch_id` 列(迁移 `a3f8c2d91e47`,**升级后首次启动自动执行**),`GET /api/batch/{batch_id}` 改为 PG 聚合查询 + `STPFile.user_id` 归属校验,删除 Redis batch key 与内存 dict 双通道;`TaskQueryService` PG 视图与 batch 聚合响应补 `progress` / `current_step`(Redis 不可用时前端仍能看到进度);② D8 清偿——型腔分模失败不再吞异常,任务标 failed 并带明确错误(已提交的几何/网格保留);③ D9 清偿——数据本体写方法只 flush,编排层分阶段原子收口(阶段 A 几何+网格、阶段 B 型腔+HTML+特征+指标+验证、完成时参数随状态一并提交),失败先 rollback 再置 failed;进度/状态更新保留即时 commit(长任务进度可见性);upload/batch/advanced 调用方补显式 commit,STPFile + ProcessingTask 原子落库消除孤儿文件记录。D11 未动(共享卷已兜正确性,留后续批次)。**测试基线**:**105 passed, 1 skipped**(新增 [tests/test_batch_status_pg.py](../tests/test_batch_status_pg.py) 4 项 + [tests/test_redis_no_fallback.py](../tests/test_redis_no_fallback.py) 3 项)。**下一步**:批次 3(API 与代码结构:`_safe_include` 失败显式化、advanced_router 拆分 + Pydantic 请求模型、配置治理,见 [ROADMAP.md](ROADMAP.md) §3.1)。) + +> 2026-09-16(**批次 1(部署正确性)完成**:① D6 清偿——分派入参 `file_path` → `stp_file_id`,处理方按 PG 元数据从 RustFS 下载源文件到任务专属临时目录(RustFS 异常时回退节点本地路径),compose 增 `uploads_data` / `html_data` 共享卷过渡兜底;② D12 清偿——新增 `AUTO_MIGRATE` 开关(默认 true 保持单机行为;多副本设 false 改部署流程单点迁移),迁移脚本与 alembic.ini 补进镜像。**连带发现并修复**:迁移目录 `alembic/` 与 alembic 包重名,应用内 `import alembic` 被遮蔽——启动期自动迁移自引入 alembic 起**从未真正生效**(异常被 init_database 吞掉只打日志),且镜像原本未打包迁移脚本;目录已改名 `migrations/`(alembic.ini + 4 处文档引用同步);③ D13 主体——Dockerfile.moldinsight 改为 conda 运行时原生执行(不再跨镜像拷贝 site-packages),基础镜像 tag 锁定;pip 全量锁文件遗留,随下次镜像构建 `pip freeze` 生成;④ compose 关键项去弱默认:`SECRET_KEY` / `ADMIN_PASSWORD` 改 `${VAR:?}` 强制显式配置(与 OPERATIONS「无默认」声明对齐),`create_admin_user` 对空口令显式报错。**测试基线**:**98 passed, 1 skipped**(新增 [tests/test_deployment_config.py](../tests/test_deployment_config.py);alembic 缺失环境 skip)。**遗留**:D13 pip 锁文件;既有问题待查——Dockerfile.celery `FROM gemold-moldinsight:latest`,而 build.sh 只构建 `gemold-backend` tag,干净机器上 build.sh 的 celery 步骤会失败。**下一步**:批次 2(任务一致性模型,见 [ROADMAP.md](ROADMAP.md) §3.1)。) + +> 2026-09-16(**批次 0(安全与诚实)完成**:① `/api/status/{task_id}` 补 JWT 鉴权 + 任务归属校验(无 token 401 / 他人或无主任务 403 / 不存在 404),归属校验收敛为 `TaskQueryService.ensure_task_access` 供 task_router 与 advanced_router 共用——技术债 [D5 清偿](TECH_DEBT.md);② 上传预检 `pythonocc_available` 从硬编码 true 改为惰性真实探测;③ bcrypt 口令治理:创建侧超 72 字节显式拒绝(此前静默截断改变有效密码),验证侧截断比较(兼容历史哈希 + 避免超长登录 500);④ `SECRET_KEY` 未配置 / `RUSTFS_*` 缺失时惰性校验抛明确错误,代码侧不再有占位弱默认。**顺带修复**:完成态任务未持久化 `analysis_metrics` 时 `/api/status` 组装视图 500(值为 None 时 `.get(key, {})` 默认值不生效)。**测试基线**:pip 无 OCC 环境 **96 passed**(新增 [tests/test_status_endpoint_auth.py](../tests/test_status_endpoint_auth.py) 8 项回归)。status 端点鉴权为接口行为变化,已同步 [API_CONTRACT.md](API_CONTRACT.md) §3.2;openapi.json 重导出仍按既有待办随下次接口变更一并执行。**下一步**:批次 1(D6 RustFS 主链路 + D12 alembic 移出 startup + D13 OCC 镜像,见 [ROADMAP.md](ROADMAP.md) §3.1)。) + > 2026-09-15(**后端设计审查完成 → 治理计划入文档**:完成 moldinsight 后端设计审查(部署 / 任务一致性 / API / 代码结构),产出治理批次计划入 [ROADMAP.md](ROADMAP.md) §3.1(批次 0–4:安全→部署→一致性→结构→架构);新识别技术债 D5–D14 入 [TECH_DEBT.md](TECH_DEBT.md) §3——含确认安全缺口 `/api/status/{task_id}` 无鉴权、主处理链路依赖节点本地文件路径(API 与 Celery worker 容器无共享卷)等。**下一步**:按批次 0 + 批次 1 的 D6(RustFS 主链路)启动实施。) > 最后更新:2026-09-15(**项目规范体系对齐 ipc-chat-cortex**——参考 `ipc-chat-cortex` 的 AGENTS.md + docs 规范重整本文档体系:① [AGENTS.md](../AGENTS.md) 重写——硬约束速览(新增:接口变更三件套 Pydantic→openapi.json→gen:api、配置只走 .env 且关键项不兜底、单数据库刻意设计)+ 代码地图逐文件化 + 开发约定映射表(改什么→同步什么文档);② 新增 [OPERATIONS.md](OPERATIONS.md)(配置来源与优先级 / 本地启动 / Compose / 运维硬性要求)与 [API_CONTRACT.md](API_CONTRACT.md)(端点总览 / 统一约定 / OpenAPI 类型生成流程);③ 本文件改为日志体,原静态内容分流到各归属文档(推荐部署模式→DEPLOYMENT §1,未完成项→ROADMAP/TECH_DEBT)。**验证**:openapi 导出命令实测可用(conda gemold 环境,unified app 76 paths);**待办**:checked-in `openapi.json`(2026-07-27,70 paths)已落后当前代码,下次接口变更时按 [API_CONTRACT.md](API_CONTRACT.md) §4 重导出并 `npm run gen:api`。) diff --git a/docs/TECH_DEBT.md b/docs/TECH_DEBT.md index c578574..5545dcc 100644 --- a/docs/TECH_DEBT.md +++ b/docs/TECH_DEBT.md @@ -25,6 +25,8 @@ - debug/history 路由补鉴权 - 任务访问控制收紧 - 无主数据不再默认放行 +- `/api/status/{task_id}` 补 JWT 鉴权与归属校验(原 D5,2026-09-16 清偿,见 D5 条目) +- bcrypt 创建口令超 72 字节显式拒绝、验证侧截断比较;`SECRET_KEY` / `RUSTFS_*` 缺失时明确报错,代码侧弱默认移除(D14 部分,2026-09-16) ### 2.2 静默失败与可用性 - `detect-undercuts` 改为基于真实 shape 分析 @@ -41,6 +43,18 @@ - 设置惰性配置校验,提升可测试性 - Generator 公共接口提取完成,补充契约测试 +### 2.5 部署正确性(2026-09-16,批次 0/1) +- `/api/status/{task_id}` 补鉴权与归属校验(原 D5) +- 主处理链路改走 RustFS:分派入参 `stp_file_id` 化,源文件按 object_key 下载;compose 共享卷过渡兜底(原 D6) +- `AUTO_MIGRATE` 开关 + 迁移脚本随镜像分发 + `alembic/`→`migrations/` 改名修复包遮蔽(原 D12) +- OCC 镜像改 conda 运行时原生执行、基础镜像 tag 锁定(D13 主体);compose 关键项去弱默认(D14 部分) + +### 2.6 任务一致性模型(2026-09-16,批次 2) +- Redis 内存回退彻底删除,PG 为任务状态单一事实源(原 D7);批量元数据入库(`processing_tasks.batch_id`,迁移 `a3f8c2d91e47`) +- 型腔生成失败任务标 failed,不再静默 completed(原 D8) +- 持久化事务边界收口:数据本体分阶段原子提交、失败先回滚再置 failed(原 D9) +- D11(HTML 双写双读)本批未动:正确性已由共享卷兜底,RustFS 单一来源留待后续批次 + 详细历史过程保留在原始技术债文档中,后续将转入归档。 --- @@ -111,74 +125,51 @@ 优先级:**P1** -### D5. `/api/status/{task_id}` 未鉴权(安全缺口) +### D5. `/api/status/{task_id}` 未鉴权(安全缺口)—— 已清偿(2026-09-16,批次 0) -现状: -- [src/moldinsight/api/task_router.py](../src/moldinsight/api/task_router.py) 的 `/api/status/{task_id}` 未挂 `get_current_active_user`,也无任务归属校验 -- 匿名可枚举任务号拉取完整分析视图(几何 / 型腔方案 / LLM 报告 / 服务器本地路径) +修复内容(保留编号以维持 D6–D14 引用稳定): +- 端点补 `Depends(get_current_active_user)`;归属校验收敛为 `TaskQueryService.ensure_task_access`,task_router 与 advanced_router 共用(advanced_router 原私有 `_ensure_task_access` 改为委托) +- 语义:无 token 401、他人/无主任务 403(无主不等于公共)、任务不存在 404 +- 回归测试:[tests/test_status_endpoint_auth.py](../tests/test_status_endpoint_auth.py) +- 接口行为变化已同步 [API_CONTRACT.md](API_CONTRACT.md) §3.2 -影响: -- 与"任务访问控制已收紧"的既有结论矛盾;任务号可经批量/历史接口关联到真实用户 -- 属确认的安全漏洞,应最先修复 +~~原现状 / 影响~~:端点未挂鉴权,匿名可枚举任务号拉取完整分析视图。 -建议: -- 补 `Depends(get_current_active_user)` 并复用 `_ensure_task_access` 归属校验 +### D6. 主处理链路依赖节点本地文件路径 —— 已清偿(2026-09-16,批次 1) -优先级:**P0** +修复内容(保留编号以维持引用稳定): +- 分派入参收敛为 `stp_file_id`(`dispatch_processing` 与 Celery 任务签名同步变更):处理方按 PG 元数据从 RustFS 下载源文件到任务专属临时目录(保留原始文件名,下游产物命名不变),任务结束即清理([processing_service.py](../src/moldinsight/services/processing_service.py) `_materialize_source_file`) +- RustFS 不可用时回退 `STPFile.file_path` 节点本地路径;compose 为 backend / celery 增加共享卷 `uploads_data` / `html_data` 作过渡兜底(HTML 产物跨容器写读同源问题一并兜住,正式修复在 D11) -### D6. 主处理链路依赖节点本地文件路径 +~~原现状 / 影响~~:worker 直读 API 节点本地路径,双容器部署必然 `FileNotFoundError`。 -现状: -- 上传保存到本地目录,任务处理直接 `load_step_file(Path(file_path))`([processing_service.py](../src/moldinsight/services/processing_service.py)) -- docker-compose 中 backend 与 moldinsight-celery 为独立容器且无共享 volume,worker 读不到 API 节点写入的本地文件 +### D7. Redis 降级为进程内 dict,多副本状态不一致 —— 已清偿(2026-09-16,批次 2) -影响: -- 双容器部署下主流程必然 `FileNotFoundError`;代码已有从 RustFS 重建几何的 [shape_loader.py](../src/moldinsight/services/shape_loader.py),主链路却未复用 +修复内容(比原建议更彻底:完全删除内存回退,而非仅限 DEBUG): +- [redis_task_manager.py](../src/shared/services/redis_task_manager.py) 删除全部 `_fallback_*` 进程内存存储:Redis 不可用时写 no-op、读返回 None(Redis 仅热缓存,任务状态事实源在 PG,缓存缺失不影响正确性) +- 批量元数据入库:`processing_tasks` 新增 `batch_id` 列(迁移 `a3f8c2d91e47`),`GET /api/batch/{batch_id}` 改为按列聚合查询 + `STPFile.user_id` 归属校验,删除 Redis batch key 与进程内 dict 双通道 +- `TaskQueryService` 的 PG 组装视图补 `progress` / `current_step`(Redis 不可用时前端轮询仍能看到进度);batch 聚合响应同步补 `current_step` -建议: -- 分派入参由 `file_path` 改为 `stp_file_id`,worker 端按 `object_key` 从 RustFS 下载后解析 +~~原现状 / 影响~~:Redis 故障时状态静默降级各进程内存,多副本互不可见、同任务不同副本读到不同状态。 -优先级:**P0** +### D8. 型腔生成失败被静默标记为 completed —— 已清偿(2026-09-16,批次 2) -### D7. Redis 降级为进程内 dict,多副本状态不一致 +修复内容: +- [processing_service.py](../src/moldinsight/services/processing_service.py) `_step_generate_cavity` 不再吞异常:分模失败直接向编排层传播 → 任务 failed(error_message 说明型腔阶段失败);已提交的几何/网格数据保留,用户可凭失败原因重新分析 +- 未采用 `completed_with_fallback`:多一个状态值会扩散到前端所有状态分支,failed + 明确错误更诚实且成本低 -现状: -- Redis 不可用时任务状态 / 批量元数据 / 任务视图缓存静默降级到各进程内存([redis_task_manager.py](../src/shared/services/redis_task_manager.py) / [batch_router.py](../src/moldinsight/api/batch_router.py) / [task_query_service.py](../src/moldinsight/services/task_query_service.py)) +~~原现状 / 影响~~:型腔失败被吞掉继续主流程,最终 completed,"完成"状态不可信。 -影响: -- 多 worker + 多 API 副本下各进程内存互相不可见:同一任务在不同副本读到不同状态 +### D9. 持久化事务边界破碎 —— 已清偿(2026-09-16,批次 2) -建议: -- PG 作为单一事实源、Redis 仅热缓存;内存回退仅限单进程 DEBUG 模式 +修复内容(进度可见性与原子性折中设计): +- **数据本体写方法只 flush 不 commit**:`save_stp_file` / `save_geometry_data` / `save_mesh_data` / `save_mold_cavity_data` / `save_html_file` / `save_features_and_recommendations` / `update_task_parameters` / `update_stp_file_analysis_summary` / `_save_analysis_metrics` / `_save_verification_metrics` +- **编排层分阶段收口**([processing_service.py](../src/moldinsight/services/processing_service.py)):阶段 A = 几何+网格(解析后确定成果,原子提交);阶段 B = 型腔+HTML+特征+指标+摘要+验证(结果包原子提交);完成时先 flush 任务参数、完成状态提交时一并落库(completed 即完整) +- **失败路径先 rollback 再置 failed**:未提交半成品回滚,失败状态单独提交,不出现"completed 但数据残缺" +- **保留即时 commit**:`update_task_status` / `update_stp_file_status`(处理中进度需跨事务对外可见,分钟级长任务不能憋在一个大事务里) +- 调用方补显式 commit:upload_router / batch_router(分派前置事务,STPFile + ProcessingTask 原子,消除孤儿文件记录)、advanced_router 导出两处 -优先级:**P1** - -### D8. 型腔生成失败被静默标记为 completed - -现状: -- `_step_generate_cavity` 异常时置 `plan_result=None` 继续主流程,最终任务标记 completed([processing_service.py](../src/moldinsight/services/processing_service.py)) - -影响: -- 核心能力失败却对外呈现"成功","完成"状态可信度低 - -建议: -- 型腔失败 → 任务 failed,或显式 `completed_with_fallback` 并前端标注 - -优先级:**P1** - -### D9. 持久化事务边界破碎 - -现状: -- [storage_integration_rustfs.py](../src/moldinsight/services/storage_integration_rustfs.py) 各方法内部自行 `session.commit()`,编排层上下文又 commit -- 型腔保存失败时几何/网格等前期数据已提交落库 - -影响: -- 失败后留下已提交的半成品数据,无对账补偿 - -建议: -- 各方法不再自提交,由编排层统一提交;明确 RustFS 与 PG 写入顺序 - -优先级:**P1** +~~原现状 / 影响~~:各存储方法内部自行 commit,型腔保存失败留半成品数据且任务仍 completed。 ### D10. OCC 全局单线程串行 + 超时重建泄漏线程 @@ -209,35 +200,26 @@ ### D12. 应用启动时自动执行 alembic 迁移 -现状: -- [init_db.py](../src/shared/database/init_db.py) 在 web 进程 startup 中执行 `alembic upgrade head` +### D12. 应用启动时自动执行 alembic 迁移 —— 已清偿(2026-09-16,批次 1) -影响: -- 多副本并发迁移有竞态,且迁移阻塞服务就绪 +修复内容: +- 新增 `AUTO_MIGRATE` 开关(settings / .env.example / compose 透传):默认 `true` 保持单机开发行为;多副本部署设 `false`,由部署流程单点执行 alembic CLI 或 `python -m shared.database.init_db` +- **连带发现并修复两个使自动迁移从未真正生效的缺陷**: + 1. 迁移目录 `alembic/` 与 alembic 包重名——应用内 `import alembic` 命中本地目录(namespace package)遮蔽真实包,启动期迁移异常被 `init_database` 吞掉只打日志;已改名 `migrations/`(alembic.ini `script_location` 与 4 处文档引用同步) + 2. 镜像未打包迁移脚本与 alembic.ini,容器内迁移必然失败——Dockerfile.base / Dockerfile.moldinsight 已补 `COPY migrations/` + `COPY alembic.ini` -建议: -- 迁移移出 web 进程,作为独立部署步骤(`AUTO_MIGRATE` 开关) - -优先级:**P1** - -### D13. PythonOCC 镜像引入方式脆弱 + 依赖无版本锁 +### D13. PythonOCC 镜像引入方式脆弱 + 依赖无版本锁(主体已清偿,锁文件遗留) 现状: -- [Dockerfile.moldinsight](../deploy/Dockerfile.moldinsight) 从 conda env 拷贝 site-packages 进 python:3.12-slim -- [requirements.txt](../requirements.txt) 全部为 `>=` 下限,无锁文件 +- ~~从 conda env 拷贝 site-packages 进 python:3.12-slim~~(2026-09-16 已修正:[Dockerfile.moldinsight](../deploy/Dockerfile.moldinsight) 改为 conda 运行时原生执行,不再跨镜像拷贝;基础镜像 tag 锁定 `continuumio/miniconda3:24.7.1-0`、`python:3.12-slim-bookworm`;tag 可用性随下次镜像构建验证) +- [requirements.txt](../requirements.txt) 全部为 `>=` 下限,无锁文件(**遗留**:首次镜像构建成功后 `pip freeze` 生成锁文件,命令已注释在 Dockerfile 内) -影响: -- slim 缺 libstdc++/libgomp 等运行时库,跨发行版拷二进制纯靠运气;构建不可复现 - -建议: -- 基础镜像改用完整 conda 环境;依赖以 pip-compile 锁文件固化 - -优先级:**P2** +优先级:**P2**(剩余锁文件部分) ### D14. 配置漂移:弱默认 / 死配置 / 重复解析 现状: -- RUSTFS_* 带 `localhost:8080` / `your-secret-key` 弱默认;compose 给 SECRET_KEY / ADMIN_PASSWORD 弱默认 +- ~~RUSTFS_* 弱默认~~(2026-09-16 代码侧已去除);~~compose 侧 SECRET_KEY / ADMIN_PASSWORD 弱默认~~(2026-09-16 已去除:改用 `${VAR:?}` 强制显式配置,`create_admin_user` 对空 ADMIN_PASSWORD 显式报错) - MAX_FILE_SIZE 配置项未被使用([file_handler.py](../src/shared/utils/file_handler.py) 硬编码 50MB) - [celery_app.py](../src/celery_app.py) 重新 load_dotenv 并手拼 REDIS URL,与 settings 两份实现 diff --git a/alembic/README b/migrations/README similarity index 100% rename from alembic/README rename to migrations/README diff --git a/alembic/env.py b/migrations/env.py similarity index 100% rename from alembic/env.py rename to migrations/env.py diff --git a/alembic/script.py.mako b/migrations/script.py.mako similarity index 100% rename from alembic/script.py.mako rename to migrations/script.py.mako diff --git a/alembic/versions/006c18c51b0d_add_product_id_to_stp_files.py b/migrations/versions/006c18c51b0d_add_product_id_to_stp_files.py similarity index 100% rename from alembic/versions/006c18c51b0d_add_product_id_to_stp_files.py rename to migrations/versions/006c18c51b0d_add_product_id_to_stp_files.py diff --git a/alembic/versions/9928d7f8c1ef_initial_schema.py b/migrations/versions/9928d7f8c1ef_initial_schema.py similarity index 100% rename from alembic/versions/9928d7f8c1ef_initial_schema.py rename to migrations/versions/9928d7f8c1ef_initial_schema.py diff --git a/migrations/versions/a3f8c2d91e47_add_batch_id_to_processing_tasks.py b/migrations/versions/a3f8c2d91e47_add_batch_id_to_processing_tasks.py new file mode 100644 index 0000000..bc06d72 --- /dev/null +++ b/migrations/versions/a3f8c2d91e47_add_batch_id_to_processing_tasks.py @@ -0,0 +1,38 @@ +"""add batch_id to processing_tasks + +批次 2(D7):批量元数据入库——processing_tasks 增加 batch_id 列, +批量任务聚合查询走 PG,替代 Redis/进程内存中的批量元数据。 + +Revision ID: a3f8c2d91e47 +Revises: 006c18c51b0d +Create Date: 2026-09-16 + +""" +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa + + +# revision identifiers, used by Alembic. +revision: str = 'a3f8c2d91e47' +down_revision: Union[str, Sequence[str], None] = '006c18c51b0d' +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.add_column( + 'processing_tasks', + sa.Column('batch_id', sa.String(length=36), nullable=True), + ) + op.create_index( + 'ix_processing_tasks_batch_id', + 'processing_tasks', + ['batch_id'], + ) + + +def downgrade() -> None: + op.drop_index('ix_processing_tasks_batch_id', table_name='processing_tasks') + op.drop_column('processing_tasks', 'batch_id') diff --git a/src/celery_tasks.py b/src/celery_tasks.py index f0d9d40..a6c1b84 100644 --- a/src/celery_tasks.py +++ b/src/celery_tasks.py @@ -7,9 +7,8 @@ logger = get_logger(__name__) @app.task(bind=True, max_retries=1, default_retry_delay=60) -def process_stp_task(self, task_id: str, file_path: str, stp_file_id: int, - process_params: dict): - """Celery 任务:异步处理 STP 文件生成模具型腔""" +def process_stp_task(self, task_id: str, stp_file_id: int, process_params: dict): + """Celery 任务:异步处理 STP 文件生成模具型腔(源文件按 stp_file_id 从 RustFS 获取)""" import asyncio async def _run(): @@ -34,7 +33,7 @@ def process_stp_task(self, task_id: str, file_path: str, stp_file_id: int, await db_manager.connect(role="celery") await processing_service.process_file_with_storage( - task_id, file_path, stp_file_id, process_params + task_id, stp_file_id, process_params ) try: diff --git a/src/moldinsight/api/advanced_router.py b/src/moldinsight/api/advanced_router.py index f34a6e4..dfc767d 100644 --- a/src/moldinsight/api/advanced_router.py +++ b/src/moldinsight/api/advanced_router.py @@ -4,7 +4,6 @@ from datetime import datetime from urllib.parse import quote from fastapi import APIRouter, Depends, HTTPException, Request -from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from shared.services.auth_service import get_current_active_user @@ -14,7 +13,6 @@ from moldinsight.services.storage_integration_rustfs import StorageIntegrationSe from moldinsight.services.task_query_service import TaskQueryService from shared.database.database import get_db_session from shared.models.database import User -from shared.models.database import ProcessingTask, STPFile from moldinsight.core.cad_exporter import CADExporter from shared.utils.logger import get_logger @@ -70,22 +68,8 @@ async def _ensure_task_access( task_id: str, user_id: int, ): - row = await db_session.execute( - select(ProcessingTask, STPFile) - .join(STPFile, ProcessingTask.stp_file_id == STPFile.id) - .where(ProcessingTask.task_id == task_id) - ) - row = row.first() - if not row: - raise HTTPException(404, "任务不存在") - - _, stp_file = row - owner_id = getattr(stp_file, "user_id", None) - if owner_id != user_id: - # 无主历史数据(owner_id is None)同样拒绝:无主不等于公共 - raise HTTPException(403, "无权访问该任务的导出文件") - - return row + # 归属校验统一走 TaskQueryService(与 /api/status 共用,含 404/403 语义) + return await TaskQueryService.ensure_task_access(db_session, task_id, user_id) def _get_export_artifacts(task_data: dict) -> dict: @@ -513,6 +497,8 @@ async def export_mold_results( task_id, {"export_artifacts": merged_artifacts}, ) + # D9:存储方法已不再自行 commit,请求侧显式提交 + await db_session.commit() await redis_task_manager.update_task( task_id, {"export_artifacts": merged_artifacts} ) @@ -543,6 +529,8 @@ async def export_mold_results( task_id, {"export_artifacts": merged_artifacts}, ) + # D9:存储方法已不再自行 commit,请求侧显式提交 + await db_session.commit() await redis_task_manager.update_task(task_id, {"export_artifacts": merged_artifacts}) TaskQueryService.invalidate_task_view(task_id) # parameters 已变更,缓存视图失效 diff --git a/src/moldinsight/api/batch_router.py b/src/moldinsight/api/batch_router.py index cee27e3..f569e85 100644 --- a/src/moldinsight/api/batch_router.py +++ b/src/moldinsight/api/batch_router.py @@ -3,17 +3,22 @@ moldinsight/api/batch_router.py — 批量分析端点 - POST /api/batch-upload 批量上传多文件,返回 batch_id + 各 task_id - GET /api/batch/{batch_id} 聚合查询批量任务进度 + +批次 2(D7):批量元数据以 PG 为单一事实源——ProcessingTask.batch_id +列聚合查询,替代此前的 Redis key + 进程内存降级存储。 """ import uuid from datetime import datetime -from typing import List, Dict, Any, Optional +from typing import List, Dict, Any from fastapi import APIRouter, UploadFile, File, Form, HTTPException, Depends +from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession +from sqlalchemy.orm import joinedload from shared.database.database import get_db_session from shared.services.auth_service import get_current_active_user -from shared.models.database import User +from shared.models.database import User, ProcessingTask, STPFile from shared.models.schemas import ProcessingStatus, create_task_info from shared.services.redis_task_manager import redis_task_manager from shared.utils.file_handler import FileHandler @@ -27,17 +32,6 @@ router = APIRouter() file_handler = FileHandler() -# ─── 批量元数据 Redis key 约定 ────────────────────────────────────── -_BATCH_KEY_PREFIX = "batch:" -_BATCH_TTL = 86400 # 24h - -# Redis 不可用时的进程内降级存储(同进程内可查,跨进程/重启不可见) -_batch_meta_memory: Dict[str, dict] = {} - - -def _batch_redis_key(batch_id: str) -> str: - return f"{_BATCH_KEY_PREFIX}{batch_id}" - @router.post("/batch-upload") async def batch_upload( @@ -91,8 +85,11 @@ async def batch_upload( ) await storage_service.create_processing_task( - db_session, task_id, stp_file.id, parameters=process_params, + db_session, task_id, stp_file.id, + parameters=process_params, batch_id=batch_id, ) + # D9:STPFile + ProcessingTask 原子提交,分派前置事务收口 + await db_session.commit() task_info = create_task_info( task_id=task_id, @@ -108,7 +105,7 @@ async def batch_upload( await redis_task_manager.set_task(task_id, task_info) # 调度处理 - dispatch_processing(task_id, str(file_path), stp_file.id, process_params) + dispatch_processing(task_id, stp_file.id, process_params) tasks.append({ "filename": file.filename, @@ -129,17 +126,6 @@ async def batch_upload( "error": str(exc), }) - # 将 batch 元数据写入 Redis;Redis 不可用时降级到进程内存储(任务状态本身有内存回退) - batch_meta = { - "batch_id": batch_id, - "user_id": current_user.id, - "created_at": str(datetime.now()), - "task_ids": [t["task_id"] for t in tasks if t.get("task_id")], - "total": len(tasks), - "params": process_params, - } - _save_batch_meta(batch_id, batch_meta) - return { "batch_id": batch_id, "total": len(tasks), @@ -148,67 +134,41 @@ async def batch_upload( } -def _save_batch_meta(batch_id: str, batch_meta: dict): - """批量元数据持久化:优先 Redis(跨进程、带 TTL),降级进程内 dict。""" - import json as _json - - if redis_task_manager.is_connected: - try: - redis_task_manager.redis_client.set( - _batch_redis_key(batch_id), - _json.dumps(batch_meta), - ex=_BATCH_TTL, - ) - return - except Exception as exc: - logger.warning(f"[BATCH] batch 元数据写 Redis 失败,降级内存: {exc}") - _batch_meta_memory[batch_id] = batch_meta - - -async def _load_batch_meta(batch_id: str) -> Optional[dict]: - """读取批量元数据,Redis 优先,内存兜底;不存在返回 None。""" - import json as _json - - if redis_task_manager.is_connected: - try: - raw = await redis_task_manager.redis_client.get(_batch_redis_key(batch_id)) - if raw: - return _json.loads(raw) - except Exception as exc: - logger.warning(f"[BATCH] batch 元数据读 Redis 失败: {exc}") - return _batch_meta_memory.get(batch_id) - - @router.get("/batch/{batch_id}") async def get_batch_status( batch_id: str, + db_session: AsyncSession = Depends(get_db_session), current_user: User = Depends(get_current_active_user), ): - """聚合查询批量任务进度""" - batch_meta = await _load_batch_meta(batch_id) - if not batch_meta: + """聚合查询批量任务进度(D7:以 PG 为单一事实源,按 batch_id 聚合;Redis 仅热缓存)""" + rows = (await db_session.execute( + select(ProcessingTask, STPFile) + .join(STPFile, ProcessingTask.stp_file_id == STPFile.id) + .where(ProcessingTask.batch_id == batch_id) + .options(joinedload(STPFile.html_file)) + .order_by(ProcessingTask.id) + )).unique().all() + + if not rows: raise HTTPException(404, "批量任务不存在或已过期") - # 权限检查 - if batch_meta.get("user_id") and batch_meta["user_id"] != current_user.id: + # 归属校验:同批任务属于同一上传用户,任一不匹配即拒绝(无主不等于公共) + if any(getattr(stp, "user_id", None) != current_user.id for _, stp in rows): raise HTTPException(403, "无权访问该批量任务") - task_ids = batch_meta.get("task_ids", []) task_statuses = [] completed = 0 failed = 0 processing = 0 + earliest_created = None - for tid in task_ids: - task_data = await redis_task_manager.get_task(tid) - if not task_data: - task_statuses.append({"task_id": tid, "status": "unknown"}) - continue - status = task_data.get("status", "unknown") - progress = task_data.get("progress", 0) - filename = task_data.get("filename", "") - error = task_data.get("error", "") - html_file = task_data.get("html_file", "") + for task, stp in rows: + status = task.status or "unknown" + + if earliest_created is None or ( + task.created_time and task.created_time < earliest_created + ): + earliest_created = task.created_time if status == ProcessingStatus.COMPLETED: completed += 1 @@ -217,19 +177,24 @@ async def get_batch_status( else: processing += 1 + html_file = "" + if stp.html_file and stp.html_file.filename: + html_file = f"/html/{stp.html_file.filename}" + task_statuses.append({ - "task_id": tid, + "task_id": task.task_id, "status": status, - "progress": progress, - "filename": filename, - "error": error, + "progress": task.progress or 0, + "current_step": task.current_step, + "filename": stp.original_filename or "", + "error": task.error_message or "", "html_file": html_file, }) - total = len(task_ids) + total = len(rows) return { "batch_id": batch_id, - "created_at": batch_meta.get("created_at"), + "created_at": earliest_created.isoformat() if earliest_created else None, "total": total, "completed": completed, "failed": failed, diff --git a/src/moldinsight/api/task_router.py b/src/moldinsight/api/task_router.py index 307fe98..f3e0d35 100644 --- a/src/moldinsight/api/task_router.py +++ b/src/moldinsight/api/task_router.py @@ -1,13 +1,13 @@ # api/v1/task_router.py from fastapi import APIRouter, HTTPException, Request, Depends -from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from moldinsight.services.task_query_service import TaskQueryService from shared.database.database import get_db_session +from shared.services.auth_service import get_current_active_user from shared.utils.logger import get_logger -from shared.models.database import ProcessingTask, STPFile +from shared.models.database import User logger = get_logger(__name__) @@ -16,15 +16,20 @@ router = APIRouter() @router.get("/status/{task_id}") @router.post("/status/{task_id}") -async def get_status(task_id: str, db_session: AsyncSession = Depends(get_db_session)): +async def get_status( + task_id: str, + db_session: AsyncSession = Depends(get_db_session), + current_user: User = Depends(get_current_active_user), +): """ - 获取任务状态 + 获取任务状态(需登录,且仅任务所有者可访问) 优先返回内存中的任务信息; 如果内存中不存在,则从 PostgreSQL + RustFS 组装一个持久化的任务视图, 结构与内存任务保持尽量一致,便于前端集中展示总结性信息。 """ try: + await TaskQueryService.ensure_task_access(db_session, task_id, current_user.id) task_view = await TaskQueryService.get_task_view(db_session, task_id) if task_view is None: raise HTTPException(404, "任务不存在") diff --git a/src/moldinsight/api/upload_router.py b/src/moldinsight/api/upload_router.py index 871fe9f..f7c0fe4 100644 --- a/src/moldinsight/api/upload_router.py +++ b/src/moldinsight/api/upload_router.py @@ -21,6 +21,19 @@ router = APIRouter() file_handler = FileHandler() +def _occ_available() -> bool: + """真实检测 PythonOCC 可用性(惰性导入,缺失时不影响本路由加载)。 + + 此前该字段硬编码 True,响应不诚实;几何处理依赖 OCC, + 不可用时任务会在处理阶段以明确错误失败。 + """ + try: + import OCC.Core.STEPControl # noqa: F401 + return True + except Exception: + return False + + @router.post("/upload") async def upload_stp( file: UploadFile = File(...), @@ -75,6 +88,9 @@ async def upload_stp( stp_file.id, parameters=process_params, ) + # D9:create_processing_task 仅 flush,STPFile + 任务记录在此一并原子提交, + # 分派前置事务收口——分派出去的任务保证在 PG 中可见 + await db_session.commit() task_info = create_task_info( task_id=task_id, @@ -89,7 +105,7 @@ async def upload_stp( task_info["file_hash"] = file_meta["sha256"] await redis_task_manager.set_task(task_id, task_info) - dispatch_processing(task_id, str(file_path), stp_file.id, process_params) + dispatch_processing(task_id, stp_file.id, process_params) return { "task_id": task_id, @@ -98,7 +114,7 @@ async def upload_stp( "file_info": { "filename": file.filename, "size": file_size, - "pythonocc_available": True, + "pythonocc_available": _occ_available(), "database_file_id": stp_file.id, "sha256": file_meta["sha256"], }, diff --git a/src/moldinsight/services/processing_service.py b/src/moldinsight/services/processing_service.py index 414b34e..e0c4054 100644 --- a/src/moldinsight/services/processing_service.py +++ b/src/moldinsight/services/processing_service.py @@ -3,14 +3,16 @@ import asyncio import os +import shutil +import tempfile import time -import traceback from collections import OrderedDict from concurrent.futures import ThreadPoolExecutor from datetime import datetime from pathlib import Path -from typing import Optional, Dict, Any, List +from typing import Optional, Dict, Any, List, Tuple +from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from moldinsight.core.stp_parser import STPParser @@ -19,11 +21,13 @@ from moldinsight.core.mesh_generator import MeshGenerator from moldinsight.core.multi_scheme_planner import MultiSchemeMoldPlanner from moldinsight.core.cad_exporter import CADExporter from moldinsight.services.storage_integration_rustfs import StorageIntegrationService +from moldinsight.storage.rustfs_storage import rustfs_manager from shared.services.redis_task_manager import redis_task_manager from moldinsight.services.material_service import MaterialService from moldinsight.services.calculation_service import CalculationService from moldinsight.services.llm_service import llm_service from shared.models.schemas import ProcessingStatus +from shared.models.database import STPFile from shared.database.database import db_manager from shared.utils.html_generator import HTMLGenerator from shared.utils.logger import get_logger @@ -71,23 +75,73 @@ class ProcessingService: loop = asyncio.get_running_loop() return await loop.run_in_executor(self._occ_executor, fn, *args) + async def _materialize_source_file(self, stp_file: STPFile) -> Tuple[Path, Optional[Path]]: + """把待处理文件落到本地磁盘,返回 (本地路径, 临时目录或 None)。 + + RustFS 为主存储:处理方按 object_key 下载到任务专属临时目录 + (文件名保留原始名——下游产物命名依赖 Path(file_path).name)。 + RustFS 不可用或对象缺失时回退 STPFile.file_path 记录的节点本地路径 + (依赖 compose 共享卷,属过渡方案);两者皆不可用则抛错置任务失败。 + """ + original_name = Path(stp_file.original_filename or "model.stp").name or "model.stp" + + if stp_file.object_key: + temp_dir = Path(tempfile.mkdtemp(prefix=f"moldinsight_{stp_file.id}_")) + try: + data = await rustfs_manager.download_file( + file_type="stp_files", object_key=stp_file.object_key + ) + target = temp_dir / original_name + target.write_bytes(data) + return target, temp_dir + except Exception as exc: + shutil.rmtree(temp_dir, ignore_errors=True) + logger.warning( + f"RustFS 源文件下载失败 (object_key={stp_file.object_key})," + f"回退节点本地路径: {exc}" + ) + + local = Path(stp_file.file_path) if stp_file.file_path else None + if local and local.exists(): + return local, None + + raise RuntimeError( + f"源文件不可用:RustFS 对象 {stp_file.object_key!r} 下载失败," + f"且节点本地路径不存在: {stp_file.file_path!r}" + ) + async def process_file_with_storage( self, task_id: str, - file_path: str, stp_file_id: int, process_params: Optional[Dict[str, Any]] = None, ): - """处理文件的后台任务 — 使用独立数据库会话""" + """处理文件的后台任务 — 使用独立数据库会话 + + 分派入参只带 stp_file_id(D6):源文件由本方法按 PG 元数据中的 + object_key 从 RustFS 获取,不再依赖分派方传入节点本地路径 + (API 与 Celery worker 容器文件系统不互通)。 + """ # 创建独立的数据库会话,避免请求范围会话关闭 async with db_manager.session() as db_session: + temp_dir: Optional[Path] = None try: + result = await db_session.execute( + select(STPFile).where(STPFile.id == stp_file_id) + ) + stp_file = result.scalar_one_or_none() + if not stp_file: + raise RuntimeError(f"STPFile 记录不存在: stp_file_id={stp_file_id}") + + source_path, temp_dir = await self._materialize_source_file(stp_file) + file_path = str(source_path) + logger.info(f"开始处理文件并生成模具型腔: {file_path}") from shared.config.settings import settings - file_size_bytes = Path(file_path).stat().st_size if Path(file_path).exists() else 0 + file_size_bytes = source_path.stat().st_size file_size_mb = max(file_size_bytes / (1024 * 1024), 1) timeout_seconds = min( max(settings.PROCESSING_TIMEOUT_BASE, int(file_size_mb * settings.PROCESSING_TIMEOUT_PER_MB)), @@ -111,12 +165,16 @@ class ProcessingService: except Exception as e: logger.error(f"模具型腔生成失败: {e}") + # D9:先丢弃未提交的数据本体,失败状态单独提交, + # 避免 failed 更新把半成品 flush 数据一起带上 + await db_session.rollback() + await self.storage_service.update_stp_file_status(db_session, stp_file_id, "failed") await self.storage_service.update_task_status( db_session, task_id, "failed", error_message=str(e) ) - # 安全更新 Redis 任务状态 + # 安全更新 Redis 任务状态(Redis 仅热缓存,写失败不影响 PG 事实) task = await redis_task_manager.get_task(task_id) if task: await redis_task_manager.update_task(task_id, { @@ -124,6 +182,10 @@ class ProcessingService: "error": str(e), "completed_at": str(datetime.now()), }) + finally: + # 任务专属临时目录必须清理,长期运行不允许残留下载副本 + if temp_dir: + shutil.rmtree(temp_dir, ignore_errors=True) async def process_file_core( self, @@ -234,6 +296,10 @@ class ProcessingService: geometry_data.get("analysis_method", "mold_cavity"), ) + # 阶段 A 提交(D9):几何 + 网格原子落库——解析后的确定成果, + # 后续型腔失败任务标 failed 时这些数据仍完整保留 + await db_session.commit() + # 7. 生成HTML可视化 await self.storage_service.update_task_status( db_session, task_id, "processing", 85, "生成可视化报告" @@ -325,6 +391,11 @@ class ProcessingService: ), ) + # 9.65 阶段 B 提交(D9):型腔 / HTML / 特征 / 指标 / 摘要 / 验证指标 + # 作为完整结果包原子落库——置 completed 前必须全部就位, + # 期间任一步失败回滚后任务标 failed,不会出现"completed 但数据残缺" + await db_session.commit() + # 9.7 FreeCAD 几何验证 stage_started = time.perf_counter() verification_result = await self._step_verify( @@ -347,11 +418,7 @@ class ProcessingService: ) stage_timings["generate_llm_report"] = round(time.perf_counter() - stage_started, 3) - # 10. 完成处理 - await self.storage_service.update_stp_file_status(db_session, stp_file_id, "completed") - await self.storage_service.update_task_status( - db_session, task_id, "completed", 100, "模具型腔生成完成" - ) + # 10. 完成处理——先 flush 任务参数,完成状态提交时一并原子落库(D9) await self.storage_service.update_task_parameters( db_session, task_id, @@ -364,6 +431,10 @@ class ProcessingService: **process_params, }, ) + await self.storage_service.update_stp_file_status(db_session, stp_file_id, "completed") + await self.storage_service.update_task_status( + db_session, task_id, "completed", 100, "模具型腔生成完成" + ) # 更新任务缓存状态(仅保留轻量摘要,完整数据由PG+RustFS持久化; # 完成态视图由 TaskQueryService 从 PG+RustFS 组装,Redis 不再存 @@ -390,6 +461,9 @@ class ProcessingService: except Exception as e: logger.error(f"模具型腔生成失败: {e}") + # D9:先丢弃未提交的数据本体再置失败(同外层说明) + await db_session.rollback() + await self.storage_service.update_stp_file_status(db_session, stp_file_id, "failed") await self.storage_service.update_task_status( db_session, task_id, "failed", error_message=str(e) @@ -471,29 +545,29 @@ class ProcessingService: async def _step_generate_cavity( self, shape, selected_material: dict, is_foam_material: bool, process_params: Dict[str, Any], - ) -> Optional[Dict[str, Any]]: - """生成多方案分模结果""" - plan_result = None - try: - if shape: - loop = asyncio.get_running_loop() - plan_result = await loop.run_in_executor( - self._occ_executor, - lambda: self.multi_scheme_planner.generate_plan( - shape=shape, - material=selected_material, - is_foam_material=is_foam_material, - process_params=process_params, - ), - ) - logger.info( - f"多方案分模完成: 生成 {len(plan_result.get('candidate_schemes', []))} 套方案" - ) - except Exception as cavity_err: - logger.warning(f"多方案分模失败,使用简化数据: {cavity_err}") - traceback.print_exc() - plan_result = None + ) -> Dict[str, Any]: + """生成多方案分模结果。 + D8:型腔是任务的核心产出,生成失败必须让任务 failed—— + 此前异常在此被吞掉置 plan_result=None 继续主流程,最终任务 + completed,"完成"状态不可信。异常直接向编排层传播。 + """ + if not shape: + raise RuntimeError("无有效几何 shape,无法生成模具型腔") + + loop = asyncio.get_running_loop() + plan_result = await loop.run_in_executor( + self._occ_executor, + lambda: self.multi_scheme_planner.generate_plan( + shape=shape, + material=selected_material, + is_foam_material=is_foam_material, + process_params=process_params, + ), + ) + logger.info( + f"多方案分模完成: 生成 {len(plan_result.get('candidate_schemes', []))} 套方案" + ) return plan_result def _cache_export_shapes(self, task_id: str, export_shapes: Dict[str, Dict[str, Any]]): @@ -727,7 +801,8 @@ class ProcessingService: ) session.add(metrics) - await session.commit() + # D9:flush 不 commit,随结果包(阶段 B)由编排层统一提交 + await session.flush() logger.info(f"分析指标保存成功: {metrics.id}") async def _save_verification_metrics(self, session: AsyncSession, stp_file_id: int, verification_result: dict): @@ -759,7 +834,8 @@ class ProcessingService: ) session.add(metrics) - await session.commit() + # D9:flush 不 commit,随结果包(阶段 B)由编排层统一提交 + await session.flush() logger.info(f"验证指标保存成功: stp_file_id={stp_file_id}") diff --git a/src/moldinsight/services/storage_integration_rustfs.py b/src/moldinsight/services/storage_integration_rustfs.py index 82341d8..ba7a51c 100644 --- a/src/moldinsight/services/storage_integration_rustfs.py +++ b/src/moldinsight/services/storage_integration_rustfs.py @@ -121,7 +121,8 @@ class StorageIntegrationService: ) session.add(stp_file) - await session.commit() + # D9:仅 flush,与 ProcessingTask 由路由层一并原子提交(避免孤儿文件记录) + await session.flush() await session.refresh(stp_file) logger.info(f"STP文件保存成功 RustFS: {stp_file.id}, 批次: {batch_id}") @@ -134,8 +135,11 @@ class StorageIntegrationService: stp_file_id: int, task_type: str = "stp_parsing", parameters: Optional[Dict[str, Any]] = None, + batch_id: Optional[str] = None, ) -> ProcessingTask: - """创建处理任务记录""" + """创建处理任务记录(D9:仅 flush 不 commit,事务由调用方收口—— + 与 STPFile 记录同批提交,避免留下无任务的孤儿文件记录;batch_id 用于批量任务聚合查询) + """ try: task = ProcessingTask( task_id=task_id, @@ -144,15 +148,15 @@ class StorageIntegrationService: status="pending", started_time=datetime.now(), parameters=parameters or {}, + batch_id=batch_id, ) - + session.add(task) - await session.commit() - await session.refresh(task) - + await session.flush() + logger.info(f"处理任务创建成功: {task_id}") return task - + except Exception as e: await session.rollback() logger.error(f"创建处理任务失败: {e}") @@ -167,7 +171,8 @@ class StorageIntegrationService: current_step: Optional[str] = None, error_message: Optional[str] = None ): - """更新任务状态""" + """更新任务状态(保留即时 commit:进度/状态需跨事务对外可见, + 处理链路中的各阶段进度依赖它落库——D9 收口仅针对数据本体写方法)""" try: update_data = { "status": status, @@ -200,7 +205,7 @@ class StorageIntegrationService: task_id: str, parameters: Dict[str, Any], ): - """合并更新任务参数,便于保存阶段耗时等元数据。""" + """合并更新任务参数,便于保存阶段耗时等元数据。(D9:flush 不 commit,事务由调用方收口)""" try: task = await session.execute( select(ProcessingTask).where(ProcessingTask.task_id == task_id) @@ -212,14 +217,14 @@ class StorageIntegrationService: merged = dict(task.parameters or {}) merged.update(parameters or {}) task.parameters = merged - await session.commit() + await session.flush() except Exception as e: await session.rollback() logger.error(f"更新任务参数失败: {e}") raise async def update_stp_file_status(self, session: AsyncSession, stp_file_id: int, status: str): - """更新STP文件状态""" + """更新STP文件状态(保留即时 commit,理由同 update_task_status)""" try: await session.execute( update(STPFile) @@ -280,7 +285,8 @@ class StorageIntegrationService: ) session.add(geometry_data) - await session.commit() + # D9:数据本体仅 flush,与网格等同阶段数据由编排层统一 commit(原子落库) + await session.flush() await session.refresh(geometry_data) logger.info(f"几何数据保存成功 RustFS: {geometry_data.id}") @@ -336,7 +342,8 @@ class StorageIntegrationService: ) session.add(mesh_data) - await session.commit() + # D9:数据本体仅 flush,与几何数据同阶段由编排层统一 commit + await session.flush() await session.refresh(mesh_data) logger.info(f"网格数据保存成功 RustFS: {mesh_data.id}") @@ -417,7 +424,8 @@ class StorageIntegrationService: ) session.add(mold_cavity) - await session.commit() + # D9:数据本体仅 flush,型腔/HTML/特征同属结果包,由编排层统一 commit + await session.flush() await session.refresh(mold_cavity) logger.info(f"模具型腔数据保存成功 RustFS: {mold_cavity.id}") @@ -466,7 +474,8 @@ class StorageIntegrationService: ) session.add(html_file) - await session.commit() + # D9:数据本体仅 flush,型腔/HTML/特征同属结果包,由编排层统一 commit + await session.flush() await session.refresh(html_file) logger.info(f"HTML文件保存成功 RustFS: {html_file.id}") @@ -504,7 +513,8 @@ class StorageIntegrationService: ) session.add(rec_record) - await session.commit() + # D9:数据本体仅 flush,型腔/HTML/特征同属结果包,由编排层统一 commit + await session.flush() logger.info(f"保存了 {len(features)} 个特征和 {len(recommendations)} 个建议") async def log_user_activity(self, session: AsyncSession, @@ -798,7 +808,8 @@ class StorageIntegrationService: .where(STPFile.id == stp_file_id) .values(**update_data) ) - await session.commit() + # D9:flush 不 commit,随结果包由编排层统一提交 + await session.flush() logger.info(f"STP文件分析摘要更新: ID {stp_file_id}") except Exception as e: diff --git a/src/moldinsight/services/task_dispatcher.py b/src/moldinsight/services/task_dispatcher.py index 8534c88..8bafbce 100644 --- a/src/moldinsight/services/task_dispatcher.py +++ b/src/moldinsight/services/task_dispatcher.py @@ -28,23 +28,28 @@ _background_tasks: set = set() _dispatch_semaphore = asyncio.Semaphore(2) -async def _run_with_limit(task_id: str, file_path: str, stp_file_id: int, process_params: dict): +async def _run_with_limit(task_id: str, stp_file_id: int, process_params: dict): async with _dispatch_semaphore: from moldinsight.services.processing_service import processing_service await processing_service.process_file_with_storage( - task_id, file_path, stp_file_id, process_params + task_id, stp_file_id, process_params ) -def dispatch_processing(task_id: str, file_path: str, stp_file_id: int, process_params: dict): - """调度 STP 处理任务:优先 Celery(进程隔离),否则 API 进程内 asyncio 后台执行。""" +def dispatch_processing(task_id: str, stp_file_id: int, process_params: dict): + """调度 STP 处理任务:优先 Celery(进程隔离),否则 API 进程内 asyncio 后台执行。 + + 入参只传 stp_file_id(D6):源文件由处理方按 PG 元数据从 RustFS 获取, + 不再跨进程传节点本地路径——API 与 Celery worker 容器文件系统不互通, + 传路径在容器化部署下必然失败。 + """ if _use_celery: - process_stp_task.delay(task_id, file_path, stp_file_id, process_params) + process_stp_task.delay(task_id, stp_file_id, process_params) logger.info(f"[DISPATCH] Celery 任务已调度: task_id={task_id}") return task = asyncio.create_task( - _run_with_limit(task_id, file_path, stp_file_id, process_params) + _run_with_limit(task_id, stp_file_id, process_params) ) _background_tasks.add(task) task.add_done_callback(_background_tasks.discard) diff --git a/src/moldinsight/services/task_query_service.py b/src/moldinsight/services/task_query_service.py index 4125a57..4edceee 100644 --- a/src/moldinsight/services/task_query_service.py +++ b/src/moldinsight/services/task_query_service.py @@ -5,6 +5,7 @@ import time from collections import OrderedDict from typing import Optional, Dict, Any, List, Tuple +from fastapi import HTTPException from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.orm import joinedload @@ -55,6 +56,30 @@ class TaskQueryService: """任务 parameters 被更新后调用(export-mold / cam 等),使缓存视图失效。""" cls._view_cache.pop(task_id, None) + @staticmethod + async def ensure_task_access( + db_session: AsyncSession, task_id: str, user_id: int + ) -> "Tuple[ProcessingTask, STPFile]": + """校验任务存在且属于指定用户:不存在 404,他人/无主任务 403(无主不等于公共)。 + + task_router(状态查询)与 advanced_router(导出/倒扣检测等)共用, + 之前只有 advanced_router 有一份私有实现,/api/status 曾因此漏鉴权。 + """ + row = await db_session.execute( + select(ProcessingTask, STPFile) + .join(STPFile, ProcessingTask.stp_file_id == STPFile.id) + .where(ProcessingTask.task_id == task_id) + ) + row = row.first() + if not row: + raise HTTPException(404, "任务不存在") + + _, stp_file = row + if getattr(stp_file, "user_id", None) != user_id: + raise HTTPException(403, "无权访问该任务") + + return row + @staticmethod async def get_task_view(db_session: AsyncSession, task_id: str) -> Optional[Dict[str, Any]]: """ @@ -128,9 +153,17 @@ class TaskQueryService: cam_preferences = processing_task.parameters.get("cam_preferences", {}) or {} task_parameters = dict(processing_task.parameters) + # 注意:analysis_metrics 键可能存在但值为 None(storage 未上传指标时), + # .get(key, {}) 的默认值对 None 不生效,必须用 or {} 兜底 + analysis_metrics = file_with_data.get("analysis_metrics") or {} + task_view = { "task_id": processing_task.task_id, "status": processing_task.status, + # D7:PG 是单一事实源——Redis 不可用时本视图即前端拿到的完整状态, + # 进度字段必须从 PG 补齐(Redis 路径的 task dict 也会带同名字段) + "progress": processing_task.progress or 0, + "current_step": processing_task.current_step, "filename": stp_file.original_filename if stp_file else "", "file_path": stp_file.file_path or "", "file_size": stp_file.file_size if stp_file else 0, @@ -154,18 +187,18 @@ class TaskQueryService: "export_artifacts": task_parameters.get("export_artifacts"), "stage_timings": task_parameters.get("stage_timings", {}), "verification": task_parameters.get("verification") - or file_with_data.get("analysis_metrics", {}).get("verification_details"), + or analysis_metrics.get("verification_details"), "llm_report": task_parameters.get("llm_report"), "analysis_result": { "geometry_data": geometry_json, "detected_features": features_json, "design_recommendations": recommendations_json, "quality_metrics": { - "volume_utilization": file_with_data.get("analysis_metrics", {}).get("volume_utilization", 0), - "topology_complexity": file_with_data.get("analysis_metrics", {}).get("topology_complexity", 0), - "wall_uniformity": file_with_data.get("analysis_metrics", {}).get("wall_uniformity", 0) + "volume_utilization": analysis_metrics.get("volume_utilization", 0), + "topology_complexity": analysis_metrics.get("topology_complexity", 0), + "wall_uniformity": analysis_metrics.get("wall_uniformity", 0) }, - "analysis_summary": file_with_data.get("analysis_metrics", {}).get("analysis_summary", "分析完成") + "analysis_summary": analysis_metrics.get("analysis_summary", "分析完成") } if geometry_json or features_json or recommendations_json else None, "error": processing_task.error_message or stp_file.error_message or None, } diff --git a/src/moldinsight/storage/rustfs_storage.py b/src/moldinsight/storage/rustfs_storage.py index c20c903..11059cc 100644 --- a/src/moldinsight/storage/rustfs_storage.py +++ b/src/moldinsight/storage/rustfs_storage.py @@ -36,6 +36,20 @@ class RustFSManager: async def connect(self, endpoint: str, access_key: str, secret_key: str, timeout: int = 30): """连接到 RustFS 服务""" + # 配置缺失时给出明确错误(settings 不再给占位默认值) + missing = [ + name for name, value in ( + ("RUSTFS_ENDPOINT", endpoint), + ("RUSTFS_ACCESS_KEY", access_key), + ("RUSTFS_SECRET_KEY", secret_key), + ) if not value + ] + if missing: + self.is_connected = False + raise ValueError( + f"RustFS 配置缺失: {', '.join(missing)}(请参照 .env.example 配置后重启)" + ) + try: # 提取端口号和主机 from urllib.parse import urlparse diff --git a/src/shared/config/settings.py b/src/shared/config/settings.py index d9ec3fd..c8f0f61 100644 --- a/src/shared/config/settings.py +++ b/src/shared/config/settings.py @@ -23,9 +23,11 @@ class Settings: self.MESH_QUALITY = os.getenv("MESH_QUALITY", "high") self.PARALLEL_PROCESSING = os.getenv("PARALLEL_PROCESSING", "true").lower() == "true" - self.RUSTFS_ENDPOINT = os.getenv("RUSTFS_ENDPOINT") or os.getenv("MINIO_ENDPOINT") or "http://localhost:8080" - self.RUSTFS_ACCESS_KEY = os.getenv("RUSTFS_ACCESS_KEY") or os.getenv("MINIO_ACCESS_KEY") or "your-access-key" - self.RUSTFS_SECRET_KEY = os.getenv("RUSTFS_SECRET_KEY") or os.getenv("MINIO_SECRET_KEY") or "your-secret-key" + # RUSTFS_* 不给代码兜底默认值(含 MINIO_* 兼容别名): + # 缺失时由 rustfs_storage.connect 抛出明确配置错误,而不是拿占位口令连库 + self.RUSTFS_ENDPOINT = os.getenv("RUSTFS_ENDPOINT") or os.getenv("MINIO_ENDPOINT") + self.RUSTFS_ACCESS_KEY = os.getenv("RUSTFS_ACCESS_KEY") or os.getenv("MINIO_ACCESS_KEY") + self.RUSTFS_SECRET_KEY = os.getenv("RUSTFS_SECRET_KEY") or os.getenv("MINIO_SECRET_KEY") self.RUSTFS_TIMEOUT = int(os.getenv("RUSTFS_TIMEOUT", "30")) self.RUSTFS_PRESIGNED_URL_EXPIRES = int(os.getenv("RUSTFS_PRESIGNED_URL_EXPIRES", "3600")) @@ -36,6 +38,10 @@ class Settings: self.DB_USER = os.getenv("DB_USER") self.DB_PASSWORD = os.getenv("DB_PASSWORD") + # 启动时是否自动执行 alembic 迁移(D12):多副本同时启动会并发迁移, + # 生产多副本应设 false,改由部署流程单点执行 alembic CLI 或本模块 __main__ + self.AUTO_MIGRATE = os.getenv("AUTO_MIGRATE", "true").lower() == "true" + self.SECRET_KEY = os.getenv("SECRET_KEY") self.ALGORITHM = os.getenv("ALGORITHM", "HS256") self.ACCESS_TOKEN_EXPIRE_MINUTES = int(os.getenv("ACCESS_TOKEN_EXPIRE_MINUTES", "1440")) diff --git a/src/shared/database/init_db.py b/src/shared/database/init_db.py index c24e65b..d584818 100644 --- a/src/shared/database/init_db.py +++ b/src/shared/database/init_db.py @@ -118,6 +118,13 @@ async def init_roles(session, perm_map): async def create_admin_user(session): """创建默认管理员""" + # compose 不再给 ADMIN_PASSWORD 弱默认(D14):缺失时显式失败, + # 而不是静默创建空口令管理员 + if not settings.ADMIN_PASSWORD: + raise RuntimeError( + "ADMIN_PASSWORD 未配置:请在 .env 中设置管理员初始密码后重启" + ) + result = await session.execute(select(User).where(User.username == settings.ADMIN_USERNAME)) existing_admin = result.scalar_one_or_none() @@ -150,7 +157,13 @@ async def init_database(keep_connected: bool = True): """初始化数据库""" try: await db_manager.connect() - await _run_alembic_migrations() + if settings.AUTO_MIGRATE: + await _run_alembic_migrations() + else: + logger.info( + "AUTO_MIGRATE=false:跳过启动期 alembic 迁移," + "schema 由部署流程单点执行(alembic CLI 或 python -m shared.database.init_db)" + ) async with db_manager.session() as session: perm_map = await init_permissions(session) diff --git a/src/shared/models/database.py b/src/shared/models/database.py index 12b6d69..7ea9e3b 100644 --- a/src/shared/models/database.py +++ b/src/shared/models/database.py @@ -282,6 +282,10 @@ class ProcessingTask(Base): id = Column(Integer, primary_key=True, index=True) task_id = Column(String(36), unique=True, index=True, nullable=False) stp_file_id = Column(Integer, ForeignKey("stp_files.id"), nullable=False, index=True) + + # 批量上传聚合 ID(批次 2:批量元数据入库——PG 为单一事实源, + # 同批任务经此列聚合查询,不再依赖 Redis/进程内存存批量元数据) + batch_id = Column(String(36), nullable=True, index=True) # 任务类型和状态 task_type = Column(String(50), default="stp_parsing") # stp_parsing, geometry_analysis, mold_generation diff --git a/src/shared/services/auth_routes.py b/src/shared/services/auth_routes.py index 5460fbe..118eeb8 100644 --- a/src/shared/services/auth_routes.py +++ b/src/shared/services/auth_routes.py @@ -212,10 +212,15 @@ async def create_user( if existing_email.scalar_one_or_none(): raise HTTPException(status_code=400, detail="邮箱已存在") + try: + hashed_password = get_password_hash(user_data.password) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + user = User( username=user_data.username, email=user_data.email, - hashed_password=get_password_hash(user_data.password), + hashed_password=hashed_password, full_name=user_data.full_name, is_active=True ) @@ -313,7 +318,10 @@ async def reset_user_password( if not user: raise HTTPException(status_code=404, detail="用户不存在") - user.hashed_password = get_password_hash(new_password) + try: + user.hashed_password = get_password_hash(new_password) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc await db_session.commit() logger.info(f"管理员 {current_user.username} 重置了用户 {user.username} 的密码") diff --git a/src/shared/services/auth_service.py b/src/shared/services/auth_service.py index f782e67..8970812 100644 --- a/src/shared/services/auth_service.py +++ b/src/shared/services/auth_service.py @@ -20,13 +20,28 @@ pwd_context = bcrypt oauth2_scheme = OAuth2PasswordBearer(tokenUrl="/api/auth/login", auto_error=False) +def _require_secret_key() -> str: + """SECRET_KEY 惰性校验:未配置时给出明确错误,而不是让 jwt.encode/decode 报晦涩 TypeError。""" + if not settings.SECRET_KEY: + raise RuntimeError("SECRET_KEY 未配置:请在 .env 中设置后重启服务(认证功能不可用)") + return settings.SECRET_KEY + + def verify_password(plain_password: str, hashed_password: str) -> bool: - return pwd_context.checkpw(plain_password.encode('utf-8'), hashed_password.encode('utf-8')) + # 比较侧按 bcrypt 语义截断到 72 字节:兼容历史上被截断存储的口令, + # 且避免 checkpw 对超长输入直接抛 ValueError(登录会变 500); + # 新口令的超长拒绝在 get_password_hash 中完成 + password_bytes = plain_password.encode('utf-8')[:72] + try: + return pwd_context.checkpw(password_bytes, hashed_password.encode('utf-8')) + except ValueError: + return False def get_password_hash(password: str) -> str: + # bcrypt 算法上限 72 字节:超长密码必须显式拒绝,静默截断会改变有效密码 if len(password.encode('utf-8')) > 72: - password = password[:72] + raise ValueError("密码长度超过 72 字节限制,请使用更短的密码") return pwd_context.hashpw(password.encode('utf-8'), pwd_context.gensalt()).decode('utf-8') @@ -37,7 +52,7 @@ def create_access_token(data: dict, expires_delta: Optional[timedelta] = None) - else: expire = datetime.utcnow() + timedelta(minutes=settings.ACCESS_TOKEN_EXPIRE_MINUTES) to_encode.update({"exp": expire}) - encoded_jwt = jwt.encode(to_encode, settings.SECRET_KEY, algorithm=settings.ALGORITHM) + encoded_jwt = jwt.encode(to_encode, _require_secret_key(), algorithm=settings.ALGORITHM) return encoded_jwt @@ -49,7 +64,7 @@ async def get_current_user( return None try: - payload = jwt.decode(token, settings.SECRET_KEY, algorithms=[settings.ALGORITHM]) + payload = jwt.decode(token, _require_secret_key(), algorithms=[settings.ALGORITHM]) username: str = payload.get("sub") if username is None: logger.warning(f"[AUTH] Token 中缺少 sub 字段") diff --git a/src/shared/services/redis_task_manager.py b/src/shared/services/redis_task_manager.py index 27657ff..e4a6bbc 100644 --- a/src/shared/services/redis_task_manager.py +++ b/src/shared/services/redis_task_manager.py @@ -1,11 +1,14 @@ # services/redis_task_manager.py -"""Redis 任务管理器 - 替代内存字典,支持 TTL 自动清理。 +"""Redis 任务管理器 - 任务状态热缓存(D7:不再有进程内存回退)。 存储格式:Redis Hash(field -> JSON 字符串)。 - update_task 走 HSET 字段级原子更新,消除旧 get->merge->set 三步竞态 (后台处理流程与导出端点并发写同一任务时丢更新); - 进度 tick 只重写变化字段,不再全量重写整个任务 blob; -- 兼容读旧 string 格式(升级前写入的在途任务),新写入一律 Hash。 +- 兼容读旧 string 格式(升级前写入的在途任务),新写入一律 Hash; +- **PG 是任务状态单一事实源**:Redis 不可用时本管理器不再降级进程内 dict + (多副本下各进程内存互相不可见,造成同一任务不同副本读到不同状态), + 而是 no-op / 返回 None——状态查询路径(TaskQueryService)自然落到 PG。 """ import json @@ -101,24 +104,6 @@ class RedisTaskManager: raise RuntimeError("Redis 未连接,无法直接访问 redis_client") return self._redis - # ---- 内存回退 ---- - _fallback_tasks: Dict[str, Dict[str, Any]] = {} - - def _fallback_set(self, task_id: str, data: Dict[str, Any]): - self._fallback_tasks[task_id] = data - - def _fallback_get(self, task_id: str) -> Optional[Dict[str, Any]]: - return self._fallback_tasks.get(task_id) - - def _fallback_delete(self, task_id: str): - self._fallback_tasks.pop(task_id, None) - - def _fallback_all(self) -> Dict[str, Dict[str, Any]]: - return dict(self._fallback_tasks) - - def _fallback_count(self) -> int: - return len(self._fallback_tasks) - # ---- 内部工具 ---- def _key(self, task_id: str) -> str: @@ -161,143 +146,116 @@ class RedisTaskManager: # ---- 公共接口 ---- async def set_task(self, task_id: str, data: Dict[str, Any], ttl: Optional[int] = None): - """整包写入任务数据(Hash,覆盖旧值,含旧 string 格式清理)""" + """整包写入任务数据(Hash,覆盖旧值,含旧 string 格式清理)。 + + Redis 不可用时 no-op:任务状态事实源在 PG,缓存缺失不影响正确性。 + """ + if not self.is_connected: + return + effective_ttl = ttl or self._ttl mapping = self._dump_mapping(data) - - if self.is_connected: - try: - key = self._key(task_id) - # DEL 先清掉可能存在的旧 string/Hash,保证覆盖语义 - pipe = self._redis.pipeline() - pipe.delete(key) - pipe.hset(key, mapping=mapping) - pipe.expire(key, effective_ttl) - await pipe.execute() - return - except Exception as e: - logger.warning(f"Redis 写入失败,回退到内存: {e}") - - self._fallback_set(task_id, self._make_serializable(data)) + try: + key = self._key(task_id) + # DEL 先清掉可能存在的旧 string/Hash,保证覆盖语义 + pipe = self._redis.pipeline() + pipe.delete(key) + pipe.hset(key, mapping=mapping) + pipe.expire(key, effective_ttl) + await pipe.execute() + except Exception as e: + logger.warning(f"Redis 写入失败(任务状态以 PG 为准): task={task_id}, {e}") async def get_task(self, task_id: str) -> Optional[Dict[str, Any]]: - """获取任务数据(Hash / 旧 string 兼容)""" - if self.is_connected: - try: - return await self._load_any(self._key(task_id)) - except Exception as e: - logger.warning(f"Redis 读取失败,回退到内存: {e}") + """获取任务数据(Hash / 旧 string 兼容)。 - return self._fallback_get(task_id) + Redis 不可用 / 未命中返回 None,调用方落到 PG 路径。 + """ + if not self.is_connected: + return None + try: + return await self._load_any(self._key(task_id)) + except Exception as e: + logger.warning(f"Redis 读取失败(任务状态以 PG 为准): task={task_id}, {e}") + return None async def update_task(self, task_id: str, updates: Dict[str, Any]): """字段级原子更新(HSET),无读改写竞态。 兼容旧 string 格式:先迁移为 Hash 再更新。 + Redis 不可用时 no-op(状态事实源在 PG)。 """ - mapping = self._dump_mapping(updates) - - if self.is_connected: - try: - key = self._key(task_id) - key_type = await self._redis.type(key) - - if key_type == "none": - logger.warning(f"任务 {task_id} 不存在,无法更新") - return - - if key_type == "string": - # 旧格式迁移:string -> Hash - legacy = await self._redis.get(key) - try: - base = json.loads(legacy) if legacy else {} - except json.JSONDecodeError: - base = {} - base.update(mapping) - pipe = self._redis.pipeline() - pipe.delete(key) - pipe.hset(key, mapping=self._dump_mapping(base)) - pipe.expire(key, self._ttl) - await pipe.execute() - return - - await self._redis.hset(key, mapping=mapping) - await self._redis.expire(key, self._ttl) - return - except Exception as e: - logger.warning(f"Redis 更新失败,回退到内存: {e}") - - # 内存回退保持读改写语义(单进程内存无并发竞态) - current = self._fallback_get(task_id) - if current is None: - logger.warning(f"任务 {task_id} 不存在,无法更新") + if not self.is_connected: return - current.update(self._make_serializable(updates)) - self._fallback_set(task_id, current) + mapping = self._dump_mapping(updates) + try: + key = self._key(task_id) + key_type = await self._redis.type(key) + + if key_type == "none": + logger.warning(f"任务 {task_id} 不存在,无法更新") + return + + if key_type == "string": + # 旧格式迁移:string -> Hash + legacy = await self._redis.get(key) + try: + base = json.loads(legacy) if legacy else {} + except json.JSONDecodeError: + base = {} + base.update(mapping) + pipe = self._redis.pipeline() + pipe.delete(key) + pipe.hset(key, mapping=self._dump_mapping(base)) + pipe.expire(key, self._ttl) + await pipe.execute() + return + + await self._redis.hset(key, mapping=mapping) + await self._redis.expire(key, self._ttl) + except Exception as e: + logger.warning(f"Redis 更新失败(任务状态以 PG 为准): task={task_id}, {e}") async def delete_task(self, task_id: str): - """删除任务(DEL 对 Hash/string 均有效)""" - if self.is_connected: - try: - await self._redis.delete(self._key(task_id)) - return - except Exception as e: - logger.warning(f"Redis 删除失败,回退到内存: {e}") - - self._fallback_delete(task_id) + """删除任务(DEL 对 Hash/string 均有效);Redis 不可用时 no-op""" + if not self.is_connected: + return + try: + await self._redis.delete(self._key(task_id)) + except Exception as e: + logger.warning(f"Redis 删除失败: task={task_id}, {e}") async def get_all_tasks(self) -> Dict[str, Dict[str, Any]]: - """获取所有任务""" - if self.is_connected: - try: - pattern = f"{self._prefix}*" - result = {} - async for key in self._redis.scan_iter(match=pattern): - task_id = key.replace(self._prefix, "") - task = await self._load_any(key) - if task: - result[task_id] = task - return result - except Exception as e: - logger.warning(f"Redis 扫描失败,回退到内存: {e}") - - return self._fallback_all() + """获取所有任务;Redis 不可用时返回空 dict(调用方需容忍)""" + if not self.is_connected: + return {} + try: + pattern = f"{self._prefix}*" + result = {} + async for key in self._redis.scan_iter(match=pattern): + task_id = key.replace(self._prefix, "") + task = await self._load_any(key) + if task: + result[task_id] = task + return result + except Exception as e: + logger.warning(f"Redis 扫描失败: {e}") + return {} async def get_task_count(self) -> int: - """获取任务总数""" - if self.is_connected: - try: - pattern = f"{self._prefix}*" - count = 0 - async for _ in self._redis.scan_iter(match=pattern): - count += 1 - return count - except Exception as e: - logger.warning(f"Redis 计数失败,回退到内存: {e}") - - return self._fallback_count() - - async def cleanup_old_tasks(self, max_age_seconds: int = 86400 * 7): - """清理过期任务(Redis 由 TTL 自动管理,内存回退需手动清理)""" - now = datetime.now() - to_delete = [] - - for task_id, task in self._fallback_tasks.items(): - completed_at = task.get("completed_at") - if completed_at: - try: - completed_dt = datetime.fromisoformat(completed_at) - if (now - completed_dt).total_seconds() > max_age_seconds: - to_delete.append(task_id) - except (ValueError, TypeError): - pass - - for task_id in to_delete: - del self._fallback_tasks[task_id] - - if to_delete: - logger.info(f"清理了 {len(to_delete)} 个过期内存任务") + """获取任务总数;Redis 不可用时返回 0(调用方需容忍)""" + if not self.is_connected: + return 0 + try: + pattern = f"{self._prefix}*" + count = 0 + async for _ in self._redis.scan_iter(match=pattern): + count += 1 + return count + except Exception as e: + logger.warning(f"Redis 计数失败: {e}") + return 0 # ---- 工具方法 ---- diff --git a/tests/test_batch_status_pg.py b/tests/test_batch_status_pg.py new file mode 100644 index 0000000..422a2be --- /dev/null +++ b/tests/test_batch_status_pg.py @@ -0,0 +1,126 @@ +"""批次 2(D7)回归测试:批量任务聚合查询以 PG 为单一事实源。 + +覆盖: +- GET /api/batch/{batch_id} 按 ProcessingTask.batch_id 聚合(此前依赖 Redis batch key + 进程内存降级) +- 归属校验:他人批次 403、不存在 404、所有者 200 + 聚合数字正确 +""" +import pytest +from fastapi import FastAPI +from httpx import AsyncClient, ASGITransport +from sqlalchemy.ext.asyncio import async_sessionmaker, AsyncSession + +from moldinsight.api.batch_router import router as batch_router +from shared.database.database import get_db_session +from shared.models.database import User, STPFile, ProcessingTask +from shared.services.auth_service import get_current_active_user + + +@pytest.fixture(scope="function") +async def batch_client(async_engine, seeded_db): + """带 batch_router 的测试应用:播种 batch-1(user 1,completed+processing)与 batch-2(user 999)。""" + session_factory = async_sessionmaker(async_engine, class_=AsyncSession, expire_on_commit=False) + + async with session_factory() as session: + stp_a = STPFile( + id=9101, user_id=1, object_key="test/batch-a.step", storage_bucket="moldinsight", + original_filename="batch-a.step", file_size=128, status="completed", + ) + stp_b = STPFile( + id=9102, user_id=1, object_key="test/batch-b.step", storage_bucket="moldinsight", + original_filename="batch-b.step", file_size=128, status="processing", + ) + stp_c = STPFile( + id=9103, user_id=999, object_key="test/batch-c.step", storage_bucket="moldinsight", + original_filename="batch-c.step", file_size=128, status="completed", + ) + task_a = ProcessingTask( + task_id="task-batch-a", stp_file_id=9101, batch_id="batch-1", + task_type="stp_parsing", status="completed", progress=100, parameters={}, + ) + task_b = ProcessingTask( + task_id="task-batch-b", stp_file_id=9102, batch_id="batch-1", + task_type="stp_parsing", status="processing", progress=40, + current_step="生成模具型腔", parameters={}, + ) + task_c = ProcessingTask( + task_id="task-batch-c", stp_file_id=9103, batch_id="batch-2", + task_type="stp_parsing", status="failed", progress=40, + error_message="处理超时", parameters={}, + ) + session.add_all([stp_a, stp_b, stp_c, task_a, task_b, task_c]) + await session.commit() + + test_app = FastAPI() + test_app.include_router(batch_router, prefix="/api") + + async def override_get_db_session(): + async with session_factory() as session: + yield session + + test_app.dependency_overrides[get_db_session] = override_get_db_session + + transport = ASGITransport(app=test_app) + async with AsyncClient(transport=transport, base_url="http://test") as ac: + yield ac, test_app + + test_app.dependency_overrides.clear() + + +def _override_user(app: FastAPI, user_id: int): + app.dependency_overrides[get_current_active_user] = lambda: User(id=user_id, username="tester") + + +@pytest.mark.asyncio +async def test_batch_status_owner_aggregates_from_pg(batch_client): + """所有者查询:聚合数字与逐任务字段来自 PG,而非 Redis。""" + ac, app = batch_client + _override_user(app, 1) + + resp = await ac.get("/api/batch/batch-1") + assert resp.status_code == 200 + body = resp.json() + + assert body["batch_id"] == "batch-1" + assert body["total"] == 2 + assert body["completed"] == 1 + assert body["processing"] == 1 + assert body["failed"] == 0 + assert body["progress_percent"] == 50.0 + + by_id = {t["task_id"]: t for t in body["tasks"]} + assert by_id["task-batch-a"]["status"] == "completed" + assert by_id["task-batch-a"]["filename"] == "batch-a.step" + assert by_id["task-batch-b"]["progress"] == 40 + assert by_id["task-batch-b"]["current_step"] == "生成模具型腔" + + +@pytest.mark.asyncio +async def test_batch_status_non_owner_is_403(batch_client): + """他人批次必须 403(按 STPFile.user_id 校验,无主不等于公共)。""" + ac, app = batch_client + _override_user(app, 2) + + resp = await ac.get("/api/batch/batch-1") + assert resp.status_code == 403 + + +@pytest.mark.asyncio +async def test_batch_status_unknown_is_404(batch_client): + ac, app = batch_client + _override_user(app, 1) + + resp = await ac.get("/api/batch/batch-missing") + assert resp.status_code == 404 + + +@pytest.mark.asyncio +async def test_batch_status_includes_error_from_pg(batch_client): + """failed 任务的 error_message 经 PG 返回(此前 error 只存在于 Redis task dict)。""" + ac, app = batch_client + _override_user(app, 999) + + resp = await ac.get("/api/batch/batch-2") + assert resp.status_code == 200 + body = resp.json() + assert body["failed"] == 1 + assert body["tasks"][0]["error"] == "处理超时" diff --git a/tests/test_deployment_config.py b/tests/test_deployment_config.py new file mode 100644 index 0000000..869a099 --- /dev/null +++ b/tests/test_deployment_config.py @@ -0,0 +1,37 @@ +"""批次 1 部署治理回归测试。 + +覆盖: +- AUTO_MIGRATE 开关:env 解析与默认值(默认 true 保持现行启动行为) +- create_admin_user:ADMIN_PASSWORD 未配置时显式报错,不创建空口令管理员 +""" +import pytest +from unittest.mock import MagicMock + +from shared.config.settings import Settings + + +def test_auto_migrate_default_true(monkeypatch): + """未配置时默认 true:保持既有单机开发行为(启动即迁移)。""" + monkeypatch.delenv("AUTO_MIGRATE", raising=False) + assert Settings().AUTO_MIGRATE is True + + +def test_auto_migrate_env_parsing(monkeypatch): + monkeypatch.setenv("AUTO_MIGRATE", "false") + assert Settings().AUTO_MIGRATE is False + monkeypatch.setenv("AUTO_MIGRATE", "true") + assert Settings().AUTO_MIGRATE is True + + +@pytest.mark.asyncio +async def test_create_admin_without_password_is_explicit_error(monkeypatch): + """ADMIN_PASSWORD 缺失必须显式失败(compose 已去弱默认),而非静默创建空口令管理员。""" + # init_db 顶层 import alembic;未安装 alembic 的环境跳过(与 OCC 测试同策略) + pytest.importorskip("alembic.config") + from shared.config.settings import settings + from shared.database.init_db import create_admin_user + + monkeypatch.setattr(settings, "ADMIN_PASSWORD", None) + # 校验发生在任何 DB 访问之前,MagicMock 会话不会被触碰 + with pytest.raises(RuntimeError, match="ADMIN_PASSWORD"): + await create_admin_user(MagicMock()) diff --git a/tests/test_redis_no_fallback.py b/tests/test_redis_no_fallback.py new file mode 100644 index 0000000..435041b --- /dev/null +++ b/tests/test_redis_no_fallback.py @@ -0,0 +1,57 @@ +"""批次 2(D7 / D9)回归测试。 + +- D7:Redis 任务管理器不再有进程内存回退——Redis 不可用时写 no-op、读返回 None, + 状态查询路径落到 PG(PG 为单一事实源) +- D9:存储服务的数据本体写方法只 flush 不 commit,事务由编排层收口 +""" +import pytest +from unittest.mock import AsyncMock, MagicMock + +from shared.services.redis_task_manager import RedisTaskManager + + +@pytest.mark.asyncio +async def test_no_memory_fallback_when_disconnected(): + """Redis 未连接:写 no-op、读 None——不得再出现进程内可见的副本。""" + mgr = RedisTaskManager() # 不 connect + assert not mgr.is_connected + + await mgr.set_task("t1", {"status": "processing"}) + assert await mgr.get_task("t1") is None + + await mgr.update_task("t1", {"status": "completed"}) # no-op,不得抛异常 + + assert await mgr.get_all_tasks() == {} + assert await mgr.get_task_count() == 0 + await mgr.delete_task("t1") # no-op,不得抛异常 + + assert await mgr.get_task("t1") is None + + +def test_fallback_storage_removed(): + """防回归:内存回退存储必须已删除,防止静默回归。""" + assert not hasattr(RedisTaskManager, "_fallback_tasks") + assert not hasattr(RedisTaskManager, "_fallback_set") + assert not hasattr(RedisTaskManager, "cleanup_old_tasks") + + +@pytest.mark.asyncio +async def test_storage_writes_flush_but_never_commit(): + """D9:数据本体写方法仅 flush;commit 由编排层/请求侧负责。""" + pytest.importorskip("minio") + from sqlalchemy.ext.asyncio import AsyncSession + from moldinsight.services.storage_integration_rustfs import StorageIntegrationService + + svc = StorageIntegrationService() + session = AsyncMock(spec=AsyncSession) + # update_task_parameters:select 返回 None(任务不存在)→ 直接 return + result_mock = MagicMock() + result_mock.scalar_one_or_none.return_value = None + session.execute.return_value = result_mock + + await svc.update_task_parameters(session, "task-x", {"a": 1}) + await svc.create_processing_task(session, "task-y", stp_file_id=1, batch_id="b-1") + + assert session.flush.await_count >= 1 + # 注意:不能用 .awaited(AsyncMock 上访问会自动创建 truthy 子 mock),用 await_count + assert session.commit.await_count == 0, "存储写方法不得自行 commit(D9 事务收口)" diff --git a/tests/test_status_endpoint_auth.py b/tests/test_status_endpoint_auth.py new file mode 100644 index 0000000..008d7de --- /dev/null +++ b/tests/test_status_endpoint_auth.py @@ -0,0 +1,139 @@ +"""批次 0 安全修复回归测试(TECH_DEBT D5 等)。 + +覆盖: +- /api/status/{task_id} 鉴权与任务归属校验(无 token 401 / 他人任务 403 / 不存在 404 / 所有者 200) +- bcrypt 72 字节上限:超长密码显式拒绝而非静默截断 +- SECRET_KEY 惰性校验:未配置时给出明确错误 +""" +import pytest +from fastapi import FastAPI +from httpx import AsyncClient, ASGITransport +from sqlalchemy.ext.asyncio import async_sessionmaker, AsyncSession + +from moldinsight.api.task_router import router as task_router +from shared.database.database import get_db_session +from shared.models.database import User, STPFile, ProcessingTask +from shared.services.auth_service import ( + get_current_active_user, + get_password_hash, + verify_password, + create_access_token, +) +from shared.config.settings import settings + + +@pytest.fixture(scope="function") +async def status_client(async_engine, seeded_db): + """带 task_router 的测试应用:播种两个任务(owner=user 1 / user 999)。""" + session_factory = async_sessionmaker(async_engine, class_=AsyncSession, expire_on_commit=False) + + async with session_factory() as session: + stp_own = STPFile( + id=9001, user_id=1, object_key="test/own.step", storage_bucket="moldinsight", + original_filename="own.step", file_size=128, status="completed", + ) + task_own = ProcessingTask( + task_id="task-owned", stp_file_id=9001, + task_type="stp_parsing", status="completed", parameters={}, + ) + stp_other = STPFile( + id=9002, user_id=999, object_key="test/other.step", storage_bucket="moldinsight", + original_filename="other.step", file_size=128, status="completed", + ) + task_other = ProcessingTask( + task_id="task-foreign", stp_file_id=9002, + task_type="stp_parsing", status="completed", parameters={}, + ) + session.add_all([stp_own, task_own, stp_other, task_other]) + await session.commit() + + test_app = FastAPI() + # 生产环境中 /api 前缀由入口层挂载时添加,测试中需显式指定才能对齐真实路径 + test_app.include_router(task_router, prefix="/api") + + async def override_get_db_session(): + async with session_factory() as session: + yield session + + test_app.dependency_overrides[get_db_session] = override_get_db_session + + transport = ASGITransport(app=test_app) + async with AsyncClient(transport=transport, base_url="http://test") as ac: + yield ac, test_app + + test_app.dependency_overrides.clear() + + +@pytest.mark.asyncio +async def test_status_without_token_is_401(status_client): + """未携带 token 访问 /api/status 必须拒绝(此前该端点完全未鉴权)。""" + ac, _ = status_client + resp = await ac.post("/api/status/task-owned") + assert resp.status_code == 401 + + +@pytest.mark.asyncio +async def test_status_owner_can_view(status_client): + ac, app = status_client + app.dependency_overrides[get_current_active_user] = lambda: User(id=1, username="tester") + + resp = await ac.post("/api/status/task-owned") + assert resp.status_code == 200 + body = resp.json() + assert body["task_id"] == "task-owned" + assert body["status"] == "completed" + + +@pytest.mark.asyncio +async def test_status_non_owner_is_403(status_client): + """他人任务(含无主任务)必须 403,不允许凭任务号枚举。""" + ac, app = status_client + app.dependency_overrides[get_current_active_user] = lambda: User(id=2, username="intruder") + + resp = await ac.post("/api/status/task-owned") + assert resp.status_code == 403 + + resp = await ac.post("/api/status/task-foreign") + assert resp.status_code == 403 + + +@pytest.mark.asyncio +async def test_status_unknown_task_is_404(status_client): + ac, app = status_client + app.dependency_overrides[get_current_active_user] = lambda: User(id=1, username="tester") + + resp = await ac.post("/api/status/task-missing") + assert resp.status_code == 404 + + +def test_password_hash_rejects_over_72_bytes(): + """bcrypt 72 字节上限:必须显式报错,不能静默截断。""" + with pytest.raises(ValueError, match="72"): + get_password_hash("a" * 73) + # 多字节字符按字节数计:25 个汉字 = 75 字节 + with pytest.raises(ValueError, match="72"): + get_password_hash("模" * 25) + + +def test_password_hash_roundtrip_at_limit(): + password = "a" * 72 + hashed = get_password_hash(password) + assert verify_password(password, hashed) + # 历史口令按 bcrypt 语义截断比较:73 字节输入截断后与 72 字节口令匹配(兼容旧数据), + # 但新口令在 get_password_hash 处已被显式拒绝,不会再产生这类哈希 + assert verify_password(password + "x", hashed) + assert not verify_password("b" * 72, hashed) + + +def test_verify_password_over_72_bytes_returns_false_not_raise(): + """超长密码登录不得抛 ValueError(否则登录接口 500),应返回 False 走正常失败路径。""" + hashed = get_password_hash("short-password") + assert not verify_password("长" * 40, hashed) # 120 字节 + assert not verify_password("a" * 73, hashed) + + +def test_create_token_without_secret_key_is_explicit_error(monkeypatch): + """SECRET_KEY 未配置时给出可读错误,而不是 jwt.encode 的晦涩 TypeError。""" + monkeypatch.setattr(settings, "SECRET_KEY", None) + with pytest.raises(RuntimeError, match="SECRET_KEY"): + create_access_token({"sub": "tester"})