From 8516c2ae9c622f0b3e0f4270fa4bd3f1b8994b7e Mon Sep 17 00:00:00 2001 From: Codex Date: Sun, 2 Aug 2026 01:41:10 +0800 Subject: [PATCH 1/3] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=EF=BC=9A=E7=BB=9F?= =?UTF-8?q?=E4=B8=80E=E7=9B=98=E8=A7=86=E9=A2=91=E5=AD=98=E5=82=A8?= =?UTF-8?q?=E4=B8=8E=E6=B0=B8=E4=B9=85=E5=88=A0=E9=99=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .env.example | 4 +- DEVELOPMENT_LOG.md | 10 + NEXT_STEPS.md | 9 + app/core/config.py | 6 +- app/main.py | 12 + app/routers/tasks.py | 96 +++++-- app/services/database_backup_service.py | 51 ++++ app/services/job_service.py | 2 + app/services/job_worker.py | 2 + app/services/pipeline_engine.py | 2 + app/services/storage_service.py | 252 +++++++++++++++-- app/services/task_lifecycle_service.py | 146 ++++++++-- app/services/task_query_service.py | 6 + app/static/js/app.js | 11 +- app/templates/new_task.html | 4 +- app/templates/system_status.html | 6 + app/templates/tasks.html | 2 +- docs/DATABASE_SCHEMA.md | 8 +- docs/DEPLOYMENT.md | 5 + docs/UI_REFERENCE.md | 10 +- scripts/purge_deleted_task_media.py | 161 +++++++++++ tests/test_database_backup_service.py | 12 + tests/test_media_storage_lifecycle.py | 345 ++++++++++++++++++++++++ tests/test_split_services.py | 16 +- tests/test_task_query_service.py | 5 +- 25 files changed, 1088 insertions(+), 95 deletions(-) create mode 100644 scripts/purge_deleted_task_media.py create mode 100644 tests/test_media_storage_lifecycle.py diff --git a/.env.example b/.env.example index cd2c69c..1bccb09 100644 --- a/.env.example +++ b/.env.example @@ -25,6 +25,8 @@ FFMPEG_TIMEOUT=600 # 如果以后要换存储盘,可以同步修改 docker-compose.yml 里的 volumes 和这里的路径。 STORAGE_ROOT=E:\直播间切片工作流存储 TASKS_DIR=E:\直播间切片工作流存储 +# 浏览器上传超过 1MB 时使用的临时目录;必须放在大容量存储盘。 +UPLOAD_TEMP_DIR=E:\直播间切片工作流存储\_临时上传 AI_DEFAULT_PROVIDER=remote AI_REQUEST_TIMEOUT_SECONDS=120 @@ -68,7 +70,7 @@ PUBLISH_SCHEDULER_INTERVAL_SECONDS=5 PUBLISH_DEFAULT_MODE=local_browser PUBLISH_JOB_STALE_MINUTES=30 PUBLISH_SCHEDULER_MAX_RETRY_COUNT=3 -PUBLISH_SCHEDULER_EXPORT_DIR= +PUBLISH_SCHEDULER_EXPORT_DIR=E:\直播间切片工作流存储\_发布包 PUBLISH_SCHEDULER_ALLOW_PUBLISH_WITHOUT_REVIEW=false PUBLISH_ENABLE_OPENCLI_FALLBACK=false # 请使用随机长字符串;start_publish_worker.ps1 会在本地 .env 缺失时自动生成。 diff --git a/DEVELOPMENT_LOG.md b/DEVELOPMENT_LOG.md index 229053e..be204d5 100644 --- a/DEVELOPMENT_LOG.md +++ b/DEVELOPMENT_LOG.md @@ -902,3 +902,13 @@ - “移出内容准备”仍是用户主动隐藏记录,“跳过任务”仍是终止当前任务;二者继续使用 `CANCELLED`,不会被普通取消逻辑误恢复。 - 兼容历史数据:数据库初始化时,只把旧版错误标记为“用户取消任务”的最后一条记录安全恢复到 `WAITING`;若同一切片和平台已经有活跃替代任务,则不会制造重复任务。 - 新增状态机、历史数据恢复和浏览器交互回归测试;专项测试 `19 passed`、完整测试 `336 passed`,Python 编译、两个 JavaScript 文件语法检查、Ruff 和差异检查均通过。 + +## 2026-08-02 E 盘统一视频存储与任务永久删除 + +- 新增 `UPLOAD_TEMP_DIR`,并在应用启动时把当前进程的 `TEMP`、`TMP` 和 Python 临时目录指向 E 盘;大于 1 MB 的浏览器上传不再先写入 C 盘系统临时目录。 +- 手动发布包默认目录从项目 `outputs/publish_packages` 调整到 `E:\直播间切片工作流存储\_发布包`;任务原片、音频、切片、字幕和封面继续统一使用 `TASKS_DIR`。 +- 任务列表“移入回收站”改为“永久删除”:只删除系统托管目录,外部 NAS / E 盘原片保留,数据库历史隐藏保留;运行中的处理和发布任务禁止删除。 +- 新增 `scripts/purge_deleted_task_media.py`,默认只预览;`--apply` 会先创建 SQLite 元数据备份,再清理已隐藏任务的 E 盘目录、发布包和精确匹配的旧版 C 盘任务目录。 +- 新增大文件 multipart 临时目录、失败上传回滚、外部原片保护、路径越界、运行中拦截、删除失败回滚、幂等删除和旧任务清理测试。 +- 已对真实数据先预演再执行清理:15 条已隐藏任务共删除 16 个托管目录,释放 `6,213,311,934` 字节(约 6.21 GB);4 条有效任务目录清理前后均完整,外部测试原片保留,清理后再次预演为 0 个残留目录。 +- 清理前 SQLite 元数据备份完整性检查为 `ok`;最终完整测试 `395 passed`,Python 编译、Ruff、JavaScript 语法和差异检查均通过。 diff --git a/NEXT_STEPS.md b/NEXT_STEPS.md index 179255f..f1d5557 100644 --- a/NEXT_STEPS.md +++ b/NEXT_STEPS.md @@ -847,3 +847,12 @@ 4. 之前由旧版“取消任务”产生、错误信息为“用户取消任务”的记录,会在数据库初始化时自动恢复;如果同一切片和平台已经存在新的活跃任务,则保留新的任务且不重复恢复。 5. 继续确认“移出内容准备”仍会隐藏任务,并可在执行记录中恢复;“跳过任务”仍保持终止状态,不会自动回到准备区。 6. 本次验收不要点击“立即发送”,不会触发抖音或 B站真实投稿。 + +## 2026-08-02 E 盘存储与永久删除验收 + +1. 重启本地后台,打开“系统状态”,确认“视频临时与导出目录”显示“E 盘就绪”,路径分别是 `_临时上传` 和 `_发布包`。 +2. 新建一个测试任务并上传视频,确认任务原片和后续切片只出现在 `E:\直播间切片工作流存储\任务名` 下;C 盘项目目录不应出现新的生产视频副本。 +3. 在任务列表点击“永久删除”,确认提示明确说明无法恢复且外部原片保留;删除后对应 E 盘任务目录和发布包应消失。 +4. 对 NAS 或 E 盘其他目录的引用任务执行删除时,只删除任务生成物,外部唯一原片必须仍然存在。 +5. 转写、切片或真实发布进行中时,删除应被拒绝并显示原因;等待任务结束后再删除。 +6. `scripts/purge_deleted_task_media.py` 默认只输出清单,只有显式带 `--apply` 才会永久清理,并在 `data/backups` 留下 SQLite 元数据备份。 diff --git a/app/core/config.py b/app/core/config.py index d1ef0df..383d705 100644 --- a/app/core/config.py +++ b/app/core/config.py @@ -60,6 +60,10 @@ class Settings: data_dir: Path = _env_path("DATA_DIR", PROJECT_ROOT / "data") storage_root: Path = _env_path("STORAGE_ROOT", EXTERNAL_STORAGE_ROOT) tasks_dir: Path = _env_path("TASKS_DIR", _env_path("STORAGE_ROOT", EXTERNAL_STORAGE_ROOT)) + upload_temp_dir: Path = _env_path( + "UPLOAD_TEMP_DIR", + _env_path("TASKS_DIR", _env_path("STORAGE_ROOT", EXTERNAL_STORAGE_ROOT)) / "_临时上传", + ) database_path: Path = _env_path( "DATABASE_PATH", _env_path("DATA_DIR", PROJECT_ROOT / "data") / "workflow.sqlite3", @@ -200,7 +204,7 @@ class Settings: publish_worker_allowed_roots: str = _env("PUBLISH_WORKER_ALLOWED_ROOTS", "") publish_scheduler_export_dir: Path = _env_path( "PUBLISH_SCHEDULER_EXPORT_DIR", - PROJECT_ROOT / "outputs" / "publish_packages", + _env_path("STORAGE_ROOT", EXTERNAL_STORAGE_ROOT) / "_发布包", ) publish_scheduler_allow_publish_without_review: bool = _env_bool( "PUBLISH_SCHEDULER_ALLOW_PUBLISH_WITHOUT_REVIEW", diff --git a/app/main.py b/app/main.py index 387a2ee..d826e7f 100644 --- a/app/main.py +++ b/app/main.py @@ -1,4 +1,6 @@ from contextlib import asynccontextmanager +import os +import tempfile from fastapi import FastAPI, Request, Response from fastapi.responses import FileResponse, JSONResponse @@ -8,6 +10,7 @@ from app.db.database import init_db from app.routers import ai_prompts, files, media, pages, publish, settings as settings_router, tasks from app.services.publish_scheduler import start_scheduler_background +from app.services.storage_service import configure_runtime_media_storage # /media 和 /static 的 Origin 白名单 @@ -43,6 +46,9 @@ def _build_allow_origin_header(origin: str) -> str: @asynccontextmanager async def lifespan(app: FastAPI): + previous_temp = tempfile.tempdir + previous_temp_env = {name: os.environ.get(name) for name in ("TEMP", "TMP")} + app.state.media_storage = configure_runtime_media_storage() init_db() scheduler = await start_scheduler_background() app.state.publish_scheduler = scheduler @@ -51,6 +57,12 @@ async def lifespan(app: FastAPI): finally: if scheduler: scheduler.stop() + tempfile.tempdir = previous_temp + for name, value in previous_temp_env.items(): + if value is None: + os.environ.pop(name, None) + else: + os.environ[name] = value app = FastAPI( diff --git a/app/routers/tasks.py b/app/routers/tasks.py index eb2f4f1..fdd239e 100644 --- a/app/routers/tasks.py +++ b/app/routers/tasks.py @@ -19,7 +19,13 @@ from app.services import task_service from app.services.ai_prompt_preset_service import update_task_ai_prompt_preset from app.services.pipeline_engine import start_auto_pipeline -from app.services.storage_service import allocate_task_dir_name, save_uploaded_video +from app.services.storage_service import ( + allocate_task_dir_name, + remove_failed_task_directory, + save_uploaded_video, + StorageSafetyError, +) +from app.services.task_lifecycle_service import TaskDeletionConflictError from app.services import job_service from app.services import job_worker @@ -67,40 +73,51 @@ async def create_upload_task( ) -> dict: task_id = uuid4().hex[:12] task_dir_name = allocate_task_dir_name(task_name, exclude_task_id=task_id) - saved_path = save_uploaded_video( - task_id, - video_file.filename or "source_video.mp4", - video_file.file, - task_dir_name=task_dir_name, - ) - payload = TaskCreate( - task_name=task_name, - source_type="upload", - platform=platform, - original_video_path=str(saved_path), - max_clip_duration=max_clip_duration, - candidate_clip_count=candidate_clip_count, - selection_profile=selection_profile, - final_clip_target=final_clip_target, - ai_preference=ai_preference, - auto_mode=auto_mode, - auto_clip_count=auto_clip_count, - auto_min_clip_seconds=auto_min_clip_seconds, - auto_max_clip_seconds=auto_max_clip_seconds, - auto_schedule_mode=auto_schedule_mode, - auto_schedule_start_at=auto_schedule_start_at, - auto_schedule_interval_hours=auto_schedule_interval_hours, - auto_schedule_daily_start_time=auto_schedule_daily_start_time, - auto_schedule_daily_end_time=auto_schedule_daily_end_time, - auto_metadata_use_ai=auto_metadata_use_ai, - ) + task_record_created = False try: + saved_path = await run_in_threadpool( + save_uploaded_video, + task_id, + video_file.filename or "source_video.mp4", + video_file.file, + task_dir_name, + ) + payload = TaskCreate( + task_name=task_name, + source_type="upload", + platform=platform, + original_video_path=str(saved_path), + max_clip_duration=max_clip_duration, + candidate_clip_count=candidate_clip_count, + selection_profile=selection_profile, + final_clip_target=final_clip_target, + ai_preference=ai_preference, + auto_mode=auto_mode, + auto_clip_count=auto_clip_count, + auto_min_clip_seconds=auto_min_clip_seconds, + auto_max_clip_seconds=auto_max_clip_seconds, + auto_schedule_mode=auto_schedule_mode, + auto_schedule_start_at=auto_schedule_start_at, + auto_schedule_interval_hours=auto_schedule_interval_hours, + auto_schedule_daily_start_time=auto_schedule_daily_start_time, + auto_schedule_daily_end_time=auto_schedule_daily_end_time, + auto_metadata_use_ai=auto_metadata_use_ai, + ) result = task_service.create_task_record(payload, task_id=task_id, task_dir_name=task_dir_name) + task_record_created = True if payload.auto_mode: result["auto_pipeline"] = start_auto_pipeline(task_id, background_tasks=background_tasks) return result except ValueError as exc: + if not task_record_created: + await run_in_threadpool(remove_failed_task_directory, task_id, task_dir_name) raise HTTPException(status_code=400, detail=str(exc)) from exc + except Exception: + if not task_record_created: + await run_in_threadpool(remove_failed_task_directory, task_id, task_dir_name) + raise + finally: + await video_file.close() @router.get("/{task_id}") @@ -131,8 +148,14 @@ async def get_ai_analysis_status(task_id: str) -> dict: async def delete_task(task_id: str) -> dict: try: return task_service.soft_delete_task(task_id) + except TaskDeletionConflictError as exc: + raise HTTPException(status_code=409, detail=str(exc)) from exc + except StorageSafetyError as exc: + raise HTTPException(status_code=409, detail=str(exc)) from exc except ValueError as exc: raise HTTPException(status_code=404, detail=str(exc)) from exc + except RuntimeError as exc: + raise HTTPException(status_code=500, detail=str(exc)) from exc @router.patch("/{task_id}/status") @@ -293,6 +316,23 @@ async def batch_update_clip_candidates( raise HTTPException(status_code=400, detail=str(exc)) from exc +@router.post("/{task_id}/clips/sync-publish") +async def sync_reviewed_clips_to_publish_center( + task_id: str, + payload: ClipCandidateBatchUpdate, +) -> dict: + try: + return await run_in_threadpool( + task_service.sync_reviewed_clips_to_publish_center, + task_id, + payload.clips, + ) + except (ValueError, FileNotFoundError) as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + except RuntimeError as exc: + raise HTTPException(status_code=500, detail=str(exc)) from exc + + @router.get("/{task_id}/clips/{clip_id}/transcript-excerpt") async def get_clip_transcript_excerpt( task_id: str, diff --git a/app/services/database_backup_service.py b/app/services/database_backup_service.py index 50d0db8..f868c51 100644 --- a/app/services/database_backup_service.py +++ b/app/services/database_backup_service.py @@ -17,6 +17,7 @@ PUBLISH_MIGRATION_JOURNAL_GLOB = f"{PUBLISH_MIGRATION_BACKUP_GLOB}-journal" PUBLISH_MIGRATION_BACKUP_COOLDOWN = timedelta(hours=24) PUBLISH_MIGRATION_BACKUP_KEEP_DAYS = 14 +MEDIA_CLEANUP_BACKUP_PREFIX = "workflow-before-media-cleanup-" BACKUP_TIMEZONE = ZoneInfo("Asia/Shanghai") @@ -262,3 +263,53 @@ def create_publish_migration_backup( backup_connection.close() if source_connection is not None: source_connection.close() + + +def create_media_cleanup_backup( + database_path: Path, + backup_dir: Path, + *, + now: datetime | None = None, +) -> Path: + """永久删除任务媒体前,原子创建一份仅包含 SQLite 元数据的备份。""" + database_path = database_path.resolve() + backup_dir = backup_dir.resolve() + now = now.astimezone(BACKUP_TIMEZONE) if now else datetime.now(BACKUP_TIMEZONE) + backup_dir.mkdir(parents=True, exist_ok=True) + + timestamp = now.strftime("%Y%m%d-%H%M%S-%f") + final_path = backup_dir / ( + f"{MEDIA_CLEANUP_BACKUP_PREFIX}{timestamp}-{os.getpid()}-{uuid4().hex[:8]}.sqlite3" + ) + temporary_path = final_path.with_name(f"{final_path.name}.tmp-{uuid4().hex}") + source_connection: sqlite3.Connection | None = None + backup_connection: sqlite3.Connection | None = None + try: + source_connection = sqlite3.connect( + f"{database_path.as_uri()}?mode=ro", + uri=True, + timeout=10, + ) + backup_connection = sqlite3.connect(str(temporary_path), timeout=10) + source_connection.backup(backup_connection) + backup_connection.close() + backup_connection = None + source_connection.close() + source_connection = None + + integrity = sqlite_quick_check(temporary_path) + if integrity != "ok": + raise BackupSafetyError(f"媒体清理前备份完整性检查失败:{integrity}") + os.replace(temporary_path, final_path) + return final_path + except Exception as exc: + if temporary_path.exists(): + temporary_path.unlink() + if isinstance(exc, BackupSafetyError): + raise + raise BackupSafetyError(f"创建媒体清理前备份失败:{exc}") from exc + finally: + if backup_connection is not None: + backup_connection.close() + if source_connection is not None: + source_connection.close() diff --git a/app/services/job_service.py b/app/services/job_service.py index 859d8d9..b5455ef 100644 --- a/app/services/job_service.py +++ b/app/services/job_service.py @@ -23,12 +23,14 @@ JOB_STATUS_RUNNING = "running" JOB_STATUS_COMPLETED = "completed" JOB_STATUS_FAILED = "failed" +JOB_STATUS_CANCELLED = "cancelled" JOB_STATUS_LABELS = { JOB_STATUS_QUEUED: "排队中", JOB_STATUS_RUNNING: "运行中", JOB_STATUS_COMPLETED: "已完成", JOB_STATUS_FAILED: "失败", + JOB_STATUS_CANCELLED: "已取消", } JOB_TYPE_LABELS = { diff --git a/app/services/job_worker.py b/app/services/job_worker.py index 2909cde..074ff5c 100644 --- a/app/services/job_worker.py +++ b/app/services/job_worker.py @@ -17,6 +17,8 @@ def execute_job(job_id: str) -> dict: job = job_service.get_job(job_id) if not job: raise ValueError(f"job 不存在:{job_id}") + if job.get("status") == job_service.JOB_STATUS_CANCELLED: + return job job_type = job.get("job_type") task_id = job.get("task_id") diff --git a/app/services/pipeline_engine.py b/app/services/pipeline_engine.py index 3bdb0f9..83c7ab8 100644 --- a/app/services/pipeline_engine.py +++ b/app/services/pipeline_engine.py @@ -143,6 +143,8 @@ def _get_task(self, task_id: str) -> dict: task = task_service.get_task(task_id, include_video_probe=False) if not task: raise ValueError("任务不存在") + if task.get("is_deleted"): + raise ValueError("任务已永久删除,已停止后续自动处理") return task def _load_auto_config(self, task: dict) -> dict: diff --git a/app/services/storage_service.py b/app/services/storage_service.py index aff4d1b..56065aa 100644 --- a/app/services/storage_service.py +++ b/app/services/storage_service.py @@ -1,7 +1,10 @@ +from dataclasses import dataclass from pathlib import Path, PureWindowsPath +import os import re import sqlite3 import shutil +import tempfile from typing import BinaryIO from uuid import uuid4 @@ -26,6 +29,62 @@ _PATH_TRAVERSAL_MARKERS = ("..", "~") +class StorageSafetyError(RuntimeError): + """存储路径不安全或不满足清理条件。""" + + +@dataclass(frozen=True) +class ManagedMediaTarget: + label: str + path: Path + + +@dataclass(frozen=True) +class TaskMediaCleanupPlan: + task_id: str + targets: tuple[ManagedMediaTarget, ...] + external_source_path: Path | None + + @property + def existing_targets(self) -> tuple[ManagedMediaTarget, ...]: + return tuple(target for target in self.targets if target.path.exists()) + + +@dataclass(frozen=True) +class TaskMediaCleanupResult: + deleted_paths: tuple[str, ...] + freed_bytes: int + external_source_preserved: bool + + +def _ensure_writable_directory(path: Path, label: str) -> Path: + try: + path.mkdir(parents=True, exist_ok=True) + probe_path = path / f".niuma-write-test-{uuid4().hex}" + probe_path.write_bytes(b"ok") + probe_path.unlink() + except OSError as exc: + raise RuntimeError(f"{label}不可用或不可写:{path};原因:{exc}") from exc + return path.resolve() + + +def configure_runtime_media_storage() -> dict[str, str]: + """准备大文件目录,并把当前应用进程的临时目录固定到存储盘。""" + tasks_dir = _ensure_writable_directory(settings.tasks_dir, "任务存储目录") + upload_temp_dir = _ensure_writable_directory(settings.upload_temp_dir, "上传临时目录") + export_dir = _ensure_writable_directory(settings.publish_scheduler_export_dir, "发布包目录") + + temp_value = str(upload_temp_dir) + os.environ["TEMP"] = temp_value + os.environ["TMP"] = temp_value + tempfile.tempdir = temp_value + return { + "tasks_dir": str(tasks_dir), + "upload_temp_dir": temp_value, + "publish_export_dir": str(export_dir), + } + + def _collect_allowed_roots() -> list[Path]: """收集所有允许访问的文件系统根目录。""" roots: list[Path] = [] @@ -107,6 +166,11 @@ def ensure_storage_root() -> Path: return settings.storage_root +def ensure_tasks_root() -> Path: + settings.tasks_dir.mkdir(parents=True, exist_ok=True) + return settings.tasks_dir + + def _storage_relative_parts(task_dir_name: str) -> tuple[str, ...]: return tuple(part for part in PureWindowsPath(task_dir_name).parts if part not in {"", "."}) @@ -159,14 +223,14 @@ def allocate_task_dir_name( base_name = sanitize_task_dir_name(task_name, fallback=exclude_task_id or "untitled") parent_parts = _storage_relative_parts(parent_dir_name or "") existing_names = _get_existing_task_dir_names(exclude_task_id=exclude_task_id) - root = ensure_storage_root().joinpath(*parent_parts) + root = ensure_tasks_root().joinpath(*parent_parts) root.mkdir(parents=True, exist_ok=True) for index in range(1, 1000): candidate_name = base_name if index == 1 else f"{base_name} ({index})" candidate_parts = (*parent_parts, candidate_name) relative_name = str(PureWindowsPath(*candidate_parts)) - candidate_path = ensure_storage_root().joinpath(*candidate_parts) + candidate_path = ensure_tasks_root().joinpath(*candidate_parts) if relative_name.lower() not in existing_names and not candidate_path.exists(): return relative_name @@ -308,10 +372,9 @@ def get_source_video_path(task: dict) -> Path | None: def save_uploaded_video(task_id: str, filename: str, file_object: BinaryIO, task_dir_name: str | None = None) -> Path: - create_task_directory(task_id, task_dir_name) - # 扩展名校验 _validate_upload_extension(filename) + create_task_directory(task_id, task_dir_name) safe_name = Path(filename or "source_video").name if not Path(safe_name).suffix: @@ -321,25 +384,176 @@ def save_uploaded_video(task_id: str, filename: str, file_object: BinaryIO, task # 流式写入 + 大小限制检查 max_size = settings.max_upload_size_bytes written = 0 - with output_path.open("wb") as target: - while True: - chunk = file_object.read(1024 * 1024) # 1MB chunks - if not chunk: - break - written += len(chunk) - if written > max_size: - # 删除已写入的部分 - try: - output_path.unlink() - except OSError: - pass - max_gb = max_size / (1024 * 1024 * 1024) - raise ValueError(f"上传文件超过大小限制({max_gb:.1f} GB)") - target.write(chunk) + try: + with output_path.open("wb") as target: + while True: + chunk = file_object.read(1024 * 1024) # 1MB chunks + if not chunk: + break + written += len(chunk) + if written > max_size: + max_gb = max_size / (1024 * 1024 * 1024) + raise ValueError(f"上传文件超过大小限制({max_gb:.1f} GB)") + target.write(chunk) + except Exception: + try: + output_path.unlink(missing_ok=True) + except OSError: + pass + raise return output_path +def remove_failed_task_directory(task_id: str, task_dir_name: str) -> None: + """仅清理本次尚未写入数据库的新任务目录。""" + if _fetch_task_dir_name(task_id): + return + task_dir = get_task_directory(task_id, task_dir_name) + tasks_root = settings.tasks_dir.resolve() + resolved = task_dir.resolve(strict=False) + try: + within_root = resolved.is_relative_to(tasks_root) + except AttributeError: # pragma: no cover - Python 3.8 兼容 + within_root = str(resolved).lower().startswith(str(tasks_root).lower() + os.sep) + if resolved == tasks_root or not within_root or task_dir.is_symlink(): + raise StorageSafetyError(f"拒绝清理不安全的任务目录:{task_dir}") + if task_dir.exists(): + shutil.rmtree(task_dir) + + +def _safe_relative_parts(value: str, label: str) -> tuple[str, ...]: + windows_path = PureWindowsPath(str(value or "").strip()) + parts = tuple(part for part in windows_path.parts if part not in {"", "."}) + if ( + not parts + or windows_path.is_absolute() + or windows_path.drive + or any(part in _PATH_TRAVERSAL_MARKERS for part in parts) + ): + raise StorageSafetyError(f"{label}包含不安全路径:{value}") + return parts + + +def _safe_managed_child(root: Path, parts: tuple[str, ...], label: str) -> Path: + resolved_root = root.resolve(strict=False) + candidate = root.joinpath(*parts) + resolved_candidate = candidate.resolve(strict=False) + try: + within_root = resolved_candidate.is_relative_to(resolved_root) + except AttributeError: # pragma: no cover - Python 3.8 兼容 + within_root = str(resolved_candidate).lower().startswith(str(resolved_root).lower() + os.sep) + if resolved_candidate == resolved_root or not within_root or candidate.is_symlink(): + raise StorageSafetyError(f"拒绝删除不安全的{label}:{candidate}") + return candidate + + +def _deduplicate_targets(targets: list[ManagedMediaTarget]) -> tuple[ManagedMediaTarget, ...]: + unique: list[ManagedMediaTarget] = [] + seen: set[str] = set() + for target in targets: + key = str(target.path.resolve(strict=False)).lower() + if key in seen: + continue + seen.add(key) + unique.append(target) + return tuple(unique) + + +def _path_is_within(path: Path, parent: Path) -> bool: + try: + return path.resolve(strict=False).is_relative_to(parent.resolve(strict=False)) + except (AttributeError, OSError, ValueError): + path_value = str(path.resolve(strict=False)).lower() + parent_value = str(parent.resolve(strict=False)).lower() + return path_value.startswith(parent_value + os.sep) + + +def build_task_media_cleanup_plan(task: dict, *, include_legacy: bool = True) -> TaskMediaCleanupPlan: + task_id = str(task.get("id") or "").strip() + task_id_parts = _safe_relative_parts(task_id, "任务 ID") + if len(task_id_parts) != 1: + raise StorageSafetyError(f"任务 ID 必须是单层目录名:{task_id}") + + task_dir_name = str(task.get("task_dir_name") or task_id) + task_parts = _safe_relative_parts(task_dir_name, "任务目录名") + task_dir = _safe_managed_child(settings.tasks_dir, task_parts, "任务目录") + targets = [ManagedMediaTarget("E 盘任务目录", task_dir)] + + export_dir = _safe_managed_child( + settings.publish_scheduler_export_dir, + task_id_parts, + "发布包目录", + ) + targets.append(ManagedMediaTarget("E 盘发布包目录", export_dir)) + + if include_legacy: + legacy_root = settings.project_root / "tasks" + legacy_values = [task_id] + if len(task_parts) == 1 and task_dir_name.lower() != task_id.lower(): + legacy_values.append(task_dir_name) + for legacy_value in legacy_values: + legacy_parts = _safe_relative_parts(legacy_value, "旧版任务目录名") + legacy_path = _safe_managed_child(legacy_root, legacy_parts, "旧版 C 盘任务目录") + targets.append(ManagedMediaTarget("旧版 C 盘任务目录", legacy_path)) + + managed_targets = _deduplicate_targets(targets) + source_path = get_source_video_path(task) + external_source_path = None + if source_path and source_path.exists(): + if not any(_path_is_within(source_path, target.path) for target in managed_targets): + external_source_path = source_path + + return TaskMediaCleanupPlan( + task_id=task_id, + targets=managed_targets, + external_source_path=external_source_path, + ) + + +def _directory_size_bytes(path: Path) -> int: + total = 0 + for child in path.rglob("*"): + try: + if child.is_file() and not child.is_symlink(): + total += child.stat().st_size + except OSError: + continue + return total + + +def task_media_cleanup_plan_size(plan: TaskMediaCleanupPlan) -> int: + return sum( + _directory_size_bytes(target.path) + for target in plan.existing_targets + if target.path.is_dir() and not target.path.is_symlink() + ) + + +def apply_task_media_cleanup_plan(plan: TaskMediaCleanupPlan) -> TaskMediaCleanupResult: + deleted_paths: list[str] = [] + freed_bytes = 0 + for target in plan.targets: + path = target.path + if not path.exists(): + continue + if path.is_symlink() or not path.is_dir(): + raise StorageSafetyError(f"拒绝删除异常的{target.label}:{path}") + size = _directory_size_bytes(path) + try: + shutil.rmtree(path) + except OSError as exc: + raise RuntimeError(f"删除{target.label}失败:{path};原因:{exc}") from exc + freed_bytes += size + deleted_paths.append(str(path)) + + return TaskMediaCleanupResult( + deleted_paths=tuple(deleted_paths), + freed_bytes=freed_bytes, + external_source_preserved=plan.external_source_path is not None, + ) + + def move_task_directory_to_trash(task_id: str, task_name: str, task_dir_name: str | None = None) -> tuple[str, Path]: current_dir_name = resolve_task_dir_name(task_id, task_dir_name) source_dir = get_task_directory(task_id, current_dir_name) diff --git a/app/services/task_lifecycle_service.py b/app/services/task_lifecycle_service.py index 7103e25..34a2323 100644 --- a/app/services/task_lifecycle_service.py +++ b/app/services/task_lifecycle_service.py @@ -8,10 +8,37 @@ from app.db.database import get_connection from app.models.task import TaskCreate, TaskStatus -from app.services.storage_service import allocate_task_dir_name, create_task_directory, validate_source_video_path +from app.services.storage_service import ( + apply_task_media_cleanup_plan, + allocate_task_dir_name, + build_task_media_cleanup_plan, + create_task_directory, + validate_source_video_path, +) from app.services.task_log_service import append_task_log +class TaskDeletionConflictError(RuntimeError): + """任务仍在执行,暂时不能删除其媒体文件。""" + + +ACTIVE_TASK_STATUSES = { + TaskStatus.CREATED.value, + TaskStatus.PREPARING_SOURCE.value, + TaskStatus.TRANSCRIBING.value, + TaskStatus.AI_ANALYZING.value, + TaskStatus.CLIP_SELECTING.value, + TaskStatus.VIDEO_CUTTING.value, + TaskStatus.METADATA_GENERATING.value, + TaskStatus.SCHEDULE_CREATING.value, + TaskStatus.PUBLISH_JOB_CREATING.value, + TaskStatus.audio_extracting.value, + TaskStatus.transcribing.value, + TaskStatus.ai_analyzing.value, + TaskStatus.cutting.value, +} + + def create_task_record(payload: TaskCreate, task_id: str | None = None, task_dir_name: str | None = None) -> dict: from app.services.task_service import _now_iso, get_status_label, STATUS_PROGRESS # noqa: F811 @@ -228,34 +255,109 @@ def update_task_selection_settings( } -def soft_delete_task(task_id: str) -> dict: +def delete_task_permanently(task_id: str) -> dict: from app.services.task_service import _now_iso, get_task # noqa: F811 task = get_task(task_id, include_video_probe=False) if not task: raise ValueError("任务不存在") - if task.get("is_deleted"): - return { - "message": "任务已隐藏,无需重复操作。", - "task_id": task_id, - "task_dir": task["task_dir"], - } - + cleanup_plan = build_task_media_cleanup_plan(task) + existing_target_count = len(cleanup_plan.existing_targets) now = _now_iso() with get_connection() as connection: - connection.execute( - """ - UPDATE tasks - SET is_deleted = 1, deleted_at = ?, updated_at = ? - WHERE id = ? - """, - (now, now, task_id), - ) - connection.commit() - - append_task_log(task_id, "任务已从列表隐藏,文件未删除") + try: + connection.execute("BEGIN IMMEDIATE") + current = connection.execute( + "SELECT status, COALESCE(is_deleted, 0) AS is_deleted FROM tasks WHERE id = ?", + (task_id,), + ).fetchone() + if not current: + raise ValueError("任务不存在") + + if not current["is_deleted"] and str(current["status"] or "") in ACTIVE_TASK_STATUSES: + raise TaskDeletionConflictError("任务正在处理,请等待处理结束后再永久删除。") + + conflicting_task = connection.execute( + """ + SELECT id + FROM tasks + WHERE id != ? AND COALESCE(is_deleted, 0) = 0 + AND LOWER(COALESCE(task_dir_name, id)) = LOWER(?) + LIMIT 1 + """, + (task_id, str(task.get("task_dir_name") or task_id)), + ).fetchone() + if conflicting_task: + raise TaskDeletionConflictError( + "该目录仍被另一条有效任务使用,已拒绝删除以避免误删。" + ) + + running_job = connection.execute( + "SELECT id FROM workflow_jobs WHERE task_id = ? AND status = 'running' LIMIT 1", + (task_id,), + ).fetchone() + if running_job: + raise TaskDeletionConflictError("任务仍有后台切片工作正在运行,请等待结束后再删除。") + + publishing_job = connection.execute( + "SELECT id FROM publish_jobs WHERE task_id = ? AND status = 'PUBLISHING' LIMIT 1", + (task_id,), + ).fetchone() + if publishing_job: + raise TaskDeletionConflictError("任务正在向平台发送视频,请等待发送结束后再删除。") + + cleanup_result = apply_task_media_cleanup_plan(cleanup_plan) + connection.execute( + """ + UPDATE workflow_jobs + SET status = 'cancelled', progress = 100, + message = '任务已永久删除,排队任务已取消', + error_message = '任务已永久删除', finished_at = ?, updated_at = ? + WHERE task_id = ? AND status = 'queued' + """, + (now, now, task_id), + ) + connection.execute( + """ + UPDATE publish_jobs + SET status = 'CANCELLED', scheduled_at = '', next_attempt_at = NULL, + error_code = 'task_deleted', error_message = '任务已永久删除', + last_error = '任务已永久删除', history_hidden = 1, + finished_at = ?, updated_at = ? + WHERE task_id = ? + AND status NOT IN ('PUBLISHED', 'EXPORTED', 'NEED_REVIEW', 'CANCELLED') + """, + (now, now, task_id), + ) + connection.execute( + """ + UPDATE tasks + SET is_deleted = 1, deleted_at = COALESCE(deleted_at, ?), updated_at = ? + WHERE id = ? + """, + (now, now, task_id), + ) + connection.commit() + except Exception: + connection.rollback() + raise + + status = "already_deleted" if task.get("is_deleted") and existing_target_count == 0 else "deleted" + freed_mb = cleanup_result.freed_bytes / (1024 * 1024) + if status == "already_deleted": + message = "任务已经永久删除,当前没有残留的任务视频文件。" + else: + message = f"任务已永久删除,共释放约 {freed_mb:.1f} MB;数据库历史记录已隐藏保留。" return { - "message": "任务已隐藏,原视频、切片和任务目录都已保留。", + "status": status, "task_id": task_id, - "task_dir": task["task_dir"], + "freed_bytes": cleanup_result.freed_bytes, + "external_source_preserved": cleanup_result.external_source_preserved, + "deleted_paths": list(cleanup_result.deleted_paths), + "message": message, } + + +def soft_delete_task(task_id: str) -> dict: + """兼容旧调用名称;实际执行永久媒体删除并保留隐藏数据库记录。""" + return delete_task_permanently(task_id) diff --git a/app/services/task_query_service.py b/app/services/task_query_service.py index d1f93ff..9cb420d 100644 --- a/app/services/task_query_service.py +++ b/app/services/task_query_service.py @@ -407,6 +407,12 @@ def get_system_status_context() -> dict: return { "storage_root": str(settings.storage_root), "storage_exists": settings.storage_root.exists(), + "tasks_dir": str(settings.tasks_dir), + "tasks_dir_exists": settings.tasks_dir.exists(), + "upload_temp_dir": str(settings.upload_temp_dir), + "upload_temp_dir_exists": settings.upload_temp_dir.exists(), + "publish_export_dir": str(settings.publish_scheduler_export_dir), + "publish_export_dir_exists": settings.publish_scheduler_export_dir.exists(), "database_path": str(settings.database_path), "database_exists": settings.database_path.exists(), "ffmpeg_path": ffmpeg_path or "未找到", diff --git a/app/static/js/app.js b/app/static/js/app.js index f491924..90c8726 100644 --- a/app/static/js/app.js +++ b/app/static/js/app.js @@ -1464,23 +1464,24 @@ document.querySelectorAll("[data-sync-publish-task]").forEach((button) => { document.querySelectorAll(".js-hide-task").forEach((button) => { button.addEventListener("click", async () => { const taskTitle = button.dataset.taskTitle || "这条任务"; - const confirmed = window.confirm(`确认把“${taskTitle}”移入 E 盘回收站吗?\n\n这会从列表隐藏任务,并把对应项目文件夹移动到 E:\\直播间切片工作流存储\\_回收站,不会删除原视频、切片文件和任务目录。`); + const confirmed = window.confirm(`确认永久删除“${taskTitle}”吗?\n\n系统会永久删除 E 盘任务目录内的原片副本、音频、转写、切片、字幕、封面和发布包,删除后无法恢复。\n\nNAS 或任务目录外的原始视频不会被删除。`); if (!confirmed) return; const originalText = button.textContent; button.disabled = true; - button.textContent = "移动中..."; + button.textContent = "删除中..."; try { const response = await fetch(`/api/tasks/${button.dataset.taskId}`, { method: "DELETE" }); const data = await response.json(); if (!response.ok) { - throw new Error(data.detail || "移入回收站失败"); + throw new Error(data.detail || "永久删除失败"); } - window.alert(data.message || "任务已移入回收站。"); + const externalNotice = data.external_source_preserved ? "\n\n任务目录外的原始视频已保留。" : ""; + window.alert(`${data.message || "任务已永久删除。"}${externalNotice}`); window.location.reload(); } catch (error) { - window.alert(`移入回收站失败:${error.message}`); + window.alert(`永久删除失败:${error.message}`); } finally { button.disabled = false; button.textContent = originalText; diff --git a/app/templates/new_task.html b/app/templates/new_task.html index 0cd86bb..16c5d3a 100644 --- a/app/templates/new_task.html +++ b/app/templates/new_task.html @@ -6,7 +6,7 @@

New Task

新建任务

-

上传本机视频后,系统会为它创建独立任务目录,并保存到 E 盘工作流存储目录。

+

上传临时文件、原片副本和后续切片都会直接保存在 E 盘工作流存储目录,不占用 C 盘视频空间。

@@ -29,7 +29,7 @@

基本信息

上传本机视频 - 选择视频后,会复制到该任务的 source 目录。 + 选择视频后,会直接写入 E 盘任务的 source 目录;上传临时文件也在 E 盘。
- + {% block extra_scripts %}{% endblock %} diff --git a/app/templates/clip_review.html b/app/templates/clip_review.html index 7ec99fa..309c3be 100644 --- a/app/templates/clip_review.html +++ b/app/templates/clip_review.html @@ -180,7 +180,7 @@

审核操作

@@ -188,8 +188,8 @@

审核操作

去字幕推送 + {% if output_clips %} - 查看本任务发送内容 {% endif %} diff --git a/app/templates/publish.html b/app/templates/publish.html index fe45fe9..f59bb91 100644 --- a/app/templates/publish.html +++ b/app/templates/publish.html @@ -12,7 +12,7 @@ {% endmacro %} {% macro content_row(job) %} -
+
@@ -22,6 +22,11 @@
裁剪片段{{ job.output_file_name or job.title }}
+ + + 已排期 + + {{ "内容完整" if job.content_complete else "缺少:" ~ (job.missing_fields|join("、")) }} @@ -49,7 +54,7 @@ - +
@@ -228,6 +233,13 @@

