diff --git a/.env.example b/.env.example index f165010..340165e 100644 --- a/.env.example +++ b/.env.example @@ -108,6 +108,14 @@ TAILSCALE_API_KEY=请填写Tailscale_API_key # 设置后不启动 cron 调度器(跑测试脚本时避免真实触发任务、占用设备) # DISABLE_SCHEDULER=1 +# 视频发布计划的素材目录(默认 data/videos)。 +# **自动化测试必须指到临时目录**,否则测试传的视频会落进用户真实数据目录。 +# DATA_VIDEO_DIR=/tmp/videos_test + +# 单次上传体积上限(MB,默认 2048)。 +# ⚠ 必须大于「应用管理」里最大的 APK(实测有 336MB 的包),否则 APK 上传会被一起卡死。 +# MAX_UPLOAD_MB=2048 + # ============================================================================== # 已废弃的历史配置(STF 已从代码层摘除,下列键代码不再使用,保留仅为兼容旧 .env) diff --git a/.gitignore b/.gitignore index 5e588ca..f478b1f 100644 --- a/.gitignore +++ b/.gitignore @@ -12,6 +12,8 @@ env/ logs/*.log data/*.log data/apks/*.apk +# 视频发布计划的素材(几百 MB 一个,绝不入库) +data/videos/ data/*.migrated # uiauto.pid 之类的运行时 PID 文件(**注释必须单独一行**:gitignore 不支持行尾注释) data/*.pid diff --git a/README.md b/README.md index 984aab6..8100149 100644 --- a/README.md +++ b/README.md @@ -47,13 +47,14 @@ | **设备池** | SQLite 清单 + adb 在线状态;支持手工添加、网段自动发现、一键重连、型号采集、启用/停用 | | **电量监控** | 后台每分钟 `dumpsys battery`(只读)采一次电量,**大屏卡片**与**监控页设备列表**都显示(按档位着色 + ⚡ 充电中);低于阈值推 webhook 告警(**充电中不报、"插着却没充电"照报**,掉档/恢复各报一次,不刷屏)。电量只放内存不落库(见 [NOTIFY.md](doc/NOTIFY.md) §3.1) | | **任务调度** | 手动 / cron 定时 / 定时启停;运行窗口;失败重试(含端口耗尽类的长退避) | -| **步骤编辑器** | 可视化拖拽编排 22 种步骤(含循环/条件/OCR/通知/录制回放),可打包成"自定义动作"复用 | +| **步骤编辑器** | 可视化拖拽编排 24 种步骤(含循环/条件/OCR/通知/录制回放),可打包成"自定义动作"复用 | | **拟人化操作** | 滑动默认走弧线轨迹、位置/幅度/时长每次抖动,且**每台设备有自己的手速与习惯**(同设备风格稳定、设备之间明显不同)——批量跑时不像同步机器人(见 [TASK_DEV.md](doc/TASK_DEV.md) §8.4) | | **录制回放** | 「录制手势」步骤 = **纯录制**:**用手指在真机上划**(读 getevent 真触屏)或在网页画面上拖,把轨迹点列原样存下来,回放时**照原路径与时长重放**——不套滑动那套方向/幅度/拟人参数 | | **动作配置 / 录制** | 「任务 → 动作配置」:① 配**步骤默认值**(新建步骤预填:滑动时长/幅度/抖动/拟人、点击超时、等待区间…)② **录制动作**——在设备画面上点/划/按键,自动翻译成步骤(点元素记选择器、划一下记方向/幅度/时长),可存成可复用动作,也可只取一个手势回填滑动步骤 | | **公共巡检** | 任务级的守护条件(**独立于步骤画布**):每 N 秒检查屏幕亮/熄、元素在不在、前台是不是某 App,命中就点亮/息屏/停本设备/推通知——通知标题正文自己写(见 [TASK_DEV.md](doc/TASK_DEV.md) §4.5) | | **去重账本** | 解决"**同一个号被做两次、有的号还没做**":任务里用「条件判断→去重」+「记为已做」,把"这个身份做过了"记进一张**所有设备共享**的账本——多台手机、反复重跑都不会重复做同一个号。「任务 → 去重记录」页能看"谁做过了、还差谁",也能删记录让它重跑(见 [TASK_DEV.md](doc/TASK_DEV.md) §4.6) | | **账号台账** | 「账号」页维护"**哪台设备上登着哪些号**"(设备号/手机号/账号名/抖音号/注册时间/卡在机内/可发视频/简介/备注):支持**从 Excel 粘贴批量导入**(表头识别、逐行结果、先预览后写入)。三个用处:① 任务的「条件判断」**直接从这里取号**,不用把几十个号手写进比对值;② 手机端 Agent 的**身份大字页**显示本机账号;③ 设备池里一眼看到每台登记了几个号(见 [DATA_MODEL.md](doc/DATA_MODEL.md) §2.10、[TASK_DEV.md](doc/TASK_DEV.md) §4.2) | +| **视频发布计划** | 「账号 → 发布计划」页:批量上传视频(文件名 `手机号_日期_编号`)与标题 txt → 自动按 (手机号,日期,编号) 配对到台账账号 → **按日期的时间线**看每天谁要发、发到哪一步。**平台只负责把素材推到手机**(推到 `/sdcard/DCIM/rp/`、触发相册刷新、把标题写进剪贴板),**抖音里怎么发由你自己在任务里写步骤**(步骤顺序:`推送发布视频` → 你的发布步骤 → `标记发布结果`;「发布计划」页顶部有**「发布任务」**块可一键建好骨架 + 看下次运行/启停,`输入文字` 那一步可以选**取值来源=计划标题**自动填文案并回读校验);发完可抓作品的**分享链接**存下来(平台不长期囤视频,链接才是长期资产,也方便后续铺评论:可一键复制、按日期/账号导出 CSV)。状态区分"**可重试**"(推送阶段失败,还没到抖音)与"**结果未知、需人工确认**"(推送之后出的岔子,绝不自动重发)(见 [DATA_MODEL.md](doc/DATA_MODEL.md) §2.11、[TASK_DEV.md](doc/TASK_DEV.md) §4.7) | | **元素抓取** | 拉取设备 UI 元素树 → 点选回填选择器;支持"点一下"与"测选择器"真机验证 | | **实时看屏** | MJPEG 实时流 + 点击/滑动/按键/文字输入;全屏监控大屏(`/wall`) | | **应用管理** | APK 上传/解析/批量安装;设备已装应用与版本查询;剪贴板注入 | @@ -141,6 +142,7 @@ auto_control/ │ ├── device_discovery.py # 网段扫描发现设备 → 待连接池 │ ├── device_battery.py # 设备电量采集(dumpsys 只读)+ 低电量告警 │ ├── ledger.py # 账号台账(CRUD / 表格粘贴解析 / 给任务取号 / 设备端账号块) +│ ├── video_plan.py # 视频发布计划(文件名解析配对 / 时间线 / 幂等占位 / 素材清理) │ ├── u2_helper.py # uiautomator2 通用辅助(等首页/安全点击等) │ ├── uiauto_helper.py # uiautodev 客户端(元素树 + XPath 建议) │ ├── ocr.py # 屏幕 OCR(RapidOCR,条件判断用) @@ -160,7 +162,7 @@ auto_control/ ├── tasks/ # 任务定义层 │ ├── base.py # BaseTask + _TASK_TYPES + register_task │ └── generic/ # 通用步骤任务(task_type=generic_steps,当前唯一类型) -│ └── task.py # STEP_TYPES(22 种步骤)+ Worker + 执行器 +│ └── task.py # STEP_TYPES(24 种步骤)+ Worker + 执行器 │ ├── mcp_server/ # MCP Server(20 个 de_* 工具,:8033) │ ├── mcp_server.py # 工具定义 + 平台登录 + 门控 @@ -286,7 +288,7 @@ self.set_progress(done=5, total=80, unit="视频", action_counts={"like": 3}, el `generic_steps` 的执行内容全在 `params.steps`(JSON 数组),由步骤编辑器产出。**没有默认步骤**——空步骤任务执行时会明确报错。 -**22 种步骤**(完整参数见 [doc/TASK_DEV.md](doc/TASK_DEV.md)): +**24 种步骤**(完整参数见 [doc/TASK_DEV.md](doc/TASK_DEV.md)): | 类别 | 步骤 | |------|------| @@ -408,7 +410,7 @@ read_log_text()`,`file` 参数走白名单,不接受任意路径。 | [doc/ARCHITECTURE.md](doc/ARCHITECTURE.md) | 架构详解(分层、装配顺序、线程模型、设备生命周期、调度链路、设计决策) | | [doc/DATA_MODEL.md](doc/DATA_MODEL.md) | 数据模型(表结构、迁移、`app_meta`、数据目录、备份覆盖清单) | | [doc/API.md](doc/API.md) | HTTP 接口全量说明 | -| [doc/TASK_DEV.md](doc/TASK_DEV.md) | 任务与步骤开发(22 种步骤、公共巡检、录制回放、去重、选择器、自定义动作) | +| [doc/TASK_DEV.md](doc/TASK_DEV.md) | 任务与步骤开发(24 种步骤、公共巡检、录制回放、去重、选择器、自定义动作) | | [doc/MCP.md](doc/MCP.md) · [doc/MCP_DESIGN.md](doc/MCP_DESIGN.md) | MCP 使用手册 / 设计文档 | | [doc/AI_CONSOLE.md](doc/AI_CONSOLE.md) · [doc/AI_TASK_GEN.md](doc/AI_TASK_GEN.md) | AI 控制台机制 / AI 建任务设计(规划) | | [doc/DEPLOY.md](doc/DEPLOY.md) | 部署与运维(含生产容器、备份导入、故障排查) | diff --git a/config.py b/config.py index 8e9cdfd..8e1e570 100644 --- a/config.py +++ b/config.py @@ -87,6 +87,14 @@ DATA_DIR = os.path.join(os.path.dirname(os.path.abspath(__file__)), "data") # APK 文件存储目录(应用管理功能) APK_DIR = os.path.join(DATA_DIR, "apks") +# 视频发布计划的素材目录(发布计划功能)。**支持 env 覆盖**,理由同下面三个: +# 自动化测试必须把素材目录指到临时位置,否则测试上传的视频会落进用户真实数据目录。 +VIDEO_DIR = _env("DATA_VIDEO_DIR", "") or os.path.join(DATA_DIR, "videos") + +# 单次请求体上限(字节)。**必须大于现有 APK 上传的最大包**(实测 336MB), +# 否则「应用管理」的 APK 上传会被一起卡死。默认 2GiB:一个视频 + 余量。 +MAX_CONTENT_LENGTH = int(_env("MAX_UPLOAD_MB", "2048")) * 1024 * 1024 + # ================== 系统数据备份(导出/导入) ================== # 导出 zip 与恢复前快照存放;导入暂存目录;待下次启动生效的恢复目录 # (运行时产物,不入 git,见 .gitignore data/backups 等) diff --git a/core/ledger.py b/core/ledger.py index fc8c8cb..b1d9a3c 100644 --- a/core/ledger.py +++ b/core/ledger.py @@ -209,6 +209,31 @@ def resolve_device(serial="", device_name=""): return name, rows, f"设备『{name}』登记 {len(rows)} 个号" +def serial_of(device_name): + """设备名 → **当前**地址(拿不到返回空串)。 + + ⚠ 为什么不能直接用台账/计划行里的 `serial` 快照:设备地址会变(换 IP、重连)。 + 快照是"录入时"的值,久了就是错的 —— 要下发操作(发布视频等)时必须现查。 + 优先设备池的**内存快照**(零 DB、已经是刷新过的),再退回 DB。 + """ + name = (device_name or "").strip() + if not name: + return "" + try: + from core import device_pool + for d in device_pool.list_devices(): + if (d.get("name") or "") == name: + return d.get("serial") or "" + except Exception: + pass + try: + with _ctx(): + dev = Device.query.filter(Device.name == name).first() + return (dev.serial if dev else "") or "" + except Exception: + return "" + + def douyin_ids(scope, serial="", device_name="", group=""): """按范围取抖音号(任务「条件判断」取号用)。返回 (去重后的纯号列表, 说明)。 diff --git a/core/models.py b/core/models.py index 45d761d..da53953 100644 --- a/core/models.py +++ b/core/models.py @@ -484,6 +484,93 @@ class DeviceAccount(db.Model): return f"" +class VideoPlan(db.Model): + """视频发布计划:账号 × 发布日期 × 编号 → 一个视频素材 + 一条标题 + 发布结果。 + + 一条记录 = **一个账号在某天要发的一个视频**(素材与计划天然 1:1,所以不拆两张表; + 但上传是分两步的 —— 先视频后标题或反过来 —— 靠 `status` 的 `pending` 态兜住)。 + + **状态机**(这是本表的灵魂,不要简化): + + | status | 含义 | + |---|---| + | `pending` | 有视频、还没标题 | + | `ready` | 素材齐,等发布日期 | + | `pushing` | 已原子占位,正在把视频推到手机 | + | `publishing` | 已推到手机,正在走抖音发布流程 | + | `done` | 发布成功(终态) | + | `failed` | 失败在 `push`/`scan` 阶段 —— 还没碰抖音,**可安全重试** | + | `unknown` | 失败在 `post`/`verify` 阶段 —— **可能已经发出去了,绝不自动重试**,要人工裁决 | + | `skipped` | 人工跳过(终态) | + + ⚠ **`failed` 与 `unknown` 必须分开**:把"不知道自己发没发"混成"知道自己没发", + 就是重复发布的来源。`stage` 记录失败发生在哪一步,是这两者互相转换的唯一依据。 + + **唯一性**:`(phone, release_date, seq)` 由服务层(`core/video_plan.py`)保证, + **不加 DB 唯一索引** —— `seq` 从 1 起、没有"空值"可言,做部分唯一索引要写三处方言适配 + (见 §唯一索引那段注释与 `_ensure_unique_indexes`),收益不匹配;违反的代价只是 + 低频人工上传产生的重复行,可见、可删。 + + **分享链接**:发布成功后抓作品的分享链接存 `share_url` —— 平台**不长期囤视频** + (存不下),链接才是长期资产,也方便后续拿它去铺评论。 + """ + __tablename__ = "video_plan" + id = db.Column(db.String(32), primary_key=True) # uuid 前 8 位 + account_id = db.Column(db.String(32), default="", index=True) # → device_account.id(不做外键) + phone = db.Column(db.String(32), default="", index=True) # 配对键(冗余存:账号删了也留痕) + device_name = db.Column(db.String(80), default="", index=True) # 设备号快照(聚合/下发免 join) + nickname = db.Column(db.String(80), default="") # 账号名称快照(时间线卡片直接显示) + douyin_id = db.Column(db.String(64), default="") # 抖音号快照 + serial = db.Column(db.String(120), default="") # 地址快照 + release_date = db.Column(db.String(10), default="", index=True) # "YYYY-MM-DD"(纯日期,等值比较) + seq = db.Column(db.Integer, default=1) # 编号,从 1 起(不用 0 表示"无") + seq_auto = db.Column(db.Boolean, default=True) # 编号是自动分配出来的(界面要提示) + title = db.Column(db.Text, default="") # 文案 + video_file = db.Column(db.String(120), default="") # 平台落盘文件名(不含绝对路径) + video_name = db.Column(db.String(200), default="") # 原始上传名(排查用) + video_size = db.Column(db.Integer, default=0) + video_sha1 = db.Column(db.String(40), default="") # 内容指纹(重复上传的判据) + status = db.Column(db.String(16), default="", index=True) + stage = db.Column(db.String(16), default="") # push/scan/post/verify(决定 failed vs unknown) + attempts = db.Column(db.Integer, default=0) # 尝试次数(超上限不再自动取) + published_at = db.Column(db.String(20), default="") + share_url = db.Column(db.String(300), default="") # 作品分享链接(发布后抓取) + link_at = db.Column(db.String(20), default="") # 抓到链接的时刻 + video_deleted_at = db.Column(db.String(20), default="") # 平台素材文件何时被清理 + push_verify = db.Column(db.String(16), default="") # 推送后的相册校验:ok=进索引 / no_index=没进 / nofile=文件不在 + push_remote = db.Column(db.String(200), default="") # 推到手机上的绝对路径(删它/排查用) + last_error = db.Column(db.String(500), default="") + note = db.Column(db.Text, default="") + created_at = db.Column(db.String(20), default="", index=True) + updated_at = db.Column(db.String(20), default="") + + __table_args__ = ( + db.Index("ix_video_plan_date_status", "release_date", "status"), + db.Index("ix_video_plan_acct_date", "account_id", "release_date"), + db.Index("ix_video_plan_phone_slot", "phone", "release_date", "seq"), + db.Index("ix_video_plan_file", "video_file"), + ) + + def to_dict(self): + return {"id": self.id, "account_id": self.account_id or "", + "phone": self.phone or "", "device_name": self.device_name or "", + "nickname": self.nickname or "", "douyin_id": self.douyin_id or "", + "serial": self.serial or "", "release_date": self.release_date or "", + "seq": int(self.seq or 1), "seq_auto": bool(self.seq_auto), + "title": self.title or "", "video_file": self.video_file or "", + "video_name": self.video_name or "", "video_size": int(self.video_size or 0), + "video_sha1": self.video_sha1 or "", "status": self.status or "", + "stage": self.stage or "", "attempts": int(self.attempts or 0), + "published_at": self.published_at or "", "share_url": self.share_url or "", + "link_at": self.link_at or "", "video_deleted_at": self.video_deleted_at or "", + "push_verify": self.push_verify or "", "push_remote": self.push_remote or "", + "last_error": self.last_error or "", "note": self.note or "", + "created_at": self.created_at or "", "updated_at": self.updated_at or ""} + + def __repr__(self): + return f"" + + class AgentConversation(db.Model): """AI 控制台会话:整个消息序列以 JSON 存在一行里(单会话几十 KB,够用)。""" __tablename__ = "agent_conversation" @@ -510,6 +597,9 @@ SCHEMA_MIGRATIONS = [ (6, "自动发现:pending_device 表新增 fingerprint 列(扫描时读取,用于提示是已有设备换了 IP)", None), (7, "去重账本:done_mark 表(跨设备幂等的「已做过」标记,唯一索引 scope_key)", None), (8, "账号台账:device_account 表(设备号/手机号/账号名称/抖音号/注册时间/卡在机内/可发视频/简介/备注)", None), + (9, "视频发布计划:video_plan 表(账号×发布日期×编号 → 素材 + 标题 + 发布状态 + 分享链接)", None), + (10, "视频发布计划:video_plan 新增 push_verify(推送后相册校验:ok/no_index/nofile)" + "与 push_remote(手机上的绝对路径,删它/排查用)", None), ] # 当前 schema 版本(备份/恢复用它判断新旧,也写进 app_meta.schema_version) diff --git a/core/notify_events.py b/core/notify_events.py index 91f8f03..d766941 100644 --- a/core/notify_events.py +++ b/core/notify_events.py @@ -111,6 +111,17 @@ EVENTS = [ ["serial", "device_name", "selector", "miss_count"], "某选择器连续 10 次未命中——任务可能显示成功但什么都没做", agg_window=0, recommend=True), + _e("task.video.published", "视频已发布", "业务", + ["serial", "device_name", "model", "phone", "nickname", "release_date", "seq", + "title", "share_url"], + "「发布视频」步骤把某个账号的视频发出去了(share_url 是作品分享链接," + "没抓到链接时为空——界面上会标「缺链接」)", agg_window=0, recommend=True), + _e("task.video.failed", "视频发布失败/结果未知", "业务", + ["serial", "device_name", "model", "phone", "nickname", "release_date", "seq", + "title", "msg"], + "「发布视频」步骤失败:推文件阶段失败=可重试;抖音那一步失败=**可能已经发出去了**," + "要到计划页看状态并人工确认(平台不会自动重发)", agg_window=0, + agg_key=_BY_SERIAL, recommend=True), # ---------------- 设备 ---------------- _e("device.online", "设备恢复在线", "设备", diff --git a/core/system_backup.py b/core/system_backup.py index 9d39273..f4240c8 100644 --- a/core/system_backup.py +++ b/core/system_backup.py @@ -59,6 +59,7 @@ TABLE_LABELS = { "task_step_log": "任务步骤明细", "done_mark": "去重记录(已做过)", "device_account": "账号台账", + "video_plan": "视频发布计划", } _STAGE_TTL = 1800 # 导入暂存有效期(秒) @@ -140,6 +141,46 @@ def _summary_info(con): return rows +def _video_summary(): + """视频素材目录的概况(写进 manifest;**不进备份包**,但要交代清楚)。 + + · `missing`:`video_plan` 引用了、但文件已经不在了的清单 —— 反向自检 + (与 `coverage_missing` 同一个思路:别让"看起来有其实没有"悄悄发生) + """ + try: + from config import VIDEO_DIR + except Exception: + return {"included": False, "error": "config 读取失败"} + count = 0 + size = 0 + try: + for root, _dirs, files in os.walk(VIDEO_DIR): + for fn in files: + if fn.startswith(".tmp_"): + continue + try: + size += os.path.getsize(os.path.join(root, fn)) + count += 1 + except OSError: + pass + except Exception: + pass + missing = [] + try: + from core.models import VideoPlan + with db.app_context(): + for r in VideoPlan.query.filter(VideoPlan.video_file != "").all(): + p = os.path.join(VIDEO_DIR, r.video_file) + if not os.path.exists(p): + missing.append({"id": r.id, "file": r.video_file}) + except Exception: + pass + return {"included": False, "dir": "data/videos", "count": count, "bytes": size, + "note": "视频素材不进备份包(体积)。恢复后需重新上传;" + "计划、发布状态与分享链接在 video_plan 表里,已备份。", + "missing": missing[:200], "missing_count": len(missing)} + + def _prune_old_exports(max_age=3600): """清理过期的导出临时 zip(下载完成后的 call_on_close 在 Windows 上可能 因文件锁删不掉,这里按时间兜底清理;保留近 1 小时的便于失败重试)。""" @@ -299,6 +340,9 @@ def create_export(include_apk=True): "include_apk": include_apk, "tables": _summary_info(con), "apks": apk_meta, + # 视频发布计划的素材**不进备份包**(几十 GB 会把备份能力搞坏), + # 但必须在 manifest 里交代清楚 —— 否则就是"以为备份了其实没有"(红线)。 + "videos": _video_summary(), } # 覆盖自检:登记在册的业务表若在快照里缺失(新增功能忘了登记 / 建表失败), # 显式告警并写进 manifest——避免"以为备份了其实没有"(备份覆盖红线)。 diff --git a/core/video_plan.py b/core/video_plan.py new file mode 100644 index 0000000..19c4e06 --- /dev/null +++ b/core/video_plan.py @@ -0,0 +1,1159 @@ +"""视频发布计划:批量上传 → 自动配对 → 时间线 → 自动发布 → 存分享链接。 + +**它解决什么**(用户场景):人员在网页上批量传视频与标题 txt,文件名里写着 +`手机号_日期_编号`,标题行写着 `标题_手机号_日期_编号`。平台把它们自动配到台账里的账号上, +排成一张发布计划;到发布日期由任务把视频推到手机相册、走抖音发布流程、回写状态, +**发布后抓作品的分享链接存进来**(平台不长期囤视频,链接才是长期资产 +—— 后续铺评论也是靠这个链接去打开视频)。 + +**三个使用方**: + 1. web「账号 → 发布计划」页(`web/video_plan_api.py`)—— 上传/时间线/单条操作/导出链接 + 2. 任务步骤「发布视频」(`tasks/generic/task.py` 的 `_exec_publish_video`)—— 取计划、占位、回写 + 3. 每日清理 job(`web_server.py`)—— 僵尸回收 + 过期行/孤儿文件清理 + +**两条不容妥协的规则**(都是"界面会显示成功但结果是错的"那类坑): + + · **原子占位**:发布前必须 `claim()` —— 带条件的 UPDATE + 受影响行数判据 + (`WHERE id=:id AND status IN ('ready','failed')`,`rowcount==1` 才算抢到)。 + 绝不能"先查后改":多台设备/任务重跑会同时通过,同一个号发两条。 + · **失败分两类**:`push`/`scan` 阶段失败 = 还没碰抖音 → `failed`(可安全重试); + `post`/`verify` 阶段失败或超时 = **可能已经发出去了** → `unknown`(**绝不自动重试**, + 人工到抖音确认后再裁决)。把"不知道自己发没发"当成"知道自己没发",就是重复发布的来源。 + +DB 访问自建 app context(照 `core/dedup.py` / `core/ledger.py` 的 `_ctx()`),调用方不用关心。 +""" +import hashlib +import os +import re +import shutil +import time +import uuid + +from config import VIDEO_DIR +from core.logger import get_logger +from core.models import db, DeviceAccount, VideoPlan + +_log = get_logger("core.video_plan") + +# ================== 常量 ================== +ST_PENDING = "pending" # 有视频、还没标题 +ST_READY = "ready" # 素材齐,等发布日期 +ST_PUSHING = "pushing" # 已占位,正在推视频到手机 +ST_PUBLISHING = "publishing" # 已推到手机,正在走抖音发布流程 +ST_DONE = "done" # 已发布(终态) +ST_FAILED = "failed" # 失败在 push/scan —— **可安全重试** +ST_UNKNOWN = "unknown" # 失败在 post/verify —— **绝不自动重试**,人工裁决 +ST_SKIPPED = "skipped" # 人工跳过(终态) +STATUSES = (ST_PENDING, ST_READY, ST_PUSHING, ST_PUBLISHING, ST_DONE, + ST_FAILED, ST_UNKNOWN, ST_SKIPPED) +# 界面上给状态配色/文案用(前端也有一份,这里给接口和后端日志用) +STATUS_LABELS = {ST_PENDING: "待标题", ST_READY: "待发布", ST_PUSHING: "推送中", + ST_PUBLISHING: "发布中", ST_DONE: "已发布", ST_FAILED: "失败(可重试)", + ST_UNKNOWN: "结果未知(需确认)", ST_SKIPPED: "已跳过"} +# 失败阶段 → 落到哪个状态(本模块最重要的一张表) +# `manual` = "已推到手机、后面由人(或用户自己写的步骤)负责" —— 再出岔子只能算**未知**: +# 人可能已经把视频发出去了,标成"可重试"就会被任务再推一遍、再发一遍。 +STAGE_MANUAL = "manual" +STAGE_STATUS = {"push": ST_FAILED, "scan": ST_FAILED, + "post": ST_UNKNOWN, "verify": ST_UNKNOWN, + STAGE_MANUAL: ST_UNKNOWN} + +VIDEO_EXTS = (".mp4", ".mov", ".mkv", ".avi", ".3gp", ".flv", ".wmv", ".m4v", ".ts", ".webm") +MAX_TITLE_LEN = 100 # 超过就拒收(**不静默截断**:截断会发出错文案,用户看不见) +MAX_ATTEMPTS = 3 # 尝试上限,超了不再自动取(人工确认后重试可清零) +ZOMBIE_MINUTES = 30 # 占位行超过这么久没动静 → 按 stage 降级(否则崩溃一次就永远发不出去) +VIDEO_KEEP_DAYS = 7 # 平台素材:抓到分享链接后再留这么多天 +PLAN_KEEP_DAYS = 90 # 过期(未发布/失败)的计划行保留这么多天 + +_PHONE_RE = re.compile(r"^1\d{10}$") +_DATE8_RE = re.compile(r"^\d{8}$") +_SEQ_RE = re.compile(r"^\d{1,3}$") +_URL_RE = re.compile(r"https?://v\.douyin\.com/[A-Za-z0-9_\-]+") +_SAFE_RE = re.compile(r"[^0-9A-Za-z_\-.一-鿿]+") + +_app = None + + +def init_app(app): + """web_server 启动时调用:绑 app(任务线程/清理 job 访问 db 要自推 context)。""" + global _app + _app = app + + +def _ctx(): + if _app is None: + raise RuntimeError("video_plan 未关联 Flask app(web_server 启动时调用 init_app)") + return _app.app_context() + + +def _now(): + return time.strftime("%Y-%m-%d %H:%M:%S") + + +def _today(): + return time.strftime("%Y-%m-%d") + + +# ================== 纯函数区(不碰 DB,好单测) ================== +def parse_date8(s): + """`20260912` → `2026-09-12`;不是合法日期返回 None。""" + s = (s or "").strip() + if not _DATE8_RE.match(s): + return None + try: + t = time.strptime(s, "%Y%m%d") + except ValueError: + return None + if not (2000 <= t.tm_year <= 2100): + return None + return f"{t.tm_year:04d}-{t.tm_mon:02d}-{t.tm_mday:02d}" + + +def parse_video_name(filename): + """视频文件名 `手机号_日期_编号` → (meta, error)。 + + meta = {"phone","date","seq"|None,"seq_given"};error 是给人看的拒收原因。 + + **按形态认,不按位置切**:手机号 11 位、日期 8 位、编号 1-3 位 —— 三个长度区间 + 互不相交,所以永远不会有一段同时匹配两种角色。按位置硬切的话,用户把顺序写反、 + 或文件名带下载站前缀时,会解析出一个**错误的手机号**;而它若恰好命中台账里 + 另一个账号,就**不报错**,静默地把 A 号的视频排给 B 号。 + """ + name = (filename or "").strip() + if not name: + return None, "文件名为空" + stem, ext = os.path.splitext(name) + if ext.lower() not in VIDEO_EXTS: + return None, f"不是视频文件(支持 {'/'.join(e.lstrip('.') for e in VIDEO_EXTS)})" + parts = [p.strip() for p in stem.split("_")] + if len(parts) < 2: + return None, "文件名格式不对,应为 手机号_日期_编号(编号可省略)" + + phone = date = seq = None + leftovers = [] + for p in parts: + if not p: + leftovers.append(p) + continue + if _PHONE_RE.match(p): + if phone is not None: + return None, f"文件名里有多个手机号({phone}、{p})" + phone = p + elif _DATE8_RE.match(p): + if date is not None: + return None, f"文件名里有多个日期({date}、{p})" + d = parse_date8(p) + if not d: + return None, f"日期 {p} 不是合法日期(应为 20260912 这种 8 位)" + date = d + elif _SEQ_RE.match(p): + if seq is not None: + return None, f"文件名里有多个编号({seq}、{p})" + seq = int(p) + else: + leftovers.append(p) + if leftovers: + return None, f"文件名里有多余片段 {'_'.join(leftovers)} —— 应为 手机号_日期_编号" + if not phone: + return None, "文件名里找不到手机号(11 位、1 开头)" + if not date: + return None, "文件名里找不到发布日期(8 位,如 20260912)" + return {"phone": phone, "date": date, "seq": seq, "seq_given": seq is not None}, "" + + +def parse_title_line(line): + """标题行 `标题内容_手机号_日期_编号` → (meta, error)。标题里可以带下划线。 + + 右锚定 + 逐级回退(**必须先试带编号的**,否则带编号的行会被当成 + "标题里含日期"而错配): + A) …_手机号_日期_编号 B) …_手机号_日期 + """ + t = (line or "").strip() + if not t: + return None, "空行" + parts = t.split("_") + phone = date = seq = None + if len(parts) >= 4 and _SEQ_RE.match(parts[-1]) \ + and _DATE8_RE.match(parts[-2]) and _PHONE_RE.match(parts[-3]): + seq = int(parts[-1]) + phone, date8, title_parts = parts[-3], parts[-2], parts[:-3] + elif len(parts) >= 3 and _DATE8_RE.match(parts[-1]) and _PHONE_RE.match(parts[-2]): + phone, date8, title_parts = parts[-2], parts[-1], parts[:-2] + else: + return None, "格式不对,应为 标题内容_手机号_日期_编号(编号可省略)" + d = parse_date8(date8) + if not d: + return None, f"日期 {date8} 不是合法日期(应为 20260912 这种 8 位)" + date = d + title = "_".join(title_parts).strip() + if not title: + return None, "缺标题内容(行首就是手机号)" + if len(title) > MAX_TITLE_LEN: + return None, f"标题太长({len(title)} 字 > {MAX_TITLE_LEN} 字上限),请改短或拆成两条" + return {"title": title, "phone": phone, "date": date, + "seq": seq, "seq_given": seq is not None}, "" + + +def parse_titles_text(text): + """标题 txt 全文 → (rows, errors)。`#` 开头的行与空行自动跳过(不算错)。""" + rows, errors = [], [] + lines = (text or "").replace("\r\n", "\n").replace("\r", "\n").split("\n") + for i, raw in enumerate(lines, start=1): + s = raw.strip() + if not s or s.startswith("#"): + continue + meta, err = parse_title_line(s) + if err: + errors.append({"line": i, "reason": err, "raw": s[:60]}) + continue + meta["_line"] = i + rows.append(meta) + return rows, errors + + +def next_free_slot(taken, start=1): + """最小的空闲编号(`taken` 是已被占用的编号集合)。""" + n = max(1, int(start or 1)) + while n in (taken or set()): + n += 1 + return n + + +def assign_slot(explicit_seq, taken): + """定一个槽位号:显式编号优先(被占则报冲突),没给编号就取最小空位。""" + taken = set(taken or set()) + if explicit_seq: + if int(explicit_seq) in taken: + return None + return int(explicit_seq) + return next_free_slot(taken) + + +def pair_slots(videos, titles, accounts_by_phone, db_taken=None): + """批量配对的**语义参考实现**(纯函数):把一批视频与一批标题配成 (账号,日期,编号)。 + + videos / titles:`parse_video_name` / `parse_title_line` 的返回(可带 "ref" 标识来源) + accounts_by_phone:{手机号: [账号 dict, ...]} + db_taken:{(手机号,日期): {已用编号}} —— 分两批上传时不能从 1 重号 + + 返回 (pairs, rejected):pairs=[{"video","title","phone","date","seq","seq_auto"}]。 + 规则:手机号在台账里 0 条或 ≥2 条 → 整组拒收;同侧编号重复 → 两条都拒收; + 没编号的视频填空位;**没编号的标题按号升序配到"还没有标题的视频"上** + (即"第 1 个视频配第 1 条标题",与逐条上传时 `import_titles()` 的行为一致 —— + 两处必须是同一套规则,不然批量路径与日常路径会给出不同结果); + 视频多标题少 → 多出的视频 title=None(待在计划页补);标题多视频少 → 多出的标题拒收。 + """ + pairs, rejected = [], [] + db_taken = db_taken or {} + + def _key(m): + return (m.get("phone") or "", m.get("date") or "") + + def _reject(side, meta, reason): + rejected.append({"side": side, "ref": meta.get("ref") or meta.get("title") + or meta.get("_line") or "", "reason": reason}) + + buckets = {} + for side, items in (("video", videos), ("title", titles)): + for m in items or []: + buckets.setdefault(_key(m), {}).setdefault(side, []).append(m) + + for key, sides in buckets.items(): + phone, date = key + accs = accounts_by_phone.get(phone) or [] + if len(accs) != 1: + reason = (f"手机号 {phone} 不在账号台账里" if not accs + else f"手机号 {phone} 在台账里对应 {len(accs)} 个账号,无法确定是哪个") + for side in ("video", "title"): + for m in sides.get(side, []): + _reject(side, m, reason) + continue + taken = set(db_taken.get(key) or set()) + # 显式编号:先占位,同侧重复 → 两条都拒收 + for side in ("video", "title"): + seen = {} + for m in sides.get(side, []): + if m.get("seq_given"): + s = int(m["seq"]) + seen.setdefault(s, []).append(m) + for s, ms in seen.items(): + if len(ms) > 1: + for m in ms: + _reject(side, m, f"编号 {s} 在这批里重复了(同侧不能有两个一样的号)") + sides[side] = [m for m in sides[side] if m not in ms] + else: + taken.add(s) + # 没编号的:填空位 + assigned = {} + for side in ("video", "title"): + for m in list(sides.get(side, [])): + if m.get("seq_given"): + continue + s = next_free_slot(taken) + taken.add(s) + assigned[id(m)] = s + vids = list(sides.get("video", [])) + tits = list(sides.get("title", [])) + # 视频 → 槽位 + vmap = {} + for m in vids: + s = int(m["seq"]) if m.get("seq_given") else assigned.get(id(m)) + vmap[s] = {"video": m, "title": None, "phone": phone, "date": date, + "seq": s, "seq_auto": not m.get("seq_given")} + # 标题 → 槽位:有编号的精确对;没编号的填进还没配到标题的槽位(按号升序) + free_slots = sorted(vmap.keys()) + for m in tits: + if m.get("seq_given"): + s = int(m["seq"]) + if s in vmap and vmap[s]["title"] is None: + vmap[s]["title"] = m + else: + _reject("title", m, f"编号 {s} 没有对应的视频(该账号 {date} 没有这条)") + else: + s = assigned.get(id(m)) + if s in vmap and vmap[s]["title"] is None: + vmap[s]["title"] = m + elif free_slots: + # 号被别的标题占了 → 退到第一个还没标题的槽位 + target = next((x for x in free_slots if vmap[x]["title"] is None), None) + if target is None: + _reject("title", m, f"没有对应视频(该账号 {date} 的视频都有标题了)") + else: + vmap[target]["title"] = m + else: + _reject("title", m, f"没有对应视频(该账号 {date} 没有视频)") + pairs.extend(vmap.values()) + pairs.sort(key=lambda p: (p["date"], p["phone"], p["seq"])) + return pairs, rejected + + +def build_caption(plan, topics=""): + """文案 = 标题 + 话题。**只在这一处定义拼接规则**(两处各配必然分叉)。 + + `topics` 支持逗号/换行分隔,不带 `#` 的自动补上。 + """ + title = (plan.get("title") or "").strip() if isinstance(plan, dict) else "" + tags = [] + for chunk in re.split(r"[,,\n]", topics or ""): + c = chunk.strip().lstrip("#").strip() + if c: + tags.append("#" + c) + caption = title + if tags: + caption = (caption + " " + " ".join(tags)).strip() if caption else " ".join(tags) + if len(caption) > MAX_TITLE_LEN + 200: + return caption, f"文案太长({len(caption)} 字)" + return caption, "" + + +def has_push_step(steps): + """这组步骤里有没有「推送发布视频」(含嵌套在条件/循环里的)——界面用它筛"发布任务"。""" + return _has_type(steps, "push_release") + + +def has_mark_step(steps): + """这组步骤里有没有「标记发布结果」。""" + return _has_type(steps, "mark_release") + + +def is_release_task(steps): + """算不算「发布任务」:有推送**或**有标记步骤。 + + 为什么要认"只有标记"的任务:`push_release` 是后加的平台步骤, + 手写的抖音发布流程(自己点相册/输入框/发布)本来就没有它 —— + 只认 push 会把这些任务漏在「发布任务」列表外,人就会觉得"这页没有能改的地方"。 + """ + return has_push_step(steps) or has_mark_step(steps) + + +def _has_type(steps, want): + for s in steps or []: + if not isinstance(s, dict): + continue + if s.get("type") == want: + return True + p = s.get("params") or {} + for k in ("children", "then", "else"): + if _has_type(p.get(k), want): + return True + return False + + +def build_push_step(label="推送发布视频(平台)"): + """「推送发布视频」步骤(骨架与「给已有任务插一步」共用这一处定义)。""" + return {"id": f"step_rel_push_{uuid.uuid4().hex[:6]}", "type": "push_release", + "label": label, + "params": {"date_source": "today", "fixed_date": "", "account": "device", + "account_phone": "", "album_dir": "/sdcard/DCIM/Camera", + "to_clipboard": True, "retry_failed": True, "max_same_run": 1}} + + +def build_mark_step(label="标记发布结果(平台)"): + """「标记发布结果」步骤(同上,共用一处定义)。""" + return {"id": f"step_rel_mark_{uuid.uuid4().hex[:6]}", "type": "mark_release", + "label": label, + "params": {"result": "published", "why": "", "capture_link": True, + "delete_phone": True, "album_dir": "/sdcard/DCIM/Camera"}} + + +def insert_release_steps(steps): + """给**已有任务**插上平台那两步:`push_release` 放最前、`mark_release` 放最后(纯函数)。 + + 位置是刻意的:推送(adb push + 触发相册刷新 + 标题写进剪贴板)必须在 + **打开抖音之前**;标记必须在**发布流程之后**。中间你原来的步骤一步不动。 + """ + out = [dict(s) if isinstance(s, dict) else s for s in (steps or [])] + if not out: + return [build_push_step(), build_mark_step()] + if not has_push_step(out): + out.insert(0, build_push_step()) + if not has_mark_step(out): + out.append(build_mark_step()) + return out + + +def build_release_steps(): + """「一键新建发布任务」的步骤骨架(纯函数,好单测)。 + + ⚠ 第一步是**亮屏**、第二步是「打开抖音」并**等「首页」出现**(`wait_home`)—— + 真机实测(2026-09-28,A01):**屏幕没亮就启动抖音,它会永远停在启动页** + (`mCurrentFocus=splash.SplashActivity`,UI 树是空的)→ 后面"点击我"必然 miss、 + "等抖音号出现"必然超时,最后显示成一句误导人的"账号不符"。 + 屏幕亮着时冷启动 5 秒就出底部栏。 + + **平台能保证的步骤都填好了**(推送 / 填标题 / 标记 + 账号校验), + 抖音里的点击步骤只给出**空选择器的占位**(标签写清楚该抓哪个), + 因为那些是混淆 id / 结构 xpath,写死了改版就废 —— 执行时会明确提示 + "selector_value 为空,请抓取元素",不会静默跳过。 + + 唯一预填了选择器的是底部「+」入口(`descriptionContains=拍摄`)与 + 「我」(`descriptionContains=我`),因为这两条我在真机上实测过、 + 且是按描述定位、跨版本相对稳。 + """ + def sid(name): + return f"step_rel_{name}_{uuid.uuid4().hex[:6]}" + + def clk(name, label, sel_type="", sel_val=""): + return {"id": sid(name), "type": "click", "label": label, + "params": {"selector_type": sel_type or "xpath", "selector_value": sel_val, + "wait_timeout": 8, "probability": 100}} + + # 命中(当前登录的就是这条计划要发的号)才走发布流程;未命中走 else 只发通知。 + # 账号校验用**现成的「条件判断」+ 候选值来源=发布计划**,平台不自作主张读账号。 + publish_flow_steps = [ + dict(build_push_step("⑤ 推送发布视频")), + clk("open", "⑥ 点击抖音底部「+」(实测可用;改版了重新抓)", + "descriptionContains", "拍摄"), + clk("album", "⑦ 点击「相册」(没命中的话用抓取元素选它)"), + clk("video", "⑧ 点击第一个视频(刚推的排最前;抓一次就行)"), + clk("next", "⑨ 点击「下一步」"), + clk("input", "⑩ 点击输入框(抖音发布页的文案框)"), + {"id": sid("title"), "type": "input_text", + "label": "⑪ 填标题(自动取这一条计划的标题,会回读校验)", + "params": {"text_source": "release_title", "release_topics": "", + "clear_first": True, "mode": "random", "texts": "", "fixed_text": ""}}, + clk("publish", "⑫ 点击「发布」"), + dict(build_mark_step("⑬ 标记发布结果(成功走这里;失败分支再放一个 result=failed)")), + ] + return [ + {"id": sid("light"), "type": "screen_on", "label": "⓪ 亮屏", + "params": {"probability": 100}}, + {"id": sid("openapp"), "type": "open_app", "label": "① 打开抖音(等「首页」出来)", + "params": {"package": "com.ss.android.ugc.aweme", "home_feature": "首页", + "wait_home": True, "probability": 100}}, + clk("me", "② 点击「我」(实测可用;改版了重新抓)", "descriptionContains", "我"), + {"id": sid("wait"), "type": "wait", "label": "③ 等「我」页加载出来", + "params": {"min": 3, "max": 5, "vary_pace": False, "probability": 100}}, + {"id": sid("check"), "type": "if_el", + "label": "④ 校验账号:当前登录的是不是这条计划要发的号", + "params": {"selector_type": "xpath", + "selector_value": "//*[contains(@text,'抖音号')]", + "timeout": 8, "cmp_op": "包含", "cmp_source": "release", + "cmp_value": "", "cmp_group": "", + "ident_type": "text", "ident_value": "", + "then": publish_flow_steps, + "else": [{"id": sid("else"), "type": "notify", + "label": "账号不符:本条跳过(不发)", + "params": {"title": "发布跳过:账号不符", + "message": "设备 {device} 当前登录的不是要发的号" + "(计划:见「账号 → 发布计划」),本条已跳过", + "level": "warning"}}]}}, + ] + + +def extract_share_url(text): + """从剪贴板/分享文本里抠抖音分享链接(发布后抓链接用)。""" + m = _URL_RE.search(text or "") + return m.group(0) if m else "" + + +def safe_name(name): + """原始文件名 → 可安全落盘的名字(**绝不信任客户端文件名**:挡住路径分隔与 ..)。""" + base = os.path.basename((name or "").replace("\\", "/")) + stem, ext = os.path.splitext(base) + stem = _SAFE_RE.sub("_", stem)[:60] or "video" + ext = _SAFE_RE.sub("", ext)[:10] or ".mp4" + return stem + ext + + +# ================== DB 服务区 ================== +def accounts_by_phone(): + """{手机号: [账号 dict, ...]}(配对用;同号多账号要能看出来)。""" + with _ctx(): + out = {} + for r in DeviceAccount.query.all(): + p = (r.phone or "").strip() + if p: + out.setdefault(p, []).append(r.to_dict()) + return out + + +def taken_slots(phone, date): + """该账号该日期已占用的编号集合(分两批上传时不能从 1 重号)。""" + with _ctx(): + rows = VideoPlan.query.filter(VideoPlan.phone == phone, + VideoPlan.release_date == date).all() + return {int(r.seq or 1) for r in rows} + + +def _video_path(video_file): + """库里的 `video_file` 是 "YYYY-MM/xxx.mp4" 这种相对路径。""" + return os.path.join(VIDEO_DIR, video_file) if video_file else "" + + +def store_video(file_storage): + """落盘 + 算 sha1。返回 (相对路径, 字节数, sha1)。 + + · **分块读** `file_storage.stream` 边写边算 sha1 —— 绝不 `read()` 整个文件进内存 + (视频几百 MB × 几个并发就是内存峰值;现有 APK 路径没这问题,是视频引进来的) + · 临时文件 + `os.replace`(学 `core/system_backup.stage_upload()`):中途断只留 `.tmp_*`, + 不会留半个正式文件;清理 job 会扫掉 `.tmp_*` + · 落盘名 `{sha1[:12]}_{安全原名}`,目录按 `YYYY-MM` 分(避免单目录堆几万个文件) + """ + now = time.localtime() + rel_dir = time.strftime("%Y-%m", now) + abs_dir = os.path.join(VIDEO_DIR, rel_dir) + os.makedirs(abs_dir, exist_ok=True) + tmp = os.path.join(VIDEO_DIR, f".tmp_{uuid.uuid4().hex[:8]}") + sha = hashlib.sha1() + size = 0 + try: + with open(tmp, "wb") as f: + while True: + chunk = file_storage.stream.read(1 << 20) + if not chunk: + break + sha.update(chunk) + size += len(chunk) + f.write(chunk) + digest = sha.hexdigest() + final_name = f"{digest[:12]}_{safe_name(file_storage.filename)}" + rel = f"{rel_dir}/{final_name}" + final = os.path.join(VIDEO_DIR, rel) + if not os.path.exists(final): + os.replace(tmp, final) + else: + os.remove(tmp) # 同一份文件重复上传:不存第二份 + return rel, size, digest + except Exception: + if os.path.exists(tmp): + try: + os.remove(tmp) + except OSError: + pass + raise + + +def upload_video(file_storage, order=None, replace=False): + """上传单个视频 → 解析 → 配对 → 落库。返回给接口用的 dict。 + + 幂等:同槽位已有行时 —— `sha1` 相同 → skip;不同 → **拒收**(除非 `replace=True`)。 + """ + name = getattr(file_storage, "filename", "") or "" + meta, err = parse_video_name(name) + if err: + return {"ok": False, "name": name, "error": err} + phone, date = meta["phone"], meta["date"] + accs = accounts_by_phone().get(phone) or [] + if not accs: + return {"ok": False, "name": name, + "error": f"手机号 {phone} 不在「账号」台账里(先在台账登记这个号,或改文件名)"} + if len(accs) > 1: + return {"ok": False, "name": name, + "error": f"手机号 {phone} 在台账里对应 {len(accs)} 个账号(" + f"{'、'.join(a['device_name'] for a in accs)}),无法确定是哪个"} + acc = accs[0] + if not acc.get("serial"): + # 台账是粘贴导入的老记录时可能没存上地址快照(导入那一刻设备池里没匹配上)—— + # 现查一次补上,别把空快照带进计划行(发布时要用"当前地址") + try: + from core import ledger as _ledger + acc = dict(acc, serial=_ledger.serial_of(acc.get("device_name") or "")) + except Exception: + pass + try: + with _ctx(): + taken = taken_slots(phone, date) + # 文件名里**写了编号就认它**(哪怕这个号已被占 —— 那正是"重复上传/要覆盖" + # 的场景,交给下面的 sha1 判据处理);没写编号才自动排到空位 + seq = int(meta["seq"]) if meta["seq_given"] else next_free_slot(taken) + old = VideoPlan.query.filter(VideoPlan.phone == phone, + VideoPlan.release_date == date, + VideoPlan.seq == seq).first() + existed = old.to_dict() if old else None + rel, size, digest = store_video(file_storage) + if existed: + if existed["video_sha1"] == digest: + # 同一份内容:库里的行用的就是同一个(内容寻址)文件, + # **绝不能删**它 —— 删了这个已引用的文件就没了 + return {"ok": True, "name": name, "action": "skip", "plan": existed, + "msg": f"{phone} {date} 编号{seq} 已经传过这个文件了(跳过)"} + if not replace: + _remove_video_file(rel) # 内容不同 → 这是刚写的新文件,没人引用,删掉 + return {"ok": False, "name": name, + "error": f"该槽位已有视频 {existed['video_name']}(编号{seq})——" + f"要换先删掉它,或在结果里点「覆盖」"} + delete_plan(existed["id"]) # 覆盖:删旧行与旧文件(旧内容若与新内容同,不会被删) + row = VideoPlan(id=uuid.uuid4().hex[:8], account_id=acc.get("id", ""), + phone=phone, device_name=acc.get("device_name", ""), + nickname=acc.get("nickname", ""), douyin_id=acc.get("douyin_id", ""), + serial=acc.get("serial", ""), release_date=date, seq=seq, + seq_auto=not meta["seq_given"], video_file=rel, + video_name=name[:200], video_size=size, video_sha1=digest, + status=ST_PENDING, created_at=_now(), updated_at=_now()) + with _ctx(): + db.session.add(row) + db.session.commit() + out = row.to_dict() + _log.info(f"发布计划新增:{phone} {date} 编号{seq} {name}({size} 字节)" + f"{'覆盖' if existed else ''}") + return {"ok": True, "name": name, "action": "replace" if existed else "add", + "plan": out, + "msg": f"已录入:{acc.get('device_name')} / {date} / 编号{seq}" + f"({'待标题' if not (existed or {}).get('title') else '取旧标题'})"} + except Exception as e: + _log.warning(f"视频上传失败: {e}") + return {"ok": False, "name": name, "error": f"保存失败:{str(e)[:120]}"} + + +def _remove_video_file(rel): + p = _video_path(rel) + if p and os.path.exists(p): + try: + os.remove(p) + except OSError: + pass + + +def import_titles(text, dry_run=False): + """标题 txt 全文 → 配到已有计划行上。返回逐行结果 + 汇总。 + + 规则:有编号的对到同编号的行;没编号的按号升序填进**还没有标题**的行。 + 对不到行(该槽位还没视频)→ 拒收并说清"请先上传视频"。 + """ + rows, errors = parse_titles_text(text) + counts = {"total": len(rows) + len(errors), "attached": 0, "skipped": 0, + "rejected": len(errors)} + out = [{"line": e["line"], "phone": "", "date": "", "seq": 0, "title": "", + "action": "reject", "error": e["reason"]} for e in errors] + with _ctx(): + accs = accounts_by_phone() + for m in rows: + phone, date = m["phone"], m["date"] + item = {"line": m.get("_line", 0), "phone": phone, "date": date, + "seq": m.get("seq") or 0, "title": m["title"], + "action": "", "error": ""} + if len(accs.get(phone) or []) != 1: + item["action"] = "reject" + item["error"] = (f"手机号 {phone} 不在账号台账里" if not accs.get(phone) + else f"手机号 {phone} 在台账里对应多个账号") + out.append(item) + counts["rejected"] += 1 + continue + q = VideoPlan.query.filter(VideoPlan.phone == phone, + VideoPlan.release_date == date) + if m["seq_given"]: + row = q.filter(VideoPlan.seq == int(m["seq"])).first() + else: + row = (q.filter(VideoPlan.title == "").order_by(VideoPlan.seq).first()) + if not row: + item["action"] = "reject" + item["error"] = (f"该账号 {date} 没有编号 {m['seq']} 的视频" + if m["seq_given"] else + f"该账号 {date} 没有待配标题的视频(请先上传视频)") + out.append(item) + counts["rejected"] += 1 + continue + item["seq"] = int(row.seq or 1) + if (row.title or "").strip() == m["title"]: + item["action"] = "skip" + item["error"] = "这条标题已经配过了" + out.append(item) + counts["skipped"] += 1 + continue + if (row.title or "").strip(): + # 已有别的标题 → 拒收(**不静默覆盖**:文案是发出去就改不了的东西; + # 要换就到计划页改标题,或先删掉这条标题) + item["action"] = "reject" + item["error"] = f"编号 {row.seq} 已经有标题了({row.title[:20]})——要换请到计划页改" + out.append(item) + counts["rejected"] += 1 + continue + item["action"] = "attach" + out.append(item) + counts["attached"] += 1 + if not dry_run: + row.title = m["title"] + if row.status == ST_PENDING: + row.status = ST_READY + row.updated_at = _now() + if not dry_run: + db.session.commit() + out.sort(key=lambda x: x.get("line") or 0) + msg = (f"标题配好 {counts['attached']} 条" + + (f",跳过 {counts['skipped']} 条" if counts["skipped"] else "") + + (f",拒收 {counts['rejected']} 条" if counts["rejected"] else "")) + return {"ok": True, "counts": counts, "rows": out, "dry_run": bool(dry_run), "msg": msg} + + +def set_title(plan_id, title): + """人工改标题(改完 `pending → ready`)。""" + title = (title or "").strip() + if len(title) > MAX_TITLE_LEN: + return None, f"标题太长(上限 {MAX_TITLE_LEN} 字)" + with _ctx(): + row = VideoPlan.query.get(plan_id) + if not row: + return None, "计划不存在" + row.title = title + if title and row.status == ST_PENDING: + row.status = ST_READY + elif not title and row.status == ST_READY: + row.status = ST_PENDING + row.updated_at = _now() + db.session.commit() + return row.to_dict(), "" + + +def set_serial(plan_id, serial): + """补/更新计划行里的地址快照(自愈用:老台账没存上快照时,别让它卡住发布)。""" + if not serial: + return + with _ctx(): + row = VideoPlan.query.get(plan_id) + if row and not row.serial: + row.serial = serial + row.updated_at = _now() + db.session.commit() + + +def get_plan(plan_id): + with _ctx(): + row = VideoPlan.query.get(plan_id) + return row.to_dict() if row else None + + +def delete_plan(plan_id, delete_file=True): + """删一行(默认连平台素材一起删)。**先删行再删文件**(反过来失败会留下指向已删文件的行)。""" + with _ctx(): + row = VideoPlan.query.get(plan_id) + if not row: + return False + rel = row.video_file + db.session.delete(row) + db.session.commit() + if delete_file and rel: + _remove_video_file(rel) + return True + + +def list_plans(date_from="", date_to="", status="", account_id="", q="", limit=500): + """时间线用:按日期区间/状态/账号/关键词取行(按日期、设备号、编号排序)。""" + with _ctx(): + query = VideoPlan.query + if date_from: + query = query.filter(VideoPlan.release_date >= date_from) + if date_to: + query = query.filter(VideoPlan.release_date <= date_to) + if status: + want = [s for s in status.split(",") if s in STATUSES] + if want: + query = query.filter(VideoPlan.status.in_(want)) + if account_id: + query = query.filter(VideoPlan.account_id == account_id) + if q: + like = f"%{q.strip()}%" + query = query.filter(db.or_(VideoPlan.phone.like(like), + VideoPlan.title.like(like), + VideoPlan.video_name.like(like), + VideoPlan.nickname.like(like), + VideoPlan.douyin_id.like(like), + VideoPlan.device_name.like(like))) + rows = (query.order_by(VideoPlan.release_date, VideoPlan.device_name, + VideoPlan.seq).limit(max(1, min(int(limit or 500), 5000))).all()) + return [r.to_dict() for r in rows] + + +def timeline(date_from="", date_to="", status="", account_id="", q="", limit=500): + """时间线数据:按日期分组的卡片 + 顶部统计 + 磁盘。""" + rows = list_plans(date_from, date_to, status, account_id, q, limit) + days = {} + today = _today() + for r in rows: + d = days.setdefault(r["release_date"], {"date": r["release_date"], + "weekday": _weekday(r["release_date"]), + "total": 0, "done": 0, "plans": []}) + d["total"] += 1 + d["done"] += 1 if r["status"] == ST_DONE else 0 + d["plans"].append(r) + return {"days": [days[k] for k in sorted(days)], "stats": stats(rows), + "disk": disk_info(), "today": today} + + +def _weekday(date_str): + try: + t = time.strptime(date_str, "%Y-%m-%d") + return "周" + "一二三四五六日"[t.tm_wday] + except Exception: + return "" + + +def stats(rows=None): + """统计。**三个口径都要给**(只给一个必然有一方是错的): + + · `today_ready` 今天要发几个号 —— 按**发布日期**归集(补发的看不出"昨天还漏着") + · `today_done` 今天完成几个号 —— 同上口径 + · `published_today` 今天实际发了几条 —— 按 `published_at`(补发也算今天) + """ + if rows is None: + rows = list_plans(limit=5000) + today = _today() + out = {"total": len(rows), "pending": 0, "ready": 0, "done": 0, "failed": 0, + "unknown": 0, "skipped": 0, "pushing": 0, "publishing": 0, + "today_ready": 0, "today_done": 0, + "today_devices": 0, "published_today": 0, "overdue": 0, "no_link": 0} + today_devs = set() + for r in rows: + st = r["status"] + if st in out: + out[st] += 1 + if r["release_date"] == today: + if st in (ST_PENDING, ST_READY, ST_PUSHING, ST_PUBLISHING, ST_FAILED): + out["today_ready"] += 1 + if r["device_name"]: + today_devs.add(r["device_name"]) + if st == ST_DONE: + out["today_done"] += 1 + elif r["release_date"] < today and st in (ST_PENDING, ST_READY, ST_FAILED): + out["overdue"] += 1 + if (r["published_at"] or "").startswith(today) and st == ST_DONE: + out["published_today"] += 1 + if st == ST_DONE and not r["share_url"]: + out["no_link"] += 1 + out["today_devices"] = len(today_devs) + return out + + +def disk_info(): + """素材占用 + 磁盘余量(视频是平台唯一无界增长的大体积数据,必须能看见)。""" + total = used = free = 0 + count = 0 + size = 0 + try: + os.makedirs(VIDEO_DIR, exist_ok=True) + du = shutil.disk_usage(VIDEO_DIR) + total, used, free = du.total, du.used, du.free + for root, _dirs, files in os.walk(VIDEO_DIR): + for f in files: + try: + size += os.path.getsize(os.path.join(root, f)) + count += 1 + except OSError: + pass + except Exception as e: + _log.warning(f"磁盘统计失败: {e}") + return {"total": total, "used": used, "free": free, "videos_bytes": size, + "videos_count": count} + + +def links(date_from="", date_to="", device="", phone=""): + """已发布作品的分享链接(导出 CSV / 界面展示用)。""" + with _ctx(): + query = VideoPlan.query.filter(VideoPlan.share_url != "") + if date_from: + query = query.filter(VideoPlan.release_date >= date_from) + if date_to: + query = query.filter(VideoPlan.release_date <= date_to) + if device: + query = query.filter(VideoPlan.device_name == device) + if phone: + query = query.filter(VideoPlan.phone == phone) + rows = query.order_by(VideoPlan.release_date, VideoPlan.device_name, + VideoPlan.seq).all() + return [r.to_dict() for r in rows] + + +# ================== 发布:取计划 / 占位 / 回写 ================== +def take_for_device(phones, date, limit=1, retry_failed=True): + """本机某天可发的计划行(按发布日期、编号升序)。`phones` = 这台设备上账号的手机号。""" + if not phones: + return [] + want = [ST_READY] + ([ST_FAILED] if retry_failed else []) + with _ctx(): + rows = (VideoPlan.query + .filter(VideoPlan.phone.in_(list(phones))) + .filter(VideoPlan.release_date <= date) + .filter(VideoPlan.status.in_(want)) + .filter(VideoPlan.attempts < MAX_ATTEMPTS) + .order_by(VideoPlan.release_date, VideoPlan.seq) + .limit(max(1, int(limit or 1))).all()) + return [r.to_dict() for r in rows] + + +def peek_for_device(phones, date): + """本机**下一条将要推送**的计划(只看不占位)——条件判断的"候选值来源=发布计划"用它。 + + 为什么用"下一条"而不是"已推送的那条":账号校验要放在**推送之前** + (号不对就不该推),那时候还没占位、也没有 `_last_release`。 + """ + rows = take_for_device(phones, date, limit=1) + return rows[0] if rows else None + + +def last_pushed_for_device(phones): + """本机最近一条"已推到手机、还没标记结果"的计划(`mark_release` 找不到上一步时的兜底)。""" + if not phones: + return None + with _ctx(): + row = (VideoPlan.query.filter(VideoPlan.phone.in_(list(phones))) + .filter(VideoPlan.status == ST_PUSHING) + .order_by(VideoPlan.updated_at.desc()).first()) + return row.to_dict() if row else None + + +def claim(plan_id, status=ST_PUSHING, stage="push"): + """**原子占位**(幂等的唯一判据):返回 True=抢到 / False=别人在处理或状态不对。 + + 带条件的 UPDATE + 受影响行数 —— 绝不能"先查后改":多台设备、任务重跑会同时通过。 + """ + with _ctx(): + sql = (f"UPDATE video_plan SET status=:st, stage=:stage, " + f"attempts=COALESCE(attempts,0)+1, updated_at=:ts " + f"WHERE id=:id AND status IN (:ready, :failed) AND COALESCE(attempts,0) < :max") + r = db.session.execute(db.text(sql), + {"st": status, "stage": stage, "ts": _now(), "id": plan_id, + "ready": ST_READY, "failed": ST_FAILED, + "max": MAX_ATTEMPTS}) + db.session.commit() + ok = (r.rowcount or 0) > 0 + if not ok: + _log.info(f"发布占位失败(别人在处理或状态不对):{plan_id}") + return ok + + +def pushable_plans(date="", limit=0, retry_failed=True): + """某天**可推送**的计划(`ready` + 可选 `failed` 且没过尝试上限),按 (设备, 编号) 升序。 + + 「一键推送」用它:一次把当天所有待推的都取出来,按设备分组由调用方串行推送。 + """ + date = date or _today() + with _ctx(): + q = VideoPlan.query.filter(VideoPlan.release_date == date) + sts = [ST_READY] + ([ST_FAILED] if retry_failed else []) + q = q.filter(VideoPlan.status.in_(sts)) + q = q.filter(VideoPlan.attempts < MAX_ATTEMPTS) + q = q.order_by(VideoPlan.device_name, VideoPlan.seq) + rows = q.limit(limit).all() if limit else q.all() + return [r.to_dict() for r in rows] + + +def mark_pushed(plan_id, remote="", verify=""): + """推送完成:状态留在 `pushing`(= 已推到手机、等人/用户的步骤去发),阶段置 `manual`。 + + 阶段置 `manual` 是为了僵尸回收时的语义正确:**再出岔子只能算"结果未知"**, + 不能算"可安全重试" —— 人可能已经发出去了。 + + `verify` = 相册校验结果(`ok` 进索引 / `no_index` 没进 / `nofile` 文件不在), + **单独存一列**而不是混进 last_error:文件推上去了就不是失败, + 但"相册里看不到它"必须能在界面上看见(后面的步骤要点它)。 + """ + with _ctx(): + row = VideoPlan.query.get(plan_id) + if not row: + return + row.status = ST_PUSHING + row.stage = STAGE_MANUAL + row.note = (row.note or "")[:200] + if remote: + row.last_error = "" # 上一轮的失败原因清掉(这条已经推到手机了) + row.push_remote = remote[:200] + if verify: + row.push_verify = verify[:16] + row.updated_at = _now() + db.session.commit() + + +def mark_done(plan_id, share_url="", msg=""): + """发布成功(抓到链接就一并存)。""" + with _ctx(): + db.session.execute(db.text( + "UPDATE video_plan SET status=:st, stage='done', published_at=:ts, " + "share_url=:url, link_at=:lat, last_error=:err, updated_at=:ts WHERE id=:id"), + {"st": ST_DONE, "ts": _now(), "url": share_url or "", + "lat": _now() if share_url else "", "err": (msg or "")[:500], "id": plan_id}) + db.session.commit() + + +def mark_failed(plan_id, stage, msg): + """按阶段落状态:push/scan → `failed`(可重试);post/verify → `unknown`(不自动重试)。""" + st = STAGE_STATUS.get(stage, ST_FAILED) + with _ctx(): + db.session.execute(db.text( + "UPDATE video_plan SET status=:st, stage=:stage, last_error=:err, " + "updated_at=:ts WHERE id=:id"), + {"st": st, "stage": stage, "err": (msg or "")[:500], "ts": _now(), + "id": plan_id}) + db.session.commit() + _log.warning(f"发布计划 {plan_id} 失败于 {stage} → {STATUS_LABELS.get(st, st)}:{msg}") + return st + + +def set_status(plan_id, status, note=""): + """人工操作:跳过 / 重试(failed→ready,attempts 清零)/ 标记完成 / 裁决 unknown。 + + · 置 `done` 时补上 `published_at`(界面上"今天实际发了几条"按它统计) + · 置 `failed`/`unknown` 时 `note` 同时记进 `last_error`(界面上要看原因) + """ + if status not in STATUSES: + return None, f"未知状态 {status}" + with _ctx(): + row = VideoPlan.query.get(plan_id) + if not row: + return None, "计划不存在" + row.status = status + if status == ST_READY: + row.attempts = 0 + row.last_error = "" + if status == ST_DONE and not row.published_at: + row.published_at = _now() + if status in (ST_FAILED, ST_UNKNOWN) and note: + row.last_error = note[:500] + if note: + row.note = note[:500] + row.updated_at = _now() + db.session.commit() + return row.to_dict(), "" + + +def resolve_unknown(plan_id, published, note=""): + """`unknown` 的人工裁决:确认真发出去了 → done;确认没发 → failed(可重试)。""" + with _ctx(): + row = VideoPlan.query.get(plan_id) + if not row: + return None, "计划不存在" + if row.status != ST_UNKNOWN: + return None, f"只有「结果未知」的行需要裁决(当前是 {STATUS_LABELS.get(row.status, row.status)})" + row.status = ST_DONE if published else ST_FAILED + if published: + row.published_at = row.published_at or _now() + else: + row.attempts = 0 + if note: + row.note = note + row.updated_at = _now() + db.session.commit() + return row.to_dict(), "" + + +# ================== 清理 ================== +def reclaim_zombies(minutes=ZOMBIE_MINUTES): + """回收卡住的占位行:按 stage 落到 failed(可重试)或 unknown(要人工)。 + + **没有这一步,崩溃一次这条计划就永远发不出去**(状态永远停在 pushing/publishing)。 + """ + cutoff = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(time.time() - minutes * 60)) + n = 0 + with _ctx(): + rows = (VideoPlan.query + .filter(VideoPlan.status.in_((ST_PUSHING, ST_PUBLISHING))) + .filter(VideoPlan.updated_at < cutoff).all()) + for r in rows: + st = STAGE_STATUS.get(r.stage or "push", ST_FAILED) + r.status = st + r.last_error = "上次执行中断(已自动回收)" + r.updated_at = _now() + n += 1 + if n: + db.session.commit() + if n: + _log.warning(f"回收卡住的发布占位 {n} 条(按阶段落到 failed/unknown)") + return n + + +def purge_plan(keep_days=PLAN_KEEP_DAYS): + """过期行清理:只删"昨天以前 + 未发布/失败/跳过"的行(**done 的行永久保留**,链接有价值)。""" + cutoff = time.strftime("%Y-%m-%d", time.localtime(time.time() - keep_days * 86400)) + n = 0 + with _ctx(): + rows = (VideoPlan.query + .filter(VideoPlan.release_date < cutoff) + .filter(VideoPlan.status.in_((ST_PENDING, ST_FAILED, + ST_UNKNOWN, ST_SKIPPED))).all()) + for r in rows: + rel = r.video_file + db.session.delete(r) + if rel: + _remove_video_file(rel) + n += 1 + if n: + db.session.commit() + if n: + _log.info(f"清理过期发布计划 {n} 条({keep_days} 天前、未发布的)") + return n + + +def purge_files(keep_days=VIDEO_KEEP_DAYS): + """删平台上的视频文件(**行保留**):只删"已发布 + 已抓到分享链接 + 超过保留期"的。 + + 安全阀:没抓到链接的**不删**(链接与文件至少留一个);pending/failed/unknown 不删。 + """ + cutoff = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(time.time() - keep_days * 86400)) + n = 0 + with _ctx(): + rows = (VideoPlan.query + .filter(VideoPlan.status == ST_DONE) + .filter(VideoPlan.share_url != "") + .filter(VideoPlan.video_file != "") + .filter(VideoPlan.published_at != "") + .filter(VideoPlan.published_at < cutoff).all()) + for r in rows: + _remove_video_file(r.video_file) + r.video_deleted_at = _now() + r.video_file = "" + n += 1 + if n: + db.session.commit() + # 孤儿文件:库里没人引用、且超过保留期的(含上传中断残留的 .tmp_*) + with _ctx(): + used = {r[0] for r in db.session.query(VideoPlan.video_file).all() if r[0]} + orphans = 0 + if os.path.isdir(VIDEO_DIR): + for root, _dirs, files in os.walk(VIDEO_DIR): + for f in files: + p = os.path.join(root, f) + rel = os.path.relpath(p, VIDEO_DIR).replace("\\", "/") + if rel in used: + continue + try: + if os.path.getmtime(p) > time.time() - keep_days * 86400: + continue + except OSError: + continue + try: + os.remove(p) + orphans += 1 + except OSError: + pass + if n or orphans: + _log.info(f"清理视频素材:{n} 个(已发布且有链接的过期素材)、孤儿文件 {orphans} 个") + return {"cleaned": n, "orphans": orphans} diff --git a/doc/API.md b/doc/API.md index 0329f0d..00f72a6 100644 --- a/doc/API.md +++ b/doc/API.md @@ -62,7 +62,9 @@ ## 2. 接口总索引 -> 共 **107 条路由**。`鉴权` 列:`—` 无、`L` 登录、`T/D/A/G` = tasks/devices/apks/logs 权限位、`Admin` 仅管理员。 +> 本表列出全部业务路由(页面路由除外;同一路径的不同方法合并成一行)。 +> `鉴权` 列:`—` 无、`L` 登录、`T/D/A/G` = tasks/devices/apks/logs 权限位、`Admin` 仅管理员。 +> 想看实际注册了多少条:`python -c "from web_server import app; print(len([r for r in app.url_map.iter_rules() if r.rule.startswith('/api')]))"`。 ### 2.1 auth(`web/auth.py`) @@ -143,6 +145,21 @@ | PUT | `/api/ledger/` | D | 改一条账号 | | DELETE | `/api/ledger/` | D | 删一条账号 | | POST | `/api/ledger/import` | D | 从表格粘贴文本批量导入(`dry_run=true` 只预览不写库) | +| POST | `/api/video_plan/upload_video` | D | 上传单个视频素材(文件名 `手机号_日期_编号`)→ 解析 + 配对 + 入计划 | +| POST | `/api/video_plan/upload_titles` | D | 上传标题(txt 文本或文件)→ 配到已有计划行 | +| GET | `/api/video_plan/timeline` | D | 发布计划时间线(按日期分组的卡片 + 统计 + 磁盘) | +| GET | `/api/video_plan/stats` | D | 今天要发/已发/待标题/过期未发 + 素材占用与磁盘余量 | +| GET/PUT/DELETE | `/api/video_plan/` | D | 单条计划:查看 / 改标题 / 删除(连素材文件) | +| GET | `/api/video_plan//video` | D | 预览平台上的素材(文件清理后 → 410) | +| POST | `/api/video_plan//push` | D | **把素材推到手机**(只推送、不发布;后台线程) | +| POST | `/api/video_plan/push_all` | D | **一键推送**:某天(默认今天)所有待发布/失败的计划 → 各自手机(每台设备一个线程、设备内串行) | +| POST | `/api/video_plan//mark` | D | 人工标记结果(`done`/`failed`/`unknown`,推完之后用) | +| POST | `/api/video_plan//skip\|retry\|resolve` | D | 跳过 / 重试 / 裁决「结果未知」 | +| GET | `/api/video_plan/links` | D | 已发布作品的分享链接(`format=csv` 导出,给铺评论用) | +| GET | `/api/video_plan/tasks` | T | 可编辑的任务列表(全部通用步骤任务 + `is_release` 标记:状态/调度/下次运行 + 完整步骤) | +| POST | `/api/video_plan/tasks` | T | 一键新建标准发布任务(骨架 15 步:亮屏 → 打开抖音 → 校验账号 → 推送 → 抖音点击 → 填标题 → 标记) | +| GET/PUT | `/api/video_plan/tasks/` | T | **发布任务的就地编辑**(就在「发布计划」页改名字/目标/时间/启停/步骤) | +| POST | `/api/video_plan/tasks//adopt` | T | 给已有任务插上平台两步(推送放最前 / 标记放最后,自己的步骤不动) | ### 2.4 admin(`web/admin_api.py`) @@ -518,6 +535,39 @@ > ⚠ 台账里的抖音号是**纯号**,只当"比对用的候选值";**不要**拿它填「去重」的身份元素 > (身份是元素原文逐字算 key,格式不同会让去重**静默失效**)。 +### 视频发布计划(「账号 → 发布计划」页) + +表与状态机见 [DATA_MODEL.md](DATA_MODEL.md) §2.11,任务侧怎么自动发布见 +[TASK_DEV.md](TASK_DEV.md) §4.7;权限 `devices`。 + +| 接口 | 请求 | 响应要点 | +|------|------|---------| +| `POST /api/video_plan/upload_video` | multipart `file`(单文件)+ `replace` | `{ok,name,action:add\|skip\|replace,plan,msg}`;解析/配对不通过 → **400 + `error`(人话原因)**。前端**逐文件串行**发(才能给每个文件一条进度与结果) | +| `POST /api/video_plan/upload_titles` | `{text, dry_run}` 或 multipart `file`(txt) | `{ok, counts:{total,attached,skipped,rejected}, rows:[{line,phone,date,seq,title,action,error}], msg}` | +| `GET /api/video_plan/timeline` | `?range=today\|week\|all&from=&to=&status=&account_id=&q=&limit=` | `{ok, days:[{date,weekday,total,done,plans:[…]}], stats, disk, today}` | +| `GET /api/video_plan/stats` | — | `{ok, stats, disk, today}`;`stats` 里**三个口径都给**:`today_ready`(今天要发几个号,按发布日期)/ `today_done`(今天完成)/ `published_today`(今天实际发的,按 `published_at`,补发也算今天) | +| `GET /api/video_plan/` · `PUT` · `DELETE` | `PUT {"title"}` | `PUT` 只改标题(改完 `pending→ready`);`DELETE` **先删行再删素材文件** | +| `GET /api/video_plan//video` | — | 素材预览(`send_file`,支持 Range);文件已被清理 → **410** | +| `POST /api/video_plan//push` | — | **把素材推到手机**(只推送、不发布):推文件 → 触发相册刷新 → 标题写进剪贴板;已推送/结果未知 → 409。发布动作由你写在任务里(见 [TASK_DEV.md](TASK_DEV.md) §4.7) | +| `POST /api/video_plan/push_all` | `{date?, retry_failed?}` | **一键推送**:某天(默认今天)`ready`/`failed` 且未超尝试上限的计划,**按设备分组**推(每台设备一个线程、设备内串行;设备间并行)。返回 `{ok,count,devices,no_device,msg}`;没绑上设备地址的条数单独在 `no_device` 里报出来。全部设备都解析不到 → 400 | +| `POST /api/video_plan//mark` | `{status: done\|failed\|unknown, why}` | 人工标记结果(推完之后用):`done` 补 `published_at`;`failed`/`unknown` 把 `why` 记进 `last_error` | +| `POST /api/video_plan//skip` | `{reason}` | 人工跳过 | +| `POST /api/video_plan//retry` | — | `failed`→`ready`,`attempts` 清零 | +| `POST /api/video_plan//resolve` | `{published, note}` | **`unknown` 的人工裁决**:确认真发出去了 → `done`;确认没发 → `failed`(可重试) | +| `GET /api/video_plan/links` | `?from=&to=&device=&phone=&format=csv` | 已发布作品的分享链接;CSV 带 BOM(Excel 打开中文不乱码) | +| `GET /api/video_plan/tasks` | — | `{ok, tasks:[{id,name,enabled,mode,cron,next_run,target,steps,step_count,schedule,retry,is_release,has_push}], release_count, other_count}`:**全部通用步骤任务**(不止发布任务 —— 只列发布任务的话,手写的抖音发布流程会一条都不显示,人会以为"这页没有能改的地方"),界面按 `is_release`(有 `push_release` **或** `mark_release`,含嵌套)分两组。带完整 `steps` 是为**就地编辑**(抓选择器);`step_count` 数的是**含嵌套**的总步数。权限 `T` | +| `POST /api/video_plan/tasks` | `{"name","target":{"mode":"all\|group\|serial","group_name","serial"},"time":"10:00","enabled"}` | 一键新建:骨架 **15 步**(亮屏 → 打开抖音(等首页) → 点「我」→ `if_el(cmp_source=release)` 校验账号 → then: `push_release` → 「+」(预填 `descriptionContains=拍摄`) → 相册/第一个视频/下一步/输入框/发布(**空选择器占位**)→ `input_text(text_source=release_title)` → `mark_release`;else: `notify` 跳过);**走 TaskManager 建**(内存+库一起更新,直接写库调度器不认)。⚠ `target` 分组键名是 **`group_name`**(写成 `group` 会静默解析出 0 台设备);分组/单设备没给具体目标 → 400。权限 `T` | +| `GET/PUT /api/video_plan/tasks/` | `PUT {name?, time?, target?, steps?, enabled?}` | **发布任务的就地编辑**(「账号 → 发布计划」页里改,不用跳任务页)。只放开这几个字段,**没传的字段原样保留**;`steps` 整块替换但 `params` 里的通用参数(去重有效期、抢占…)**深合并保留**。校验:空名字/坏时间/`steps` 不是数组/什么都不传 → 400,未知任务 → 404。权限 `T`。底层仍是 `TaskManager.update_job()`(内存+库一起改) | +| `POST /api/video_plan/tasks//adopt` | — | 给**已有任务**插上平台那两步:`push_release` 放最前、`mark_release` 放最后(**中间你自己的步骤一步不动**;位置理由见 `core.video_plan.insert_release_steps`)。只对「通用步骤」任务;没有步骤 → 400,已经有这两步 → 409,未知任务 → 404。权限 `T` | + +**逐文件上传的幂等**:同一 `(手机号,日期,编号)` 已有行时 —— 文件内容相同 → `action=skip`(不重复落盘); +内容不同 → **拒收**,前端给「覆盖」按钮(二次确认后带 `replace=1` 重传)。 +标题同理:同槽位已有别的标题 → 拒收(不静默覆盖,文案是发出去就改不了的东西)。 + +> ⚠ 上传体积:`MAX_CONTENT_LENGTH`(默认 2GiB,`.env` 的 `MAX_UPLOAD_MB` 可调)**必须大于 +> 「应用管理」最大的 APK**(实测 336MB),否则 APK 上传会被一起卡死;超限返回 **JSON 413**。 +> 上反代时 nginx 那侧还要 `client_max_body_size`(见 [DEPLOY.md](DEPLOY.md) §6)。 + --- ## 6. 设备池与自动发现 diff --git a/doc/ARCHITECTURE.md b/doc/ARCHITECTURE.md index 02cad63..57d83b6 100644 --- a/doc/ARCHITECTURE.md +++ b/doc/ARCHITECTURE.md @@ -86,7 +86,7 @@ | 阶段 | 位置 | 做了什么 | |------|------|---------| | **E** | `:64-67` | `context.init(...)`;`register_blueprints(app)`(10 个蓝图);`agent_api.set_app(app)`(供后台线程推 app context) | -| **F** | 同上附近 | **第二个独立 APScheduler**:`CronTrigger(hour=3, minute=47)` 挂经验库巡检、`hour=4, minute=13` 挂步骤明细清理(`_purge_step_log`)、`hour=4, minute=23` 挂去重记录清理(`_purge_done_mark`,只清 `day`/`hours` 桶)——两个清理任务都自建 app context;失败仅 warning | +| **F** | 同上附近 | **第二个独立 APScheduler**:`CronTrigger(hour=3, minute=47)` 挂经验库巡检、`hour=4, minute=13` 挂步骤明细清理(`_purge_step_log`)、`hour=4, minute=23` 挂去重记录清理(`_purge_done_mark`,只清 `day`/`hours` 桶)、`hour=4, minute=41` 挂发布计划清理(`_purge_video_plan`:僵尸占位回收 + 过期行)、`hour=4, minute=47` 挂素材文件清理(`_purge_video_files`)——时间刻意错开(都是删数据的大操作),各任务自建 app context;失败仅 warning。**僵尸回收是必须有的一步**:占位后崩溃的行没有它会永远卡在"推送中",再也不会被取到 | | **G** | `__main__` | `_ensure_uiauto_running()`(拉起 uiautodev:20242,写 `data/uiauto.pid`,`atexit` 清理)→ `_preconnect_pool_devices()`(后台并发 connect 池内网络设备)→ `_purge_step_log_async()`(后台清理超期步骤明细)→ `_run_server()`(候选端口依次 bind:`0.0.0.0:18050` → `127.0.0.1:18050` → `127.0.0.1:18051..18055`);退出时 `notifier.shutdown()` + `step_log.shutdown()` + `mgr.shutdown()` + `device_discovery.shutdown()` + `device_battery.shutdown()` + 停 uiautodev | > ⚠️ **阶段 A~F 在 import 期就会起线程/调度器**,只有 uiautodev 拉起与预连接在 `__main__` 分支。以 WSGI 方式 import 本模块会得到"半个启动"的进程——本地调试请直接 `python web_server.py`。 @@ -280,7 +280,7 @@ connecting ──获取设备──▶ u2 连接 ──▶ running ──▶ set | 顶级 Tab | `data-tab` | 权限 | 子分栏 | |---------|-----------|------|--------| | 监控 | `monitor` | 登录即可 | — | -| 账号 | `account` | `devices` | — (账号台账:列表 / 增删改 / 从表格粘贴导入,见 [DATA_MODEL.md](DATA_MODEL.md) §2.10) | +| 账号 | `account` | `devices` | `ledger` 账号台账(列表/增删改/粘贴导入)/ `plan` **发布计划**(上传视频与标题、时间线、推送到手机、导出链接 + **发布任务就地编辑**,见 [DATA_MODEL.md](DATA_MODEL.md) §2.10 §2.11) | | 任务 | `tasks` | 登录即可(写操作需 `tasks`) | `plan` 任务计划 / `actions` 自定义动作 / `actioncfg` 动作配置 / `dedup` 去重记录 | | 日志 | `logs` | `logs` | — | | 用户 | `users` | `admin` | — | @@ -298,6 +298,7 @@ connecting ──获取设备──▶ u2 连接 ──▶ running ──▶ set | `markdown.js` | 轻量 Markdown 渲染(`renderMarkdown`) | 无 CDN 依赖;**先转义再套标记** | | `list.js` | 统一列表组件:搜索 + 分页 + 排序 | 状态注册表 `_LIST_PAGERS`;`setListPager` **不重置**页码/搜索/排序(避免轮询刷新打断用户) | | `ledger.js` | 账号台账页(列表 / 单条增删改 / 从表格粘贴导入 + 逐行结果 / 设备维度弹窗)+ `ledgerCell()` 供设备池那一列 | 服务层 `core/ledger.py`;导入默认**先预览再写**(`dry_run`) | +| `release.js` | 发布计划页(时间线按日期分组 / 多选视频**串行**上传 + 逐文件进度与结果 / 标题 txt 上传 / 推送到手机 + **一键推送**(当天全部)、相册校验徽章与裁决 / 链接复制与导出 CSV / **发布任务的就地编辑**) | 服务层 `core/video_plan.py`;上传前把队列**按文件名排序**并让用户确认解析结果(顺序 = "无编号按顺序配"的依据)。就地编辑只放常用字段(名字/目标/时间/启停/逐步参数 + 抓取元素),更复杂的编排走「打开通用编辑器」→ `openTaskModal`;步骤路径里 `then/else/children` 在**步骤的 params 下**(`_rpWalk`),走错这一层会把分支当数组下标用 | | `monitor.js` | 监控页:设备表(含**电量列**,按 `battery.tier` 上色)、批量操作、异常汇总、任务概况卡片(含覆盖设备 chip) | 5s 轮询 + 脏检查(签名不变不重渲染) | | `editor.js` | 步骤编辑器(拖拽 / 参数表单 / 条件分支 / 元素抓取 / 测试此步骤)+ `saveTask` | 最大的前端文件;`_stepEditor` 单例;设备选择卡 `_devCard` 供「抓取元素」「测试此步骤」共用(**名字优先**,型号·地址作副标题) | | `tasks.js` | 任务 Tab:任务 CRUD / 调度解析 / 运行窗口 + 自定义动作 | | diff --git a/doc/DATA_MODEL.md b/doc/DATA_MODEL.md index ddfd4e6..f2618b1 100644 --- a/doc/DATA_MODEL.md +++ b/doc/DATA_MODEL.md @@ -16,7 +16,7 @@ | SQLite 回退时的 PRAGMA | `journal_mode=WAL`、`busy_timeout=5000`、`synchronous=NORMAL`(监听器按连接类型守卫,MySQL 连接不会执行) | | 建表方式 | **唯一真相是模型**:`db.create_all()`(建缺表)+ `_sync_columns()`(补缺列);`SCHEMA_MIGRATIONS` 只作版本账本与数据回填 | | 库位置 | MySQL:由 `DB_HOST/DB_NAME` 指定;SQLite:`data/users.db` | -| 当前 schema 版本 | `app_meta.schema_version = 6` | +| 当前 schema 版本 | `app_meta.schema_version = 10`(迁移清单见 `core/models.py` 的 `SCHEMA_MIGRATIONS`:建表/补列以模型为准,这里只作版本账本) | **环境与库的绑定**(防混库,见 [DEPLOY.md](DEPLOY.md) §2.2) @@ -27,7 +27,7 @@ 启动时校验「`.env` 声明」与「库名」「库中登记的 `app_meta.deployment_env`」三方一致,不符**拒绝启动**。 -**表清单(16 张,全部是 `core/models.py` 里的 ORM 模型)** +**表清单(17 张,全部是 `core/models.py` 里的 ORM 模型)** | # | 表 | 用途 | |---|---|------| @@ -47,6 +47,7 @@ | 14 | `task_step_log` | 任务步骤明细(每次步骤执行一条,见 §2.8;**唯一有无界增长风险的表**,靠保留期清理) | | 15 | `done_mark` | 去重账本:跨设备"已做过"标记(见 §2.9;本清单此前漏列,2026-09-24 补上) | | 16 | `device_account` | 账号台账:一台设备上登录着哪些账号(见 §2.10;任务「条件判断」的取号来源、设备端身份页显示用) | +| 17 | `video_plan` | 视频发布计划:账号 × 发布日期 × 编号 → 素材 + 标题 + 发布状态 + 分享链接(见 §2.11) | > 2026-09-13 之前,`app_meta` 与 4 张 `agent_*` 表是各模块里的裸 `CREATE TABLE` > (不进模型层)。迁 MySQL 时那批 SQL 的 `AUTOINCREMENT`/`TEXT DEFAULT ''`/`TEXT PRIMARY KEY` @@ -234,6 +235,55 @@ 这张表**自动进整库备份**(§6 派生规则)。 +### 2.11 `video_plan` — 视频发布计划(一个账号在某天要发的一个视频) + +「账号 → 发布计划」页维护;服务层 `core/video_plan.py`,接口 `/api/video_plan/*`(见 [API.md](API.md) §2.14), +任务步骤「发布视频」按它自动发布(见 [TASK_DEV.md](TASK_DEV.md) §4.7)。 +**素材文件**落在 `data/videos/YYYY-MM/`(文件名 `{sha1[:12]}_{安全原名}`,内容寻址),**不进整库备份**(§6)。 + +| 列 | 类型 | 说明 | +|----|------|------| +| `id` | String(32) PK | uuid 前 8 位 | +| `account_id` / `phone` | String(32) | 台账行 id(**不做外键**)/ 配对键(冗余存:账号删了也留痕) | +| `device_name` / `nickname` / `douyin_id` / `serial` | String | 账号与设备快照(时间线卡片直接显示,免 join) | +| `release_date` | String(10) index | **`2026-09-12`** 纯日期(等值比较走索引) | +| `seq` / `seq_auto` | Integer / Boolean | 编号从 **1** 起(不用 0 表示"无");`seq_auto`=号是自动分配的 | +| `title` | Text | 文案(标题 txt 配对写入;**空标题不会被发布**) | +| `video_file` / `video_name` / `video_size` / `video_sha1` | 各自 | 落盘相对路径 / 原始名 / 大小 / 内容指纹(重复上传判据) | +| `status` | String(16) index | 见下方状态机 | +| `stage` | String(16) | 失败发生在哪一步:`push`/`scan`/`post`/`verify` | +| `attempts` | Integer | 尝试次数(上限 3,超了不再自动取) | +| `published_at` / `share_url` / `link_at` / `video_deleted_at` | String | 发布时刻 / **作品分享链接** / 抓到链接的时刻 / 素材何时清理 | +| `push_verify` / `push_remote` | String(16)/String(200) | **推送后的相册校验**:`ok`=已进相册索引 / `no_index`=文件在但没进索引(相册里可能看不到)/ `nofile`=文件不在;`push_remote` = 推到手机上的绝对路径(删它、排查用)。**文件推上去了 ≠ 相册里点得到它**,所以单独存一列而不是混进 `last_error` | +| `last_error` / `note` / `created_at` / `updated_at` | | 失败原因 / 人工备注 / 时间 | + +**状态机**(本表的灵魂,别简化): + +``` +pending 有视频、还没标题 ready 素材齐,等推送 +pushing 已推到手机(等人/用户的步骤去发) done 已发布(终态) +skipped 人工跳过(终态) failed 推送阶段就失败 —— 还没到抖音,**可安全重试** +unknown 推送之后出的岔子 —— **可能已经发出去了,绝不自动重试**,要人工裁决 +``` + +> **平台只负责把素材推到手机**(`push_release` 步骤 / 计划页「推送到手机」): +> 推文件 → 触发相册刷新 → 把标题写进手机剪贴板。**抖音里怎么发由用户在任务画布上自己写**, +> 最后放一个 `mark_release`「标记发布结果」回写这里的状态(`published`→`done`、`failed`、`unknown`)。 +> 这样抖音改版时用户改自己的步骤即可,不用等平台发版。 + +> ⚠ **`failed` 与 `unknown` 必须分开**:把"不知道自己发没发"混成"知道自己没发", +> 就是重复发布的来源。`stage` 是两者互相转换的唯一依据(`core/video_plan.STAGE_STATUS`)。 +> 界面上的 `unknown` 卡片标橙 + 硬提示,人工到抖音确认后点「已发出 / 未发出」裁决。 + +**唯一性**:`(phone, release_date, seq)` 由**服务层**保证,**不加 DB 唯一索引** —— +`seq` 从 1 起、没有"空值"可言,做部分唯一索引要写三处方言适配 + MySQL 生成列(§4.2), +收益不匹配;违反的代价只是低频人工上传产生的重复行(可见、可删)。 + +**索引**:`(release_date,status)` 时间线主查询 · `(account_id,release_date)` 账号视角 · +`(phone,release_date,seq)` 配对/幂等 · `(video_file)` 清理时反查引用。 + +这张表**自动进整库备份**(§6 派生规则);**素材文件不进备份包**,见 §6 与 §7。 + --- ## 3. 非模型表 @@ -368,6 +418,11 @@ SUMMARY_TABLES = tuple(sorted(t.name for t in db.metadata.tables.values())) - **导入侧**:备份里出现未登记的表(排除 `sqlite_` 前缀)→ 预览告警(字段 `extra_tables`) - `REQUIRED_TABLES = (app_meta, user, task_job, device_group)`:缺任一直接拒绝导入(这四项是"判定这是不是本平台备份"的最小集合,故意手工维护) +> **素材文件不进备份包**:`create_export()` 只打包 `data/apks/*.apk`,**不含 `data/videos/`** +> (几十 GB 会把"数据库备份"这个核心能力搞坏)。备份的 `manifest.json` 里有 +> `videos.included=false` + 数量/字节数,备份预览页也会提示"素材需另外备份 `data/videos/`" +> —— **不能让人以为备份了**。恢复后素材要重新上传(计划与发布状态、分享链接都在表里,已备份)。 + > **红线**:新增持久化表时**同时补 `TABLE_LABELS` 的中文标签**并更新 [DEPLOY.md](DEPLOY.md) §数据备份。 > 覆盖清单本身不再需要手工登记(已由 metadata 派生)。历史教训:`agent_action` 曾漏登记, > 导致"动作库看起来没备份"(数据其实在快照里,只是清单没列)。 @@ -380,6 +435,7 @@ SUMMARY_TABLES = tuple(sorted(t.name for t in db.metadata.tables.values())) |------|------|--------| | `data/users.db`(+`-wal`/`-shm`) | SQLite 主库(**仅回退模式用**;连 MySQL 时这些文件不被读写,可留作历史归档) | 否 | | `data/apks/*.apk` | 上传的 APK | 否 | +| `data/videos/YYYY-MM/*` | 视频发布计划的素材(几百 MB 一个;**不进整库备份**,见 §6) | 否 | | `data/backups/` | 导出临时 zip、`pre_restore_*.zip`(导入前安全网)、`restore_failed_*` | 否 | | `data/restore_staging//` | 导入暂存(TTL 1800s 自动清理) | 否 | | `data/restore_pending/` | 待生效恢复任务(重启时单事务消费) | 否 | diff --git a/doc/DEPLOY.md b/doc/DEPLOY.md index ec610b3..d2fb4f6 100644 --- a/doc/DEPLOY.md +++ b/doc/DEPLOY.md @@ -212,12 +212,16 @@ tail -20 logs/web.log # 无 ERROR/Traceback SUMMARY_TABLES = tuple(sorted(t.name for t in db.metadata.tables.values())) ``` -当前 16 张表:`app_meta` / `user` / `device_group` / `task_job` / `custom_action` / +当前 17 张表:`app_meta` / `user` / `device_group` / `task_job` / `custom_action` / `apk_file` / `device` / `pending_device` / `agent_conversation` / `agent_experience` / `experience_audit` / `agent_action` / `device_install_log` / `task_step_log` / `done_mark` / -`device_account`(账号台账)。 +`device_account`(账号台账)/ `video_plan`(视频发布计划)。 完整说明见 [DATA_MODEL.md](DATA_MODEL.md) §6。 +> ⚠ **视频素材(`data/videos/`)不进备份包**(几十 GB 会把备份搞坏):manifest 里有 +> `videos.included=false` 与数量/字节数,预览页也会提示 —— **素材要另外备份**。 +> 计划、发布状态与分享链接都在 `video_plan` 表里,随备份一起走。 + > ⚠️ `task_step_log`(任务步骤明细)是会持续增长的表:它按 `KEEP_DAYS` > (默认 14 天,见 `core/step_log.py`)自动清理,但备份包里会带上保留期内的全部行。 > 设备多、任务密时导出 zip 会明显变大——需要更小的包就调小那个常量。 @@ -263,7 +267,13 @@ SUMMARY_TABLES = tuple(sorted(t.name for t in db.metadata.tables.values())) - **HTTPS/80 端口**:用 Nginx 反代 ```nginx -location / { proxy_pass http://127.0.0.1:18050; proxy_set_header Host $host; } +location / { + proxy_pass http://127.0.0.1:18050; + proxy_set_header Host $host; + # ⚠ 必须加:nginx 默认 client_max_body_size 只有 1m, + # 不加的话「应用管理」传 APK(有 336MB 的包)和视频素材上传都会 413。 + client_max_body_size 2048m; +} ``` --- diff --git a/doc/DEVELOPMENT.md b/doc/DEVELOPMENT.md index f495066..50b0a36 100644 --- a/doc/DEVELOPMENT.md +++ b/doc/DEVELOPMENT.md @@ -140,6 +140,8 @@ MCP_ALLOW_WRITE=1 MCP_PLATFORM_USER=admin MCP_PLATFORM_PASS=<密码> \ | `WEB_HOST` / `WEB_PORT` | `0.0.0.0` / `18050` | 硬编码 | 18050 避开 Windows 动态端口段 | | `DATA_DIR` / `APK_DIR` | `data/` / `data/apks/` | 代码计算 | | | `BACKUP_DIR` / `RESTORE_STAGING_DIR` / `RESTORE_PENDING_DIR` | `data/backups` / `data/restore_staging` / `data/restore_pending` | 代码计算 | 备份相关 | +| `VIDEO_DIR` | `data/videos` | `DATA_VIDEO_DIR` | 视频发布计划的素材目录(**不进整库备份**;测试要指到临时目录) | +| `MAX_CONTENT_LENGTH` | 2 GiB | `MAX_UPLOAD_MB` | 单次上传体积上限;**必须大于最大的 APK(实测 336MB)**,否则 APK 上传会被卡死 | | `DEPLOY_ENV` | `dev` | `.env` 可覆盖 | 声明这套配置连哪个环境的库;与库名绑定(dev→`auto_control_dev`,prod→`auto_control`),不符拒绝启动 | | `DB_HOST` / `DB_PORT` / `DB_USER` / `DB_PASSWORD` / `DB_NAME` | 空 / `3306` / 空 / 空 / 空 | `.env` 可覆盖 | MySQL 目标;`DB_HOST` 为空则回退 SQLite(**生产禁止静默回退**,需 `DB_ALLOW_SQLITE_FALLBACK=1`) | | `DB_CHARSET` / `DB_COLLATION` | `utf8mb4` / `utf8mb4_bin` | `.env` 可覆盖 | 排序规则必须用 `_bin`(逐码点比较,等价 SQLite 的大小写敏感语义) | @@ -198,7 +200,7 @@ MCP_ALLOW_WRITE=1 MCP_PLATFORM_USER=admin MCP_PLATFORM_PASS=<密码> \ | 页面结构 / 样式 / 引入脚本 | `templates/admin/monitor.html` | | 公共工具(API/权限/Toast/Tab) | `static/admin/base.js` | | 列表分页排序 | `static/admin/list.js` | -| 各功能域逻辑 | `static/admin/{monitor,editor,tasks,tools,apps,admin,agent,system,markdown,dedup,ledger}.js` | +| 各功能域逻辑 | `static/admin/{monitor,editor,tasks,tools,apps,admin,agent,system,markdown,dedup,ledger,release}.js` | 新增 JS 模块:建文件 → 在 `monitor.html` 里按依赖顺序加 ` + diff --git a/web/__init__.py b/web/__init__.py index 9f15e35..216833c 100644 --- a/web/__init__.py +++ b/web/__init__.py @@ -6,7 +6,8 @@ tasks — 任务计划/分组/自定义动作/元素抓取/步骤测试 admin — 用户管理/日志 tools — adb 终端/剪贴板注入/应用版本 - devices — 设备池管理 + devices — 设备池管理 + 账号台账 + video_plan — 视频发布计划(上传/时间线/链接导出/单条发布) apks — 应用管理 tailscale — Tailscale 管理 system — 系统数据备份导出/导入 @@ -24,6 +25,7 @@ def register_blueprints(app): from .admin_api import bp as admin_bp from .tools_api import bp as tools_bp from .devices_api import bp as devices_bp + from .video_plan_api import bp as video_plan_bp from .apks_api import bp as apks_bp from .tailscale_api import bp as tailscale_bp from . import agent_api as _agent_mod @@ -33,6 +35,6 @@ def register_blueprints(app): from .device_agent_api import bp as devagent_bp from .notify_api import bp as notify_bp for bp in (auth_bp, monitor_bp, tasks_bp, admin_bp, tools_bp, - devices_bp, apks_bp, tailscale_bp, agent_bp, system_bp, + devices_bp, video_plan_bp, apks_bp, tailscale_bp, agent_bp, system_bp, devagent_bp, notify_bp): app.register_blueprint(bp) diff --git a/web/video_plan_api.py b/web/video_plan_api.py new file mode 100644 index 0000000..9d1f019 --- /dev/null +++ b/web/video_plan_api.py @@ -0,0 +1,504 @@ +"""视频发布计划 API(「账号 → 发布计划」页)。 + +数据在 `video_plan` 表、服务层 `core/video_plan.py`;任务侧怎么自动发布见 +`tasks/generic/task.py` 的 `_exec_publish_video`(本文件的「立即发布」调的是同一段逻辑)。 + +红线(与 `core/video_plan.py` 顶部一致):发布前必须原子占位;失败要分 +"推文件阶段(可重试)"与"抖音阶段(结果未知、绝不自动重试)"两类。 +""" +import csv +import io + +from flask import Blueprint, jsonify, request, send_file + +from core import video_plan as vp +from core.logger import get_logger +from web.auth import perm_required, PERM_DEVICES, PERM_TASKS + +_log = get_logger("web") +bp = Blueprint("video_plan", __name__) + + +@bp.route("/api/video_plan/upload_video", methods=["POST"]) +@perm_required(PERM_DEVICES) +def api_vp_upload_video(): + """上传单个视频(前端**逐文件**发,才能给每个文件一条结果与进度)。 + + 文件名规则 `手机号_日期_编号`(编号可省略);解析出的手机号必须**恰好**命中台账一个账号。 + """ + f = request.files.get("file") + if not f or not f.filename: + return jsonify({"ok": False, "error": "没收到文件"}), 400 + replace = str(request.form.get("replace") or "") in ("1", "true", "yes") + order = request.form.get("order") + result = vp.upload_video(f, order=order, replace=replace) + if not result.get("ok"): + return jsonify(result), 400 + return jsonify(result) + + +@bp.route("/api/video_plan/upload_titles", methods=["POST"]) +@perm_required(PERM_DEVICES) +def api_vp_upload_titles(): + """上传标题:`{text, dry_run}`;也接受 multipart 的 txt 文件(字段名 file)。 + + 每行 `标题内容_手机号_日期_编号`(编号可省略、标题里可以有下划线);`#` 开头的行忽略。 + """ + if request.files.get("file"): + raw = request.files["file"].stream.read() + try: + text = raw.decode("utf-8") + except UnicodeDecodeError: + text = raw.decode("gbk", errors="replace") + dry = str(request.form.get("dry_run") or "") in ("1", "true", "yes") + else: + data = request.json or {} + text = data.get("text") or "" + dry = bool(data.get("dry_run")) + if not text.strip(): + return jsonify({"ok": False, "error": "内容是空的"}), 400 + result = vp.import_titles(text, dry_run=dry) + return jsonify(result) + + +@bp.route("/api/video_plan/timeline") +@perm_required(PERM_DEVICES) +def api_vp_timeline(): + """时间线:?range=today|week|all &from=&to=&status=&account_id=&q=&limit=""" + rng = (request.args.get("range") or "").strip() + today = vp._today() + date_from = (request.args.get("from") or "").strip() + date_to = (request.args.get("to") or "").strip() + if rng == "today" and not date_from: + date_from = date_to = today + elif rng == "week" and not date_from: + import time as _t + date_from = today + date_to = _t.strftime("%Y-%m-%d", _t.localtime(_t.time() + 7 * 86400)) + data = vp.timeline(date_from, date_to, + status=(request.args.get("status") or "").strip(), + account_id=(request.args.get("account_id") or "").strip(), + q=(request.args.get("q") or "").strip(), + limit=int(request.args.get("limit") or 500)) + return jsonify({"ok": True, **data}) + + +@bp.route("/api/video_plan/stats") +@perm_required(PERM_DEVICES) +def api_vp_stats(): + """顶部统计(今天要发/已发/待标题/过期未发 + 素材占用与磁盘余量)。""" + return jsonify({"ok": True, "stats": vp.stats(), "disk": vp.disk_info(), + "today": vp._today()}) + + +@bp.route("/api/video_plan/", methods=["GET", "PUT", "DELETE"]) +@perm_required(PERM_DEVICES) +def api_vp_one(pid): + if request.method == "GET": + row = vp.get_plan(pid) + if not row: + return jsonify({"ok": False, "error": "计划不存在"}), 404 + return jsonify({"ok": True, "plan": row}) + if request.method == "DELETE": + if not vp.delete_plan(pid): + return jsonify({"ok": False, "error": "计划不存在"}), 404 + return jsonify({"ok": True, "msg": "已删除(素材文件一并删除)"}) + data = request.json or {} + if "title" in data: + row, err = vp.set_title(pid, data.get("title") or "") + if err: + return jsonify({"ok": False, "error": err}), 400 if "太长" in err else 404 + return jsonify({"ok": True, "msg": "标题已改", "plan": row}) + return jsonify({"ok": False, "error": "没有要改的字段"}), 400 + + +@bp.route("/api/video_plan//video") +@perm_required(PERM_DEVICES) +def api_vp_video(pid): + """预览平台上的素材(文件被清理后返回 410)。""" + import os + from flask import send_file as _send + row = vp.get_plan(pid) + if not row: + return jsonify({"ok": False, "error": "计划不存在"}), 404 + path = vp._video_path(row.get("video_file")) + if not path or not os.path.exists(path): + return jsonify({"ok": False, "error": "素材文件已清理"}), 410 + return _send(path, conditional=True) + + +@bp.route("/api/video_plan//push", methods=["POST"]) +@perm_required(PERM_DEVICES) +def api_vp_push(pid): + """把这一条的素材**推到手机**(只推送,不在手机上发布)。 + + 发布动作由你写在任务里(推送步骤后面接你自己的步骤),这里只是"人工先推一条"。 + """ + import threading + row = vp.get_plan(pid) + if not row: + return jsonify({"ok": False, "error": "计划不存在"}), 404 + # 地址现查(设备名 → 当前地址):**不能用计划行里的快照** —— 设备换 IP 后快照就是错的, + # 而且台账是粘贴导入的老记录可能压根没存上快照(实测踩过:设备池里明明有 A01, + # 却因为快照为空点不动按钮)。快照只作为兜底。 + from core import ledger + serial = ledger.serial_of(row.get("device_name") or "") or (row.get("serial") or "") + if not serial: + return jsonify({"ok": False, + "error": f"找不到设备『{row.get('device_name') or '?'}』的地址 ——" + f"确认「工具 → 设备池」里有这台设备且已连接"}), 400 + if row.get("device_name") and not row.get("serial"): + vp.set_serial(pid, serial) # 顺手把空快照补上(自愈,免得下次又踩) + if row["status"] in (vp.ST_PUSHING, vp.ST_PUBLISHING): + return jsonify({"ok": False, "error": "这条已经推到手机上了,别重复推"}), 409 + if row["status"] == vp.ST_UNKNOWN: + return jsonify({"ok": False, + "error": "这条上次结果未知(可能已经发出去了)—— 先去抖音确认," + "再用「已发出 / 未发出」裁决"}), 409 + from tasks.generic.publish_flow import run_push + t = threading.Thread(target=run_push, args=(serial, pid), daemon=True) + t.start() + _log.info(f"推送素材到手机 video_plan {pid}({row['phone']} {row['release_date']})") + return jsonify({"ok": True, + "msg": "已开始推送到手机(推完状态会变成「已推送到手机」," + "之后你在手机上发;发完点「标记已发布」)"}) + + +def _serial_for_plan(row): + """这条计划要推给哪台设备:**设备名现查**(设备换 IP 后快照是错的),快照兜底。""" + from core import ledger + return (ledger.serial_of(row.get("device_name") or "") or (row.get("serial") or "")).strip() + + +@bp.route("/api/video_plan/push_all", methods=["POST"]) +@perm_required(PERM_DEVICES) +def api_vp_push_all(): + """**一键推送**:把某天(默认今天)所有「待发布/失败」的计划推到各自的手机上。 + + 并发口径:**每台设备一个线程、设备内串行**(同一个手机同时推两个文件会互相干扰: + 两条 adb push 抢带宽 + 两条 `content call scan_file` 抢索引),设备之间并行。 + + 立刻返回受理条数(推送是后台跑的)——进度与结果看时间线的状态徽章/ + 「已推送·相册可见」标记。 + """ + import threading + from core import ledger + data = request.json or {} + date = (data.get("date") or "").strip() or vp._today() + rows = vp.pushable_plans(date, retry_failed=bool(data.get("retry_failed", True))) + if not rows: + return jsonify({"ok": True, "count": 0, "devices": 0, + "msg": f"{date} 没有待推送的计划(有标题的才推)"}) + by_dev, no_dev = {}, [] + for row in rows: + serial = (ledger.serial_of(row.get("device_name") or "") or (row.get("serial") or "")).strip() + (by_dev.setdefault(serial, []) if serial else no_dev).append(row) + by_dev.pop("", None) # 兜底:空 key 不该出现 + if not by_dev: + return jsonify({"ok": False, "error": f"{len(rows)} 条计划都没绑上设备地址 ——" + f"到「工具 → 设备池」确认设备在线"}), 400 + + from tasks.generic.publish_flow import push_one + + def _push_device(serial, plans): + for row in plans: + try: + push_one(serial, row) + except Exception as e: + _log.warning(f"一键推送 {row.get('id')} 异常: {e}") + + for serial, plans in by_dev.items(): + threading.Thread(target=_push_device, args=(serial, plans), daemon=True).start() + _log.info(f"一键推送 {date}:{len(rows)} 条 → {len(by_dev)} 台设备" + f"(没绑设备的 {len(no_dev)} 条已跳过)") + return jsonify({"ok": True, "count": len(rows), "devices": len(by_dev), + "no_device": len(no_dev), + "msg": f"已开始推送 {len(rows)} 条({len(by_dev)} 台设备,每台串行)——" + f"推完状态会变成「已推送·相册可见/相册未见」,刷新看结果"}) + + +@bp.route("/api/video_plan//mark", methods=["POST"]) +@perm_required(PERM_DEVICES) +def api_vp_mark(pid): + """人工标记结果:`{status: done|failed|unknown, why}`(推送到手机之后用)。""" + data = request.json or {} + st = (data.get("status") or "").strip() + if st not in (vp.ST_DONE, vp.ST_FAILED, vp.ST_UNKNOWN): + return jsonify({"ok": False, "error": "status 只能是 done / failed / unknown"}), 400 + why = (data.get("why") or "").strip() + row, err = vp.set_status(pid, st, note=why) + if err: + return jsonify({"ok": False, "error": err}), 400 + _log.info(f"人工标记 video_plan {pid} → {st}") + return jsonify({"ok": True, "msg": "已标记", "plan": row}) + + +@bp.route("/api/video_plan//", methods=["POST"]) +@perm_required(PERM_DEVICES) +def api_vp_action(pid, action): + """单条操作:skip 跳过 / retry 重试(failed→ready)/ resolve 裁决 unknown。""" + data = request.json or {} + if action == "skip": + row, err = vp.set_status(pid, vp.ST_SKIPPED, note=data.get("reason") or "") + msg = "已跳过" + elif action == "retry": + row, err = vp.set_status(pid, vp.ST_READY, note="人工重试") + msg = "已放回待发布(尝试次数已清零)" + elif action == "resolve": + row, err = vp.resolve_unknown(pid, bool(data.get("published")), + note=data.get("note") or "") + msg = "已记为「已发出」" if data.get("published") else "已记为「未发出」(可重试)" + else: + return jsonify({"ok": False, "error": f"未知操作 {action}"}), 400 + if err: + return jsonify({"ok": False, "error": err}), 400 + return jsonify({"ok": True, "msg": msg, "plan": row}) + + +def _norm_target(target): + """规范化任务目标。**分组键必须是 `group_name`** —— `TaskJob.resolve_serials()` + 只认这个键,写成 `group` 会静默解析出 0 台设备(任务"跑完了"但谁也没跑)。""" + t = dict(target or {}) + mode = (t.get("mode") or "all").strip() + out = {"mode": mode} + if mode == "group": + out["group_name"] = (t.get("group_name") or t.get("group") or "").strip() + elif mode == "serial": + out["serial"] = (t.get("serial") or "").strip() + return out + + +def _target_error(target): + if target["mode"] not in ("all", "group", "serial"): + return "目标只能是 全部/分组/单设备" + if target["mode"] == "group" and not target.get("group_name"): + return "选了「设备分组」但没指定分组名" + if target["mode"] == "serial" and not target.get("serial"): + return "选了「单台设备」但没指定设备" + return "" + + +def _parse_hhmm(text): + try: + h, m = [int(x) for x in str(text).split(":")] + assert 0 <= h < 24 and 0 <= m < 60 + return h, m, "" + except Exception: + return 0, 0, "时间格式应为 HH:MM" + + +def _count_steps(steps): + """步骤总数(**含嵌套**:条件判断的 then/else、循环块里的)—— 界面上给人看的数字, + 只数顶层会让人以为任务很短(骨架顶层 4 步、实际 14 步)。""" + n = 0 + for s in steps or []: + if not isinstance(s, dict): + continue + n += 1 + p = s.get("params") or {} + for k in ("then", "else", "children"): + n += _count_steps(p.get(k) or []) + return n + + +def _job_brief(j, mgr): + """任务摘要(列表用)+ 步骤(就地编辑用)。""" + steps = (j.params or {}).get("steps") or [] + try: + nxt = mgr.next_run_of(j) + nxt = nxt.strftime("%Y-%m-%d %H:%M") if nxt else "" + except Exception: + nxt = "" + sch = j.schedule or {} + return {"id": j.id, "name": j.name, "enabled": bool(j.enabled), + "mode": sch.get("mode", "once"), "cron": sch.get("cron", ""), + "next_run": nxt, "target": j.target, "steps": steps, + "step_count": _count_steps(steps), + "retry": j.retry or {}, + "schedule": sch} + + +@bp.route("/api/video_plan/tasks") +@perm_required(PERM_DEVICES) +def api_vp_tasks(): + """**可编辑的任务列表** —— 「发布计划」页顶部那一块。 + + 返回**全部通用步骤任务**,每条带 `is_release`(含「推送/标记发布结果」= 发布任务): + - 只列发布任务的话,手写的抖音发布流程(自己点相册/输入框/发布,没有平台的 + `push_release` 步骤)就会**一条都不显示** —— 人会以为"这页没有能改的地方" + - 界面上按 `is_release` 分成两组:上面「发布任务」,下面「其它任务」也能就地改步骤 + + 带上完整 `steps`:这一页要**就地编辑**这些步骤(抓选择器),不用再跳「任务」页。 + """ + from web import context + release, other = [], [] + try: + jobs = list((context.mgr.jobs or {}).values()) + except Exception: + jobs = [] + for j in jobs: + # 只有通用步骤任务有步骤可编(其它任务类型的 params 是另一套) + if getattr(j, "task_type", "generic_steps") != "generic_steps": + continue + b = _job_brief(j, context.mgr) + b["is_release"] = vp.is_release_task((j.params or {}).get("steps") or []) + b["has_push"] = vp.has_push_step((j.params or {}).get("steps") or []) + (release if b["is_release"] else other).append(b) + release.sort(key=lambda x: x["name"]) + other.sort(key=lambda x: x["name"]) + return jsonify({"ok": True, "tasks": release + other, + "release_count": len(release), "other_count": len(other)}) + + +@bp.route("/api/video_plan/tasks", methods=["POST"]) +@perm_required(PERM_TASKS) +def api_vp_task_create(): + """一键新建「标准发布任务」:账号校验 → 推送 → 抖音点击(占位待你抓)→ 填标题 → 标记。 + + 请求:{"name": "发布视频", "target": {"mode":"all"|"group"|"serial", …}, + "time": "10:00", "enabled": true} + **走 TaskManager 建**(内存 + 库一起更新)—— 直接写库的话调度器不认。 + """ + from web import context + data = request.json or {} + name = (data.get("name") or "发布视频").strip()[:120] + target = _norm_target(data.get("target")) + err = _target_error(target) + if err: + return jsonify({"ok": False, "error": err}), 400 + h, m, err = _parse_hhmm((data.get("time") or "10:00").strip()) + if err: + return jsonify({"ok": False, "error": err}), 400 + schedule = {"mode": "cron", "cron": f"{m} {h} * * *", "stop_cron": "", + "window": {"start": "", "end": ""}} + params = {"steps": vp.build_release_steps(), + "dedup_reset": "day", "dedup_hours": 6, + "preempt": False, "skip_offline": True, "max_duration": 0} + try: + job = context.mgr.add_job(name=name, task_type="generic_steps", target=target, + params=params, schedule=schedule, + retry={"max_attempts": 1, "delay": 60}, + enabled=bool(data.get("enabled", True))) + except Exception as e: + _log.warning(f"新建发布任务失败: {e}") + return jsonify({"ok": False, "error": f"创建失败: {str(e)[:120]}"}), 500 + _log.info(f"新建发布任务 {job.name}({job.id}),每天 " + f"{h:02d}:{m:02d},目标 {target}") + return jsonify({"ok": True, "msg": f"已创建「{name}」(每天 {h:02d}:{m:02d} 跑)——" + f"就在下面这一块把抖音那几步的选择器抓一下就能用了", + "job_id": job.id, "task": _job_brief(job, context.mgr)}) + + +@bp.route("/api/video_plan/tasks//adopt", methods=["POST"]) +@perm_required(PERM_TASKS) +def api_vp_task_adopt(job_id): + """给**已有任务**插上平台那两步(`push_release` 放最前、`mark_release` 放最后)。 + + 手写的抖音发布流程常缺这两步:不插的话计划页的视频永远不会被推到手机、 + 发布结果也不会回写状态。**中间你自己的步骤一步不动**(位置理由见 + `core.video_plan.insert_release_steps`)。 + """ + from web import context + job = (context.mgr.jobs or {}).get(job_id) + if not job: + return jsonify({"ok": False, "error": "任务不存在"}), 404 + if getattr(job, "task_type", "generic_steps") != "generic_steps": + return jsonify({"ok": False, "error": "只有「通用步骤」任务能插平台发布步骤"}), 400 + steps = (job.params or {}).get("steps") or [] + if not steps: + return jsonify({"ok": False, "error": "这个任务没有步骤,先去「任务」页把抖音流程画好"}), 400 + if vp.has_push_step(steps) and vp.has_mark_step(steps): + return jsonify({"ok": False, "error": "这个任务已经有平台发布步骤了"}), 409 + params = dict(job.params or {}) + params["steps"] = vp.insert_release_steps(steps) + try: + context.mgr.update_job(job_id, params=params) + except Exception as e: + _log.warning(f"给任务 {job_id} 插平台发布步骤失败: {e}") + return jsonify({"ok": False, "error": f"保存失败: {str(e)[:120]}"}), 500 + _log.info(f"给任务 {job.name}({job_id}) 插上平台发布步骤") + return jsonify({"ok": True, "msg": "已插入「推送发布视频」(最前)与「标记发布结果」(最后)——" + "你原来的步骤没动;到「账号 → 发布计划」里核对一下顺序", + "task": _job_brief(job, context.mgr)}) + + +@bp.route("/api/video_plan/tasks/", methods=["GET", "PUT"]) +@perm_required(PERM_TASKS) +def api_vp_task_one(job_id): + """就地编辑发布任务(「账号 → 发布计划」页里,不用跳「任务」页)。 + + 只放开这一页用得上的字段:`name` / `time`(HH:MM) / `target` / `steps` / `enabled`。 + **没传的字段原样保留**(深合并:`params` 里还有去重有效期、抢占等通用参数, + 整块覆盖会把它们抹掉)。底层仍是 `context.mgr.update_job()` —— 内存 + 库一起改, + 否则调度器不认(直接写库的坑见 `core/task_manager` 的 `_save_jobs`)。 + """ + from web import context + job = (context.mgr.jobs or {}).get(job_id) + if not job: + return jsonify({"ok": False, "error": "任务不存在"}), 404 + if request.method == "GET": + return jsonify({"ok": True, "task": _job_brief(job, context.mgr)}) + data = request.json or {} + fields = {} + if "name" in data: + nm = (data.get("name") or "").strip()[:120] + if not nm: + return jsonify({"ok": False, "error": "任务名不能为空"}), 400 + fields["name"] = nm + if "target" in data: + target = _norm_target(data.get("target")) + err = _target_error(target) + if err: + return jsonify({"ok": False, "error": err}), 400 + fields["target"] = target + if "enabled" in data: + fields["enabled"] = bool(data.get("enabled")) + if "time" in data: + h, m, err = _parse_hhmm((data.get("time") or "").strip()) + if err: + return jsonify({"ok": False, "error": err}), 400 + sch = dict(job.schedule or {}) + sch["mode"] = "cron" + sch["cron"] = f"{m} {h} * * *" + sch.setdefault("stop_cron", "") + sch.setdefault("window", {"start": "", "end": ""}) + fields["schedule"] = sch + if "steps" in data: + steps = data.get("steps") + if not isinstance(steps, list): + return jsonify({"ok": False, "error": "steps 必须是数组"}), 400 + params = dict(job.params or {}) + params["steps"] = steps + fields["params"] = params + if not fields: + return jsonify({"ok": False, "error": "没有要改的字段"}), 400 + try: + context.mgr.update_job(job_id, **fields) + except Exception as e: + _log.warning(f"改发布任务 {job_id} 失败: {e}") + return jsonify({"ok": False, "error": f"保存失败: {str(e)[:120]}"}), 500 + _log.info(f"改发布任务 {job.name}({job_id}):{'/'.join(fields)}") + return jsonify({"ok": True, "msg": "已保存", "task": _job_brief(job, context.mgr)}) + + +@bp.route("/api/video_plan/links") +@perm_required(PERM_DEVICES) +def api_vp_links(): + """已发布作品的分享链接(界面展示 / `format=csv` 导出,给后续铺评论用)。""" + rows = vp.links(date_from=(request.args.get("from") or "").strip(), + date_to=(request.args.get("to") or "").strip(), + device=(request.args.get("device") or "").strip(), + phone=(request.args.get("phone") or "").strip()) + if (request.args.get("format") or "") == "csv": + buf = io.StringIO() + w = csv.writer(buf) + w.writerow(["发布日期", "设备号", "手机号", "账号名称", "抖音号", + "编号", "标题", "分享链接", "发布时间"]) + for r in rows: + w.writerow([r["release_date"], r["device_name"], r["phone"], r["nickname"], + r["douyin_id"], r["seq"], r["title"], r["share_url"], + r["published_at"]]) + data = ("" + buf.getvalue()).encode("utf-8") # BOM:Excel 打开中文不乱码 + return send_file(io.BytesIO(data), mimetype="text/csv", as_attachment=True, + download_name=f"video_links_{vp._today()}.csv") + return jsonify({"ok": True, "links": rows, "count": len(rows)}) diff --git a/web_server.py b/web_server.py index 6b865d1..11442ef 100644 --- a/web_server.py +++ b/web_server.py @@ -35,6 +35,7 @@ app.config["SECRET_KEY"] = _web_secret app.config["TEMPLATES_AUTO_RELOAD"] = True # 数据库目标与连接参数统一由 core/db_config 装配(.env 的 DEPLOY_ENV / DB_* 决定)。 # 这里就把配置校验做完:连不上、环境与库名不匹配 → 直接拒绝启动,不要带着错配置跑起来。 +from config import MAX_CONTENT_LENGTH as _MAX_UPLOAD from core import db_config try: _DB_URI = db_config.build_db_uri() @@ -45,11 +46,26 @@ except db_config.DBConfigError as _e: app.config["SQLALCHEMY_DATABASE_URI"] = _DB_URI app.config["SQLALCHEMY_ENGINE_OPTIONS"] = db_config.engine_options(_DB_URI) app.config["SQLALCHEMY_TRACK_MODIFICATIONS"] = False +# 单次请求体上限(视频发布计划要传几百 MB 的视频素材)。 +# 默认 2GiB:**必须大于现有 APK 上传的最大包(实测 336MB)**,否则会把「应用管理」的 +# APK 上传一起卡死。原来是没设的(Werkzeug 默认无上限)—— 现在有了上限, +# 超限必须回 **JSON**(见下面 413 处理器),否则前端 JSON.parse 会炸成 +# "解析响应失败",用户看不出是文件太大。 +app.config["MAX_CONTENT_LENGTH"] = _MAX_UPLOAD login_manager = LoginManager(app) login_manager.login_view = "auth.login" # 蓝图化后路由前缀 auth +@app.errorhandler(413) +def _payload_too_large(e): + """上传超限要说人话(默认返回 HTML,前端 JSON.parse 会炸成"解析响应失败")。""" + from flask import jsonify as _jsonify + mb = _MAX_UPLOAD // (1024 * 1024) + return _jsonify({"ok": False, + "error": f"文件太大:单次上传上限 {mb} MB(可在 .env 里调 MAX_UPLOAD_MB)"}), 413 + + def _consume_pending_restore(): """消费「待生效的备份恢复」(若存在 data/restore_pending/users.db)。 @@ -107,6 +123,9 @@ dedup.init_app(app) # 账号台账:绑 app 供任务线程(条件判断取号)/ Web / 设备端身份页自推 context from core import ledger ledger.init_app(app) +# 视频发布计划:绑 app 供任务线程(发布步骤)/ Web / 清理 job 自推 context +from core import video_plan +video_plan.init_app(app) mgr = TaskManager(app=app) # 设备换地址(认领/迁址)后同步分组与任务的引用:它们的"内存副本"在 TaskManager 里 # (调度用内存对象),device_pool 只改库不改内存不生效,所以在此注册回调打通 @@ -154,6 +173,35 @@ def _purge_done_mark(): _log.warning(f"去重记录清理失败(不影响主服务): {e}") +def _purge_video_plan(): + """发布计划:僵尸占位回收 + 过期行清理(每日 04:41)。 + + 僵尸回收是**必须有的一步**:占位(pushing/publishing)后崩溃的话,没有它 + 这条计划永远停在"发布中",再也不会被取到 —— 表现成"这个号就是不发"。 + """ + try: + from core.video_plan import reclaim_zombies, purge_plan + with app.app_context(): + reclaim_zombies() + purge_plan() + except Exception as e: + _log.warning(f"发布计划清理失败(不影响主服务): {e}") + + +def _purge_video_files(): + """发布计划:平台素材文件清理(每日 04:47)。 + + **只删"已发布 + 已抓到分享链接 + 超过保留期"的**:没抓到链接的一律留着 + (链接与素材至少留一个),未发布/失败/未知的更不删。 + """ + try: + from core.video_plan import purge_files + with app.app_context(): + purge_files() + except Exception as e: + _log.warning(f"发布素材清理失败(不影响主服务): {e}") + + # 经验库每日 AI 巡检(凌晨 03:47):评审标记疑似问题经验,删除只走人工确认。 # 巡检线程在 run_experience_audit 内自建 app context,不依赖这里。 try: @@ -167,6 +215,10 @@ try: _audit_sched.add_job(_purge_step_log, CronTrigger(hour=4, minute=13)) # 去重记录清理:只清 day/hours 桶(all 永不清理,清了等于"只做一次"失效) _audit_sched.add_job(_purge_done_mark, CronTrigger(hour=4, minute=23)) + # 发布计划:僵尸占位回收 + 过期行清理(04:41)、孤儿素材清理(04:47)。 + # 时间刻意与 04:13/04:23 错开 15 分钟以上:都是删数据的大操作,别挤同一分钟。 + _audit_sched.add_job(_purge_video_plan, CronTrigger(hour=4, minute=41)) + _audit_sched.add_job(_purge_video_files, CronTrigger(hour=4, minute=47)) _audit_sched.start() _log.info("经验库每日巡检已注册(03:47 Asia/Shanghai)") except Exception as e: