From c1bd53d832028b740b6facb813f2d0352bd7b77d Mon Sep 17 00:00:00 2001 From: Misaka Date: Sun, 2 Aug 2026 14:14:30 +0800 Subject: [PATCH] feat(tasks): record trigger mode, target date and force in task history --- inbound_verify/cli/server.py | 13 ++++++- inbound_verify/state_store.py | 73 ++++++++++++++++++++++++++++------- 2 files changed, 70 insertions(+), 16 deletions(-) diff --git a/inbound_verify/cli/server.py b/inbound_verify/cli/server.py index 821de4e..182f1e5 100644 --- a/inbound_verify/cli/server.py +++ b/inbound_verify/cli/server.py @@ -128,7 +128,10 @@ def _enqueue_fetch(site, kind): if not _in_active_window(cfg["active_start"], cfg["active_end"]): return # 不在激活时段,跳过本次 fire try: - tid = state_store.create_task_if_idle(site, kind) + target_date = state_store.resolve_target_date(site, kind) + tid = state_store.create_task_if_idle( + site, kind, trigger="auto", target_date=target_date + ) if tid is None: return # 上一次同类任务还没跑完,跳过避免堆积 task_queue.put((tid, {"site": site, "kind": kind})) @@ -225,7 +228,13 @@ def create_task(req: TaskRequest): raise HTTPException( status_code=400, detail="百世固定下载当天,不支持指定日期" ) - task_id = state_store.create_task(req.site, req.kind) + task_id = state_store.create_task( + req.site, + req.kind, + trigger="manual", + target_date=state_store.resolve_target_date(req.site, req.kind, req.date), + force=req.force, + ) spec = {"site": req.site, "kind": req.kind, "force": req.force} if req.date: spec["date"] = req.date diff --git a/inbound_verify/state_store.py b/inbound_verify/state_store.py index 2cfc669..0327c96 100644 --- a/inbound_verify/state_store.py +++ b/inbound_verify/state_store.py @@ -7,7 +7,7 @@ import os import sqlite3 -from datetime import datetime +from datetime import datetime, timedelta from inbound_verify.paths import STATE_DB_PATH @@ -70,9 +70,22 @@ def init_db(): status TEXT, started_at TEXT, finished_at TEXT, - error TEXT + error TEXT, + trigger TEXT NOT NULL DEFAULT '', + target_date TEXT NOT NULL DEFAULT '', + force INTEGER NOT NULL DEFAULT 0 ) """) + # 旧库迁移:补触发方式/目标日期/强制重下三列(新库已含;重复添加抛 OperationalError,忽略) + for _col, _typedef in [ + ("trigger", "TEXT NOT NULL DEFAULT ''"), + ("target_date", "TEXT NOT NULL DEFAULT ''"), + ("force", "INTEGER NOT NULL DEFAULT 0"), + ]: + try: + conn.execute(f"ALTER TABLE task_history ADD COLUMN {_col} {_typedef}") + except sqlite3.OperationalError: + pass conn.execute(""" CREATE TABLE IF NOT EXISTS site_config ( site TEXT PRIMARY KEY, @@ -357,6 +370,27 @@ def set_offset(site, kind, offset): return offset +def resolve_target_date(site, kind, date=None): + """计算一条任务的目标下载日期(YYYY-MM-DD,供任务日志展示 / 重试回放)。 + 有 date 用 date;否则按站点偏移推算(与 runtime._record_business_date 同源): + expected → 应到偏移;actual → 实到偏移;百世 undelivered → 当天; + 4 站 undelivered → 跟随应到偏移。__compare__ 无数据概念,返回 ''。""" + if site == "__compare__": + return "" + if date: + return date + today = datetime.now().date() + if kind == "expected": + return (today - timedelta(days=get_offset(site, "expected"))).strftime( + "%Y-%m-%d" + ) + if kind == "actual": + return (today - timedelta(days=get_offset(site, "actual"))).strftime("%Y-%m-%d") + if site == "百世": + return today.strftime("%Y-%m-%d") + return (today - timedelta(days=get_offset(site, "expected"))).strftime("%Y-%m-%d") + + def set_schedule(site, enabled, time_str): """【DEPRECATED】旧"每日单时点定时"——已被 fetch_schedule 的周期+激活时段模式取代。 保留死代码以防外部残留调用;新代码请用 set_fetch_schedule。""" @@ -539,26 +573,28 @@ def get_site_settings(site): # ============================ 任务历史 ============================ -def create_task(site, kind): - """新建一条 pending 任务,返回其 id。""" +def create_task(site, kind, trigger="manual", target_date="", force=False): + """新建一条 pending 任务(手动触发),返回其 id。trigger='manual'/'auto'; + target_date 为该任务的目标下载日期(YYYY-MM-DD,可为 '');force 是否强制重下。""" now = _now() with sqlite3.connect(STATE_DB_PATH, timeout=3.0) as conn: cur = conn.execute( - "INSERT INTO task_history (site, kind, status, started_at, finished_at, error) " - "VALUES (?, ?, ?, ?, '', '')", - (site, kind, TASK_PENDING, now), + "INSERT INTO task_history " + "(site, kind, status, started_at, finished_at, error, trigger, target_date, force) " + "VALUES (?, ?, ?, ?, '', '', ?, ?, ?)", + (site, kind, TASK_PENDING, now, trigger, target_date, 1 if force else 0), ) conn.commit() return cur.lastrowid -def create_task_if_idle(site, kind): +def create_task_if_idle(site, kind, trigger="auto", target_date=""): """周期调度专用:若该 (site,kind) 已有 pending/running 任务则返回 None(跳过本次周期), 否则建一条 pending 任务返回其 id。单连接内 check-then-insert,靠 SQLite 写锁把竞态压到忽略不计。 与 create_task 的区别:手动触发(POST /tasks)用 create_task(用户点的必建);周期 job 用本函数 ——上一次还没跑完时跳过,避免同 (site,kind) 任务堆积。手动建的任务会让紧随其后的周期 fire - 判到 inflight 而跳过,天然互斥。""" + 判到 inflight 而跳过,天然互斥。trigger='auto';target_date 为目标下载日期(YYYY-MM-DD)。""" now = _now() with sqlite3.connect(STATE_DB_PATH, timeout=3.0) as conn: row = conn.execute( @@ -569,9 +605,10 @@ def create_task_if_idle(site, kind): if row: return None cur = conn.execute( - "INSERT INTO task_history (site, kind, status, started_at, finished_at, error) " - "VALUES (?, ?, ?, ?, '', '')", - (site, kind, TASK_PENDING, now), + "INSERT INTO task_history " + "(site, kind, status, started_at, finished_at, error, trigger, target_date, force) " + "VALUES (?, ?, ?, ?, '', '', ?, ?, 0)", + (site, kind, TASK_PENDING, now, trigger, target_date), ) conn.commit() return cur.lastrowid @@ -598,7 +635,8 @@ def get_task(task_id): """返回单条任务 dict,不存在返回 None。""" with sqlite3.connect(STATE_DB_PATH, timeout=3.0) as conn: row = conn.execute( - "SELECT id, site, kind, status, started_at, finished_at, error " + "SELECT id, site, kind, status, started_at, finished_at, error, " + "trigger, target_date, force " "FROM task_history WHERE id=?", (task_id,), ).fetchone() @@ -612,6 +650,9 @@ def get_task(task_id): "started_at": row[4], "finished_at": row[5], "error": row[6], + "trigger": row[7], + "target_date": row[8], + "force": bool(row[9]), } @@ -619,7 +660,8 @@ def list_tasks(limit=20): """返回最近 limit 条任务(按 id 倒序)。""" with sqlite3.connect(STATE_DB_PATH, timeout=3.0) as conn: rows = conn.execute( - "SELECT id, site, kind, status, started_at, finished_at, error " + "SELECT id, site, kind, status, started_at, finished_at, error, " + "trigger, target_date, force " "FROM task_history ORDER BY id DESC LIMIT ?", (limit,), ).fetchall() @@ -632,6 +674,9 @@ def list_tasks(limit=20): "started_at": r[4], "finished_at": r[5], "error": r[6], + "trigger": r[7], + "target_date": r[8], + "force": bool(r[9]), } for r in rows ]