{{ group.task_name }}

+

Tasks

抖音任务清单

未排期任务也会保留在清单中,勾选后可批量设置时间。
计划时间视频平台账号状态操作
@@ -339,4 +351,4 @@

全部执行记录

{% endblock %} -{% block extra_scripts %}{% endblock %} +{% block extra_scripts %}{% endblock %} diff --git a/app/templates/task_detail.html b/app/templates/task_detail.html index d64b613..dcac926 100644 --- a/app/templates/task_detail.html +++ b/app/templates/task_detail.html @@ -10,28 +10,18 @@

任务详情 · {{ task.title }}

任务 ID:{{ task.id }} 写入时间:{{ task.created_at }} 主题:{{ task.title }} - 当前状态:{{ task.status_label }} + 当前状态:{{ task.status_label }}

{% set header_transcript_progress = task.transcript_progress or {} %} -
+
{% if task.auto_mode %} - {% if task.status.startswith("FAILED_") %} - - {% elif task.status in ["pending_review", "completed", "completed_with_errors", "failed"] %} - - {% elif task.status in ["READY_TO_PUBLISH", "COMPLETED"] %} - 前往发送中心 - {% else %} - - {% endif %} - {% if task.candidate_count > 0 %} - 检查候选片段 - {% endif %} - {% if task.output_clip_count > 0 %} - - 发送中心 · {{ publish_link_state.label }} - {% endif %} + + + 前往发送中心 + + 检查候选片段 + {% elif task.transcript_exists %} @@ -61,7 +51,7 @@

