"""视频发布计划:批量上传 → 自动配对 → 时间线 → 自动发布 → 存分享链接。 **它解决什么**(用户场景):人员在网页上批量传视频与标题 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) 标题内容_手机号_日期 ② 手机号在前(左锚定,跟视频文件名同序 —— "文件名去掉扩展名 + 标题"): C) 手机号_日期_编号_标题 D) 手机号_日期_标题 为什么**先右后左**:右锚定是主格式(标题里带下划线也不会错配)。 左锚定的四种里先试带编号的 C(`手机号_日期_1_标题` —— 那个 `1` 当编号, 不当标题的开头),保证同一行永远只有一种解释。 什么时候两种都不会误判:右锚定要求"倒数第 3/2 段是手机号",左锚定要求 "第 1 段是手机号、第 2 段是日期" —— 一行同时满足两边时(例如 `手机号_日期_1`) 右锚定先命中,且标题为空会被下面的校验拒掉,不会静默发出一条错文案。 """ t = (line or "").strip() if not t: return None, "空行" parts = t.split("_") phone = date = seq = None title_parts = 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] # ② 左锚定(手机号在前,跟视频文件名同序):同样先试带编号的 elif len(parts) >= 4 and _PHONE_RE.match(parts[0]) and _DATE8_RE.match(parts[1]) \ and _SEQ_RE.match(parts[2]): phone, date8, seq = parts[0], parts[1], int(parts[2]) title_parts = parts[3:] elif len(parts) >= 3 and _PHONE_RE.match(parts[0]) and _DATE8_RE.match(parts[1]): # `手机号_日期_1` 这种"只有编号、没有标题"要**拒掉**,不能被当成标题"1" # (静默生成一条标题是"1"的文案,比拒收危险得多;标题真是纯数字就用"标题在前"的写法) if len(parts) == 3 and _SEQ_RE.match(parts[2]): return None, ("只有手机号/日期/编号、没有标题文字 ——" "标题是纯数字的话请用「标题_手机号_日期」写法") phone, date8, title_parts = parts[0], 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)。空行跳过;`#` 开头**只有认不出标题行时**才算注释。 ⚠ 顺序是"先当标题解析、再当注释":`#中秋快乐_19286456013_20260929` 是一条**正常标题** (`#话题` 本来就长这样),过去的写法会把整行当注释**静默丢掉**。现在只有像 `# 这是注释` 这种解析不出手机号/日期的行才按注释忽略。 """ 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: continue meta, err = parse_title_line(s) if err: if s.startswith("#"): # 认不出 + 以 # 开头 → 当注释,静默跳过 continue 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}