From f2bef1f2f7a2f1bd50dd342977577dee6a845705 Mon Sep 17 00:00:00 2001 From: Nixevol Date: Mon, 18 May 2026 18:21:20 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=96=B0=E5=A2=9E=E8=BF=9C=E7=A8=8B?= =?UTF-8?q?=E6=95=B0=E6=8D=AE=E8=87=AA=E5=8A=A8=E5=8C=96=E5=A4=84=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- README.md | 3 + app/api/routers/config.py | 12 +- app/api/routers/remote.py | 129 +++++++++++++ app/config.py | 77 ++++++++ app/main.py | 3 +- app/services/remote_download.py | 217 ++++++++++++++++++++++ docs/project_context.md | 9 + frontend/src/components/FileWorkflow.vue | 44 ++++- frontend/src/components/SettingsPanel.vue | 160 +++++++++++++++- frontend/src/styles.css | 35 ++++ frontend/src/types.ts | 13 ++ requirements.txt | 1 + 12 files changed, 698 insertions(+), 5 deletions(-) create mode 100644 app/api/routers/remote.py create mode 100644 app/services/remote_download.py diff --git a/README.md b/README.md index 5d389c4..40bb437 100644 --- a/README.md +++ b/README.md @@ -89,6 +89,7 @@ Vite 会把 `/api` 和 `/health` 代理到 `http://localhost:9081`。 `Configure.json` 主要包含: - `MySQL_DBInfo`:MySQL 连接信息 +- `RemoteData`:FTP/SFTP 远程数据源配置,用于递归下载目录后自动处理 - `SheetFilter`:Excel Sheet 过滤规则 - `ExtractField`:字段抽取映射配置 @@ -114,6 +115,8 @@ Docker 镜像构建采用多阶段流程:先在 Node 阶段执行前端构建 - `POST /api/login`:登录 - `POST /api/change-password`:修改密码 - `POST /api/upload`:上传文件 +- `POST /api/remote/test`:测试 FTP/SFTP 远程数据源 +- `POST /api/remote/start`:从远程目录下载数据并自动启动处理 - `POST /api/process/start`:启动处理 - `POST /api/process/status`:查询处理状态 - `GET /api/history`:历史记录 diff --git a/app/api/routers/config.py b/app/api/routers/config.py index 5a36ed0..1a75eba 100644 --- a/app/api/routers/config.py +++ b/app/api/routers/config.py @@ -6,7 +6,7 @@ from fastapi import APIRouter, Body, HTTPException, UploadFile, File from fastapi.responses import FileResponse from app import state -from app.config import CONFIG_FILE +from app.config import CONFIG_FILE, RemoteDataConfig router = APIRouter(tags=["config"]) @@ -39,6 +39,13 @@ async def update_mysql_config( return {"success": True, "message": "数据库配置已更新", "update": state.config.update} +@router.post("/api/config/remote") +async def update_remote_config(config: dict[str, Any] = Body(...)): + state.config.remote_data = RemoteDataConfig.from_dict(config) + state.config.save() + return {"success": True, "message": "远程数据配置已更新", "update": state.config.update} + + @router.post("/api/config/sheet-filter") async def update_sheet_filter(filters: list[str] = Body(...)): state.config.sheet_filter = filters @@ -97,3 +104,6 @@ def _apply_config_data(data: dict[str, Any]) -> None: if "ExtractField" in data: state.config.extract_fields = data["ExtractField"] if isinstance(data["ExtractField"], list) else [] + remote_data = data.get("RemoteData") + if isinstance(remote_data, dict): + state.config.remote_data = RemoteDataConfig.from_dict(remote_data) diff --git a/app/api/routers/remote.py b/app/api/routers/remote.py new file mode 100644 index 0000000..dc6aeac --- /dev/null +++ b/app/api/routers/remote.py @@ -0,0 +1,129 @@ +from datetime import datetime +from pathlib import Path +from threading import Thread +from typing import Any + +from fastapi import APIRouter, Body, HTTPException + +from app import state +from app.config import CACHE_DIR, RemoteDataConfig +from app.processor import DataProcessor, ProcessLogger +from app.services.remote_download import RemoteDataDownloader + + +router = APIRouter(tags=["remote"]) + + +@router.post("/api/remote/test") +async def test_remote_connection(config: dict[str, Any] | None = Body(None)): + remote_config = RemoteDataConfig.from_dict(config) if config else state.config.remote_data + try: + RemoteDataDownloader(remote_config).test_connection() + return {"success": True, "message": "远程服务器连接成功"} + except Exception as exc: + return {"success": False, "message": f"远程服务器连接失败: {exc}"} + + +@router.post("/api/remote/start") +async def start_remote_processing(): + if state.global_task_lock["locked"]: + raise HTTPException(status_code=409, detail="已有任务在运行,请等待当前任务完成") + + remote_config = state.config.remote_data.normalized() + if not remote_config.enabled: + raise HTTPException(status_code=400, detail="请先启用远程数据配置") + + task_id = datetime.now().strftime("%Y%m%d_%H%M%S") + work_dir = CACHE_DIR / task_id + work_dir.mkdir(parents=True, exist_ok=True) + + logs: list[str] = [] + + def log_callback(message: str) -> None: + logs.append(message) + state.processing_tasks[task_id] = {"logs": logs.copy(), "status": "processing"} + + logger = ProcessLogger(log_file=work_dir / "log.txt", callback=log_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.global_task_lock.update( + { + "locked": True, + "task_id": task_id, + "stage": "downloading", + "started_at": datetime.now().isoformat(), + } + ) + + thread = Thread( + target=_run_remote_processing, + args=(task_id, work_dir, remote_config, logger), + daemon=True, + ) + thread.start() + + return { + "success": True, + "message": "远程下载处理任务已启动", + "task_id": task_id, + "stage": "downloading", + } + + +def _run_remote_processing( + task_id: str, + work_dir: Path, + remote_config: RemoteDataConfig, + logger: ProcessLogger, +) -> None: + 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) + logger.success( + f"远程下载完成,共 {download_result.file_count} 个文件," + f"{_format_bytes(download_result.total_bytes)}" + ) + + if download_result.file_count == 0: + raise RuntimeError("远程目录中未下载到任何文件") + + state.history_manager.update(task_id, file_count=download_result.file_count) + state.global_task_lock["stage"] = "processing" + + processor = DataProcessor(state.config, work_dir, logger) + result = processor.process() + status = "completed" if result.get("success") else "failed" + state.history_manager.update( + task_id, + status=status, + elapsed_time=result.get("elapsed_time", 0), + error=result.get("error"), + result_tables=["4G_结果表", "5G_结果表"], + ) + state.processing_tasks[task_id] = { + "logs": state.history_manager.get_logs(task_id), + "status": status, + } + except Exception as exc: + 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), + "status": "failed", + } + finally: + state.reset_task_lock() + + +def _format_bytes(size: int) -> str: + value = float(size) + for unit in ("B", "KB", "MB", "GB"): + if value < 1024: + return f"{value:.1f} {unit}" + value /= 1024 + return f"{value:.1f} TB" diff --git a/app/config.py b/app/config.py index 05bbfbe..2fc2401 100644 --- a/app/config.py +++ b/app/config.py @@ -23,10 +23,82 @@ class MySQLConfig: dbname: str = "CapacityReport" +@dataclass +class RemoteDataConfig: + enabled: bool = False + protocol: str = "sftp" + host: str = "" + port: int = 22 + user: str = "" + passwd: str = "" + remote_dir: str = "/" + passive: bool = True + timeout: int = 30 + + def normalized(self) -> "RemoteDataConfig": + protocol = self.protocol.lower().strip() + if protocol not in {"ftp", "sftp"}: + protocol = "sftp" + + port = self.port or (22 if protocol == "sftp" else 21) + return RemoteDataConfig( + enabled=bool(self.enabled), + protocol=protocol, + host=self.host.strip(), + port=port, + user=self.user.strip(), + passwd=self.passwd, + remote_dir=(self.remote_dir or "/").strip() or "/", + passive=bool(self.passive), + timeout=max(int(self.timeout or 30), 1), + ) + + def to_dict(self, include_password: bool = False) -> Dict[str, Any]: + data = { + "enabled": self.enabled, + "protocol": self.protocol, + "host": self.host, + "port": self.port, + "user": self.user, + "remote_dir": self.remote_dir, + "passive": self.passive, + "timeout": self.timeout, + } + if include_password: + data["passwd"] = self.passwd + return data + + @classmethod + def from_dict(cls, data: Dict[str, Any] | None) -> "RemoteDataConfig": + data = data or {} + protocol = str(data.get("protocol", "sftp")).lower() + default_port = 22 if protocol == "sftp" else 21 + try: + port = int(data.get("port") or default_port) + except (TypeError, ValueError): + port = default_port + try: + timeout = int(data.get("timeout") or 30) + except (TypeError, ValueError): + timeout = 30 + return cls( + enabled=bool(data.get("enabled", False)), + protocol=protocol, + host=str(data.get("host", "")), + port=port, + user=str(data.get("user", "")), + passwd=str(data.get("passwd", "")), + remote_dir=str(data.get("remote_dir", "/")), + passive=bool(data.get("passive", True)), + timeout=timeout, + ).normalized() + + @dataclass class AppConfig: update: str = "" mysql: MySQLConfig = field(default_factory=MySQLConfig) + remote_data: RemoteDataConfig = field(default_factory=RemoteDataConfig) sheet_filter: List[str] = field(default_factory=list) extract_fields: List[Dict[str, Any]] = field(default_factory=list) @@ -47,10 +119,12 @@ class AppConfig: passwd=mysql_data.get("passwd", ""), dbname=mysql_data.get("dbname", "CapacityReport") ) + remote_config = RemoteDataConfig.from_dict(data.get("RemoteData")) return cls( update=data.get("Update", ""), mysql=mysql_config, + remote_data=remote_config, sheet_filter=data.get("SheetFilter", []), extract_fields=data.get("ExtractField", []) ) @@ -68,6 +142,7 @@ class AppConfig: "passwd": self.mysql.passwd, "dbname": self.mysql.dbname }, + "RemoteData": self.remote_data.normalized().to_dict(include_password=True), "SheetFilter": self.sheet_filter, "ExtractField": self.extract_fields } @@ -85,6 +160,7 @@ class AppConfig: "user": self.mysql.user, "dbname": self.mysql.dbname }, + "remote_data": self.remote_data.normalized().to_dict(), "sheet_filter": self.sheet_filter, "extract_fields": self.extract_fields } @@ -100,6 +176,7 @@ class AppConfig: "passwd": self.mysql.passwd, "dbname": self.mysql.dbname }, + "remote_data": self.remote_data.normalized().to_dict(include_password=True), "sheet_filter": self.sheet_filter, "extract_fields": self.extract_fields } diff --git a/app/main.py b/app/main.py index 5de7a6d..b16b5d5 100644 --- a/app/main.py +++ b/app/main.py @@ -8,7 +8,7 @@ from fastapi.middleware.cors import CORSMiddleware from fastapi.responses import FileResponse, JSONResponse from app import state -from app.api.routers import auth, cache, config, database, health, history, script, service, tasks, upload +from app.api.routers import auth, cache, config, database, health, history, remote, script, service, tasks, upload from app.auth import verify_jwt_token from app.config import BASE_DIR @@ -65,6 +65,7 @@ def register_routes(app: FastAPI) -> None: auth.router, health.router, upload.router, + remote.router, tasks.router, history.router, database.router, diff --git a/app/services/remote_download.py b/app/services/remote_download.py new file mode 100644 index 0000000..59cc2df --- /dev/null +++ b/app/services/remote_download.py @@ -0,0 +1,217 @@ +from __future__ import annotations + +import stat +from dataclasses import dataclass +from ftplib import FTP +from pathlib import Path +from typing import Callable + +from app.config import RemoteDataConfig + + +LogFn = Callable[[str], None] + + +@dataclass +class RemoteDownloadResult: + file_count: int = 0 + total_bytes: int = 0 + + +class RemoteDownloadError(RuntimeError): + pass + + +class RemoteDataDownloader: + def __init__(self, config: RemoteDataConfig, logger: LogFn | None = None): + self.config = config.normalized() + self.logger = logger + + def test_connection(self) -> None: + self._validate_config() + if self.config.protocol == "ftp": + with self._ftp_client() as ftp: + ftp.cwd(self.config.remote_dir) + return + + ssh = self._sftp_ssh_client() + try: + with ssh.open_sftp() as sftp: + sftp.stat(self.config.remote_dir) + finally: + ssh.close() + + def download_to(self, destination: Path) -> RemoteDownloadResult: + self._validate_config() + destination.mkdir(parents=True, exist_ok=True) + if self.config.protocol == "ftp": + return self._download_ftp(destination) + return self._download_sftp(destination) + + def _validate_config(self) -> None: + if self.config.protocol not in {"ftp", "sftp"}: + raise RemoteDownloadError("远程协议只支持 FTP 或 SFTP") + if not self.config.host: + raise RemoteDownloadError("请填写远程服务器地址") + if not self.config.user: + raise RemoteDownloadError("请填写远程服务器用户名") + if not self.config.remote_dir: + raise RemoteDownloadError("请填写远程数据目录") + + def _log(self, message: str) -> None: + if self.logger: + self.logger(message) + + def _ftp_client(self) -> FTP: + ftp = FTP() + ftp.connect(self.config.host, self.config.port, timeout=self.config.timeout) + ftp.login(self.config.user, self.config.passwd) + ftp.set_pasv(self.config.passive) + return ftp + + def _download_ftp(self, destination: Path) -> RemoteDownloadResult: + result = RemoteDownloadResult() + with self._ftp_client() as ftp: + self._log(f"已连接 FTP: {self.config.host}:{self.config.port}") + self._download_ftp_dir(ftp, self.config.remote_dir, destination, result) + return result + + def _download_ftp_dir( + self, + ftp: FTP, + remote_dir: str, + local_dir: Path, + result: RemoteDownloadResult, + ) -> None: + local_dir.mkdir(parents=True, exist_ok=True) + entries = self._list_ftp_entries(ftp, remote_dir) + + for name, entry_type, size in entries: + if name in {".", ".."}: + continue + + remote_path = self._join_remote_path(remote_dir, name) + local_path = local_dir / name + + if entry_type == "dir": + self._download_ftp_dir(ftp, remote_path, local_path, result) + continue + + if entry_type == "unknown" and self._ftp_is_dir(ftp, remote_path): + self._download_ftp_dir(ftp, remote_path, local_path, result) + continue + + self._download_ftp_file(ftp, remote_path, local_path, result, size) + + def _list_ftp_entries(self, ftp: FTP, remote_dir: str) -> list[tuple[str, str, int]]: + try: + return [ + ( + name, + facts.get("type", "unknown"), + int(facts.get("size", "0") or 0), + ) + for name, facts in ftp.mlsd(remote_dir) + ] + except Exception: + names = ftp.nlst(remote_dir) + entries = [] + for name in names: + clean_name = Path(name.replace("\\", "/")).name + if clean_name: + entries.append((clean_name, "unknown", 0)) + return entries + + def _ftp_is_dir(self, ftp: FTP, remote_path: str) -> bool: + current = ftp.pwd() + try: + ftp.cwd(remote_path) + return True + except Exception: + return False + finally: + try: + ftp.cwd(current) + except Exception: + pass + + def _download_ftp_file( + self, + ftp: FTP, + 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}") + with local_path.open("wb") as file: + ftp.retrbinary(f"RETR {remote_path}", file.write) + result.file_count += 1 + result.total_bytes += expected_size or local_path.stat().st_size + + def _sftp_ssh_client(self): + try: + import paramiko + except ImportError as exc: + raise RemoteDownloadError("SFTP 功能需要安装 paramiko,请执行 uv pip install -r requirements.txt") from exc + + ssh = paramiko.SSHClient() + ssh.set_missing_host_key_policy(paramiko.AutoAddPolicy()) + ssh.connect( + hostname=self.config.host, + port=self.config.port, + username=self.config.user, + password=self.config.passwd, + timeout=self.config.timeout, + banner_timeout=self.config.timeout, + auth_timeout=self.config.timeout, + ) + return ssh + + def _download_sftp(self, destination: Path) -> RemoteDownloadResult: + result = RemoteDownloadResult() + ssh = self._sftp_ssh_client() + try: + with ssh.open_sftp() as sftp: + self._log(f"已连接 SFTP: {self.config.host}:{self.config.port}") + self._download_sftp_path(sftp, self.config.remote_dir, destination, result) + finally: + ssh.close() + return result + + def _download_sftp_path( + self, + sftp, + remote_path: str, + local_path: Path, + result: RemoteDownloadResult, + ) -> None: + attrs = sftp.stat(remote_path) + if stat.S_ISDIR(attrs.st_mode): + local_path.mkdir(parents=True, exist_ok=True) + for item in sftp.listdir_attr(remote_path): + if item.filename in {".", ".."}: + continue + self._download_sftp_path( + sftp, + self._join_remote_path(remote_path, item.filename), + local_path / item.filename, + result, + ) + return + + 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) + + @staticmethod + def _join_remote_path(parent: str, child: str) -> str: + parent = parent.replace("\\", "/").rstrip("/") + if not parent: + return child + if parent == "/": + return f"/{child}" + return f"{parent}/{child}" diff --git a/docs/project_context.md b/docs/project_context.md index aa213d9..8aff0f7 100644 --- a/docs/project_context.md +++ b/docs/project_context.md @@ -1,5 +1,14 @@ # 项目上下文记录 +## 2026-05-18:新增 FTP/SFTP 远程自动化处理 + +- `Configure.json` 新增 `RemoteData` 配置,包含启用状态、协议、主机、端口、用户名、密码、远程目录、FTP 被动模式和超时时间;`app/config.py` 会兼容旧配置并在保存时写回该配置块。 +- 新增 `app/services/remote_download.py`,FTP 使用标准库 `ftplib`,SFTP 使用 `paramiko`,会递归下载远程目录下的全部文件和文件夹到本地任务缓存目录。 +- 新增 `app/api/routers/remote.py`,`POST /api/remote/test` 用于测试远程连接,`POST /api/remote/start` 会创建历史任务、下载远程数据并复用 `DataProcessor` 完成现有处理流程。 +- `frontend/src/components/SettingsPanel.vue` 新增“远程数据源”配置卡片,支持 FTP/SFTP 切换、保存和测试连接;`frontend/src/components/FileWorkflow.vue` 新增“远程下载并处理”入口。 +- 远程任务启动后会把全局任务阶段设置为 `downloading`,前端处理进度页会显示“远程下载中...”,下载完成后再切换为既有数据处理流程。 +- 新增依赖 `paramiko`,当前 `.venv` 已执行 `uv pip install -r requirements.txt`;已执行 `.venv\Scripts\python.exe -m compileall app` 和 `npm run build`,均通过;未启动浏览器或 headless Chrome。 + ## 2026-05-18:完善服务重启交互和运行时兼容 - `frontend/src/AppShell.vue` 的重启按钮点击后会先弹出确认框,确认后显示全屏“正在重启服务”遮罩和旋转加载动画,避免用户重复操作。 diff --git a/frontend/src/components/FileWorkflow.vue b/frontend/src/components/FileWorkflow.vue index a101f71..1c80439 100644 --- a/frontend/src/components/FileWorkflow.vue +++ b/frontend/src/components/FileWorkflow.vue @@ -40,6 +40,17 @@ /> +
+
+

