From f8a5269284e5090c12f0fd0ceca59cb4c2946a7d Mon Sep 17 00:00:00 2001 From: Nixevol Date: Tue, 19 May 2026 16:08:14 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E4=BC=98=E5=8C=96=E5=A4=84=E7=90=86?= =?UTF-8?q?=E8=BF=9B=E5=BA=A6=E5=92=8C=E6=97=A5=E5=BF=97=E8=B7=9F=E9=9A=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/api/routers/remote.py | 30 +++++++++-- app/api/routers/tasks.py | 38 ++++++++++++-- app/processor.py | 18 ++++++- docs/project_context.md | 7 +++ frontend/src/components/FileWorkflow.vue | 66 +++++++++++++++++++++--- frontend/src/types.ts | 1 + 6 files changed, 142 insertions(+), 18 deletions(-) diff --git a/app/api/routers/remote.py b/app/api/routers/remote.py index 07432d0..21b6a40 100644 --- a/app/api/routers/remote.py +++ b/app/api/routers/remote.py @@ -38,15 +38,25 @@ async def start_remote_processing(): work_dir.mkdir(parents=True, exist_ok=True) logs: list[str] = [] + current_stage = "downloading" def log_callback(message: str) -> None: logs.append(message) - state.processing_tasks[task_id] = {"logs": logs.copy(), "status": "processing"} + _set_task_stage(task_id, current_stage, logs) - logger = ProcessLogger(log_file=work_dir / "log.txt", callback=log_callback) + def stage_callback(stage: str) -> None: + nonlocal current_stage + current_stage = stage + _set_task_stage(task_id, current_stage, logs) + + logger = ProcessLogger( + log_file=work_dir / "log.txt", + callback=log_callback, + stage_callback=stage_callback, + ) state.history_manager.create(work_dir, 0, record_id=task_id) state.history_manager.update(task_id, status="processing") - state.processing_tasks[task_id] = {"logs": [], "status": "processing"} + state.processing_tasks[task_id] = {"logs": [], "status": "processing", "stage": current_stage} state.global_task_lock.update( { "locked": True, @@ -71,6 +81,16 @@ async def start_remote_processing(): } +def _set_task_stage(task_id: str, stage: str, logs: list[str], status: str = "processing") -> None: + state.processing_tasks[task_id] = { + "logs": logs.copy(), + "status": status, + "stage": stage, + } + if state.global_task_lock["task_id"] == task_id: + state.global_task_lock["stage"] = stage + + def _run_remote_processing( task_id: str, work_dir: Path, @@ -93,7 +113,7 @@ def _run_remote_processing( raise RuntimeError("远程目录中未下载到任何文件") state.history_manager.update(task_id, file_count=download_result.file_count) - state.global_task_lock["stage"] = "processing" + logger.set_stage("extracting") processor = DataProcessor(state.config, work_dir, logger) result = processor.process() @@ -115,6 +135,7 @@ def _run_remote_processing( state.processing_tasks[task_id] = { "logs": state.history_manager.get_logs(task_id), "status": status, + "stage": status, } except Exception as exc: logger.error(f"远程自动化任务失败: {exc}") @@ -122,6 +143,7 @@ def _run_remote_processing( state.processing_tasks[task_id] = { "logs": state.history_manager.get_logs(task_id), "status": "failed", + "stage": "failed", } finally: try: diff --git a/app/api/routers/tasks.py b/app/api/routers/tasks.py index 57c7261..fb487f3 100644 --- a/app/api/routers/tasks.py +++ b/app/api/routers/tasks.py @@ -38,7 +38,7 @@ async def get_global_task_status(): return { "has_active": True, "task_id": task_id, - "stage": "processing", + "stage": active_tasks[task_id].get("stage", "processing"), "logs": active_tasks[task_id].get("logs", []), } @@ -88,14 +88,24 @@ async def start_processing(task_id: str = Body(..., embed=True)): raise HTTPException(status_code=400, detail="工作目录不存在") logs: list[str] = [] + current_stage = "processing" def log_callback(message: str) -> None: logs.append(message) - state.processing_tasks[task_id] = {"logs": logs.copy(), "status": "processing"} + _set_task_stage(task_id, current_stage, logs) - logger = ProcessLogger(log_file=work_dir / "log.txt", callback=log_callback) + def stage_callback(stage: str) -> None: + nonlocal current_stage + current_stage = stage + _set_task_stage(task_id, current_stage, logs) + + logger = ProcessLogger( + log_file=work_dir / "log.txt", + callback=log_callback, + stage_callback=stage_callback, + ) state.history_manager.update(task_id, status="processing") - state.processing_tasks[task_id] = {"logs": [], "status": "processing"} + state.processing_tasks[task_id] = {"logs": [], "status": "processing", "stage": current_stage} state.global_task_lock.update( { "locked": True, @@ -116,7 +126,12 @@ async def get_processing_status(task_id: str = Body(..., embed=True)): if task_id in state.processing_tasks: task_info = state.processing_tasks[task_id] logs = task_info.get("logs") or state.history_manager.get_logs(task_id) - return {"task_id": task_id, "status": task_info["status"], "logs": logs} + return { + "task_id": task_id, + "status": task_info["status"], + "stage": task_info.get("stage"), + "logs": logs, + } record = state.history_manager.get(task_id) if not record: @@ -125,6 +140,7 @@ async def get_processing_status(task_id: str = Body(..., embed=True)): return { "task_id": task_id, "status": record.status, + "stage": record.status, "logs": state.history_manager.get_logs(task_id), "elapsed_time": record.elapsed_time, "error": record.error, @@ -140,6 +156,16 @@ def _task_finished(task_id: str) -> bool: return bool(task_info and task_info.get("status") in {"completed", "failed"}) +def _set_task_stage(task_id: str, stage: str, logs: list[str], status: str = "processing") -> None: + state.processing_tasks[task_id] = { + "logs": logs.copy(), + "status": status, + "stage": stage, + } + if state.global_task_lock["task_id"] == task_id: + state.global_task_lock["stage"] = stage + + def _run_processing(task_id: str, work_dir: Path, logger: ProcessLogger) -> None: try: processor = DataProcessor(state.config, work_dir, logger) @@ -155,12 +181,14 @@ def _run_processing(task_id: str, work_dir: Path, logger: ProcessLogger) -> None state.processing_tasks[task_id] = { "logs": state.history_manager.get_logs(task_id), "status": status, + "stage": status, } except Exception as exc: state.history_manager.update(task_id, status="failed", error=str(exc)) state.processing_tasks[task_id] = { "logs": state.history_manager.get_logs(task_id), "status": "failed", + "stage": "failed", } finally: try: diff --git a/app/processor.py b/app/processor.py index fb2acc8..d9e59d2 100644 --- a/app/processor.py +++ b/app/processor.py @@ -22,10 +22,16 @@ from app.database import DatabaseManager class ProcessLogger: """处理日志记录器""" - def __init__(self, log_file: Optional[Path] = None, callback: Optional[Callable[[str], None]] = None): + def __init__( + self, + log_file: Optional[Path] = None, + callback: Optional[Callable[[str], None]] = None, + stage_callback: Optional[Callable[[str], None]] = None, + ): self.logs: List[str] = [] self.log_file = log_file self.callback = callback + self.stage_callback = stage_callback # 如果指定了日志文件,确保目录存在 if self.log_file: self.log_file.parent.mkdir(parents=True, exist_ok=True) @@ -62,6 +68,10 @@ class ProcessLogger: def success(self, message: str): self.log(message, "SUCCESS") + def set_stage(self, stage: str): + if self.stage_callback: + self.stage_callback(stage) + def get_logs(self) -> List[str]: return self.logs.copy() @@ -195,24 +205,30 @@ class DataProcessor: try: # 1. 解压 ZIP 文件 + self.logger.set_stage("extracting") self._unzip_files() # 2. 处理 Excel 文件(并行) + self.logger.set_stage("converting") self._process_excel_files_parallel() # 3. 处理 CSV 文件并上传到数据库(高性能批量插入) + self.logger.set_stage("importing") self._process_csv_files() # 4. 执行 SQL 脚本 + self.logger.set_stage("scripting") self._execute_sql_script() elapsed = round(time.time() - start_time, 2) + self.logger.set_stage("completed") self.logger.success(f"处理完成!总耗时: {elapsed} 秒") self.results["success"] = True self.results["elapsed_time"] = elapsed except Exception as e: + self.logger.set_stage("failed") self.logger.error(f"处理失败: {str(e)}") self.results["success"] = False self.results["error"] = str(e) diff --git a/docs/project_context.md b/docs/project_context.md index 3e6dad8..7c9cad2 100644 --- a/docs/project_context.md +++ b/docs/project_context.md @@ -1,5 +1,12 @@ # 项目上下文记录 +## 2026-05-19:优化处理进度阶段显示和日志跟随 + +- `ProcessLogger` 新增轻量阶段回调,`DataProcessor.process()` 会在远程下载后依次上报 `extracting`、`converting`、`importing`、`scripting`、`completed/failed` 阶段。 +- `/api/process/status` 和 `/api/task/status` 返回当前 `stage`,本地上传处理和远程下载处理都通过 `state.processing_tasks` 与全局任务锁同步阶段,前端轮询即可实时显示“远程下载中 / 解压数据中 / 上传数据中 / 运行脚本中”等状态。 +- `frontend/src/components/FileWorkflow.vue` 的处理进度卡片新增“保持最新 Log”勾选框,勾选后新日志到达会自动滚动到日志底部;当前任务提示不再显示原始阶段码,改为中文阶段文本。 +- 已执行 `.venv\Scripts\python.exe -m compileall app`、`uvx --offline ruff check .` 和 `npm run build`,均通过。 + ## 2026-05-19:清理后端冗余代码和未用依赖 - `app/database.py` 移除未使用的 SQLAlchemy 连接池、`engine` 属性、`dispose()` 空释放路径和未引用的 `delete_rows()`;数据库访问统一保留现有 PyMySQL 上下文连接。 diff --git a/frontend/src/components/FileWorkflow.vue b/frontend/src/components/FileWorkflow.vue index cd9308f..aa4e7e3 100644 --- a/frontend/src/components/FileWorkflow.vue +++ b/frontend/src/components/FileWorkflow.vue @@ -126,6 +126,7 @@ {{ processStatusText }}
+ 保持最新 Log 刷新 @@ -134,12 +135,12 @@
- 当前任务:{{ activeTask.task_id }} / {{ activeTask.stage || 'processing' }} + 当前任务:{{ activeTask.task_id }} / {{ currentStageText }} {{ taskStatus.error }} -
+
{{ logText }}
@@ -149,7 +150,7 @@