diff --git a/README.md b/README.md index 5010826..a5d0dfe 100644 --- a/README.md +++ b/README.md @@ -6,6 +6,7 @@ CapacityReport 用于导入每周容量报表数据,按 `Configure.json` 的 - Excel/CSV/ZIP 数据导入与自动解压、转换、入库。 - FTP/SFTP 远程数据源配置、连接测试、远程下载并处理。 +- 远程自动调度:按远程 ZIP 文件名日期检查目标自然周 7 天数据,就绪后自动下载并处理。 - MySQL 数据表查看、清空、删除、CSV/XLSX 导出。 - SQL 脚本在线查看、保存和执行。 - 处理历史、日志查看、历史原始数据打包下载。 @@ -219,6 +220,26 @@ Server Portable 和桌面版需要在目标系统原生构建:Windows 包在 W - `SheetFilter`:Excel Sheet 过滤规则。 - `ExtractField`:字段映射配置。 +`RemoteData.auto_scheduler` 用于远程自动调度: + +```json +{ + "enabled": false, + "check_interval_hours": 1, + "expected_directories": ["4G/FDD", "4G/900", "5G/2.6", "5G/700"], + "week_offset": 0 +} +``` + +- `enabled`:是否启用自动调度。 +- `check_interval_hours`:检查间隔,最小 1 小时。 +- `expected_directories`:相对 `remote_dir` 的预期数据目录;为空时按远程 ZIP 实际所在目录检测。 +- `week_offset`:`0` 表示上周自然周,`-1` 表示上上周。 + +自动调度开启后,系统会强制开启 `RemoteData.enabled` 和 `auto_delete_source`。调度器每轮先检查 `cache/auto_scheduler/ready.flag`;如果标识存在则直接触发远程下载并处理。没有标识时,会按 ZIP 文件名中的第一个时间戳判断每个目录是否覆盖目标自然周 7 天;全部就绪后写入标识,下一个检查周期再启动处理。处理成功并完成远程源文件清理后会删除标识;处理失败或源文件清理失败会保留标识,后续自动重试。 + +无论手动上传还是远程下载,处理流程都会按文件名日期对每个目录只保留最近 7 天文件。文件名支持 `XXX_YYYYMMDDHHMM_YYYYMMDDHHMM` 和 `XXX_YYYYMMDDHHMM` 两类格式,数据日期始终取第一个时间戳。 + 登录密码保存在本地 `auth.ini`,该文件不应提交到版本库。 API Token 保存在本地 `api_tokens.json`,包含 HMAC 哈希、显示用前后缀和完整 Token,登录后可在列表中重复复制。该文件属于运行时数据,不应提交到版本库。 @@ -254,6 +275,8 @@ API Token 与登录态一样可访问业务 API,包括文件上传、远程下 - `POST /api/upload`:上传文件 - `POST /api/remote/test`:测试 FTP/SFTP 连接 - `POST /api/remote/start`:远程下载并处理 +- `GET /api/remote/scheduler/status`:查询远程自动调度状态 +- `POST /api/remote/scheduler/trigger`:手动触发一次自动调度检查 - `POST /api/process/start`:启动本地处理 - `POST /api/process/status`:查询处理状态 - `GET /api/license/status`:查询授权状态 diff --git a/app/api/routers/remote.py b/app/api/routers/remote.py index 264ad69..9242c0e 100644 --- a/app/api/routers/remote.py +++ b/app/api/routers/remote.py @@ -1,7 +1,7 @@ -from datetime import datetime +from datetime import date, datetime from pathlib import Path from threading import Thread -from typing import Any +from typing import Any, Callable, Iterable from fastapi import APIRouter, Body, HTTPException @@ -27,6 +27,29 @@ async def test_remote_connection(config: dict[str, Any] | None = Body(None)): @router.post("/api/remote/start") async def start_remote_processing(): + return start_remote_processing_job(source="manual") + + +@router.get("/api/remote/scheduler/status") +def get_scheduler_status(): + if state.auto_scheduler is None: + return {"enabled": False, "running": False, "message": "自动调度器未启动"} + return state.auto_scheduler.get_status() + + +@router.post("/api/remote/scheduler/trigger") +def trigger_scheduler_check(): + if state.auto_scheduler is None: + raise HTTPException(status_code=503, detail="自动调度器未启动") + return state.auto_scheduler.check_and_run(manual=True) + + +def start_remote_processing_job( + *, + source: str = "manual", + on_finish: Callable[[str, str], None] | None = None, + target_dates: Iterable[date] | None = None, +) -> dict[str, Any]: if state.global_task_lock["locked"]: raise HTTPException(status_code=409, detail="已有任务在运行,请等待当前任务完成") @@ -70,14 +93,14 @@ async def start_remote_processing(): thread = Thread( target=_run_remote_processing, - args=(task_id, work_dir, app_config, remote_config, logger), + args=(task_id, work_dir, app_config, remote_config, logger, source, on_finish, target_dates), daemon=True, ) thread.start() return { "success": True, - "message": "远程下载处理任务已启动", + "message": "自动调度远程下载处理任务已启动" if source == "scheduler" else "远程下载处理任务已启动", "task_id": task_id, "stage": "downloading", } @@ -99,14 +122,18 @@ def _run_remote_processing( app_config: AppConfig, remote_config: RemoteDataConfig, logger: ProcessLogger, + source: str = "manual", + on_finish: Callable[[str, str], None] | None = None, + target_dates: Iterable[date] | None = None, ) -> None: + final_status = "failed" try: logger.info( f"开始远程下载,协议: {remote_config.protocol.upper()}," f"服务器: {remote_config.host}:{remote_config.port},目录: {remote_config.remote_dir}" ) downloader = RemoteDataDownloader(remote_config, logger.info) - download_result = downloader.download_to(work_dir) + download_result = downloader.download_to(work_dir, target_dates=target_dates) logger.success( f"远程下载完成,共 {download_result.file_count} 个文件," f"{_format_bytes(download_result.total_bytes)}" @@ -127,8 +154,12 @@ def _run_remote_processing( try: deleted_count = downloader.delete_source_files(download_result.remote_files) logger.success(f"远程源文件清理完成,共删除 {deleted_count} 个文件,目录已保留") + final_status = status except Exception as exc: logger.warning(f"远程源文件清理失败,数据处理结果已保留: {exc}") + final_status = "source_cleanup_failed" + else: + final_status = status state.history_manager.update( task_id, @@ -145,7 +176,10 @@ def _run_remote_processing( } except Exception as exc: error_detail = exc.to_detail() if isinstance(exc, LicenseError) else None - logger.error(f"远程自动化任务失败: {exc}") + if source == "scheduler": + logger.error(f"自动调度任务失败: {exc}") + else: + logger.error(f"远程自动化任务失败: {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), @@ -160,6 +194,8 @@ def _run_remote_processing( except Exception as exc: print(f"自动清理处理历史失败: {exc}") state.reset_task_lock() + if on_finish: + on_finish(task_id, final_status) def _format_bytes(size: int) -> str: diff --git a/app/config.py b/app/config.py index 44f28d9..74c1dbc 100644 --- a/app/config.py +++ b/app/config.py @@ -34,6 +34,59 @@ class MySQLConfig: dbname: str = "CapacityReport" +@dataclass +class AutoSchedulerConfig: + enabled: bool = False + check_interval_hours: int = 1 + expected_directories: List[str] = field(default_factory=list) + week_offset: int = 0 + + def normalized(self) -> "AutoSchedulerConfig": + try: + check_interval_hours = int(self.check_interval_hours) + except (TypeError, ValueError): + check_interval_hours = 1 + try: + week_offset = int(self.week_offset) + except (TypeError, ValueError): + week_offset = 0 + + directories = [] + seen = set() + for directory in self.expected_directories or []: + normalized = str(directory).replace("\\", "/").strip().strip("/") + if normalized and normalized not in seen: + directories.append(normalized) + seen.add(normalized) + + return AutoSchedulerConfig( + enabled=bool(self.enabled), + check_interval_hours=max(check_interval_hours, 1), + expected_directories=directories, + week_offset=week_offset, + ) + + def to_dict(self) -> Dict[str, Any]: + normalized = self.normalized() + return { + "enabled": normalized.enabled, + "check_interval_hours": normalized.check_interval_hours, + "expected_directories": normalized.expected_directories, + "week_offset": normalized.week_offset, + } + + @classmethod + def from_dict(cls, data: Dict[str, Any] | None) -> "AutoSchedulerConfig": + data = data or {} + directories = data.get("expected_directories", []) + return cls( + enabled=bool(data.get("enabled", False)), + check_interval_hours=data.get("check_interval_hours", 1), + expected_directories=directories if isinstance(directories, list) else [], + week_offset=data.get("week_offset", 0), + ).normalized() + + @dataclass class RemoteDataConfig: enabled: bool = False @@ -46,6 +99,7 @@ class RemoteDataConfig: passive: bool = True timeout: int = 30 auto_delete_source: bool = False + auto_scheduler: AutoSchedulerConfig = field(default_factory=AutoSchedulerConfig) def normalized(self) -> "RemoteDataConfig": protocol = self.protocol.lower().strip() @@ -53,8 +107,9 @@ class RemoteDataConfig: protocol = "sftp" port = self.port or (22 if protocol == "sftp" else 21) + scheduler = self.auto_scheduler.normalized() return RemoteDataConfig( - enabled=bool(self.enabled), + enabled=bool(self.enabled) or scheduler.enabled, protocol=protocol, host=self.host.strip(), port=port, @@ -63,7 +118,8 @@ class RemoteDataConfig: remote_dir=(self.remote_dir or "/").strip() or "/", passive=bool(self.passive), timeout=max(int(self.timeout or 30), 1), - auto_delete_source=bool(self.auto_delete_source), + auto_delete_source=bool(self.auto_delete_source) or scheduler.enabled, + auto_scheduler=scheduler, ) def to_dict(self, include_password: bool = False) -> Dict[str, Any]: @@ -77,6 +133,7 @@ class RemoteDataConfig: "passive": self.passive, "timeout": self.timeout, "auto_delete_source": self.auto_delete_source, + "auto_scheduler": self.auto_scheduler.normalized().to_dict(), } if include_password: data["passwd"] = self.passwd @@ -106,6 +163,7 @@ class RemoteDataConfig: passive=bool(data.get("passive", True)), timeout=timeout, auto_delete_source=bool(data.get("auto_delete_source", False)), + auto_scheduler=AutoSchedulerConfig.from_dict(data.get("auto_scheduler")), ).normalized() diff --git a/app/main.py b/app/main.py index 4004569..d4ce670 100644 --- a/app/main.py +++ b/app/main.py @@ -1,4 +1,5 @@ import argparse +from contextlib import asynccontextmanager from pathlib import Path import uvicorn @@ -25,6 +26,7 @@ from app.api.routers import ( from app.auth import extract_access_token, resolve_access_context, resolve_login_context from app.config import BASE_DIR from app.services.api_tokens import touch_token_usage +from app.services.auto_scheduler import AutoScheduler APP_VERSION = "3.0.0" @@ -139,6 +141,14 @@ OPENAPI_OPERATION_DOCS = { "summary": "远程下载并处理", "description": "从已配置的 FTP/SFTP 目录递归下载源数据,然后自动执行完整处理流程。", }, + ("get", "/api/remote/scheduler/status"): { + "summary": "查询远程自动调度状态", + "description": "返回自动调度启用状态、目标周、就绪标识、下次检查时间和各远程目录的日期覆盖情况。", + }, + ("post", "/api/remote/scheduler/trigger"): { + "summary": "手动触发自动调度检查", + "description": "立即执行一次远程目录就绪检查;如果已存在就绪标识,会直接触发远程下载并处理。", + }, ("post", "/api/history"): { "summary": "查询处理历史", "description": "按最近时间返回处理历史记录。", @@ -249,6 +259,12 @@ OPENAPI_OPERATION_DOCS = { "passive": True, "timeout": 30, "auto_delete_source": False, + "auto_scheduler": { + "enabled": False, + "check_interval_hours": 1, + "expected_directories": ["4G/FDD", "4G/900", "5G/2.6", "5G/700"], + "week_offset": 0, + }, }, }, ("post", "/api/config/history-retention"): { @@ -312,6 +328,17 @@ OPENAPI_OPERATION_DOCS = { } +@asynccontextmanager +async def app_lifespan(app: FastAPI): + state.auto_scheduler = AutoScheduler() + state.auto_scheduler.start() + try: + yield + finally: + if state.auto_scheduler is not None: + state.auto_scheduler.stop() + + def create_app() -> FastAPI: app = FastAPI( title="CapacityReport", @@ -320,6 +347,7 @@ def create_app() -> FastAPI: docs_url=None, redoc_url=None, openapi_url=None, + lifespan=app_lifespan, ) app.openapi = lambda: custom_openapi(app) # type: ignore[method-assign] diff --git a/app/processor.py b/app/processor.py index 95802a0..a88a46f 100644 --- a/app/processor.py +++ b/app/processor.py @@ -18,6 +18,7 @@ from concurrent.futures import ThreadPoolExecutor, as_completed from app.config import AppConfig, SQL_SCRIPT from app.database import DatabaseManager +from app.utils.file_dates import DirectoryDateSelection, select_recent_items_by_directory class ProcessLogger: @@ -246,7 +247,7 @@ class DataProcessor: def _unzip_files(self): """解压所有 ZIP 文件(支持中文文件名)""" self.logger.info("正在解压 ZIP 文件...") - zip_files = list(self.work_dir.rglob("*.zip")) + zip_files = self._filter_recent_files(list(self.work_dir.rglob("*.zip")), "ZIP") zip_count = 0 for zip_file in zip_files: @@ -338,6 +339,43 @@ class DataProcessor: for ext in extensions: for file in directory.rglob(f"*{ext}"): yield file + + def _filter_recent_files(self, files: list[Path], label: str, root: Path | None = None) -> list[Path]: + if not files: + return files + + base = (root or self.work_dir).resolve() + + def parent_key(file_path: Path) -> str: + try: + parent = file_path.parent.resolve().relative_to(base) + except ValueError: + parent = file_path.parent + parent_text = str(parent).replace("\\", "/") + return "" if parent_text == "." else parent_text + + selected, summaries = select_recent_items_by_directory( + files, + parent_key=parent_key, + name_key=lambda file_path: file_path.name, + ) + self._log_recent_file_selection(label, summaries) + return sorted(selected) + + def _log_recent_file_selection(self, label: str, summaries: list[DirectoryDateSelection]) -> None: + skipped_total = sum(summary.skipped_count for summary in summaries) + if not skipped_total: + return + + for summary in summaries: + if not summary.skipped_count or not summary.start_date or not summary.max_date: + continue + self.logger.info( + f"{label}目录 {summary.directory or '.'}: 仅处理 " + f"{summary.start_date.isoformat()} 至 {summary.max_date.isoformat()} " + f"的 {summary.selected_count}/{summary.total_count} 个文件," + f"跳过 {summary.skipped_count} 个旧文件" + ) def _process_single_excel(self, excel_file: Path, sheet_filter: set) -> int: """处理单个 Excel 文件(用于并行)""" @@ -368,7 +406,7 @@ class DataProcessor: def _process_excel_files_parallel(self): """并行处理 Excel 文件""" self.logger.info("正在并行处理 Excel 文件...") - excel_files = list(self._scan_files(self.work_dir, ['.xlsx', '.xls'])) + excel_files = self._filter_recent_files(list(self._scan_files(self.work_dir, ['.xlsx', '.xls'])), "Excel") self.logger.info(f"找到 {len(excel_files)} 个 Excel 文件") if not excel_files: @@ -794,7 +832,7 @@ class DataProcessor: self.db.drop_table(table_name) # 处理该目录下的所有 CSV - csv_files = list(self._scan_files(subdir, ['.csv'])) + csv_files = self._filter_recent_files(list(self._scan_files(subdir, ['.csv'])), "CSV", root=subdir) self.logger.info(f"找到 {len(csv_files)} 个 CSV 文件") total_rows = 0 diff --git a/app/services/auto_scheduler.py b/app/services/auto_scheduler.py new file mode 100644 index 0000000..735d960 --- /dev/null +++ b/app/services/auto_scheduler.py @@ -0,0 +1,408 @@ +from __future__ import annotations + +import json +import threading +from collections import defaultdict +from dataclasses import dataclass +from datetime import date, datetime, timedelta +from typing import Any + +from fastapi import HTTPException + +from app import state +from app.config import CACHE_DIR, AutoSchedulerConfig, RemoteDataConfig +from app.services.remote_download import RemoteDataDownloader, RemoteFileInfo +from app.utils.file_dates import extract_file_date, required_week_days + + +READY_DIR = CACHE_DIR / "auto_scheduler" +READY_FLAG = READY_DIR / "ready.flag" +DISABLED_CHECK_SECONDS = 60 +STARTUP_CHECK_SECONDS = 5 +FAILURE_RESULTS = {"scan_failed", "trigger_failed", "source_cleanup_failed", "failed", "invalid_flag"} + + +@dataclass(frozen=True) +class DirectoryReadyStatus: + directory: str + ready: bool + found_days: list[date] + missing_days: list[date] + file_count: int + error: str | None = None + + def to_dict(self) -> dict[str, Any]: + return { + "ready": self.ready, + "found_days": [item.isoformat() for item in self.found_days], + "missing_days": [item.isoformat() for item in self.missing_days], + "found_count": len(self.found_days), + "required_count": len(self.found_days) + len(self.missing_days), + "file_count": self.file_count, + "error": self.error, + } + + +class AutoScheduler: + def __init__(self) -> None: + self._stop_event = threading.Event() + self._check_lock = threading.Lock() + self._status_lock = threading.Lock() + self._thread: threading.Thread | None = None + self._running = False + self._last_check_at: datetime | None = None + self._next_check_at: datetime | None = None + self._last_result = "not_started" + self._last_message = "自动调度器尚未检查" + self._directory_status: dict[str, dict[str, Any]] = {} + self._task_running = False + self._task_id: str | None = None + self._failure_count = 0 + + def start(self) -> None: + if self._thread and self._thread.is_alive(): + return + self._running = True + self._stop_event.clear() + self._set_next_check(datetime.now() + timedelta(seconds=STARTUP_CHECK_SECONDS)) + self._thread = threading.Thread(target=self._run_loop, name="auto-scheduler", daemon=True) + self._thread.start() + + def stop(self) -> None: + self._running = False + self._stop_event.set() + if self._thread and self._thread.is_alive(): + self._thread.join(timeout=5) + + def get_status(self) -> dict[str, Any]: + config = state.current_config().remote_data.normalized() + scheduler = config.auto_scheduler.normalized() + target_days = required_week_days(scheduler.week_offset) + ready_flag = self._read_ready_flag() + + with self._status_lock: + return { + "enabled": scheduler.enabled, + "running": self._running, + "check_interval_hours": scheduler.check_interval_hours, + "expected_directories": scheduler.expected_directories, + "week_offset": scheduler.week_offset, + "auto_delete_source": config.auto_delete_source, + "next_check_at": self._format_dt(self._next_check_at), + "last_check_at": self._format_dt(self._last_check_at), + "last_result": self._last_result, + "last_message": self._last_message, + "failure_count": self._failure_count, + "task_running": self._task_running, + "task_id": self._task_id, + "ready_flag": ready_flag, + "target_week": { + "start": target_days[0].isoformat(), + "end": target_days[-1].isoformat(), + "days": [item.isoformat() for item in target_days], + }, + "directory_status": self._directory_status, + } + + def check_and_run(self, manual: bool = False) -> dict[str, Any]: + if not self._check_lock.acquire(blocking=False): + return self._finish_check("busy", "自动调度器正在检查中", manual) + + try: + return self._check_and_run_locked(manual) + finally: + self._check_lock.release() + + def _run_loop(self) -> None: + while not self._stop_event.wait(self._seconds_until_next_check()): + self.check_and_run(manual=False) + + def _check_and_run_locked(self, manual: bool) -> dict[str, Any]: + now = datetime.now() + remote_config = state.current_config().remote_data.normalized() + scheduler = remote_config.auto_scheduler.normalized() + self._last_check_at = now + + if not scheduler.enabled: + self._set_next_check(now + timedelta(seconds=DISABLED_CHECK_SECONDS)) + return self._finish_check("disabled", "自动调度未启用", manual) + + self._set_next_check(now + timedelta(hours=scheduler.check_interval_hours)) + + if not remote_config.enabled: + return self._finish_check("remote_disabled", "远程数据源未启用", manual) + + if state.global_task_lock["locked"]: + return self._finish_check("task_running", "已有任务在运行,本轮自动调度跳过", manual) + + ready_flag = self._read_ready_flag() + if ready_flag["exists"]: + if ready_flag.get("invalid"): + self._clear_ready_flag() + return self._finish_check("invalid_flag", "就绪标识格式错误,已清除,本轮不触发处理", manual) + target_dates = self._target_dates_from_flag(ready_flag, scheduler) + return self._trigger_processing(target_dates, ready_flag, manual) + + target_days = required_week_days(scheduler.week_offset) + directory_status = self._check_remote_ready(remote_config, scheduler, target_days) + ready = bool(directory_status) and all(item.ready for item in directory_status.values()) + self._set_directory_status(directory_status) + + error_count = sum(1 for item in directory_status.values() if item.error) + if error_count: + return self._finish_check( + "scan_failed", + f"远程目录扫描失败,{error_count}/{len(directory_status)} 个目录无法访问或扫描失败", + manual, + ) + + if ready: + self._mark_ready(target_days, directory_status) + return self._finish_check( + "marked_ready", + "远程数据已满足目标周 7 天,已写入就绪标识,下次检查将自动处理", + manual, + ) + + ready_count = sum(1 for item in directory_status.values() if item.ready) + return self._finish_check( + "waiting", + f"远程数据未就绪,{ready_count}/{len(directory_status)} 个目录满足目标周 7 天", + manual, + ) + + def _check_remote_ready( + self, + remote_config: RemoteDataConfig, + scheduler: AutoSchedulerConfig, + target_days: list[date], + ) -> dict[str, DirectoryReadyStatus]: + downloader = RemoteDataDownloader(remote_config) + expected_directories = scheduler.expected_directories + if expected_directories: + return { + directory: self._directory_ready_status( + directory, + self._safe_list_remote_zip_files(downloader, directory), + target_days, + ) + for directory in expected_directories + } + + files = self._safe_list_remote_zip_files(downloader, None) + if files is None: + return { + ".": DirectoryReadyStatus( + directory=".", + ready=False, + found_days=[], + missing_days=target_days, + file_count=0, + error="远程目录扫描失败", + ) + } + + grouped: dict[str, list[RemoteFileInfo]] = defaultdict(list) + for remote_file in files: + grouped[remote_file.parent or "."].append(remote_file) + + if not grouped: + return { + ".": DirectoryReadyStatus( + directory=".", + ready=False, + found_days=[], + missing_days=target_days, + file_count=0, + error="远程目录未找到 ZIP 文件", + ) + } + + return { + directory: self._directory_ready_status(directory, directory_files, target_days) + for directory, directory_files in sorted(grouped.items(), key=lambda item: item[0]) + } + + def _safe_list_remote_zip_files( + self, + downloader: RemoteDataDownloader, + directory: str | None, + ) -> list[RemoteFileInfo] | None: + try: + return downloader.list_remote_zip_files(directory) + except Exception: + return None + + def _directory_ready_status( + self, + directory: str, + files: list[RemoteFileInfo] | None, + target_days: list[date], + ) -> DirectoryReadyStatus: + if files is None: + return DirectoryReadyStatus( + directory=directory, + ready=False, + found_days=[], + missing_days=target_days, + file_count=0, + error="远程目录不存在或无法访问", + ) + + required = set(target_days) + found = { + file_date + for remote_file in files + if (file_date := extract_file_date(remote_file.name)) in required + } + missing = [item for item in target_days if item not in found] + return DirectoryReadyStatus( + directory=directory, + ready=not missing, + found_days=sorted(found), + missing_days=missing, + file_count=len(files), + ) + + def _trigger_processing( + self, + target_dates: list[date], + ready_flag: dict[str, Any], + manual: bool, + ) -> dict[str, Any]: + from app.api.routers.remote import start_remote_processing_job + + try: + result = start_remote_processing_job( + source="scheduler", + on_finish=self._on_processing_finish, + target_dates=target_dates, + ) + except HTTPException as exc: + return self._finish_check("trigger_failed", str(exc.detail), manual) + except Exception as exc: + return self._finish_check("trigger_failed", f"自动调度触发失败: {exc}", manual) + + task_id = result.get("task_id") + with self._status_lock: + self._task_running = True + self._task_id = str(task_id) if task_id else None + self._last_result = "triggered" + self._last_message = "已根据就绪标识触发远程下载并处理" + return { + "success": True, + "result": "triggered", + "message": "已根据就绪标识触发远程下载并处理", + "task_id": task_id, + "ready_flag": ready_flag, + "status": self.get_status(), + } + + def _mark_ready( + self, + target_days: list[date], + directory_status: dict[str, DirectoryReadyStatus], + ) -> None: + READY_DIR.mkdir(parents=True, exist_ok=True) + payload = { + "ready_at": datetime.now().isoformat(timespec="seconds"), + "week_start": target_days[0].isoformat(), + "week_end": target_days[-1].isoformat(), + "target_dates": [item.isoformat() for item in target_days], + "directories": { + name: item.to_dict() + for name, item in directory_status.items() + }, + } + READY_FLAG.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8") + + def _clear_ready_flag(self) -> None: + try: + READY_FLAG.unlink(missing_ok=True) + except OSError: + pass + + def _read_ready_flag(self) -> dict[str, Any]: + if not READY_FLAG.exists(): + return {"exists": False} + try: + data = json.loads(READY_FLAG.read_text(encoding="utf-8")) + if isinstance(data, dict): + return {"exists": True, **data} + except Exception as exc: + return {"exists": True, "invalid": True, "error": str(exc)} + return {"exists": True, "invalid": True, "error": "ready.flag 格式错误"} + + def _target_dates_from_flag( + self, + ready_flag: dict[str, Any], + scheduler: AutoSchedulerConfig, + ) -> list[date]: + raw_dates = ready_flag.get("target_dates") + if isinstance(raw_dates, list): + dates: list[date] = [] + for raw_date in raw_dates: + try: + dates.append(date.fromisoformat(str(raw_date))) + except ValueError: + continue + if dates: + return sorted(dates) + return required_week_days(scheduler.week_offset) + + def _on_processing_finish(self, task_id: str, status: str) -> None: + if status == "completed": + self._clear_ready_flag() + message = "自动调度任务处理成功,已清除就绪标识" + result = "completed" + else: + message = "自动调度任务未完成,就绪标识已保留,下次检查会重试" + result = status or "failed" + + with self._status_lock: + self._task_running = False + self._task_id = task_id + self._last_result = result + self._last_message = message + self._failure_count = 0 if result == "completed" else self._failure_count + 1 + + def _finish_check(self, result: str, message: str, manual: bool) -> dict[str, Any]: + with self._status_lock: + self._last_result = result + self._last_message = message + if result in FAILURE_RESULTS: + self._failure_count += 1 + elif result not in {"busy", "task_running"}: + self._failure_count = 0 + if result not in {"triggered", "task_running"}: + self._task_running = False + + return { + "success": result not in FAILURE_RESULTS, + "manual": manual, + "result": result, + "message": message, + "status": self.get_status(), + } + + def _set_directory_status(self, directory_status: dict[str, DirectoryReadyStatus]) -> None: + with self._status_lock: + self._directory_status = { + name: item.to_dict() + for name, item in directory_status.items() + } + + def _set_next_check(self, value: datetime) -> None: + with self._status_lock: + self._next_check_at = value + + def _seconds_until_next_check(self) -> float: + with self._status_lock: + next_check_at = self._next_check_at + if not next_check_at: + return STARTUP_CHECK_SECONDS + return max((next_check_at - datetime.now()).total_seconds(), 0) + + @staticmethod + def _format_dt(value: datetime | None) -> str | None: + return value.isoformat(timespec="seconds") if value else None diff --git a/app/services/remote_download.py b/app/services/remote_download.py index a489b95..f10102a 100644 --- a/app/services/remote_download.py +++ b/app/services/remote_download.py @@ -1,12 +1,15 @@ from __future__ import annotations import stat +from collections import defaultdict from dataclasses import dataclass, field +from datetime import date from ftplib import FTP from pathlib import Path from typing import Callable, Iterable from app.config import RemoteDataConfig +from app.utils.file_dates import extract_file_date, select_recent_items_by_directory LogFn = Callable[[str], None] @@ -19,6 +22,15 @@ class RemoteDownloadResult: remote_files: list[str] = field(default_factory=list) +@dataclass(frozen=True) +class RemoteFileInfo: + path: str + relative_path: str + parent: str + name: str + size: int = 0 + + class RemoteDownloadError(RuntimeError): pass @@ -42,13 +54,71 @@ class RemoteDataDownloader: finally: ssh.close() - def download_to(self, destination: Path) -> RemoteDownloadResult: + def download_to(self, destination: Path, target_dates: Iterable[date] | None = None) -> RemoteDownloadResult: self._validate_config() destination.mkdir(parents=True, exist_ok=True) + zip_files = self.list_remote_zip_files() + date_filter = set(target_dates or []) + if date_filter: + selected_files = self._select_files_by_dates(zip_files, date_filter) + return self._download_selected_files(destination, selected_files) + + if zip_files: + selected_files, summaries = select_recent_items_by_directory( + zip_files, + parent_key=lambda item: item.parent, + name_key=lambda item: item.name, + ) + for summary in summaries: + if summary.max_date and summary.start_date and summary.skipped_count: + self._log( + f"远程目录 {summary.directory or '.'}: 仅下载 " + f"{summary.start_date.isoformat()} 至 {summary.max_date.isoformat()} " + f"的 {summary.selected_count}/{summary.total_count} 个 ZIP 文件," + f"跳过 {summary.skipped_count} 个旧文件" + ) + return self._download_selected_files(destination, selected_files) + if self.config.protocol == "ftp": return self._download_ftp(destination) return self._download_sftp(destination) + def _select_files_by_dates(self, zip_files: list[RemoteFileInfo], target_dates: set[date]) -> list[RemoteFileInfo]: + selected_files: list[RemoteFileInfo] = [] + grouped: dict[str, list[RemoteFileInfo]] = defaultdict(list) + for remote_file in zip_files: + grouped[remote_file.parent].append(remote_file) + + for parent, files in sorted(grouped.items(), key=lambda item: item[0]): + selected = [ + remote_file + for remote_file in files + if extract_file_date(remote_file.name) in target_dates + ] + selected_files.extend(selected) + skipped_count = len(files) - len(selected) + if skipped_count: + first_day = min(target_dates).isoformat() + last_day = max(target_dates).isoformat() + self._log( + f"远程目录 {parent or '.'}: 仅下载调度目标日期 " + f"{first_day} 至 {last_day} 的 {len(selected)}/{len(files)} 个 ZIP 文件," + f"跳过 {skipped_count} 个非目标日期文件" + ) + + return selected_files + + def list_remote_zip_files(self, directory: str | None = None) -> list[RemoteFileInfo]: + self._validate_config() + remote_dir = ( + self._join_remote_path(self.config.remote_dir, directory.strip("/")) + if directory + else self.config.remote_dir + ) + if self.config.protocol == "ftp": + return self._list_ftp_zip_files(remote_dir) + return self._list_sftp_zip_files(remote_dir) + def delete_source_files(self, remote_files: Iterable[str] | None = None) -> int: self._validate_config() source_files = list(remote_files or []) @@ -164,6 +234,59 @@ class RemoteDataDownloader: result.total_bytes += expected_size or local_path.stat().st_size result.remote_files.append(remote_path) + def _download_selected_files(self, destination: Path, remote_files: list[RemoteFileInfo]) -> RemoteDownloadResult: + result = RemoteDownloadResult() + if not remote_files: + return result + + if self.config.protocol == "ftp": + with self._ftp_client() as ftp: + self._log(f"已连接 FTP: {self.config.host}:{self.config.port}") + for remote_file in remote_files: + self._download_ftp_file( + ftp, + remote_file.path, + destination / remote_file.relative_path, + result, + remote_file.size, + ) + return result + + ssh = self._sftp_ssh_client() + try: + with ssh.open_sftp() as sftp: + self._log(f"已连接 SFTP: {self.config.host}:{self.config.port}") + for remote_file in remote_files: + self._download_sftp_file( + sftp, + remote_file.path, + destination / remote_file.relative_path, + result, + remote_file.size, + ) + finally: + ssh.close() + return result + + def _list_ftp_zip_files(self, remote_dir: str) -> list[RemoteFileInfo]: + with self._ftp_client() as ftp: + return self._collect_ftp_zip_files(ftp, remote_dir) + + def _collect_ftp_zip_files(self, ftp: FTP, remote_dir: str) -> list[RemoteFileInfo]: + files: list[RemoteFileInfo] = [] + for name, entry_type, size in self._list_ftp_entries(ftp, remote_dir): + if name in {".", ".."}: + continue + + remote_path = self._join_remote_path(remote_dir, name) + if entry_type == "dir" or (entry_type == "unknown" and self._ftp_is_dir(ftp, remote_path)): + files.extend(self._collect_ftp_zip_files(ftp, remote_path)) + continue + + if name.lower().endswith(".zip"): + files.append(self._remote_file_info(remote_path, size)) + return files + def _delete_ftp_source_files(self) -> int: with self._ftp_client() as ftp: self._log(f"开始清理 FTP 源文件: {self.config.remote_dir}") @@ -246,13 +369,50 @@ class RemoteDataDownloader: ) return + self._download_sftp_file(sftp, remote_path, local_path, result, int(getattr(attrs, "st_size", 0) or 0)) + + def _download_sftp_file( + self, + sftp, + remote_path: str, + local_path: Path, + result: RemoteDownloadResult, + expected_size: int = 0, + ) -> None: local_path.parent.mkdir(parents=True, exist_ok=True) self._log(f"下载: {remote_path}") sftp.get(remote_path, str(local_path)) result.file_count += 1 - result.total_bytes += int(getattr(attrs, "st_size", 0) or local_path.stat().st_size) + result.total_bytes += expected_size or local_path.stat().st_size result.remote_files.append(remote_path) + def _list_sftp_zip_files(self, remote_dir: str) -> list[RemoteFileInfo]: + ssh = self._sftp_ssh_client() + try: + with ssh.open_sftp() as sftp: + return self._collect_sftp_zip_files(sftp, remote_dir) + finally: + ssh.close() + + def _collect_sftp_zip_files(self, sftp, remote_path: str) -> list[RemoteFileInfo]: + attrs = sftp.stat(remote_path) + if not stat.S_ISDIR(attrs.st_mode): + name = Path(remote_path.replace("\\", "/")).name + if name.lower().endswith(".zip"): + return [self._remote_file_info(remote_path, int(getattr(attrs, "st_size", 0) or 0))] + return [] + + files: list[RemoteFileInfo] = [] + for item in sftp.listdir_attr(remote_path): + if item.filename in {".", ".."}: + continue + child_path = self._join_remote_path(remote_path, item.filename) + if stat.S_ISDIR(item.st_mode): + files.extend(self._collect_sftp_zip_files(sftp, child_path)) + elif item.filename.lower().endswith(".zip"): + files.append(self._remote_file_info(child_path, int(getattr(item, "st_size", 0) or 0))) + return files + def _delete_sftp_source_files(self) -> int: ssh = self._sftp_ssh_client() try: @@ -301,3 +461,22 @@ class RemoteDataDownloader: if parent == "/": return f"/{child}" return f"{parent}/{child}" + + def _remote_file_info(self, remote_path: str, size: int = 0) -> RemoteFileInfo: + normalized_root = self.config.remote_dir.replace("\\", "/").rstrip("/") + normalized_path = remote_path.replace("\\", "/") + if normalized_root and normalized_root != "/" and normalized_path.startswith(f"{normalized_root}/"): + relative_path = normalized_path[len(normalized_root) + 1 :] + else: + relative_path = normalized_path.lstrip("/") + relative = Path(relative_path) + parent = str(relative.parent).replace("\\", "/") + if parent == ".": + parent = "" + return RemoteFileInfo( + path=remote_path, + relative_path=relative_path, + parent=parent, + name=relative.name, + size=size, + ) diff --git a/app/state.py b/app/state.py index c78d684..f02cf6b 100644 --- a/app/state.py +++ b/app/state.py @@ -8,6 +8,7 @@ config = AppConfig.load() history_manager = HistoryManager() processing_tasks: dict[str, dict[str, Any]] = {} upload_sessions: dict[str, dict[str, Any]] = {} +auto_scheduler: Any | None = None global_task_lock: dict[str, Any] = { "locked": False, diff --git a/app/utils/file_dates.py b/app/utils/file_dates.py new file mode 100644 index 0000000..ed1534c --- /dev/null +++ b/app/utils/file_dates.py @@ -0,0 +1,108 @@ +"""Filename date parsing and per-directory date filtering.""" +from __future__ import annotations + +import re +from dataclasses import dataclass +from datetime import date, datetime, timedelta +from pathlib import Path +from typing import Callable, Iterable, TypeVar + + +FILENAME_DATE_RE = re.compile(r"(? date | None: + """Extract the first timestamp date from a filename.""" + match = FILENAME_DATE_RE.search(Path(str(filename)).name) + if not match: + return None + return parse_file_timestamp(match.group(1)) + + +def parse_file_timestamp(value: str) -> date | None: + fmt = "%Y%m%d%H%M%S" if len(value) == 14 else "%Y%m%d%H%M" + try: + return datetime.strptime(value, fmt).date() + except ValueError: + return None + + +def required_week_days(week_offset: int = 0, today: date | None = None) -> list[date]: + """Return Monday-Sunday dates for the target natural week. + + week_offset=0 means last week, -1 means the week before last. + """ + base_day = today or date.today() + this_monday = base_day - timedelta(days=base_day.weekday()) + target_monday = this_monday - timedelta(days=7 * (1 - week_offset)) + return [target_monday + timedelta(days=offset) for offset in range(7)] + + +def select_recent_items_by_directory( + items: Iterable[T], + *, + parent_key: Callable[[T], str], + name_key: Callable[[T], str], + days: int = 7, +) -> tuple[list[T], list[DirectoryDateSelection]]: + """Select the latest N natural days of dated files in each directory. + + If a directory has no dated files at all, all files in that directory are + kept to preserve compatibility with legacy/manual data. + """ + safe_days = max(int(days), 1) + groups: dict[str, list[T]] = {} + for item in items: + groups.setdefault(parent_key(item), []).append(item) + + selected: list[T] = [] + summaries: list[DirectoryDateSelection] = [] + for directory, group_items in sorted(groups.items(), key=lambda pair: pair[0]): + dated_items = [(item, extract_file_date(name_key(item))) for item in group_items] + valid_dates = [file_date for _, file_date in dated_items if file_date] + if not valid_dates: + selected.extend(group_items) + summaries.append( + DirectoryDateSelection( + directory=directory, + total_count=len(group_items), + selected_count=len(group_items), + skipped_count=0, + undated_count=len(group_items), + ) + ) + continue + + max_date = max(valid_dates) + start_date = max_date - timedelta(days=safe_days - 1) + group_selected = [ + item + for item, file_date in dated_items + if file_date is not None and start_date <= file_date <= max_date + ] + selected.extend(group_selected) + summaries.append( + DirectoryDateSelection( + directory=directory, + total_count=len(group_items), + selected_count=len(group_selected), + skipped_count=len(group_items) - len(group_selected), + undated_count=sum(1 for _, file_date in dated_items if file_date is None), + max_date=max_date, + start_date=start_date, + ) + ) + + return selected, summaries diff --git a/docs/auto_scheduler_plan.md b/docs/auto_scheduler_plan.md new file mode 100644 index 0000000..5611d02 --- /dev/null +++ b/docs/auto_scheduler_plan.md @@ -0,0 +1,366 @@ +# 远程数据自动调度功能设计方案 + +## 一、需求概述 + +### 1.1 背景 + +当前系统需要每周一手动点击"远程下载并处理"来跑上周数据(周一到周日,7天数据)。由于服务器时间不准确,无法使用定时任务在固定时间运行。需要一个基于**数据就绪状态**的自动调度机制。 + +### 1.2 核心需求 + +1. **文件时间过滤**:处理时只处理最近7天的文件 +2. **自动就绪检测**:每小时检查远程目录,判断上周7天数据是否齐全 +3. **自动触发处理**:数据就绪后自动下载并处理 +4. **状态标识管理**:使用标识文件避免重复触发 + +--- + +## 二、文件时间解析规则 + +### 2.1 文件名格式 + +| 格式 | 示例 | 数据日期 | +|------|------|----------| +| `XXX_YYYYMMDDHHMM_YYYYMMDDHHMM` | `CapacityReportData2.6_202605110000_202605120000.zip` | 第一个时间戳 `202605110000` → 2026-05-11 | +| `XXX_YYYYMMDDHHMM` | `CapacityReportData2.6_202605110000.zip` | 时间戳 `202605110000` → 2026-05-11 | + +### 2.2 解析逻辑 + +复用项目已有的正则 `_ZIP_DATE_RE = re.compile(r"(? date | None: + """从文件名提取数据日期(取第一个匹配的时间戳)""" + match = _ZIP_DATE_RE.search(filename) + if not match: + return None + timestamp = match.group(1) # 12位: YYYYMMDDHHMM 或 14位: YYYYMMDDHHMMSS + fmt = "%Y%m%d%H%M%S" if len(timestamp) == 14 else "%Y%m%d%H%M" + try: + return datetime.strptime(timestamp, fmt).date() + except ValueError: + return None +``` + +--- + +## 三、自动调度机制 + +### 3.1 工作流程 + +``` +每小时定时器触发 + │ + ▼ +┌─────────────────────┐ +│ 检查就绪标识文件 │ +│ (ready.flag) │ +└─────────────────────┘ + │ + ├── 标识存在 ──► 跳过检测,直接触发处理 + │ + ▼ +┌─────────────────────┐ +│ 连接 FTP/SFTP │ +│ 遍历远程目录 │ +└─────────────────────┘ + │ + ▼ +┌─────────────────────────────────────┐ +│ 对每个 expected_directory: │ +│ 1. 列出所有 ZIP 文件名 │ +│ 2. 从文件名提取数据日期 │ +│ 3. 检查是否覆盖上周7天(周一~周日)│ +└─────────────────────────────────────┘ + │ + ├── 未就绪 ──► 记录日志,等待下次检查 + │ + ▼ +┌─────────────────────┐ +│ 写入就绪标识文件 │ +│ (ready.flag) │ +└─────────────────────┘ + │ + ▼ +┌─────────────────────┐ +│ 触发远程下载处理 │ +│ (复用现有流程) │ +└─────────────────────┘ + │ + ▼ +┌─────────────────────┐ +│ 处理成功后 │ +│ 1. 删除远程源文件 │ +│ 2. 删除就绪标识 │ +└─────────────────────┘ + │ + ▼ + 继续每小时检查 +``` + +### 3.2 目标日期范围(严格按自然周) + +**永远是上周一到上周日**(自然周,7天)。无论今天是周几,程序都会自动计算正确的上周一和上周日。 + +```python +def get_target_week_range(week_offset: int = 0) -> tuple[date, date]: + """获取目标周的周一到周日。 + week_offset=0 表示上周,-1 表示上上周 + """ + today = date.today() + this_monday = today - timedelta(days=today.weekday()) + target_monday = this_monday - timedelta(days=7 * (1 - week_offset)) + target_sunday = target_monday + timedelta(days=6) + return target_monday, target_sunday +``` + +示例(假设 week_offset=0): + +| 今天 | 本周一 | 上周一 | 上周日 | 检查范围 | +|------|--------|--------|--------|----------| +| 周一 5/19 | 5/19 | 5/12 | 5/18 | 5/12~5/18 | +| 周三 5/21 | 5/19 | 5/12 | 5/18 | 5/12~5/18 | +| 周日 5/25 | 5/19 | 5/12 | 5/18 | 5/12~5/18 | + +**不会**把非自然周的7天范围(如上周五到本周四)当作目标。`week_offset` 参数仅用于补跑场景:如果上周处理失败,可设为 `-1` 检查上上周的周一到周日。 + +### 3.3 就绪判断 + +一个目录的"就绪"条件:该目录下所有 ZIP 文件提取出的数据日期集合,**完全覆盖**目标周(上周一~上周日)的全部 7 天。 + +**关键**: +1. 不依赖系统日期,而是从文件名中提取日期 +2. 文件名只取**第一个**时间戳作为数据日期 +3. 例如文件 `CapacityReportData2.6_202605120000_202605130000.zip` → 数据日期为 2026-05-12 +4. 如果某个日期没有对应的 ZIP 文件,该目录判定为未就绪 + +### 3.4 标识文件 + +- **位置**:`cache/auto_scheduler/ready.flag` +- **格式**:JSON + +```json +{ + "ready_at": "2026-05-19T10:00:00", + "week_start": "2026-05-12", + "week_end": "2026-05-18", + "directories": { + "4G/FDD": {"found_days": 7, "required_days": 7}, + "4G/900": {"found_days": 7, "required_days": 7}, + "5G/2.6": {"found_days": 7, "required_days": 7}, + "5G/700": {"found_days": 7, "required_days": 7} + } +} +``` + +--- + +## 四、配置项设计 + +### 4.1 新增配置(`Configure.json` → `RemoteData` 节点下) + +```json +{ + "RemoteData": { + "enabled": true, + "protocol": "sftp", + "host": "127.0.0.1", + "port": 2022, + "user": "nixevol", + "remote_dir": "/CapacityReportData", + "passive": true, + "timeout": 30, + "auto_delete_source": true, + "passwd": "242520", + + "auto_scheduler": { + "enabled": false, + "check_interval_hours": 1, + "expected_directories": [ + "4G/FDD", + "4G/900", + "5G/2.6", + "5G/700" + ], + "week_offset": 0 + } + } +} +``` + +### 4.2 字段说明 + +| 字段 | 类型 | 默认值 | 说明 | +|------|------|--------|------| +| `enabled` | bool | `false` | 是否启用自动调度 | +| `check_interval_hours` | int | `1` | 检查间隔(小时),最小1 | +| `expected_directories` | list[str] | `[]` | 需要检测就绪的子目录路径列表(相对于 `remote_dir`)。为空时按实际 ZIP 所在目录逐个检测 | +| `week_offset` | int | `0` | 周偏移,`0` = 检查上周,`-1` = 检查上上周 | + +### 4.3 约束规则 + +1. **自动调度开启时,`auto_delete_source` 强制为 `true`**(前端灰显、后端校验) +2. `expected_directories` 为空时按远程 ZIP 实际所在目录逐个检测 +3. `check_interval_hours` 最小值为 1 + +--- + +## 五、代码改动计划 + +### 5.1 后端 + +#### `app/config.py` +- `RemoteDataConfig` 新增 4 个字段 + `AutoSchedulerConfig` 子配置类 +- 更新 `from_dict()`、`to_dict()`、`to_file_dict()` 序列化逻辑 +- 新增约束校验:`auto_scheduler.enabled = true` 时强制 `auto_delete_source = true` + +#### `app/services/remote_download.py` +- 新增 `list_remote_zip_files(directory: str | None) -> list[RemoteFileInfo]`:递归列出指定远程目录下所有 ZIP 文件路径、相对路径、所在目录和大小(只扫描,不下载) +- 远程下载会先按每个实际目录筛选最近 7 天 ZIP;自动调度触发时按 ready flag 中的目标日期精确下载,避免数据就绪后远程目录又出现新文件时误下载非目标周数据。 + +#### `app/services/auto_scheduler.py`(**新建**) +- `AutoScheduler` 类 + - `start()` / `stop()`:启动/停止后台定时线程 + - `check_and_run()`:检查 → 就绪 → 触发处理 + - `_is_ready()`:检查所有目录就绪状态 + - `_mark_ready()` / `_clear_ready()`:标识文件管理 + - `_trigger_processing()`:复用 `remote.py` 中的 `_run_remote_processing` 逻辑 + - `get_status()` → dict:返回调度器当前状态 + +#### `app/main.py` +- 应用启动时初始化 `AutoScheduler` 并 start +- 应用关闭时 stop + +#### `app/api/routers/remote.py` +- 新增 `GET /api/remote/scheduler/status`:查询调度状态 +- 新增 `POST /api/remote/scheduler/trigger`:手动触发一次检测 +- 修改 `POST /api/remote/start`:自动调度期间禁止手动触发(避免冲突) + +### 5.2 前端 + +#### `frontend/src/components/SettingsPanel.vue` +在"远程数据源"配置区域新增"自动调度"折叠卡片: +- 启用开关 +- 检查间隔(小时)输入框 +- 预期目录列表(支持增删) +- 周偏移选择(默认上周) +- 状态面板:下次检查时间、上次检查结果、就绪状态、各目录检测结果 + +--- + +## 六、日志与监控 + +### 6.1 日志示例 + +``` +[10:00:00] 自动调度:开始检查远程目录就绪状态 +[10:00:01] 自动调度:连接 SFTP 127.0.0.1:2022 +[10:00:02] 自动调度:检查 4G/FDD → 找到 5/7 天文件(缺少 2026-05-13, 2026-05-14) +[10:00:03] 自动调度:检查 4G/900 → 找到 7/7 天文件 ✓ +[10:00:04] 自动调度:检查 5G/2.6 → 找到 7/7 天文件 ✓ +[10:00:05] 自动调度:检查 5G/700 → 找到 6/7 天文件(缺少 2026-05-18) +[10:00:06] 自动调度:数据未就绪(2/4 目录满足),等待下次检查 +... +[11:00:00] 自动调度:开始检查远程目录就绪状态 +[11:00:05] 自动调度:所有目录数据就绪 (4/4),写入标识文件 +[11:00:06] 自动调度:触发远程下载处理任务 +[11:05:00] 自动调度:任务处理成功 +[11:05:01] 自动调度:远程源文件已清理 +[11:05:02] 自动调度:已清除就绪标识,继续监控 +``` + +### 6.2 状态接口响应 + +```json +GET /api/remote/scheduler/status + +{ + "enabled": true, + "running": true, + "next_check_at": "2026-05-19T12:00:00", + "last_check_at": "2026-05-19T11:00:00", + "last_result": "triggered", + "failure_count": 0, + "task_running": false, + "ready_flag": { + "exists": false + }, + "target_week": { + "start": "2026-05-12", + "end": "2026-05-18" + }, + "directory_status": { + "4G/FDD": {"found_days": ["2026-05-12"], "found_count": 7, "required_count": 7, "ready": true}, + "4G/900": {"found_days": ["2026-05-12"], "found_count": 7, "required_count": 7, "ready": true}, + "5G/2.6": {"found_days": ["2026-05-12"], "found_count": 7, "required_count": 7, "ready": true}, + "5G/700": {"found_days": ["2026-05-12"], "found_count": 7, "required_count": 7, "ready": true} + } +} +``` + +--- + +## 七、异常处理 + +| 场景 | 处理方式 | +|------|----------| +| SFTP/FTP 连接失败 | 记录错误日志,等下次周期重试 | +| 连续失败 ≥3 次 | 在前端调度状态面板显示红色失败状态 | +| 处理失败 | **保留就绪标识**,下次检查时直接触发处理(不重新检测) | +| 远程目录不存在 | 跳过该目录,记录警告,其他目录正常检测 | +| 目录内无 ZIP 文件 | 该目录判定为未就绪 | +| 手动触发远程处理 | 自动调度运行中禁止手动触发,反之亦然(互斥锁复用 `global_task_lock`) | + +--- + +## 八、测试要点 + +1. **文件名解析** + - `XXX_202605110000_202605120000.zip` → 2026-05-11 ✓ + - `XXX_202605110000.zip` → 2026-05-11 ✓ + - `no_date_here.zip` → None(跳过) + - `XXX_20260511000012345.zip` → 14位匹配 → 2026-05-11 ✓ + +2. **就绪检测** + - 所有目录都覆盖7天 → 就绪 + - 某目录缺少1天 → 未就绪 + - `expected_directories` 为空 → 只检测根目录 + - 子目录不存在 → 跳过并记录警告 + +3. **调度流程** + - 检测 → 写标识 → 触发处理 → 成功 → 删标识 → 循环 + - 处理失败 → 保留标识 → 下次直接触发 + +4. **配置联动** + - 开启自动调度 → `auto_delete_source` 自动变为 true + - 修改配置 → 调度器热重载 + +5. **互斥** + - 自动调度运行中 → 手动触发被拒绝 + - 手动任务运行中 → 自动调度跳过本轮 + +--- + +## 九、实施步骤与工时估算 + +| 步骤 | 内容 | 预估工时 | +|------|------|----------| +| 1 | `app/config.py`:扩展 `RemoteDataConfig`,新增 `AutoSchedulerConfig` | 20min | +| 2 | `app/services/remote_download.py`:新增 `list_remote_zip_names()` | 30min | +| 3 | `app/services/auto_scheduler.py`:新建调度器核心逻辑 | 2h | +| 4 | `app/api/routers/remote.py`:新增状态查询和手动触发接口 | 30min | +| 5 | `app/main.py`:集成调度器启动/停止 | 10min | +| 6 | `frontend/src/components/SettingsPanel.vue`:自动调度配置 UI | 1.5h | +| 7 | `frontend/src/api/client.ts`:新增 API 调用方法 | 10min | +| 8 | 联调测试 | 1h | +| **合计** | | **约6h** | + +--- + +## 十、注意事项 + +1. **不改动现有处理流程**:自动调度只是在现有"远程下载并处理"之上加了一层检测和触发机制 +2. **标识文件是运行时数据**:`cache/auto_scheduler/ready.flag` 不应提交到版本库 +3. **`week_offset` 用途**:如果某周处理失败需要下周补跑,可将 `week_offset` 设为 `-1` 来检查上上周的数据 +4. **配置热重载**:修改自动调度配置后,调度器应在下一次检查周期使用新配置(不需要重启服务) +5. **前端 `auto_delete_source` 联动**:当自动调度开启时,前端应将"处理成功后删除源文件"开关灰显为强制开启状态 diff --git a/docs/project_context.md b/docs/project_context.md index 02528dd..2d6b0ef 100644 --- a/docs/project_context.md +++ b/docs/project_context.md @@ -1,5 +1,16 @@ # 项目上下文记录 +## 2026-05-25:新增远程自动调度和每目录 7 天处理窗口 + +- 新增 `app/utils/file_dates.py`,统一解析文件名中的第一个 `YYYYMMDDHHMM` 或 `YYYYMMDDHHMMSS` 时间戳,并提供按目录筛选最近 7 个自然日文件的工具;本地手动上传和远程下载后的 ZIP、Excel、CSV 处理都会按所在目录只保留最近 7 天文件,未携带日期且同目录没有任何可识别日期时保留兼容。 +- `app/services/remote_download.py` 增加远程 ZIP 清单扫描和筛选下载:普通远程处理按每个远程目录下载最近 7 天 ZIP;自动调度触发时按 ready flag 中记录的目标日期精确下载,若没有匹配 ZIP 不会回退全量下载,避免误处理新旧混杂数据。 +- `app/config.py` 增加 `RemoteData.auto_scheduler` 配置,包含 `enabled`、`check_interval_hours`、`expected_directories` 和 `week_offset`。自动调度开启时后端会强制 `RemoteData.enabled=True` 和 `auto_delete_source=True`,前端也同步灰显并强制打开相关开关。 +- 新增 `app/services/auto_scheduler.py` 后台线程:应用启动后按配置间隔检查 FTP/SFTP 目录,使用本机当前日期计算目标自然周(`week_offset=0` 为上周,`-1` 为上上周),但文件覆盖情况完全以 ZIP 文件名日期为准;全部目录覆盖 7 天后写入 `cache/auto_scheduler/ready.flag`,下一轮检测再触发远程下载并处理。 +- 调度成功且远程源文件清理成功后会删除 ready flag;处理失败、触发失败或源文件清理失败会保留 ready flag 供下轮重试。扫描失败会记录 `scan_failed`,连续失败次数通过状态接口返回,前端会显示红色失败状态。 +- `app/api/routers/remote.py` 新增 `/api/remote/scheduler/status` 和 `/api/remote/scheduler/trigger`,并将远程处理启动逻辑抽成 `start_remote_processing_job()` 供手动按钮和调度器复用。 +- `frontend/src/components/SettingsPanel.vue` 在远程数据源配置中加入自动调度区域:启用开关、检查间隔、目标周期、预期目录维护、调度状态、刷新状态和立即检查。配置下载/上传会随 `RemoteData` 一起携带自动调度配置。 +- 已验证:文件名日期解析、目标周计算、每目录 7 天筛选、调度开启强制删除源文件配置、调度日期精确筛选逻辑均通过临时 Python 片段;`.venv\Scripts\python.exe -m compileall app` 与 `npm run build` 通过,前端构建仅保留既有 Vite 大 chunk 警告。 + ## 2026-05-25:修复 API Token 指定日期输入不可见 - `frontend/src/components/ApiTokenManager.vue` 中创建/编辑 Token 的到期日期控件从 `n-date-picker` 改为原生 `input[type=date]`,避免日期选择组件在弹窗内出现占位但输入框不可见的问题。 diff --git a/frontend/src/components/SettingsPanel.vue b/frontend/src/components/SettingsPanel.vue index 7124dde..72d96ea 100644 --- a/frontend/src/components/SettingsPanel.vue +++ b/frontend/src/components/SettingsPanel.vue @@ -84,12 +84,15 @@ - + - + @@ -142,6 +145,111 @@

自动化执行会递归下载该目录下的全部文件和文件夹到本地缓存,再按现有处理流程入库和执行脚本;自动删除源文件只会在处理成功后删除远程文件,保留目录结构。

+ + 自动调度 + + + + + + + + + + + + + + + + + + +
+ + + + 添加 + + +
+ + {{ directory }} + + + 按实际 ZIP 目录检测 + +
+
+
+

自动调度会按文件名日期检查目标自然周 7 天;开启后会强制启用远程自动化和处理成功后删除源文件。

+ +
+
+ 调度状态 + + {{ schedulerStatusLabel }} + +
+
+
+ 目标周 + {{ schedulerTargetWeekText }} +
+
+ 下次检查 + {{ schedulerStatus?.next_check_at || '-' }} +
+
+ 就绪标识 + {{ schedulerStatus?.ready_flag?.exists ? '已存在' : '无' }} +
+
+

{{ schedulerStatus?.last_message || '自动调度状态尚未加载' }}

+
+
+ {{ row.name }} + + {{ row.ready ? '就绪' : `缺 ${row.missing_days.length} 天` }} + +
+
+ + + 刷新状态 + + + 立即检查 + + +