文件
wangchuanli 3751dffef9 feat: 新增备份恢复与公网加固
- 新增备份管理页与 API:在线快照、自动周期备份、按份数清理、下载、一键恢复(恢复前自动兜底)
- 新增 /profile/export,普通用户可导出本人全部数据(不含 Cookie 明文)
- 修复 X-Forwarded-For 可伪造导致三道 IP 防线失效,统一走 client_ip() 取客户端地址
- 取消 admin123 硬编码默认口令,留空则生成随机初始口令并仅打印一次
- .dockerignore 排除 backups/ 并加构建期断言,防止密钥随镜像分发
- 新增会话版本号,改密/停用/删除及恢复备份后其他会话立即失效
- 新增容器资源上限、采集跨度硬顶 31 天、重操作最小间隔与并发 409
- 新增访问日志、HSTS 条件下发、口令黑名单、验证码抗模板匹配、instance.json 0600
- 版本号升至 1.4.0,同步更新 README、SECURITY、.env.example 与 compose 配置
2026-09-18 08:46:34 +08:00

212 行
8.0 KiB
Python

此文件含有模棱两可的 Unicode 字符
此文件含有可能会与其他字符混淆的 Unicode 字符。 如果您是想特意这样的,可以安全地忽略该警告。 使用 Escape 按钮显示他们。
# -*- coding: utf-8 -*-
# SPDX-License-Identifier: MIT
# Copyright (c) 2026 Wang Chuanli
"""进程内调度器(手写,不依赖 APScheduler)。
为什么不引 APScheduler:
* 需求只是「每天固定几个时刻跑一次」,一个线程 + 20 秒轮询足够
* 需要「启动补跑」(程序没开的时候错过了时刻,开机后要补上)
* 需要和 CLI 共享同一把文件锁,避免两处同时采集
多用户
------
调度配置(开关 / 时刻 / 补跑 / slot:* 簿记)都是**个人级**设置,
所以 tick 会遍历所有启用状态的账号,各自判断有没有到期槽位。
好处是「A 想 9 点采、B 想 21 点采」互不影响;代价是串行执行 ——
这是刻意的,SQLite 单写者不允许并发采集。
单实例保证:
* Flask 的 reloader 会 fork 两个进程 → 只在 WERKZEUG_RUN_MAIN 里启动
* 多进程部署时用环境变量 WB_DISABLE_SCHEDULER=1 关掉除一个之外的所有实例
* 真正防重复采集靠 collect._Lock(文件锁),即使两个调度线程同时触发也只会跑一个
"""
import logging
import os
import threading
from datetime import datetime, timedelta
from . import backup, collect, db
log = logging.getLogger("wb.scheduler")
SLOT_PREFIX = "slot:" # settings 键:slot:09:00 -> 最近执行的日期(按 user_id 存)
def parse_times(raw):
"""把 '09:00,17:00' 解析成 ['09:00','17:00'];非法片段直接丢弃。
对外公开(config.normalize_setting 用它做写入校验),所以不要改名。
"""
out = []
for part in str(raw or "").replace(";", ",").replace(",", ",").split(","):
p = part.strip()
if not p:
continue
bits = p.split(":")
try:
hh = int(bits[0])
mm = int(bits[1]) if len(bits) > 1 else 0
except ValueError:
continue
if 0 <= hh <= 23 and 0 <= mm <= 59:
out.append("%02d:%02d" % (hh, mm))
return sorted(set(out))
# 兼容旧名(早期版本用 _parse_times)
_parse_times = parse_times
def slots(conn, uid=0):
return parse_times(db.get_setting(conn, "schedule_times", uid=uid))
def last_run_of_slot(conn, uid, slot):
return db.get_setting(conn, SLOT_PREFIX + slot, "", uid)
def mark_slot(conn, uid, slot, day):
db.set_setting(conn, SLOT_PREFIX + slot, day, uid)
def next_run_at(conn, uid=0, now=None):
"""下一次计划执行时间(仅按配置推算,不含补跑)。"""
if not db.get_bool(conn, "schedule_enabled", True, uid):
return None
now = now or datetime.now()
best = None
for s in slots(conn, uid):
hh, mm = map(int, s.split(":"))
cand = now.replace(hour=hh, minute=mm, second=0, microsecond=0)
if cand <= now:
cand += timedelta(days=1)
if best is None or cand < best:
best = cand
return best
def due_slots(conn, uid=0, now=None):
"""返回此刻应当执行的槽位列表(含启动补跑)。"""
if not db.get_bool(conn, "schedule_enabled", True, uid):
return []
now = now or datetime.now()
today = now.strftime("%Y-%m-%d")
# 用 get_int 兜底:catch_up_grace_hours 在后台是自由文本框,
# 历史上填成 "12h" 会让这里 int() 抛 ValueError,把 /tasks 打成 500。
grace_hours = db.get_int(conn, "catch_up_grace_hours", 12, uid)
grace = timedelta(hours=max(1, grace_hours))
catch_up = db.get_bool(conn, "catch_up", True, uid)
out = []
for s in slots(conn, uid):
hh, mm = map(int, s.split(":"))
when = now.replace(hour=hh, minute=mm, second=0, microsecond=0)
if when > now:
continue # 还没到点
if last_run_of_slot(conn, uid, s) == today:
continue # 今天这个槽位已跑过
if when < now - grace and catch_up:
continue # 错过太久,不补(避免开机狂刷)
if when < now - timedelta(seconds=90) and not catch_up:
continue # 未开启补跑,只认刚到的点
out.append(s)
return out
class Scheduler:
def __init__(self, interval=20):
self.interval = interval
self._stop = threading.Event()
self._thread = None
# ---- 生命周期 ----
def start(self):
if self._thread and self._thread.is_alive():
return False
self._stop.clear()
self._thread = threading.Thread(target=self._loop, name="wb-scheduler", daemon=True)
self._thread.start()
log.info("调度器已启动,轮询间隔 %ss", self.interval)
return True
def stop(self):
self._stop.set()
if self._thread:
self._thread.join(timeout=5)
@property
def running(self):
return bool(self._thread and self._thread.is_alive())
# ---- 主循环 ----
def _loop(self):
while not self._stop.is_set():
try:
self.tick()
except Exception as e: # 任何异常都不能让线程死掉
log.exception("调度 tick 出错:%s", e)
self._stop.wait(self.interval)
def tick(self, now=None):
conn = db.thread_conn()
now = now or datetime.now()
today = now.strftime("%Y-%m-%d")
for u in db.active_users(conn):
uid = u["id"]
try:
pending = due_slots(conn, uid, now)
except Exception as e: # 单个账号配置坏了不能拖垮其他人
log.error("账号 #%s(%s) 读取调度配置失败:%s", uid, u["username"], e)
continue
for slot in pending:
scheduled = now.replace(hour=int(slot[:2]), minute=int(slot[3:]),
second=0, microsecond=0)
trigger = "startup" if now - scheduled > timedelta(minutes=5) else "schedule"
log.info("触发采集:账号 %s 槽位 %s(%s)", u["username"], slot, trigger)
# 先占位,避免采集失败被无限重试打爆云端
mark_slot(conn, uid, slot, today)
if not db.secret_state(conn, "cookie", uid)["set"]:
log.info("跳过:账号 %s 还没配置自己的 Cookie", u["username"])
continue
try:
r = collect.run_sync(trigger=trigger, uid=uid)
log.info("采集完成:%s → %s", u["username"], r["message"])
except collect.Busy as e:
log.warning("跳过(%s)", e)
except db.SecretUnreadable as e:
log.error("账号 %s 的 Cookie 解不开:%s", u["username"], e)
except Exception as e:
log.error("账号 %s 采集失败:%s", u["username"], e)
# ---- 自动备份(实例级,与具体账号无关,所以放在账号循环之外)----
# 有采集在跑就跳过,等下一轮:备份会整库读一遍,没必要和采集抢磁盘。
try:
if os.path.exists(collect.LOCK_PATH):
log.debug("有采集在跑,本次跳过自动备份")
else:
backup.maybe_auto(conn, now)
except Exception as e: # 备份失败不能拖累调度本身
log.exception("自动备份出错:%s", e)
return True
_scheduler = None
def get_scheduler(interval=20):
global _scheduler
if _scheduler is None:
_scheduler = Scheduler(interval=interval)
return _scheduler
def start_from_app(app):
"""由 create_app 调用。遵守 reloader 与显式禁用开关。"""
if os.environ.get("WB_DISABLE_SCHEDULER") == "1":
app.logger.info("WB_DISABLE_SCHEDULER=1,调度器未启动")
return None
if app.debug and os.environ.get("WERKZEUG_RUN_MAIN") != "true":
return None # reloader 的父进程不启动,避免跑两份
sch = get_scheduler()
sch.start()
return sch