任务详情 · {{ task.title }}

data-running="{{ 'true' if task.status in ['CREATED', 'PREPARING_SOURCE', 'TRANSCRIBING', 'AI_ANALYZING', 'CLIP_SELECTING', 'VIDEO_CUTTING', 'METADATA_GENERATING', 'SCHEDULE_CREATING', 'PUBLISH_JOB_CREATING'] else 'false' }}" > 全自动模式已接管这个任务。 -

系统会自动推进音频提取、转写、AI 选片和视频切割;页面在处理中会自动刷新。只有失败或历史任务中断时才需要点击重试/继续。

+

系统会自动推进音频提取、转写、AI 选片和视频切割;状态概览与运行日志每 3 秒自动更新,不会刷新整张页面。只有失败或历史任务中断时才需要点击重试/继续。

{% endif %} @@ -78,31 +68,31 @@

基础信息

视频来源
{{ task.source_type_label }}
视频时长
{{ task.duration }}
视频大小
{{ task.video_size }}
-
候选片段
{{ task.candidate_count }} 条
+
候选片段
{{ task.candidate_count }} 条
最长切片
{{ task.max_clip_duration }} 分钟
创建时间
{{ task.created_at }}
-
更新时间
{{ task.updated_at }}
+
更新时间
{{ task.updated_at }}
源文件状态
{{ "可访问" if task.source_exists else "不可访问 / 未选择" }}
-
切片输出
{{ task.output_clip_count }} 条
+
切片输出
{{ task.output_clip_count }} 条
任务目录
{{ task.task_dir }}
-
+

Status

状态概览

- {{ task.status_label }} + {{ task.status_label }}
- +
- {{ task.progress }}% + {{ task.progress }}%
    {% for step in workflow_steps %} -
  1. +
  2. {{ step.index }} {{ step.name }}
  3. @@ -113,6 +103,7 @@

    状态概览

    当前阶段 异常 / 待处理 +

    正在连接任务进度…

@@ -286,7 +277,7 @@

运行日志

待刷新 -
点击 AI 分析后,这里会显示最新运行日志。
+
任务开始后,这里会自动显示最新运行日志。

diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 34d8457..46dd52f 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -4,7 +4,9 @@ ### 1.1 架构形态 -当前 v1.5.0 保持 **FastAPI 单体应用 + SQLite**。视频、AI、页面和调度器仍在同一个应用中;只有必须使用宿主系统 Chrome 的真实发布动作由 Windows Worker 执行,不引入 Redis、Celery 或微服务。 +当前 v2.0.0 保持 **FastAPI 单体应用 + SQLite + Windows 发布 Worker**。视频、AI、页面和调度器仍在同一个应用中;只有必须使用宿主系统 Chrome 的真实发布动作由 Windows Worker 执行,不引入 Redis、Celery 或微服务。 + +v2.0 的架构目标不是云端多租户,而是把一台 Windows 电脑上的长视频生产与发布链路做完整、可恢复、可审计。SQLite 是唯一业务事实来源,E 盘任务目录保存大文件,浏览器 Profile 和平台登录态只保留在本机且不进入 Git。 ```text ┌─────────────────────────────────────────────────────────┐ @@ -243,21 +245,26 @@ Worker 会把容器内 `/workspace/tasks/...` 映射到宿主 `.env` 的 `TASKS_ ```text 新建任务表单 → 上传视频 / 选择 NAS 路径 +→ E 盘任务目录与安全文件边界 → FFmpeg 提取音频 → 转写(火山引擎远程 / faster-whisper 本地) -→ AI 候选片段分析(DeepSeek / Ollama) -→ 候选片段人工审核(启用/禁用/编辑时间) -→ FFmpeg 自动切割 → 05_clips/ +→ AI 候选片段分析(DeepSeek / OpenAI-compatible / Ollama) +→ 候选片段人工审核(启用/禁用/编辑标题、摘要与时间) +→ 保存审核结果并按需生成安全的新切片版本 +→ FFmpeg 自动切割 + 文案 + 封面帧 → 全自动模式跳过字幕生成/烧录 -→ 发送中心内容准备与排期 +→ 发送中心内容准备、北京时间排期月历与执行记录 → Scheduler + Windows Worker 真实投稿(抖音 + B站) +→ PUBLISHED / FAILED / NEED_REVIEW 可追溯终态 ``` +任务详情通过轻量 `live-status` 接口每 3 秒局部更新,不重新加载整页;发送中心重新切片时保留旧执行证据,只让当前激活切片进入新的内容准备和排期。 + --- ## 8. 架构演进路线 -### 8.1 当前阶段:P2-1(已完成) +### 8.1 当前阶段:v2.0 本地生产闭环(已完成) - FastAPI 单体应用 - SQLite 单文件数据库 @@ -265,19 +272,22 @@ Worker 会把容器内 `/workspace/tasks/...` 映射到宿主 `.env` 的 `TASKS_ - FFmpeg 同步本地处理 - 本地/远程 AI Provider - 抖音/B站统一真实发布、人工复核和显式手动导出 +- 全自动任务状态轮询、失败续跑和片段审核同步 +- 内容准备、跨午夜排期、最晚排期续接、月历详情与执行记录 +- E 盘统一生产存储、外部原片保护和托管产物安全删除 - 代码检查与 CI 流程 -### 8.2 短期演进(P2-2 ~ P2-3) +### 8.2 v2.0 后续重点 -**目标**:不改变单体形态,增强本地可靠性和安全边界。 +**目标**:不扩大单用户本地范围,优先用真实素材与真实账号完成灰度验收并提高可靠性。 | 方向 | 具体措施 | | --- | --- | -| **数据库** | SQLite 保持不变,增加 WAL 模式、备份脚本 | -| **Job Worker** | 引入本地后台 Job Queue(`threading` / `asyncio`),将耗时任务(转写、AI 分析、切割)异步化,不阻塞 HTTP 请求 | -| **安全增强** | `LOCAL_ADMIN_TOKEN` 鉴权强化,敏感操作确认对话框,操作日志记录 | -| **存储** | 支持 NAS 路径作为存储根目录,`STORAGE_ROOT` 可配置为网络路径 | -| **错误恢复** | 任务失败后可从中断点重试,而非从头开始 | +| **真实灰度发布** | 抖音、B站各用单条低风险素材验证当前页面选择器、成功证据和人工复核路径 | +| **长期稳定性** | 连续运行 Scheduler、Docker Watcher 和 Windows Worker,观察中断恢复与日志完整性 | +| **数据保护** | 定期验证 SQLite 备份、E 盘空间、外部原片保护和永久删除清单 | +| **内容质量** | 用真实长视频继续校准综艺笑点与通用模式的候选质量、文案和封面时间点 | +| **平台维护** | 平台页面改版后更新 Publisher 选择器,不通过绕过验证或静默重传维持“成功率” | ### 8.3 中期演进(P3) @@ -315,7 +325,7 @@ Worker 会把容器内 `/workspace/tasks/...` 映射到宿主 `.env` 的 `TASKS_ | 暂不做的 | 原因 | | --- | --- | | **SaaS 多租户** | 当前是个人本地工具,不需要租户隔离和计费系统 | -| **真实全自动发布** | 平台有验证码、风控、登录失效,全自动不可行也不安全 | +| **绕过验证的无人值守发布** | 平台有验证码、风控、登录失效;v2.0 只在登录有效且平台无需人工确认时自动执行 | | **强依赖云部署** | 首版定位 Windows 本地工具,不应强制要求云服务器 | | **移动端 App** | 核心工作流依赖 FFmpeg 和大文件处理,不适合移动端 | | **实时直播流处理** | 当前是录播后处理,实时流需要完全不同的技术栈 | diff --git a/docs/DATABASE_SCHEMA.md b/docs/DATABASE_SCHEMA.md index 17666d1..b56194b 100644 --- a/docs/DATABASE_SCHEMA.md +++ b/docs/DATABASE_SCHEMA.md @@ -408,7 +408,7 @@ data/workflow.sqlite3 ## 2026-06-09 v1.2 补充说明 -- 以上 v1.2 说明仅是历史记录。v1.5.0 已由 `PublishScheduler` 执行到期任务,并通过 Windows Worker 调用抖音/B站 Publisher。 +- 以上 v1.2 说明仅是历史记录。v2.0.0 已由 `PublishScheduler` 执行到期任务,并通过 Windows Worker 调用抖音/B站 Publisher。 - 平台发送不绕过验证码、登录失效、风控和人工确认;结果不确定写 `NEED_REVIEW`。 - 代码中仍存在兼容性 `clips` 子目录(`TASK_SUBDIRECTORIES` 同时包含 `clips` 和 `05_clips`),新任务的正式输出目录是 `05_clips`。旧 `clips` 目录为兼容保留,不建议删除。 # 2026-06-23:v1.4.0 定时发送字段 @@ -417,7 +417,7 @@ data/workflow.sqlite3 - 旧字段继续兼容:`output_clip_id` 等同于 `clip_id`,`description` 等同于 `caption`,`tags` 等同于 `hashtags`,`video_file_path` 等同于 `video_path`,`provider_response` 兼容 `publish_result`,`retry_count` 兼容 `attempt_count`。 - 发布状态使用:`DRAFT`、`SCHEDULED`、`WAITING`、`PUBLISHING`、`PUBLISHED`、`FAILED`、`CANCELLED`、`NEED_REVIEW`。 - 调度器只扫描 `status = SCHEDULED` 且 `scheduled_at <= 当前时间` 的任务;`NEED_REVIEW`、`CANCELLED`、`PUBLISHED` 不会自动发布。 -- 该 2026-06-23 版本曾默认使用 `manual_export`,2026-07-11 曾改为 `opencli_publish`;v1.5.0 当前默认是 `local_browser`。发布包导出成功写 `EXPORTED` 且不写 `published_at`;只有平台确认提交成功才写 `PUBLISHED` 和 `published_at`。 +- 该 2026-06-23 版本曾默认使用 `manual_export`,2026-07-11 曾改为 `opencli_publish`;v2.0.0 当前默认是 `local_browser`。发布包导出成功写 `EXPORTED` 且不写 `published_at`;只有平台确认提交成功才写 `PUBLISHED` 和 `published_at`。 - 没有 `scheduled_at` 的旧手动发送任务迁移为 `WAITING`,避免被自动调度器误执行。 ## 2026-07-27:取消发送状态兼容 diff --git a/docs/DEPLOYMENT.md b/docs/DEPLOYMENT.md index f4dca5a..38f6162 100644 --- a/docs/DEPLOYMENT.md +++ b/docs/DEPLOYMENT.md @@ -2,38 +2,38 @@ ## 1. 当前部署模式 -v1.3 支持两种部署方式: +v2.0.0 支持两种运行方式。日常使用推荐 Docker Desktop;真实抖音 / B站投稿无论使用哪种方式,都必须由 Windows 主机上的 Chrome Worker 执行。 -### 方式 A:Windows 本地直接运行(推荐日常使用) +### 方式 A:Docker Desktop + Windows Worker(推荐日常使用) ``` 你的 Windows 电脑 -├── Python 3.12(系统安装) -├── FFmpeg(系统安装,需在 PATH 中) -├── 项目代码(任意目录) -├── .venv(Python 虚拟环境) +├── Docker Desktop 运行 niuma-studio +│ └── FastAPI + SQLite Scheduler + FFmpeg +├── NiuMa Studio Docker Watcher +│ └── 容器运行时自动启动 Windows Chrome Worker +├── 系统 Chrome 独立账号目录 └── 浏览器打开 http://127.0.0.1:8001 ``` -**适用场景**:日常使用、开发调试、单机处理。 +**适用场景**:日常处理、排期和真实灰度发布。平时不需要手动打开 PowerShell。 -### 方式 B:Docker 容器运行(推荐测试/隔离环境) +### 方式 B:Windows 本地 Python 直接运行(开发 / 诊断) ``` 你的 Windows 电脑 -├── Docker Desktop -├── 项目代码(任意目录) -├── 容器 niuma-studio -│ ├── Python 3.12 + FFmpeg(容器内预装) -│ └── uvicorn 监听 8001 端口 +├── Python 3.12 + .venv +├── FFmpeg(系统安装,需在 PATH 中) +├── uvicorn 监听 8001 端口 +├── scripts/publish_host_worker.py(真实投稿时) └── 浏览器打开 http://127.0.0.1:8001 ``` -**适用场景**:不想装 Python/FFmpeg、测试环境隔离、CI 验证。 +**适用场景**:代码开发、自动化测试、发布 Worker 诊断。 --- -## 2. 方式 A 详细步骤:Windows 本地直接运行 +## 2. Windows 本地 Python 方式详细步骤 ### 2.1 环境要求 @@ -105,7 +105,7 @@ http://127.0.0.1:8001 --- -## 3. 方式 B 详细步骤:Docker 运行 +## 3. Docker 方式详细步骤 ### 3.1 环境要求 @@ -113,7 +113,7 @@ http://127.0.0.1:8001 ### 3.2 配置环境变量 -与方式 A 相同,先 `copy .env.example .env` 并填写配置。 +与本地 Python 方式相同,先 `copy .env.example .env` 并填写配置。 ### 3.3 构建并启动 @@ -147,7 +147,7 @@ docker compose down - **存储目录**:`docker-compose.yml` 默认将 `E:\直播间切片工作流存储` 挂载到容器内 `/workspace/tasks`。如果你的存储目录在其他位置,请修改 `docker-compose.yml` 中的 `volumes` 配置。 - **代码热更新**:`app/` 和 `prompts/` 目录以 volume 方式挂载,修改代码后容器自动重载。 - **Ollama 连接**:如果 Ollama 在宿主机运行,容器内通过 `http://host.docker.internal:11434/v1` 访问。 -- **opencli 桥接**:容器内通过 `http://host.docker.internal:8765` 访问宿主机上的 opencli 桥接服务。 +- **Windows 发布 Worker**:容器内通过 `http://host.docker.internal:8765` 访问宿主机受 Token 保护的 Chrome Worker;旧 opencli 仅是默认关闭的兼容模式。 --- @@ -157,9 +157,9 @@ docker compose down | 局限 | 说明 | 影响 | | --- | --- | --- | -| **单进程** | API 和视频处理在同一进程 | 处理大视频时页面可能卡住(请求阻塞) | -| **无后台队列** | 没有独立 Worker 进程 | 转写、AI 分析、切割都在请求线程中同步执行 | -| **单机存储** | 任务产物必须在本地磁盘 | 不能跨机器共享任务数据 | +| **单业务进程** | API、全自动流水线和 Scheduler 在同一个 FastAPI 进程 | 进程重启会中断正在处理的步骤,需从失败点重试或续跑 | +| **无外部任务队列** | 使用 FastAPI BackgroundTasks 和轻量 SQLite Job 记录,没有 Redis / Celery | 适合个人单机,不适合多机并发处理 | +| **单机主存储** | 生产任务目录由当前 Windows 主机 / Docker 挂载负责 | 可引用 NAS 原片,但不支持多台处理机共同写同一任务 | | **无负载均衡** | 不支持多实例部署 | 只能一个人用,不能横向扩展 | | **无 HTTPS** | 只有 HTTP | 只适合本地使用,不要暴露到公网 | | **单用户** | 没有登录和用户隔离 | 谁打开浏览器都能操作所有任务 | diff --git a/docs/PROJECT_GUIDE.md b/docs/PROJECT_GUIDE.md index 4e13cf2..a8aaa48 100644 --- a/docs/PROJECT_GUIDE.md +++ b/docs/PROJECT_GUIDE.md @@ -2,26 +2,47 @@ 这份文档给不熟悉代码和终端的新手使用。你只需要按顺序做,不需要理解每一行命令背后的原理。 +适用版本:牛马片场 `2.0.0`。 + ## 1. 项目是做什么的 项目中文名:牛马片场。 -它的目标是把一条本地直播录像、综艺访谈或长视频素材变成一组可审核、可切割、可继续加字幕和发送的短视频候选片段。 +它的目标是把一条本地直播录像、综艺访谈或长视频素材,变成一组可审核、可切割、可准备文案和封面、可排期并发送到抖音或 B站的短视频内容。 -当前 MVP 主要流程: +当前 v2.0 完整流程: ```text -上传本地视频 --> 创建任务 +上传本地视频 / 选择 NAS 或本地已有视频 +-> 创建独立任务目录 -> 提取音频 --> 本地语音转写 --> AI 分析候选片段 +-> 火山引擎或 faster-whisper 转写 +-> DeepSeek / OpenAI-compatible 或 Ollama 分析候选片段 -> 人工审核片段 --> 自动切割输出短视频 +-> 自动切割并准备标题、简介、话题和封面 +-> 发送中心核对账号与内容 +-> 立即发送或按北京时间排期 +-> Windows Chrome Worker 投稿抖音 / B站 +-> 保存成功、失败或人工复核记录 ``` 针对《康熙来了》类综艺,可以在新建任务或任务详情里选择“综艺笑点优先”。该模式先按重叠窗口找笑点,再补齐前后文,最后全局去重和评分;候选池默认 12 条,但只会默认启用最多 5 条 A 级内容,质量不足时不会凑数。现有直播和通用长视频继续使用“通用内容价值”模式。 +### 1.1 现在已经能做什么 + +- 素材、音频、转写、AI 结果、切片和封面按任务保存到 E 盘存储目录。 +- 全自动任务详情会每 3 秒更新进度与日志,失败后可以重试或继续。 +- 片段审核保存后,可一键生成最新切片并同步到发送中心。 +- 发送中心按任务分组管理抖音 / B站内容,支持内容补齐、封面、排期月历、跨午夜时间窗和续接最晚排期。 +- 立即发送与定时发送共用同一套 Scheduler 和 Windows Worker,执行记录不会因为重新切片而被覆盖。 + +### 1.2 现在仍需要人工做什么 + +- 首次使用抖音或 B站账号时,在系统 Chrome 独立窗口完成登录、二维码、短信或平台要求的验证。 +- 真实发送前逐条核对视频、标题、简介、话题、封面、账号、可见范围和北京时间。 +- 遇到验证码、滑块、登录失效、平台风控或发布结果不确定时,到平台创作者中心人工确认。 +- 项目不会绕过平台限制,也不会在结果不确定时自动重复上传。 + ## 2. 当前项目目录 ```text @@ -148,13 +169,16 @@ http://127.0.0.1:8001/health 启动项目后,按这个顺序检查: -1. 打开首页,确认工作台能显示。 -2. 打开“新建任务”,上传一个本地视频。 -3. 进入任务详情页,确认任务信息和处理按钮能显示。 -4. 生成转写后,确认详情页能看到转写预览。 -5. 点击远程 AI 或本地 AI 分析,生成候选片段。 -6. 打开片段审核页,勾选或修改候选片段。 -7. 触发切割,确认输出片段记录能展示。 +1. 打开首页,确认左侧显示 `v2.0 本地高光生产版`。 +2. 打开“新建任务”,上传一个短测试视频,保持“全自动流程”开启。 +3. 进入任务详情页,确认状态概览和运行日志每 3 秒局部更新,页面不会自动跳回顶部。 +4. 转写与 AI 分析完成后,打开片段审核页,勾选或修改候选片段。 +5. 点击“保存并同步发送中心”,确认需要时会生成最新切片,并自动进入当前任务的内容准备区。 +6. 在发送中心核对抖音 / B站的视频、标题、简介、话题和封面。 +7. 只使用测试任务预览排期,确认月历日期详情和任务清单按北京时间从早到晚排列。 +8. 如果要测试真实发布,必须先确认 Windows Worker 正常、平台账号已登录,并只发送一条低风险测试内容。 + +> 只检查生产链路时,不要点击“立即发送”。排期预览不会触发真实投稿,确认排期后任务会等待调度器到点执行。 ## 8. 命令行测试 diff --git a/docs/ROADMAP.md b/docs/ROADMAP.md index daa3be4..cc6f8c6 100644 --- a/docs/ROADMAP.md +++ b/docs/ROADMAP.md @@ -8,7 +8,7 @@ - **P3**:分布式与多机协作 - **P4+**:多用户与平台化 -当前版本:**v1.3.0**(分支整合稳定版) +当前版本:**v2.0.0**(Windows 本地高光生产闭环) --- @@ -52,53 +52,53 @@ - [x] CI 流程(GitHub Actions) - [x] Git 安全规范文档 ---- +### v1.3 ~ v1.5 — 自动流水线与真实发布架构 -## 进行中 +- [x] 全自动任务流水线、失败步骤重试和历史任务续跑 +- [x] SQLite Scheduler、Publisher Registry 与北京时间 / UTC 统一 +- [x] Windows Chrome Worker、抖音 Publisher 与 B站 Publisher +- [x] `PUBLISHED` / `FAILED` / `NEED_REVIEW` 保守终态和执行事件 +- [x] 内容准备、排期计划、执行记录和显式 `manual_export` -### P2-2 — 架构演进准备 +### v2.0 — 分支收拢与本地生产闭环 -- [ ] 架构文档更新([ARCHITECTURE.md](ARCHITECTURE.md)) -- [ ] 路线图文档(本文档) -- [ ] 部署文档([DEPLOYMENT.md](DEPLOYMENT.md)) -- [ ] 明确短期/中期/长期演进方向 -- [ ] 明确暂不做的边界 +- [x] E 盘统一生产存储、临时上传、发布包与安全永久删除 +- [x] 综艺笑点优先 / 通用内容价值两种 AI 选片模式 +- [x] AI 封面时间点、历史内容封面补齐与双平台复用 +- [x] 片段审核保存、按需重新切片并同步发送中心 +- [x] 任务详情每 3 秒局部更新进度、日志与可用操作 +- [x] 发送内容按原始任务分组,已排期状态清晰可见 +- [x] 跨午夜排期、续接平台最晚排期、日期详情和时间升序 +- [x] 取消发送返回内容准备,旧执行证据继续保留 +- [x] 项目说明、架构、流程、路线图和新手文档统一到 2.0 --- -## 计划中 - -### P2-3 — 本地 Job Worker - -**目标**:将耗时任务异步化,不阻塞 HTTP 请求。 +## 进行中 -| 任务 | 说明 | -| --- | --- | -| 引入 `asyncio` 后台任务 | 转写、AI 分析、切割在后台执行 | -| 任务进度实时推送 | SSE 或轮询,前端可看到实时进度 | -| 任务队列管理 | 支持任务排队、取消、重试 | -| 并发控制 | 限制同时进行的 FFmpeg/转写/AI 任务数 | +### v2.0 真实环境验收 -### P2-4 — 本地安全增强 +- [ ] 抖音使用单条低风险素材完成登录、上传、表单填写、封面、提交和成功证据验收 +- [ ] B站使用单条低风险素材完成同等灰度验收 +- [ ] 登录失效、验证码、风控和结果不确定场景都能进入 `NEED_REVIEW` +- [ ] 连续运行 Docker Watcher、Scheduler 与 Windows Worker,验证中断恢复和日志完整性 +- [ ] 用真实长视频继续校准候选质量、文案、封面时间点和排期体验 -**目标**:加固单用户本地使用的安全边界。 +--- -| 任务 | 说明 | -| --- | --- | -| 操作确认对话框 | 删除、清空回收站等敏感操作需二次确认 | -| 操作日志 | 记录关键操作到 `audit_log` 表(本地 SQLite) | -| Token 管理 | `LOCAL_ADMIN_TOKEN` 过期/刷新机制 | -| 文件校验 | 上传文件 MIME 类型和大小校验 | +## 计划中 -### P2-5 — 存储灵活性 +### v2.x — 本地可靠性与安全增强 -**目标**:让存储根目录可以灵活配置,支持 NAS。 +**目标**:保持单用户 Windows 本地架构,优先提高可靠性,不急于引入分布式组件。 | 任务 | 说明 | | --- | --- | -| NAS 路径验证 | 启动时检查 `STORAGE_ROOT` 是否可达 | -| 网络路径支持 | UNC 路径(`\\NAS\share`)读写测试 | -| 空间检查 | 新建任务前检查磁盘剩余空间 | +| 数据备份 | 增加面向新手的一键 SQLite 备份、完整性检查和恢复说明 | +| 磁盘预警 | 新建任务前检查 E 盘剩余空间,避免长视频处理中途失败 | +| 处理并发 | 明确 FFmpeg、转写和 AI 任务的本机并发上限与排队提示 | +| 操作追踪 | 为永久删除、人工确认发布结果等高风险动作补充本地审计记录 | +| 平台维护 | 平台页面变化时更新选择器和成功证据,不绕过平台验证 | --- diff --git a/docs/TASK_FLOW.md b/docs/TASK_FLOW.md index 858020e..261c1f5 100644 --- a/docs/TASK_FLOW.md +++ b/docs/TASK_FLOW.md @@ -36,7 +36,7 @@ pending_video ← 任务已创建,尚未上传视频 - `completed` / `completed_with_errors` 代表"自动切割阶段结束",不是平台发布完成。 - 字幕和发布是独立于主任务状态的后续工作流。 -## 3. v1.3.0 全自动任务状态流 +## 3. v2.0.0 全自动任务状态流 `auto_mode=true` 的任务使用独立的大写状态,不破坏原有手动流程: @@ -78,7 +78,9 @@ FAILED_PUBLISH_JOB_CREATING - 自动选片数量读取 `tasks.candidate_clip_count`,时长上限读取 `tasks.max_clip_duration`;旧自动数量和最小/最大秒数只保留兼容,不再参与新任务决策。 - 切片输出仍写入 `05_clips/`,并写入 `output_clip`;单个切片失败不会阻断其他成功切片生成文案和发布任务。 - `SCHEDULE_CREATING` 当前表示整理发送队列,不再自动计算发布时间。 -- 发布任务先以 `WAITING` / `NEED_REVIEW` 创建;用户在发送中心批量设置时间后进入 `SCHEDULED`,再由 v1.4.0 调度器执行。 +- 发布任务先以 `WAITING` / `NEED_REVIEW` 创建;用户在发送中心批量设置时间后进入 `SCHEDULED`,再由当前 `PublishScheduler` 到点执行。 +- 任务详情调用 `GET /api/tasks/{task_id}/live-status`,处理中每 3 秒局部更新状态、进度、10 步时间线、日志、候选/输出数量和可用操作;完成或失败后停止轮询。 +- 网络短暂中断只显示重试提示,不会整页刷新,也不会清空正在编辑但尚未保存的 AI Prompt。 ## 4. 失败流转 @@ -150,7 +152,8 @@ audio/source.wav ## 8. 自动切割阶段 -- 用户在片段审核页点击"生成切片"后,任务进入 `cutting`。 +- 用户在片段审核页点击“生成切片”时,页面会先保存当前启用状态、标题、摘要和出入点,再让任务进入 `cutting`。 +- 点击“保存并同步发送中心”会先比较当前启用候选与激活成片;有变化、文件缺失或数量不一致时生成新版本,完全一致时只执行幂等同步。 - 所有启用片段都切割成功时,任务进入 `completed`。 - 至少一个片段成功、同时存在失败片段时,任务进入 `completed_with_errors`。 - 所有片段都失败时,任务进入 `failed`。 @@ -214,6 +217,8 @@ output_clip 生成成功 ### 排期与立即发送 - 浏览器提交北京时间 `start_at_local`;后端按 `Asia/Shanghai` 应用每日开始/结束窗口,跨日后顺延到次日开始时间。 +- 排期支持续接当前平台最晚未来任务;抖音和 B站独立计算,选中但尚未保存的任务不会被当成已有排期重复计算。 +- 排期月历可展开某天全部任务;当天详情和主任务清单都按北京时间升序,未排期任务保留在已排期任务之后。 - `scheduled_at` 统一存 UTC ISO 8601,API 同时返回 `scheduled_at_utc` 与 `scheduled_at_local`。 - 自动调度只读取到期的 `SCHEDULED`;`NEED_REVIEW` 即使有时间也不能执行。 - “立即发送”允许 `DRAFT`、`WAITING`、`SCHEDULED`,只把 `scheduled_at` 更新为当前 UTC 并唤醒 Scheduler;不直接调用 opencli 或平台页面。 diff --git a/docs/UI_REFERENCE.md b/docs/UI_REFERENCE.md index cd2caf8..1142b5d 100644 --- a/docs/UI_REFERENCE.md +++ b/docs/UI_REFERENCE.md @@ -1,5 +1,21 @@ # UI 参考说明 +## v2.0 界面总览 + +- 左侧主导航保持“工作台、新建任务、任务列表、片段审核、字幕推送、发送中心、系统状态”七个入口,侧栏状态统一显示“v2.0 本地高光生产版”。 +- 任务详情是处理过程的事实入口:全自动任务每 3 秒局部更新状态、10 步进度、日志、候选数、输出数和重试 / 继续 / 前往发送中心操作,不整页刷新。 +- 片段审核负责最终选片与出入点确认;“保存并同步发送中心”把保存、必要时重新切片和内容同步合并成一个清楚动作。 +- 发送中心固定分为“内容准备、排期计划、执行记录”:内容准备核对素材和文案,排期计划管理北京时间与月历,执行记录保留成功、失败、取消、导出和人工复核证据。 +- v2.0 继续使用 Apple 风格浅色界面、蓝色主强调、轻量玻璃卡片和充足留白;状态颜色必须服务于“处理中、完成、警告 / 人工处理”,不能只做装饰。 + +## 2026-08-03 更新:排期日期详情与时间升序 + +- `/publish` 的“排期计划”月历保留每格前两条预览和“另有 N 条”,有排期的日期格整体增加悬停、键盘焦点和蓝色选中状态,空日期保持不可操作。 +- 点击日期后,在月历网格下方展开浅蓝玻璃卡片,标题显示日期、当前平台和任务数量;完整列表按北京时间升序展示时间、视频标题、账号、状态与“定位到任务”。 +- 点击当天详情中的任务会滚动并高亮下方原任务,所有发送、清除排期和取消操作仍只保留在主任务行,避免出现两套状态不同步的操作入口。 +- 主任务清单使用时间线语义:有效排期在前并从早到晚排列,未排期统一置后且保留原顺序;动态调整后原地重排。 +- 小屏下当天详情压缩为“时间 + 任务信息”两列,隐藏重复的定位提示;不新增前端框架、接口或数据库字段。 + ## 2026-08-02 更新:续接当前平台最晚排期 - 排期抽屉在“第 1 条发布时间”下方新增次级按钮“接在当前平台最晚排期后”,保持手动日期时间输入框可编辑。 @@ -8,6 +24,13 @@ - 每日窗口默认显示 `07:00 → 00:00`,两个时间控件仍可修改;跨午夜时段的视觉说明明确 00:00 代表次日午夜。 - 按钮、提示和输入框继续使用现有浅色抽屉、蓝色强调和紧凑表单样式,不引入新的前端框架。 +## 2026-08-02 更新:内容准备已排期视觉标记 + +- `/publish` 内容准备卡片继续保留已排期任务,但 `SCHEDULED` 卡片右上角新增蓝色“已排期 + 具体北京时间”胶囊标签。 +- 已排期卡片使用轻蓝渐变底色、蓝色描边、左侧强调线和克制阴影;普通未排期卡片继续使用原白色玻璃卡片,形成清楚但不突兀的状态区分。 +- 已排期卡片底部按钮显示“调整排期”,未排期时显示“加入发布计划”;新增、清除或取消排期后,标签、卡片颜色和按钮文案原地同步。 +- 标签复用现有 `SCHEDULED` 状态与北京时间,不调整页面结构、不新增数据库字段,也不改变抖音/B站发布边界。 + ## 2026-08-01 更新:综艺笑点优先选片与反馈 - 新建任务页把“候选池上限”和“最终启用目标”拆成两个控件,并新增“通用内容价值 / 综艺笑点优先”模式选择;弱集少选的规则在控件下直接说明。 @@ -442,3 +465,10 @@ v1.3.0 起,新建任务页可勾选“新建后自动跑完整流水线”。 - 确认提示会说明视频、标题、简介、话题和封面均被保留;操作期间按钮禁用,避免重复请求。 - 操作成功后页面自动切换到“内容准备”,展开任务所在分组并定位到返回的卡片;排期时间清空为“未排期”,用户可直接修改内容或重新设置排期。 - “取消发送并返回准备”“移出内容准备”“跳过任务”保持三种独立语义:前者回到准备区,移出操作隐藏任务,跳过操作终止任务。 + +## 2026-08-02 更新:任务详情进度局部自动刷新 + +- 全自动任务详情页的“状态概览”卡片每 3 秒原地更新进度条、百分比和 10 个地铁路线节点,不刷新整张页面。 +- 顶部当前状态、基础信息中的候选/输出数量和右侧运行日志同步更新;切回浏览器标签页时立即读取一次最新状态。 +- 状态卡底部显示“自动刷新中 / 状态已更新 / 自动更新暂时中断”等轻量提示,沿用蓝色、绿色和黄色的现有状态语义。 +- 流程完成后原地显示发送中心和片段检查入口;失败时对应步骤变黄并显示重试入口。用户正在编辑的 AI Prompt、页面滚动位置和其他未保存内容保持不变。 diff --git a/scripts/opencli_host_bridge.py b/scripts/opencli_host_bridge.py index 7aa9628..1042bbf 100644 --- a/scripts/opencli_host_bridge.py +++ b/scripts/opencli_host_bridge.py @@ -142,7 +142,7 @@ def log_message(self, format: str, *args) -> None: def main() -> None: - # 兼容旧启动命令,但实际启动 v1.5 的受保护发布 Worker。 + # 兼容旧启动命令,但实际启动 v2.0 的受保护发布 Worker。 from scripts.publish_host_worker import main as worker_main worker_main() diff --git a/scripts/publish_host_worker.py b/scripts/publish_host_worker.py index ac944ca..151effd 100644 --- a/scripts/publish_host_worker.py +++ b/scripts/publish_host_worker.py @@ -148,7 +148,7 @@ def _resolve_media_path(raw_value: str, *, required: bool) -> str: def create_worker_app(token: str | None = None) -> FastAPI: worker_token = str(token if token is not None else settings.publish_worker_token) - worker = FastAPI(title="NiuMa Studio Publish Worker", version="1.5.0") + worker = FastAPI(title="NiuMa Studio Publish Worker", version="2.0.0") def require_token(authorization: str = Header(default="")) -> None: if not worker_token: diff --git a/scripts/start_docker_opencli.ps1 b/scripts/start_docker_opencli.ps1 index 0348f35..f18a395 100644 --- a/scripts/start_docker_opencli.ps1 +++ b/scripts/start_docker_opencli.ps1 @@ -7,7 +7,7 @@ $ErrorActionPreference = "Stop" $ProjectRoot = Resolve-Path (Join-Path $PSScriptRoot "..") Set-Location $ProjectRoot -Write-Host '此兼容脚本现在会启动 v1.5 Windows Chrome 发布 Worker。' +Write-Host '此兼容脚本现在会启动 v2.0 Windows Chrome 发布 Worker。' & (Join-Path $PSScriptRoot 'start_publish_worker.ps1') -Port $BridgePort Write-Host '正在刷新 Docker 服务:http://127.0.0.1:8001' diff --git a/tests/test_auto_pipeline.py b/tests/test_auto_pipeline.py index 76cbf57..bdf226a 100644 --- a/tests/test_auto_pipeline.py +++ b/tests/test_auto_pipeline.py @@ -1,4 +1,4 @@ -"""v1.3.0 全自动任务流水线最小测试。""" +"""v2.0.0 全自动任务流水线回归测试。""" from __future__ import annotations @@ -16,7 +16,8 @@ from app.services.auto_publish_service import create_auto_publish_jobs from app.services.pipeline_engine import PipelineEngine, build_schedule_times from app.services.storage_service import get_artifact_paths -from app.services.task_lifecycle_service import create_task_record +from app.services.task_lifecycle_service import create_task_record, update_task_status +from app.services.task_log_service import append_task_log from app.services.task_service import get_task @@ -460,3 +461,81 @@ def test_auto_resume_endpoint_starts_from_clip_selection(monkeypatch): assert starter.call_args.args[0] == task["id"] assert starter.call_args.kwargs["start_step"] == TaskStatus.CLIP_SELECTING assert starter.call_args.kwargs["background_tasks"] is not None + + +def test_live_status_endpoint_tracks_running_auto_pipeline(): + task = _create_auto_task("test-auto-live-running") + update_task_status(task["id"], TaskStatus.AI_ANALYZING) + append_task_log(task["id"], "实时状态测试日志") + + with TestClient(app) as client: + response = client.get(f"/api/tasks/{task['id']}/live-status", headers=_headers()) + + assert response.status_code == 200 + payload = response.json() + assert payload["status"] == TaskStatus.AI_ANALYZING.value + assert payload["status_label"] == "AI 分析中" + assert payload["progress"] == 45 + assert payload["is_running"] is True + assert payload["should_poll"] is True + assert payload["actions"]["primary"] == "processing" + assert len(payload["workflow_steps"]) == 10 + assert payload["workflow_steps"][3]["state"] == "current" + assert any("实时状态测试日志" in line for line in payload["log_lines"]) + + +def test_live_status_endpoint_returns_failure_and_retry_action(): + task = _create_auto_task("test-auto-live-failed") + update_task_status(task["id"], TaskStatus.FAILED_AI_ANALYZING, "模型服务暂时不可用") + + with TestClient(app) as client: + response = client.get(f"/api/tasks/{task['id']}/live-status", headers=_headers()) + + assert response.status_code == 200 + payload = response.json() + assert payload["should_poll"] is False + assert payload["error_message"] == "模型服务暂时不可用" + assert payload["actions"]["primary"] == "retry" + assert payload["workflow_steps"][3]["state"] == "warning" + + +def test_live_status_endpoint_returns_completed_actions(monkeypatch): + task = _create_auto_task("test-auto-live-completed") + update_task_status(task["id"], TaskStatus.COMPLETED) + monkeypatch.setattr("app.services.task_service.count_clip_candidates", lambda _task_id: 3) + monkeypatch.setattr("app.services.task_service.count_output_clips", lambda _task_id: 2) + + with TestClient(app) as client: + response = client.get(f"/api/tasks/{task['id']}/live-status", headers=_headers()) + + assert response.status_code == 200 + payload = response.json() + assert payload["progress"] == 100 + assert payload["should_poll"] is False + assert payload["counts"] == {"candidates": 3, "outputs": 2} + assert payload["actions"] == {"primary": "publish", "review": True, "publish": True} + assert all(step["state"] == "done" for step in payload["workflow_steps"]) + + +def test_live_status_endpoint_returns_404_for_missing_task(): + with TestClient(app) as client: + response = client.get("/api/tasks/test-auto-missing/live-status", headers=_headers()) + + assert response.status_code == 404 + assert response.json()["detail"] == "任务不存在" + + +def test_task_detail_live_status_frontend_uses_partial_refresh(): + template = (settings.project_root / "app" / "templates" / "task_detail.html").read_text(encoding="utf-8") + script = (settings.project_root / "app" / "static" / "js" / "app.js").read_text(encoding="utf-8") + live_script = script[ + script.index("function renderTaskLiveStatus"): + script.index("async function pollAiAnalysisStatus") + ] + + assert "data-task-live-overview" in template + assert "data-live-task-actions" in template + assert "/live-status" in live_script + assert "TASK_LIVE_STATUS_INTERVAL_MS = 3000" in script + assert 'document.addEventListener("visibilitychange"' in live_script + assert "window.location.reload" not in live_script diff --git a/tests/test_clip_review_publish_sync.py b/tests/test_clip_review_publish_sync.py new file mode 100644 index 0000000..a4cdc36 --- /dev/null +++ b/tests/test_clip_review_publish_sync.py @@ -0,0 +1,72 @@ +from __future__ import annotations + +from unittest.mock import Mock + +from app.main import app +from app.models.task import ClipCandidateBatchItem +from app.services import task_service + + +def _payload() -> list[ClipCandidateBatchItem]: + return [ + ClipCandidateBatchItem( + id="clip-review-sync-001", + title="值得发送的片段", + start_time="00:01:00", + end_time="00:02:00", + enabled=True, + summary="完整笑点片段", + ) + ] + + +def test_review_sync_regenerates_when_active_outputs_do_not_match(monkeypatch) -> None: + save = Mock(return_value={"changed_count": 0, "message": "已保存", "clips": [], "task": {}}) + cut = Mock( + return_value={ + "publish_sync": { + "status": "ok", + "message": "发送中心同步完成。", + "link_state": {"linked_count": 1}, + } + } + ) + monkeypatch.setattr(task_service, "update_clip_candidates_batch", save) + monkeypatch.setattr(task_service, "_active_outputs_match_enabled_candidates", lambda _task_id: False) + monkeypatch.setattr(task_service, "process_task_video_cuts", cut) + + result = task_service.sync_reviewed_clips_to_publish_center("task-review-sync", _payload()) + + assert result["regenerated"] is True + assert result["link_state"]["linked_count"] == 1 + assert "重新生成最新切片" in result["message"] + save.assert_called_once() + cut.assert_called_once_with("task-review-sync") + + +def test_review_sync_reuses_matching_outputs_when_review_is_unchanged(monkeypatch) -> None: + save = Mock(return_value={"changed_count": 0, "message": "已保存", "clips": [], "task": {}}) + sync = Mock( + return_value={ + "status": "ok", + "message": "发送中心同步完成。", + "link_state": {"linked_count": 1}, + } + ) + monkeypatch.setattr(task_service, "update_clip_candidates_batch", save) + monkeypatch.setattr(task_service, "_active_outputs_match_enabled_candidates", lambda _task_id: True) + monkeypatch.setattr("app.services.publish_service.sync_task_publish_jobs", sync) + + result = task_service.sync_reviewed_clips_to_publish_center("task-review-sync", _payload()) + + assert result["regenerated"] is False + assert "无需重复生成" in result["message"] + sync.assert_called_once_with( + "task-review-sync", + prefer_subtitled=False, + restore_removed=True, + ) + + +def test_review_sync_route_is_available() -> None: + assert "/api/tasks/{task_id}/clips/sync-publish" in app.openapi()["paths"] diff --git a/tests/test_publish_center_browser.py b/tests/test_publish_center_browser.py index eabcd52..be5cc76 100644 --- a/tests/test_publish_center_browser.py +++ b/tests/test_publish_center_browser.py @@ -74,6 +74,7 @@ def test_publish_center_schedule_preview_confirm_and_export(monkeypatch, tmp_pat init_db() _cleanup() douyin_jobs = [_seed_job(tmp_path, index) for index in range(1, 11)] + unscheduled = _seed_job(tmp_path, 0) first = douyin_jobs[0] newest = douyin_jobs[-1] bilibili = _seed_job(tmp_path, 11, "bilibili") @@ -92,11 +93,23 @@ def test_publish_center_schedule_preview_confirm_and_export(monkeypatch, tmp_pat connection.commit() generated_cover = tmp_path / "browser-batch-cover.jpg" generated_cover.write_bytes(b"fake-cover") + preserved_cover = tmp_path / "browser-preserved-cover.jpg" + preserved_cover.write_bytes(b"existing-cover") + with get_connection() as connection: + connection.execute( + "UPDATE publish_jobs SET cover_file_path = ? WHERE id = ?", + (str(preserved_cover), unscheduled), + ) + connection.commit() def fake_backfill_covers(platform=None): with get_connection() as connection: rows = connection.execute( - "SELECT id FROM publish_jobs WHERE task_id LIKE ? AND status = 'WAITING' AND platform = ?", + """ + SELECT id FROM publish_jobs + WHERE task_id LIKE ? AND status = 'WAITING' AND platform = ? + AND TRIM(COALESCE(cover_file_path, '')) = '' + """, (f"{PREFIX}%", platform), ).fetchall() connection.execute( @@ -104,6 +117,7 @@ def fake_backfill_covers(platform=None): UPDATE publish_jobs SET cover_mode = 'time', cover_time_seconds = 30, cover_file_path = ? WHERE task_id LIKE ? AND status = 'WAITING' AND platform = ? + AND TRIM(COALESCE(cover_file_path, '')) = '' """, (str(generated_cover), f"{PREFIX}%", platform), ) @@ -123,6 +137,7 @@ def fake_backfill_covers(platform=None): monkeypatch.setattr(publish_service, "backfill_missing_publish_covers", fake_backfill_covers) future_start = datetime.now() + timedelta(days=2) future_day = future_start.strftime("%Y-%m-%d") + future_day_label = f"{future_start.year} 年 {future_start.month} 月 {future_start.day} 日" following_day = (future_start + timedelta(days=1)).strftime("%Y-%m-%d") port = _free_port() server = uvicorn.Server( @@ -223,6 +238,9 @@ def fake_backfill_covers(platform=None): assert page.locator('[name="daily_start_time"]').input_value() == "07:00" assert page.locator('[name="daily_end_time"]').input_value() == "00:00" page.locator("[data-use-latest-schedule]").click() + page.locator("[data-latest-schedule-note]").filter( + has_text=f"本次第 1 条:{future_day} 22:00" + ).wait_for() assert page.locator('[name="start_at_local"]').input_value() == f"{future_day}T22:00" assert page.locator("[data-latest-schedule-note]").inner_text() == ( f"当前最晚:{future_day} 19:00;本次第 1 条:{future_day} 22:00" @@ -263,6 +281,37 @@ def fake_backfill_covers(platform=None): assert f"{future_day} 06:00" in page.locator( f'[data-publish-row][data-section="schedule"][data-job-id="{first}"] [data-row-schedule]' ).inner_text() + visible_test_rows = page.locator( + f'[data-publish-row][data-section="schedule"][data-job-id^="{PREFIX}"]:visible' + ) + assert visible_test_rows.evaluate_all("rows => rows.map((row) => row.dataset.jobId)") == [ + *douyin_jobs, + unscheduled, + ] + + calendar_day = page.locator(f'[data-calendar-date="{future_day}"]') + assert calendar_day.locator(".calendar-job-chip").all_inner_texts() == [ + "06:00 浏览器测试片段 1", + "09:00 浏览器测试片段 2", + ] + assert calendar_day.locator(".calendar-job-more").inner_text() == "另有 4 条" + calendar_day.click() + day_detail = page.locator("[data-calendar-day-detail]") + assert day_detail.is_visible() + assert day_detail.locator("[data-calendar-day-title]").inner_text() == future_day_label + assert "6 条排期" in day_detail.locator("[data-calendar-day-summary]").inner_text() + assert day_detail.locator(".publish-calendar-day-item time").all_inner_texts() == [ + "06:00", "09:00", "12:00", "15:00", "18:00", "21:00", + ] + day_detail.locator(f'[data-calendar-detail-job="{first}"]').click() + assert "is-calendar-focus" in page.locator( + f'[data-publish-row][data-section="schedule"][data-job-id="{first}"]' + ).get_attribute("class") + page.locator("[data-calendar-day-close]").click() + assert day_detail.is_hidden() + calendar_day.focus() + calendar_day.press("Enter") + assert day_detail.is_visible() dialogs = [] page.on("dialog", lambda dialog: (dialogs.append(dialog.message), dialog.accept())) @@ -281,6 +330,13 @@ def fake_backfill_covers(platform=None): assert "已取消发送并返回内容准备" in page.locator("#send-center-message").inner_text() page.locator('[data-center-tab="schedule"]').click() + assert page.locator( + f'[data-publish-row][data-section="schedule"][data-job-id^="{PREFIX}"]:visible' + ).evaluate_all("rows => rows.map((row) => row.dataset.jobId)") == [ + *douyin_jobs[:-1], + newest, + unscheduled, + ] page.locator(f'[data-publish-row][data-section="schedule"][data-job-id="{first}"] [data-publish-now]').click() assert any("抖音" in message for message in dialogs) page.locator("#send-center-message").filter(has_text="统一调度").wait_for() diff --git a/tests/test_publish_task_grouping.py b/tests/test_publish_task_grouping.py index 0184718..2e2be48 100644 --- a/tests/test_publish_task_grouping.py +++ b/tests/test_publish_task_grouping.py @@ -155,6 +155,37 @@ def test_publish_page_renders_task_identity_without_full_source_path(tmp_path: P assert "移出内容准备" in html +def test_scheduled_content_card_has_clear_schedule_marker(tmp_path: Path) -> None: + task_id, _ = _insert_task(tmp_path, "已排期标记任务", created_at=_time(-5)) + job_id, _, _ = _insert_clip_job( + tmp_path, + task_id, + platform="douyin", + status="SCHEDULED", + created_at=_time(-4), + ) + + response = TestClient(app).get("/publish") + html = response.text + card_start = html.index(f'data-job-id="{job_id}"') + card_end = html.index("
", card_start) + card_html = html[card_start:card_end] + + assert response.status_code == 200 + assert 'data-status="SCHEDULED"' in card_html + assert "data-content-schedule" in card_html + assert "已排期" in card_html + assert "调整排期" in card_html + + +def test_schedule_marker_updates_without_page_reload() -> None: + script = (settings.project_root / "app" / "static" / "js" / "publish-center.js").read_text(encoding="utf-8") + + assert 'row.classList.toggle("is-scheduled", isScheduled)' in script + assert "contentScheduleBadge.hidden = !isScheduled" in script + assert 'isScheduled ? "调整排期" : "加入发布计划"' in script + + def test_dismiss_keeps_files_and_other_platform_and_blocks_recreation(tmp_path: Path) -> None: task_id, _ = _insert_task(tmp_path, "安全移出任务", created_at=_time(-5)) douyin_job, clip_id, clip_path = _insert_clip_job(