远程自动化

+

从已配置的 FTP/SFTP 目录递归下载数据,然后自动开始处理。

+
+ + + 远程下载并处理 + +
+

已选择文件

@@ -133,6 +144,7 @@ import { computed, onBeforeUnmount, onMounted, ref } from 'vue'; import { useMessage } from 'naive-ui'; import { + CloudDownloadOutline, CloudUploadOutline, DocumentTextOutline, FolderOpenOutline, @@ -160,7 +172,9 @@ interface DroppedFile { interface UploadResponse { success: boolean; task_id: string; - file_count: number; + file_count?: number; + message?: string; + stage?: string; } const validExtensions = new Set(['.zip', '.xlsx', '.xls', '.csv']); @@ -170,6 +184,7 @@ const folderInput = ref(null); const files = ref([]); const uploadProgress = ref(0); const working = ref(false); +const remoteStarting = ref(false); const isDragging = ref(false); const taskStatus = ref(null); const activeTask = ref(null); @@ -194,6 +209,7 @@ const statusClass = computed(() => { const processStatusText = computed(() => { if (taskStatus.value?.status === 'completed') return '处理完成'; if (taskStatus.value?.status === 'failed') return '处理失败'; + if (activeTask.value?.stage === 'downloading') return '远程下载中...'; if (activeTask.value?.stage === 'uploading') return '上传中...'; return '处理中...'; }); @@ -396,7 +412,7 @@ async function uploadAndStart() { item.status = 'uploaded'; item.progress = 100; }); - message.success(`上传完成:${result.file_count} 个文件`); + message.success(`上传完成:${result.file_count ?? files.value.length} 个文件`); await apiPost('/api/process/start', { task_id: result.task_id }); taskStatus.value = { task_id: result.task_id, status: 'processing', logs: ['任务已提交,等待处理日志...'] }; startPolling(result.task_id); @@ -412,6 +428,30 @@ async function uploadAndStart() { } } +async function startRemoteProcessing() { + if (working.value || remoteStarting.value) return; + + remoteStarting.value = true; + try { + const result = await apiPost('/api/remote/start'); + message.success(result.message || '远程自动化任务已启动'); + files.value = []; + uploadProgress.value = 0; + activeTask.value = { + has_active: true, + task_id: result.task_id, + stage: result.stage || 'downloading', + started_at: new Date().toISOString() + }; + taskStatus.value = { task_id: result.task_id, status: 'processing', logs: ['远程下载任务已提交,等待处理日志...'] }; + startPolling(result.task_id); + } catch (error) { + message.error(error instanceof Error ? error.message : '启动远程自动化任务失败'); + } finally { + remoteStarting.value = false; + } +} + function updateUploadingFiles(progress: number) { files.value.forEach(item => { if (item.status === 'uploading') { diff --git a/frontend/src/components/SettingsPanel.vue b/frontend/src/components/SettingsPanel.vue index 947d02f..cada288 100644 --- a/frontend/src/components/SettingsPanel.vue +++ b/frontend/src/components/SettingsPanel.vue @@ -46,6 +46,70 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + +

自动化执行会递归下载该目录下的全部文件和文件夹到本地缓存,再按现有处理流程入库和执行脚本。

+
+ + +
+

匹配这些关键词的 Sheet 将被跳过处理

@@ -193,7 +257,7 @@ import { useMessage, type SelectOption } from 'naive-ui'; import { CloseOutline, CloudDownloadOutline, CloudUploadOutline } from '@vicons/ionicons5'; import { apiGet, apiPost, downloadGet, upload } from '../api/client'; -import type { ApiMessage, AppConfig } from '../types'; +import type { ApiMessage, AppConfig, RemoteDataConfig } from '../types'; import { resetPageHeader, setPageHeader } from '../composables/pageHeader'; interface ExtractFieldConfig { @@ -208,7 +272,9 @@ const configInput = ref(null); const configUpdate = ref(''); const loading = ref(false); const testingDb = ref(false); +const testingRemote = ref(false); const savingMysql = ref(false); +const savingRemote = ref(false); const savingSheetFilter = ref(false); const savingExtractFields = ref(false); const changingPassword = ref(false); @@ -226,6 +292,18 @@ const mysqlForm = reactive({ dbname: '' }); +const remoteForm = reactive({ + enabled: false, + protocol: 'sftp', + host: '', + port: 22, + user: '', + passwd: '', + remote_dir: '/', + passive: true, + timeout: 30 +}); + const passwordForm = reactive({ current_password: '', new_password: '', @@ -238,6 +316,10 @@ const fieldTypeOptions: SelectOption[] = [ { label: '小数', value: 'float' }, { label: '日期时间', value: 'datetime' } ]; +const remoteProtocolOptions: SelectOption[] = [ + { label: 'SFTP', value: 'sftp' }, + { label: 'FTP', value: 'ftp' } +]; const visibleFieldMappings = computed(() => { const keyword = fieldSearch.value.trim().toLowerCase(); @@ -275,6 +357,7 @@ async function loadConfig() { mysqlForm.user = config.mysql.user; mysqlForm.passwd = config.mysql.passwd || ''; mysqlForm.dbname = config.mysql.dbname; + Object.assign(remoteForm, normalizeRemoteConfig(config.remote_data)); sheetFilters.value = [...config.sheet_filter]; extractFields.value = normalizeExtractFields(config.extract_fields); resetExtractInputs(); @@ -285,6 +368,15 @@ async function loadConfig() { } } +function updateRemoteProtocol(value: string) { + if (value === 'sftp' && (!remoteForm.port || remoteForm.port === 21)) { + remoteForm.port = 22; + } + if (value === 'ftp' && (!remoteForm.port || remoteForm.port === 22)) { + remoteForm.port = 21; + } +} + async function saveMysql() { if (!mysqlForm.host || !mysqlForm.user || !mysqlForm.dbname) { message.warning('请填写完整数据库配置'); @@ -306,6 +398,57 @@ async function saveMysql() { } } +function getRemotePayload(): RemoteDataConfig { + return { + ...remoteForm, + protocol: remoteForm.protocol, + host: remoteForm.host.trim(), + port: remoteForm.port || (remoteForm.protocol === 'sftp' ? 22 : 21), + user: remoteForm.user.trim(), + passwd: remoteForm.passwd || '', + remote_dir: remoteForm.remote_dir.trim() || '/', + passive: remoteForm.passive, + timeout: remoteForm.timeout || 30 + }; +} + +function validateRemoteForm(): boolean { + if (!remoteForm.host || !remoteForm.user || !remoteForm.remote_dir) { + message.warning('请填写完整远程数据源配置'); + return false; + } + return true; +} + +async function saveRemote() { + if (!validateRemoteForm()) return; + + savingRemote.value = true; + try { + const result = await apiPost('/api/config/remote', getRemotePayload()); + configUpdate.value = result.update || configUpdate.value; + message.success(result.message || '远程数据源配置已保存'); + } catch (error) { + message.error(error instanceof Error ? error.message : '保存远程数据源配置失败'); + } finally { + savingRemote.value = false; + } +} + +async function testRemote() { + if (!validateRemoteForm()) return; + + testingRemote.value = true; + try { + const result = await apiPost('/api/remote/test', getRemotePayload()); + message[result.success ? 'success' : 'error'](result.message || '远程连接测试完成'); + } catch (error) { + message.error(error instanceof Error ? error.message : '远程连接测试失败'); + } finally { + testingRemote.value = false; + } +} + async function testDatabase() { testingDb.value = true; try { @@ -457,6 +600,21 @@ async function uploadConfigFile(event: Event) { } } +function normalizeRemoteConfig(config: RemoteDataConfig | undefined): RemoteDataConfig { + const protocol = config?.protocol === 'ftp' ? 'ftp' : 'sftp'; + return { + enabled: Boolean(config?.enabled), + protocol, + host: config?.host || '', + port: config?.port || (protocol === 'sftp' ? 22 : 21), + user: config?.user || '', + passwd: config?.passwd || '', + remote_dir: config?.remote_dir || '/', + passive: config?.passive ?? true, + timeout: config?.timeout || 30 + }; +} + function normalizeExtractFields(fields: Array>): ExtractFieldConfig[] { return fields.map(field => { const extract = Array.isArray(field.Extract) ? uniqueStrings(field.Extract) : []; diff --git a/frontend/src/styles.css b/frontend/src/styles.css index 9fa68ed..1eb7f53 100644 --- a/frontend/src/styles.css +++ b/frontend/src/styles.css @@ -593,6 +593,36 @@ select { font-size: 13px !important; } +.remote-run-card { + display: flex; + align-items: center; + justify-content: space-between; + gap: 16px; + margin-top: 16px; + padding: 16px 20px; + background: var(--td-bg-color-container); + border: 1px solid var(--td-border-color-light); + border-radius: var(--td-radius-default); +} + +.remote-run-copy { + min-width: 0; +} + +.remote-run-copy h4 { + margin: 0 0 4px; + color: var(--td-text-color-primary); + font-size: 15px; + font-weight: 600; +} + +.remote-run-copy p { + margin: 0; + color: var(--td-text-color-secondary); + font-size: 13px; + line-height: 20px; +} + .file-list { margin-top: 20px; } @@ -1521,6 +1551,11 @@ select { justify-content: flex-end; } + .remote-run-card { + align-items: stretch; + flex-direction: column; + } + .card-actions, .field-card-title-row, .extract-list-header { diff --git a/frontend/src/types.ts b/frontend/src/types.ts index cdc91e2..b408901 100644 --- a/frontend/src/types.ts +++ b/frontend/src/types.ts @@ -57,10 +57,23 @@ export interface AppConfig { passwd?: string; dbname: string; }; + remote_data: RemoteDataConfig; sheet_filter: string[]; extract_fields: Array>; } +export interface RemoteDataConfig { + enabled: boolean; + protocol: 'ftp' | 'sftp'; + host: string; + port: number; + user: string; + passwd?: string; + remote_dir: string; + passive: boolean; + timeout: number; +} + export interface TableInfo { name: string; columns: Array<{ diff --git a/requirements.txt b/requirements.txt index 4ffa271..7c8deda 100644 --- a/requirements.txt +++ b/requirements.txt @@ -12,6 +12,7 @@ cryptography pandas openpyxl chardet +paramiko # Process manager supervisor