diff --git a/pipeline_service/work_env.py b/pipeline_service/work_env.py index b432bc3..70d2b52 100644 --- a/pipeline_service/work_env.py +++ b/pipeline_service/work_env.py @@ -204,6 +204,19 @@ def _validate_remote(remote_config: dict) -> str: # 目录迁移 # ═══════════════════════════════════════════════════════════ +async def _remote_mkdir(env: dict, remote_dir: str) -> bool: + """远程 mkdir -p(相对路径可能多级,rsync 不会递归建父目录)。失败返回 False。""" + try: + target = _ssh_target(env) + ssh_cmd = ["ssh"] + _ssh_common_args(env) + [target, "mkdir -p " + shlex.quote(remote_dir)] + proc = await asyncio.create_subprocess_exec( + *ssh_cmd, stdout=asyncio.subprocess.DEVNULL, stderr=asyncio.subprocess.DEVNULL) + await asyncio.wait_for(proc.communicate(), timeout=30) + return proc.returncode == 0 + except Exception: + return False + + async def migrate_work_dir(account_name: str, from_mode: str, to_mode: str, env: dict) -> dict: """把工作目录内容从旧环境迁移到新环境。 @@ -231,6 +244,7 @@ async def migrate_work_dir(account_name: str, from_mode: str, to_mode: str, env: remote_dir = env.get("remote_dir", "") if not remote_dir: return {"ok": False, "error": "缺少远程目录 remote_dir"} + await _remote_mkdir(env, _resolve_remote_dir(remote_dir)) target = _ssh_target(env) + ":" + _resolve_remote_dir(remote_dir) cmd = ["rsync", "-az", "--delete", "-e", "ssh " + " ".join(_ssh_common_args(env)),