Files
EveryPublish/sau_backend.py
T

1308 lines
50 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import asyncio
import base64
import os
import sqlite3
import subprocess
import threading
import time
import uuid
from datetime import datetime
from functools import wraps
from pathlib import Path
from queue import Queue
from flask_cors import CORS
from myUtils.auth import check_cookie
from flask import Flask, request, jsonify, Response, render_template, send_from_directory
from werkzeug.utils import secure_filename
from conf import BASE_DIR
from myUtils.postVideo import (
post_video_tencent,
post_video_DouYin,
post_video_ks,
post_video_xhs,
post_video_bilibili,
post_video_baijiahao,
post_video_alipay,
post_video_weibo,
post_video_hupu,
post_video_youtube,
post_video_tiktok,
post_note_xiaohongshu,
post_note_douyin,
post_note_kuaishou,
post_note_tencent,
)
from uploader.baijiahao_uploader.main import baijiahao_setup
from uploader.alipay_uploader.main import alipay_setup
from uploader.weibo_uploader.main import weibo_setup
from uploader.hupu_uploader.main import hupu_setup
from uploader.youtube_uploader.main import youtube_setup
from uploader.douyin_uploader.main import douyin_setup
from uploader.ks_uploader.main import ks_setup
from uploader.xiaohongshu_uploader.main import xiaohongshu_setup
from uploader.tencent_uploader.main import tencent_setup
from uploader.bilibili_uploader.runtime import ensure_biliup_binary
from loguru import logger
from utils.settings import get_setting, set_setting, get_bool
# ===== 发布任务监控(异步发布 + 进度 + 日志) =====
PUBLISH_TASKS = {}
CHECK_TASKS = {}
_publish_thread_local = threading.local()
_publish_lock = threading.Lock()
SUCCESS_KEYWORDS = ["发布成功", "发布完成", "published success", "video published", "已发布", "发布完毕", "上传成功"]
def _logs_indicate_success(logs):
return any(any(k in log for k in SUCCESS_KEYWORDS) for log in logs)
def _update_task_progress(task, text, level):
ts = datetime.now().strftime("%H:%M:%S")
task["logs"].append(ts + " " + text)
if len(task["logs"]) > 800:
task["logs"] = task["logs"][-800:]
if any(k in text for k in ["发布成功", "published success", "发布完成", "已发布", "video published"]):
task["step"] = "发布成功"
task["progress"] = 100
elif any(k in text for k in ["发布", "publish", "点击发布", "提交发布"]):
task["step"] = "发布中"
task["progress"] = max(task["progress"], 85)
elif any(k in text for k in ["封面", "cover"]):
task["step"] = "处理封面"
task["progress"] = max(task["progress"], 70)
elif any(k in text for k in ["标题", "话题", "标签", "title", "tag"]):
task["step"] = "填写信息"
task["progress"] = max(task["progress"], 60)
elif any(k in text for k in ["上传", "upload", "搬运", "选择文件"]):
task["step"] = "上传中"
task["progress"] = max(task["progress"], 30)
elif any(k in text for k in ["失败", "错误", "error", "超时", "timeout"]):
task["step"] = "出现异常"
def _publish_log_sink(message):
try:
task = getattr(_publish_thread_local, "publish_task", None)
if not task:
return
record = message.record
text = str(record.get("message", ""))
level = record.get("level", {})
level_name = getattr(level, "name", "INFO") if level else "INFO"
_update_task_progress(task, text, level_name)
except Exception:
pass
logger.add(_publish_log_sink, level="INFO", format="{message}")
def _dispatch_publish(type, title, file_list, account_list, tags, category, enableTimer, videos_per_day, daily_times, start_days, thumbnail_path, desc, productLink, productTitle, is_draft, is_image):
match type:
case 1:
if is_image:
post_note_xiaohongshu(title, file_list, tags, account_list, desc)
else:
post_video_xhs(title, file_list, tags, account_list, category, enableTimer, videos_per_day, daily_times, start_days)
case 2:
if is_image:
post_note_tencent(title, file_list, tags, account_list, desc, is_draft)
else:
post_video_tencent(title, file_list, tags, account_list, category, enableTimer, videos_per_day, daily_times, start_days, is_draft)
case 3:
if is_image:
post_note_douyin(title, file_list, tags, account_list, desc)
else:
post_video_DouYin(title, file_list, tags, account_list, category, enableTimer, videos_per_day, daily_times, start_days, thumbnail_path, productLink, productTitle)
case 4:
if is_image:
post_note_kuaishou(title, file_list, tags, account_list, desc)
else:
post_video_ks(title, file_list, tags, account_list, category, enableTimer, videos_per_day, daily_times, start_days)
case 5:
post_video_bilibili(title, file_list, tags, account_list, category, enableTimer, videos_per_day, daily_times, start_days, desc)
case 6:
post_video_baijiahao(title, file_list, tags, account_list, category, enableTimer, videos_per_day, daily_times, start_days, thumbnail_path, desc)
case 7:
post_video_alipay(title, file_list, tags, account_list, category, enableTimer, videos_per_day, daily_times, start_days, thumbnail_path, desc)
case 8:
post_video_weibo(title, file_list, tags, account_list, category, enableTimer, videos_per_day, daily_times, start_days, thumbnail_path, desc)
case 9:
post_video_hupu(title, file_list, tags, account_list, category, enableTimer, videos_per_day, daily_times, start_days, thumbnail_path, desc)
case 10:
post_video_youtube(title, file_list, tags, account_list, category, enableTimer, videos_per_day, daily_times, start_days, thumbnail_path, desc)
case 11:
post_video_tiktok(title, file_list, tags, account_list, category, enableTimer, videos_per_day, daily_times, start_days, thumbnail_path, desc)
case _:
raise ValueError("不支持的平台类型: " + str(type))
def _run_publish_task(task, type, title, file_list, account_list, tags, category, enableTimer, videos_per_day, daily_times, start_days, thumbnail_path, desc, productLink, productTitle, is_draft, is_image):
_publish_thread_local.publish_task = task
task["status"] = "running"
task["step"] = "启动浏览器"
task["progress"] = 5
try:
_dispatch_publish(type, title, file_list, account_list, tags, category, enableTimer, videos_per_day, daily_times, start_days, thumbnail_path, desc, productLink, productTitle, is_draft, is_image)
task["status"] = "success"
task["step"] = "发布完成"
task["progress"] = 100
except Exception as e:
if _logs_indicate_success(task["logs"]):
task["status"] = "success"
task["step"] = "发布成功(平台已确认)"
task["progress"] = 100
else:
task["status"] = "failed"
task["error"] = str(e)
task["step"] = "发布失败"
finally:
_publish_thread_local.publish_task = None
def _create_publish_task(data):
file_list = data.get('fileList', [])
account_list = data.get('accountList', [])
type = data.get('type')
title = data.get('title')
tags = data.get('tags')
category = data.get('category')
if category == 0:
category = None
productLink = data.get('productLink', '')
productTitle = data.get('productTitle', '')
thumbnail_path = data.get('thumbnail', '')
desc = data.get('desc', '')
enableTimer = data.get('enableTimer')
is_draft = data.get('isDraft', False)
videos_per_day = data.get('videosPerDay')
daily_times = data.get('dailyTimes')
start_days = data.get('startDays')
if not file_list:
return None, "文件列表不能为空"
if not account_list:
return None, "账号列表不能为空"
if not type:
return None, "平台类型不能为空"
if not title:
return None, "标题不能为空"
_IMG_EXTS = {".jpg", ".jpeg", ".png", ".webp", ".bmp", ".gif"}
is_image = bool(file_list) and all(Path(f).suffix.lower() in _IMG_EXTS for f in file_list)
if is_image and type not in (1, 2, 3, 4):
return None, "图文发布仅支持小红书/抖音/快手,请选择支持图文的平台"
task_id = str(uuid.uuid4())
task = {
"id": task_id,
"status": "pending",
"step": "排队中",
"progress": 0,
"logs": [],
"error": None,
"type": type,
"title": title,
"created_at": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
}
with _publish_lock:
PUBLISH_TASKS[task_id] = task
thread = threading.Thread(
target=_run_publish_task,
args=(task, type, title, file_list, account_list, tags, category, enableTimer, videos_per_day, daily_times, start_days, thumbnail_path, desc, productLink, productTitle, is_draft, is_image),
daemon=True,
)
thread.start()
return task_id, None
active_queues = {}
app = Flask(__name__)
#允许所有来源跨域访问
CORS(app)
# 限制上传文件大小为160MB
app.config['MAX_CONTENT_LENGTH'] = 160 * 1024 * 1024
# 获取当前目录(假设 index.html 和 assets 在这里)
current_dir = os.path.dirname(os.path.abspath(__file__))
# 处理所有静态资源请求(未来打包用)
@app.route('/assets/<filename>')
def custom_static(filename):
return send_from_directory(os.path.join(current_dir, 'assets'), filename)
# 处理 favicon.ico 静态资源(未来打包用)
@app.route('/favicon.ico')
def favicon():
return send_from_directory(os.path.join(current_dir, 'assets'), 'vite.svg')
@app.route('/vite.svg')
def vite_svg():
return send_from_directory(os.path.join(current_dir, 'assets'), 'vite.svg')
# (未来打包用)
@app.route('/')
def index(): # put application's code here
return send_from_directory(current_dir, 'index.html')
@app.route('/upload', methods=['POST'])
def upload_file():
if 'file' not in request.files:
return jsonify({
"code": 400,
"data": None,
"msg": "No file part in the request"
}), 400
file = request.files['file']
if file.filename == '':
return jsonify({
"code": 400,
"data": None,
"msg": "No selected file"
}), 400
try:
# 保存文件到指定位置
uuid_v1 = uuid.uuid1()
print(f"UUID v1: {uuid_v1}")
safe_name = secure_filename(file.filename)
if not safe_name:
return jsonify({"code": 400, "data": None, "msg": "Invalid filename"}), 400
filepath = Path(BASE_DIR / "videoFile" / f"{uuid_v1}_{safe_name}")
file.save(filepath)
return jsonify({"code":200,"msg": "File uploaded successfully", "data": f"{uuid_v1}_{safe_name}"}), 200
except Exception as e:
return jsonify({"code":500,"msg": str(e),"data":None}), 500
@app.route('/getFile', methods=['GET'])
def get_file():
# 获取 filename 参数
filename = request.args.get('filename')
if not filename:
return jsonify({"code": 400, "msg": "filename is required", "data": None}), 400
# 防止路径穿越攻击
if '..' in filename or filename.startswith('/'):
return jsonify({"code": 400, "msg": "Invalid filename", "data": None}), 400
# 拼接完整路径
file_path = str(Path(BASE_DIR / "videoFile"))
# 返回文件
return send_from_directory(file_path,filename)
@app.route('/uploadSave', methods=['POST'])
def upload_save():
if 'file' not in request.files:
return jsonify({
"code": 400,
"data": None,
"msg": "No file part in the request"
}), 400
file = request.files['file']
if file.filename == '':
return jsonify({
"code": 400,
"data": None,
"msg": "No selected file"
}), 400
# 获取表单中的自定义文件名(可选)
custom_filename = request.form.get('filename', None)
if custom_filename:
filename = secure_filename(custom_filename + "." + file.filename.split('.')[-1])
else:
filename = secure_filename(file.filename)
if not filename:
return jsonify({"code": 400, "data": None, "msg": "Invalid filename"}), 400
try:
# 生成 UUID v1
uuid_v1 = uuid.uuid1()
print(f"UUID v1: {uuid_v1}")
# 构造文件名和路径
final_filename = f"{uuid_v1}_{filename}"
filepath = Path(BASE_DIR / "videoFile" / f"{uuid_v1}_{filename}")
# 保存文件
file.save(filepath)
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
cursor = conn.cursor()
cursor.execute('''
INSERT INTO file_records (filename, filesize, file_path)
VALUES (?, ?, ?)
''', (filename, round(float(os.path.getsize(filepath)) / (1024 * 1024),2), final_filename))
conn.commit()
print("✅ 上传文件已记录")
return jsonify({
"code": 200,
"msg": "File uploaded and saved successfully",
"data": {
"filename": filename,
"filepath": final_filename
}
}), 200
except Exception as e:
print(f"Upload failed: {e}")
return jsonify({
"code": 500,
"msg": f"upload failed: {e}",
"data": None
}), 500
@app.route('/getFiles', methods=['GET'])
def get_all_files():
try:
# 使用 with 自动管理数据库连接
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
conn.row_factory = sqlite3.Row # 允许通过列名访问结果
cursor = conn.cursor()
# 查询所有记录
cursor.execute("SELECT * FROM file_records")
rows = cursor.fetchall()
# 将结果转为字典列表,并提取UUID
data = []
for row in rows:
row_dict = dict(row)
# 从 file_path 中提取 UUID (文件名的第一部分,下划线前)
if row_dict.get('file_path'):
file_path_parts = row_dict['file_path'].split('_', 1) # 只分割第一个下划线
if len(file_path_parts) > 0:
row_dict['uuid'] = file_path_parts[0] # UUID 部分
else:
row_dict['uuid'] = ''
else:
row_dict['uuid'] = ''
data.append(row_dict)
return jsonify({
"code": 200,
"msg": "success",
"data": data
}), 200
except Exception as e:
return jsonify({
"code": 500,
"msg": str("get file failed!"),
"data": None
}), 500
@app.route("/getAccounts", methods=['GET'])
def getAccounts():
"""快速获取所有账号信息,不进行cookie验证"""
try:
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
conn.row_factory = sqlite3.Row
cursor = conn.cursor()
cursor.execute('''
SELECT * FROM user_info''')
rows = cursor.fetchall()
rows_list = [list(row) for row in rows]
print("\n📋 当前数据表内容(快速获取):")
for row in rows:
print(row)
return jsonify(
{
"code": 200,
"msg": None,
"data": rows_list
}), 200
except Exception as e:
print(f"获取账号列表时出错: {str(e)}")
return jsonify({
"code": 500,
"msg": f"获取账号列表失败: {str(e)}",
"data": None
}), 500
@app.route("/getValidAccounts",methods=['GET'])
async def getValidAccounts():
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
cursor = conn.cursor()
cursor.execute('''
SELECT * FROM user_info''')
rows = cursor.fetchall()
rows_list = [list(row) for row in rows]
print("\n📋 当前数据表内容:")
for row in rows:
print(row)
for row in rows_list:
flag = await check_cookie(row[1],row[2])
if not flag:
row[4] = 0
cursor.execute('''
UPDATE user_info
SET status = ?
WHERE id = ?
''', (0,row[0]))
conn.commit()
print("✅ 用户状态已更新")
for row in rows:
print(row)
return jsonify(
{
"code": 200,
"msg": None,
"data": rows_list
}),200
@app.route('/deleteFile', methods=['GET'])
def delete_file():
file_id = request.args.get('id')
if not file_id or not file_id.isdigit():
return jsonify({
"code": 400,
"msg": "Invalid or missing file ID",
"data": None
}), 400
try:
# 获取数据库连接
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
conn.row_factory = sqlite3.Row
cursor = conn.cursor()
# 查询要删除的记录
cursor.execute("SELECT * FROM file_records WHERE id = ?", (file_id,))
record = cursor.fetchone()
if not record:
return jsonify({
"code": 404,
"msg": "File not found",
"data": None
}), 404
record = dict(record)
# 获取文件路径并删除实际文件
file_path = Path(BASE_DIR / "videoFile" / record['file_path'])
if file_path.exists():
try:
file_path.unlink() # 删除文件
print(f"✅ 实际文件已删除: {file_path}")
except Exception as e:
print(f"⚠️ 删除实际文件失败: {e}")
# 即使删除文件失败,也要继续删除数据库记录,避免数据不一致
else:
print(f"⚠️ 实际文件不存在: {file_path}")
# 删除数据库记录
cursor.execute("DELETE FROM file_records WHERE id = ?", (file_id,))
conn.commit()
return jsonify({
"code": 200,
"msg": "File deleted successfully",
"data": {
"id": record['id'],
"filename": record['filename']
}
}), 200
except Exception as e:
return jsonify({
"code": 500,
"msg": str("delete failed!"),
"data": None
}), 500
@app.route('/deleteAccount', methods=['GET'])
def delete_account():
account_id = request.args.get('id')
if not account_id or not account_id.isdigit():
return jsonify({
"code": 400,
"msg": "Invalid or missing account ID",
"data": None
}), 400
account_id = int(account_id)
try:
# 获取数据库连接
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
conn.row_factory = sqlite3.Row
cursor = conn.cursor()
# 查询要删除的记录
cursor.execute("SELECT * FROM user_info WHERE id = ?", (account_id,))
record = cursor.fetchone()
if not record:
return jsonify({
"code": 404,
"msg": "account not found",
"data": None
}), 404
record = dict(record)
# 删除关联的cookie文件
if record.get('filePath'):
cookie_file_path = Path(BASE_DIR / "cookiesFile" / record['filePath'])
if cookie_file_path.exists():
try:
cookie_file_path.unlink()
print(f"✅ Cookie文件已删除: {cookie_file_path}")
except Exception as e:
print(f"⚠️ 删除Cookie文件失败: {e}")
# 删除数据库记录
cursor.execute("DELETE FROM user_info WHERE id = ?", (account_id,))
conn.commit()
return jsonify({
"code": 200,
"msg": "account deleted successfully",
"data": None
}), 200
except Exception as e:
return jsonify({
"code": 500,
"msg": f"delete failed: {str(e)}",
"data": None
}), 500
# SSE 登录接口
@app.route('/login')
def login():
# 1 小红书 2 视频号 3 抖音 4 快手
type = request.args.get('type')
# 账号名
id = request.args.get('id')
# 模拟一个用于异步通信的队列
status_queue = Queue()
active_queues[id] = status_queue
def on_close():
print(f"清理队列: {id}")
del active_queues[id]
# 启动异步任务线程
thread = threading.Thread(target=run_async_function, args=(type,id,status_queue), daemon=True)
thread.start()
response = Response(sse_stream(status_queue,), mimetype='text/event-stream')
response.headers['Cache-Control'] = 'no-cache'
response.headers['X-Accel-Buffering'] = 'no' # 关键:禁用 Nginx 缓冲
response.headers['Content-Type'] = 'text/event-stream'
response.headers['Connection'] = 'keep-alive'
return response
@app.route('/postVideo', methods=['POST'])
def postVideo():
data = request.get_json()
if not data:
return jsonify({"code": 400, "msg": "请求数据不能为空", "data": None}), 400
task_id, err = _create_publish_task(data)
if err:
return jsonify({"code": 400, "msg": err, "data": None}), 400
return jsonify({"code": 200, "msg": "发布任务已提交", "data": {"taskId": task_id}}), 200
@app.route('/postVideoBatch', methods=['POST'])
def postVideoBatch():
data_list = request.get_json()
if not isinstance(data_list, list):
return jsonify({"code": 400, "msg": "Expected a JSON array", "data": None}), 400
task_ids = []
for data in data_list:
task_id, err = _create_publish_task(data)
if task_id:
task_ids.append(task_id)
if not task_ids:
return jsonify({"code": 400, "msg": "没有可提交的发布任务", "data": None}), 400
return jsonify({"code": 200, "msg": None, "data": {"taskIds": task_ids}}), 200
@app.route('/uploadCookie', methods=['POST'])
def upload_cookie():
try:
if 'file' not in request.files:
return jsonify({
"code": 400,
"msg": "没有找到Cookie文件",
"data": None
}), 400
file = request.files['file']
if file.filename == '':
return jsonify({
"code": 400,
"msg": "Cookie文件名不能为空",
"data": None
}), 400
if not file.filename.endswith('.json'):
return jsonify({
"code": 400,
"msg": "Cookie文件必须是JSON格式",
"data": None
}), 400
# 获取账号信息
account_id = request.form.get('id')
platform = request.form.get('platform')
if not account_id or not platform:
return jsonify({
"code": 400,
"msg": "缺少账号ID或平台信息",
"data": None
}), 400
# 从数据库获取账号的文件路径
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
conn.row_factory = sqlite3.Row
cursor = conn.cursor()
cursor.execute('SELECT filePath FROM user_info WHERE id = ?', (account_id,))
result = cursor.fetchone()
if not result:
return jsonify({
"code": 500,
"msg": "账号不存在",
"data": None
}), 404
# 保存上传的Cookie文件到对应路径
cookie_file_path = Path(BASE_DIR / "cookiesFile" / result['filePath'])
cookie_file_path.parent.mkdir(parents=True, exist_ok=True)
file.save(str(cookie_file_path))
# 更新数据库中的账号信息(可选,比如更新更新时间)
# 这里可以根据需要添加额外的处理逻辑
return jsonify({
"code": 200,
"msg": "Cookie文件上传成功",
"data": None
}), 200
except Exception as e:
print(f"上传Cookie文件时出错: {str(e)}")
return jsonify({
"code": 500,
"msg": f"上传Cookie文件失败: {str(e)}",
"data": None
}), 500
# Cookie文件下载API
@app.route('/downloadCookie', methods=['GET'])
def download_cookie():
try:
file_path = request.args.get('filePath')
if not file_path:
return jsonify({
"code": 500,
"msg": "缺少文件路径参数",
"data": None
}), 400
# 验证文件路径的安全性,防止路径遍历攻击
cookie_file_path = Path(BASE_DIR / "cookiesFile" / file_path).resolve()
base_path = Path(BASE_DIR / "cookiesFile").resolve()
if not cookie_file_path.is_relative_to(base_path):
return jsonify({
"code": 500,
"msg": "非法文件路径",
"data": None
}), 400
if not cookie_file_path.exists():
return jsonify({
"code": 500,
"msg": "Cookie文件不存在",
"data": None
}), 404
# 返回文件
return send_from_directory(
directory=str(cookie_file_path.parent),
path=cookie_file_path.name,
as_attachment=True
)
except Exception as e:
print(f"下载Cookie文件时出错: {str(e)}")
return jsonify({
"code": 500,
"msg": f"下载Cookie文件失败: {str(e)}",
"data": None
}), 500
# ===== 新增:CLI 主线平台登录桥接 =====
def _qrcode_to_queue(status_queue):
def callback(payload):
try:
if isinstance(payload, dict):
data_url = payload.get("image_data_url") or ""
image_path = payload.get("image_path") or ""
status_text = payload.get("status_text") or ""
if status_text:
status_queue.put(status_text)
if data_url:
status_queue.put(data_url)
elif image_path and os.path.exists(image_path):
with open(image_path, "rb") as fobj:
b64 = base64.b64encode(fobj.read()).decode()
ext = Path(image_path).suffix.lstrip(".").lower() or "png"
status_queue.put("data:image/" + ext + ";base64," + b64)
else:
status_queue.put(str(payload))
elif isinstance(payload, str):
status_queue.put(payload)
else:
status_queue.put(str(payload))
except Exception:
pass
return callback
def _insert_account(platform_type, account_name, cookie_filename):
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
cursor = conn.cursor()
cursor.execute(
"INSERT INTO user_info (type, filePath, userName, status) VALUES (?, ?, ?, ?)",
(int(platform_type), cookie_filename, account_name, 1),
)
conn.commit()
def _bilibili_login(account_file, status_queue):
try:
binary = ensure_biliup_binary(force_check=False)
except Exception:
return False
workdir = account_file.parent
qrcode_path = workdir / "qrcode.png"
status_queue.put("正在启动 B 站扫码登录,请稍候...")
proc = subprocess.Popen(
[str(binary), "-u", str(account_file), "login"],
cwd=str(workdir),
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
)
last_mtime = 0
for _ in range(600):
if qrcode_path.exists():
try:
mtime = qrcode_path.stat().st_mtime
except Exception:
mtime = 0
if mtime != last_mtime:
last_mtime = mtime
try:
b64 = base64.b64encode(qrcode_path.read_bytes()).decode()
status_queue.put("data:image/png;base64," + b64)
except Exception:
pass
if proc.poll() is not None:
break
time.sleep(1)
try:
proc.wait(timeout=60)
except Exception:
proc.kill()
return proc.returncode == 0 and account_file.exists()
async def _tiktok_login_coro(account_file, status_queue):
"""TikTok 交互式登录:弹浏览器让用户手动登录,轮询 sessionid cookie。"""
from playwright.async_api import async_playwright
try:
async with async_playwright() as p:
browser = await p.chromium.launch(headless=False)
context = await browser.new_context()
page = await context.new_page()
await page.goto("https://www.tiktok.com/login?lang=en", wait_until="domcontentloaded")
for _ in range(600):
try:
cookies = await context.cookies()
except Exception:
cookies = []
if any(c.get("name") == "sessionid" and c.get("value") for c in cookies):
await asyncio.sleep(2)
await context.storage_state(path=str(account_file))
await browser.close()
return {"success": True}
await asyncio.sleep(1)
await browser.close()
return {"success": False}
except Exception:
return {"success": False}
def _login_new_platform(platform_type, account_name, status_queue):
uuid_v1 = str(uuid.uuid1())
account_file = Path(BASE_DIR / "cookiesFile" / (uuid_v1 + ".json"))
account_file.parent.mkdir(exist_ok=True)
cookie_filename = uuid_v1 + ".json"
if int(platform_type) == 5:
ok = _bilibili_login(account_file, status_queue)
if ok:
_insert_account(platform_type, account_name, cookie_filename)
status_queue.put("200")
else:
status_queue.put("500")
return
login_headless = (get_setting("login_mode", "qr") == "qr")
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
cb = _qrcode_to_queue(status_queue)
try:
status_queue.put("正在打开浏览器登录,请稍候...")
t = int(platform_type)
if t == 1:
result = loop.run_until_complete(xiaohongshu_setup(str(account_file), handle=True, return_detail=True, qrcode_callback=cb, headless=login_headless))
elif t == 2:
# 视频号(微信)qrconnect 登录组件在 headless 下 iframe JS 不执行、扫码回调失败,
# 必须强制有头模式弹出浏览器,保证微信登录组件正常工作(二维码仍会同步到前端)
status_queue.put("已打开浏览器,请在弹出的浏览器窗口中扫码登录视频号(微信登录不支持无头模式)")
result = loop.run_until_complete(tencent_setup(str(account_file), handle=True, return_detail=True, qrcode_callback=cb, headless=False))
elif t == 3:
result = loop.run_until_complete(douyin_setup(str(account_file), handle=True, return_detail=True, qrcode_callback=cb, headless=login_headless))
elif t == 4:
result = loop.run_until_complete(ks_setup(str(account_file), handle=True, return_detail=True, qrcode_callback=cb, headless=login_headless))
elif t == 6:
result = loop.run_until_complete(baijiahao_setup(str(account_file), handle=True, return_detail=True, qrcode_callback=cb, headless=login_headless))
elif t == 7:
result = loop.run_until_complete(alipay_setup(str(account_file), handle=True, return_detail=True, qrcode_callback=cb, headless=login_headless))
elif t == 8:
result = loop.run_until_complete(weibo_setup(str(account_file), handle=True, return_detail=True, qrcode_callback=cb, headless=login_headless))
elif t == 9:
result = loop.run_until_complete(hupu_setup(str(account_file), handle=True, return_detail=True, qrcode_callback=cb, headless=login_headless))
elif t == 10:
status_queue.put("已打开浏览器,请在弹出的 Chrome 窗口中登录 Google / YouTube 账号")
result = loop.run_until_complete(youtube_setup(str(account_file), handle=True, return_detail=True, headless=False))
elif t == 11:
status_queue.put("已打开浏览器,请在 Chrome 窗口中登录 TikTok(手机号/邮箱/Google/扫码等任意方式)")
result = loop.run_until_complete(_tiktok_login_coro(str(account_file), status_queue))
else:
result = {"success": False}
if isinstance(result, dict) and result.get("success"):
_insert_account(platform_type, account_name, cookie_filename)
status_queue.put("200")
else:
status_queue.put("500")
except Exception:
status_queue.put("500")
finally:
try:
loop.close()
except Exception:
pass
# 包装函数:在线程中运行异步函数
def run_async_function(type,id,status_queue):
match type:
case '1':
_login_new_platform('1', id, status_queue)
case '2':
_login_new_platform('2', id, status_queue)
case '3':
_login_new_platform('3', id, status_queue)
case '4':
_login_new_platform('4', id, status_queue)
case '5':
_login_new_platform('5', id, status_queue)
case '6':
_login_new_platform('6', id, status_queue)
case '7':
_login_new_platform('7', id, status_queue)
case '8':
_login_new_platform('8', id, status_queue)
case '9':
_login_new_platform('9', id, status_queue)
case '10':
_login_new_platform('10', id, status_queue)
case '11':
_login_new_platform('11', id, status_queue)
# SSE 流生成器函数
def sse_stream(status_queue):
while True:
if not status_queue.empty():
msg = status_queue.get()
yield f"data: {msg}\n\n"
else:
# 避免 CPU 占满
time.sleep(0.1)
@app.route('/publish/status', methods=['GET'])
def publish_status():
task_id = request.args.get('taskId')
with _publish_lock:
task = PUBLISH_TASKS.get(task_id)
if not task:
return jsonify({"code": 404, "msg": "任务不存在", "data": None}), 404
return jsonify({"code": 200, "msg": None, "data": task}), 200
@app.route('/publish/tasks', methods=['GET'])
def publish_tasks():
with _publish_lock:
tasks = list(PUBLISH_TASKS.values())
tasks = sorted(tasks, key=lambda t: t.get("created_at", ""), reverse=True)
return jsonify({"code": 200, "msg": None, "data": tasks}), 200
# 平台能力矩阵 —— 对齐 CLI 主线(docs/CLI.md / README 能力表)
@app.route('/platforms', methods=['GET'])
def platforms():
data = [
{"key": "douyin", "name": "抖音", "type": 3, "video": True, "note": True, "schedule": True, "cli": True, "web": True, "theme": "danger", "desc": "主线重构最完整,支持图文/定时"},
{"key": "kuaishou", "name": "快手", "type": 4, "video": True, "note": True, "schedule": True, "cli": True, "web": True, "theme": "warning", "desc": "浏览器自动化,支持图文/定时"},
{"key": "xiaohongshu", "name": "小红书", "type": 1, "video": True, "note": True, "schedule": True, "cli": True, "web": True, "theme": "danger", "desc": "浏览器自动化,支持图文/定时"},
{"key": "tencent", "name": "视频号", "type": 2, "video": True, "note": True, "schedule": True, "cli": True, "web": True, "theme": "success", "desc": "浏览器自动化,支持定时"},
{"key": "bilibili", "name": "B站", "type": 5, "video": True, "note": False, "schedule": True, "cli": True, "web": True, "theme": "primary", "desc": "封装 biliup,支持定时"},
{"key": "baijiahao", "name": "百家号", "type": 6, "video": True, "note": False, "schedule": False, "cli": True, "web": True, "theme": "primary", "desc": "浏览器自动化"},
{"key": "alipay", "name": "支付宝生活号", "type": 7, "video": True, "note": False, "schedule": False, "cli": True, "web": True, "theme": "primary", "desc": "浏览器自动化,支持生活号视频"},
{"key": "weibo", "name": "微博", "type": 8, "video": True, "note": False, "schedule": False, "cli": True, "web": True, "theme": "primary", "desc": "浏览器自动化,标题最多 30 字"},
{"key": "hupu", "name": "虎扑", "type": 9, "video": True, "note": False, "schedule": False, "cli": True, "web": True, "theme": "primary", "desc": "浏览器自动化,标题 4–40 字"},
{"key": "youtube", "name": "YouTube", "type": 10, "video": True, "note": False, "schedule": False, "cli": True, "web": True, "theme": "primary", "desc": "Studio 自动化,支持播放列表/可见性"},
{"key": "tiktok", "name": "TikTok", "type": 11, "video": True, "note": False, "schedule": True, "cli": False, "web": True, "theme": "default", "desc": "交互式登录 + Chrome 版发布,已接入 Web 后端"},
]
return jsonify({"code": 200, "msg": None, "data": data})
# ===================== Open API(对外开放) =====================
def _ensure_api_keys_table():
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
conn.execute("CREATE TABLE IF NOT EXISTS api_keys (id INTEGER PRIMARY KEY AUTOINCREMENT, key TEXT NOT NULL UNIQUE, name TEXT NOT NULL, enabled INTEGER DEFAULT 1, created_at TEXT)")
conn.commit()
def _ensure_default_api_key():
_ensure_api_keys_table()
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
cur = conn.cursor()
cur.execute("SELECT COUNT(*) FROM api_keys")
if cur.fetchone()[0] == 0:
default_key = "sk-sau-" + uuid.uuid4().hex[:24]
cur.execute("INSERT INTO api_keys (key, name, enabled, created_at) VALUES (?, ?, 1, ?)", (default_key, "默认密钥", datetime.now().strftime("%Y-%m-%d %H:%M:%S")))
conn.commit()
print("已创建默认 API Key:", default_key)
def _get_valid_api_keys():
_ensure_default_api_key()
_ensure_api_keys_table()
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
cur = conn.cursor()
cur.execute("SELECT key FROM api_keys WHERE enabled = 1")
return {row[0] for row in cur.fetchall()}
def require_api_key(f):
@wraps(f)
def wrapper(*args, **kwargs):
auth = request.headers.get("X-API-Key") or ""
if not auth:
bearer = request.headers.get("Authorization", "") or ""
if bearer.startswith("Bearer "):
auth = bearer[7:].strip()
if not auth or auth not in _get_valid_api_keys():
return jsonify({"code": 401, "msg": "无效或缺失的 API Key", "data": None}), 401
return f(*args, **kwargs)
return wrapper
@app.route('/open/keys', methods=['GET'])
def list_api_keys():
_ensure_default_api_key()
_ensure_api_keys_table()
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
conn.row_factory = sqlite3.Row
cur = conn.cursor()
cur.execute("SELECT id, key, name, enabled, created_at FROM api_keys ORDER BY id DESC")
rows = [dict(r) for r in cur.fetchall()]
return jsonify({"code": 200, "msg": None, "data": rows}), 200
@app.route('/open/keys', methods=['POST'])
def create_api_key():
data = request.get_json() or {}
name = data.get('name', '新密钥')
key = "sk-sau-" + uuid.uuid4().hex[:24]
_ensure_api_keys_table()
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
conn.execute("INSERT INTO api_keys (key, name, enabled, created_at) VALUES (?, ?, 1, ?)", (key, name, datetime.now().strftime("%Y-%m-%d %H:%M:%S")))
conn.commit()
return jsonify({"code": 200, "msg": "密钥已创建", "data": {"key": key}}), 200
@app.route('/open/keys/<int:key_id>', methods=['DELETE'])
def delete_api_key(key_id):
_ensure_api_keys_table()
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
conn.execute("DELETE FROM api_keys WHERE id = ?", (key_id,))
conn.commit()
return jsonify({"code": 200, "msg": "密钥已删除", "data": None}), 200
@app.route('/open/v1/platforms', methods=['GET'])
@require_api_key
def open_platforms():
data = [
{"key": "douyin", "name": "抖音", "type": 3, "video": True, "note": True, "schedule": True},
{"key": "kuaishou", "name": "快手", "type": 4, "video": True, "note": True, "schedule": True},
{"key": "xiaohongshu", "name": "小红书", "type": 1, "video": True, "note": True, "schedule": True},
{"key": "tencent", "name": "视频号", "type": 2, "video": True, "note": False, "schedule": True},
{"key": "bilibili", "name": "B站", "type": 5, "video": True, "note": False, "schedule": True},
{"key": "baijiahao", "name": "百家号", "type": 6, "video": True, "note": False, "schedule": False},
{"key": "alipay", "name": "支付宝生活号", "type": 7, "video": True, "note": False, "schedule": False},
{"key": "weibo", "name": "微博", "type": 8, "video": True, "note": False, "schedule": False},
{"key": "hupu", "name": "虎扑", "type": 9, "video": True, "note": False, "schedule": False},
{"key": "youtube", "name": "YouTube", "type": 10, "video": True, "note": False, "schedule": False},
{"key": "tiktok", "name": "TikTok", "type": 11, "video": True, "note": False, "schedule": True},
]
return jsonify({"code": 200, "msg": None, "data": data}), 200
@app.route('/open/v1/accounts', methods=['GET'])
@require_api_key
def open_accounts():
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
conn.row_factory = sqlite3.Row
cur = conn.cursor()
cur.execute("SELECT id, type, filePath, userName, status FROM user_info ORDER BY id DESC")
rows = [dict(r) for r in cur.fetchall()]
return jsonify({"code": 200, "msg": None, "data": rows}), 200
@app.route('/open/v1/materials', methods=['GET'])
@require_api_key
def open_materials():
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
conn.row_factory = sqlite3.Row
cur = conn.cursor()
cur.execute("SELECT * FROM file_records ORDER BY id DESC")
rows = [dict(r) for r in cur.fetchall()]
return jsonify({"code": 200, "msg": None, "data": rows}), 200
@app.route('/open/v1/materials', methods=['POST'])
@require_api_key
def open_upload_material():
if 'file' not in request.files:
return jsonify({"code": 400, "msg": "缺少文件字段 file", "data": None}), 400
file = request.files['file']
if file.filename == '':
return jsonify({"code": 400, "msg": "文件名为空", "data": None}), 400
custom_filename = request.form.get('filename', None)
if custom_filename:
filename = secure_filename(custom_filename + "." + file.filename.split('.')[-1])
else:
filename = secure_filename(file.filename)
if not filename:
return jsonify({"code": 400, "msg": "非法文件名", "data": None}), 400
try:
uuid_v1 = uuid.uuid1()
final_filename = str(uuid_v1) + "_" + filename
filepath = Path(BASE_DIR / "videoFile" / final_filename)
file.save(filepath)
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
conn.execute("INSERT INTO file_records (filename, filesize, file_path) VALUES (?, ?, ?)", (filename, round(float(os.path.getsize(filepath)) / (1024 * 1024), 2), final_filename))
conn.commit()
return jsonify({"code": 200, "msg": "上传成功", "data": {"filename": filename, "filePath": final_filename}}), 200
except Exception as e:
return jsonify({"code": 500, "msg": str(e), "data": None}), 500
@app.route('/open/v1/materials/<int:file_id>', methods=['DELETE'])
@require_api_key
def open_delete_material(file_id):
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
conn.row_factory = sqlite3.Row
cur = conn.cursor()
cur.execute("SELECT * FROM file_records WHERE id = ?", (file_id,))
record = cur.fetchone()
if not record:
return jsonify({"code": 404, "msg": "素材不存在", "data": None}), 404
record = dict(record)
file_path = Path(BASE_DIR / "videoFile" / record['file_path'])
if file_path.exists():
try:
file_path.unlink()
except Exception:
pass
cur.execute("DELETE FROM file_records WHERE id = ?", (file_id,))
conn.commit()
return jsonify({"code": 200, "msg": "删除成功", "data": None}), 200
@app.route('/open/v1/publish', methods=['POST'])
@require_api_key
def open_publish():
data = request.get_json()
if not data:
return jsonify({"code": 400, "msg": "请求数据不能为空", "data": None}), 400
task_id, err = _create_publish_task(data)
if err:
return jsonify({"code": 400, "msg": err, "data": None}), 400
return jsonify({"code": 200, "msg": "发布任务已提交", "data": {"taskId": task_id}}), 200
@app.route('/open/v1/publish/<task_id>', methods=['GET'])
@require_api_key
def open_publish_status(task_id):
with _publish_lock:
task = PUBLISH_TASKS.get(task_id)
if not task:
return jsonify({"code": 404, "msg": "任务不存在", "data": None}), 404
return jsonify({"code": 200, "msg": None, "data": task}), 200
# ===================== 系统全局配置 =====================
@app.route('/settings', methods=['GET'])
def get_settings():
data = {
"login_mode": get_setting("login_mode", "qr"),
"publish_headless": get_bool("publish_headless", True),
}
return jsonify({"code": 200, "msg": None, "data": data}), 200
@app.route('/settings', methods=['POST'])
def update_settings():
data = request.get_json() or {}
if "login_mode" in data:
set_setting("login_mode", data["login_mode"])
if "publish_headless" in data:
set_setting("publish_headless", "true" if data["publish_headless"] else "false")
return jsonify({"code": 200, "msg": "配置已保存", "data": None}), 200
# ===================== 账号校验(单账号/批量,异步) =====================
def _run_check_task(task):
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
try:
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
for aid in task["ids"]:
cur = conn.cursor()
cur.execute("SELECT type, filePath, userName FROM user_info WHERE id = ?", (aid,))
row = cur.fetchone()
if not row:
task["results"][str(aid)] = {"valid": None, "name": str(aid)}
task["done"] += 1
continue
ptype, fp, name = row
try:
valid = loop.run_until_complete(check_cookie(ptype, fp))
except Exception:
valid = False
task["results"][str(aid)] = {"valid": valid, "name": name}
cur.execute("UPDATE user_info SET status = ? WHERE id = ?", (1 if valid else 0, aid))
conn.commit()
task["done"] += 1
task["logs"].append("账号「" + name + "」校验" + ("有效" if valid else "失效"))
task["status"] = "done"
except Exception as e:
task["status"] = "error"
task["error"] = str(e)
finally:
try:
loop.close()
except Exception:
pass
@app.route('/checkAccounts', methods=['POST'])
def check_accounts():
data = request.get_json() or {}
ids = data.get('ids', [])
if not ids:
return jsonify({"code": 400, "msg": "请选择要校验的账号", "data": None}), 400
ids = [int(i) for i in ids]
task_id = str(uuid.uuid4())
task = {"id": task_id, "ids": ids, "total": len(ids), "done": 0, "status": "running", "results": {}, "logs": [], "error": None}
with _publish_lock:
CHECK_TASKS[task_id] = task
threading.Thread(target=_run_check_task, args=(task,), daemon=True).start()
return jsonify({"code": 200, "msg": "校验任务已提交", "data": {"taskId": task_id}}), 200
@app.route('/checkAccounts/status', methods=['GET'])
def check_account_status():
task_id = request.args.get('taskId')
task = CHECK_TASKS.get(task_id)
if not task:
return jsonify({"code": 404, "msg": "任务不存在", "data": None}), 404
return jsonify({"code": 200, "msg": None, "data": task}), 200
@app.route('/deleteAccounts', methods=['POST'])
def delete_accounts():
data = request.get_json() or {}
ids = data.get('ids', [])
if not ids:
return jsonify({"code": 400, "msg": "请选择要删除的账号", "data": None}), 400
ids = [int(i) for i in ids]
deleted = 0
with sqlite3.connect(Path(BASE_DIR / "db" / "database.db")) as conn:
cur = conn.cursor()
for aid in ids:
cur.execute("SELECT filePath FROM user_info WHERE id = ?", (aid,))
row = cur.fetchone()
if not row:
continue
cookie_file = Path(BASE_DIR / "cookiesFile" / row[0])
if cookie_file.exists():
try:
cookie_file.unlink()
except Exception:
pass
cur.execute("DELETE FROM user_info WHERE id = ?", (aid,))
deleted += 1
conn.commit()
return jsonify({"code": 200, "msg": "已删除 " + str(deleted) + " 个账号", "data": {"deleted": deleted}}), 200
if __name__ == '__main__':
app.run(host='0.0.0.0' ,port=5409)