|
|
@@ -46,6 +46,28 @@ CREATE TABLE IF NOT EXISTS account_events (
|
|
|
|
|
|
CREATE INDEX IF NOT EXISTS idx_events_email ON account_events(email);
|
|
|
CREATE INDEX IF NOT EXISTS idx_events_ts ON account_events(ts);
|
|
|
+
|
|
|
+CREATE TABLE IF NOT EXISTS tasks (
|
|
|
+ task_id TEXT PRIMARY KEY,
|
|
|
+ mode TEXT NOT NULL, -- full | pay_only
|
|
|
+ status TEXT NOT NULL, -- queued | running | success | failed | cancelled
|
|
|
+ stage TEXT,
|
|
|
+ attempts INTEGER DEFAULT 0,
|
|
|
+ max_attempts INTEGER DEFAULT 3,
|
|
|
+ params_json TEXT,
|
|
|
+ result_json TEXT,
|
|
|
+ last_error TEXT,
|
|
|
+ email TEXT,
|
|
|
+ plan_type TEXT,
|
|
|
+ cpa_file_name TEXT,
|
|
|
+ created_at INTEGER NOT NULL,
|
|
|
+ updated_at INTEGER NOT NULL,
|
|
|
+ started_at INTEGER,
|
|
|
+ finished_at INTEGER
|
|
|
+);
|
|
|
+
|
|
|
+CREATE INDEX IF NOT EXISTS idx_tasks_status ON tasks(status);
|
|
|
+CREATE INDEX IF NOT EXISTS idx_tasks_created ON tasks(created_at);
|
|
|
"""
|
|
|
|
|
|
|
|
|
@@ -211,3 +233,97 @@ def _event_row(row: sqlite3.Row) -> dict:
|
|
|
except Exception:
|
|
|
d["payload"] = None
|
|
|
return d
|
|
|
+
|
|
|
+
|
|
|
+# ------------------------- tasks -------------------------
|
|
|
+
|
|
|
+def create_task(task_id: str, mode: str, params: dict, max_attempts: int = 3) -> dict:
|
|
|
+ init_db()
|
|
|
+ now = _now_ms()
|
|
|
+ with _conn() as c:
|
|
|
+ c.execute(
|
|
|
+ """
|
|
|
+ INSERT INTO tasks (task_id, mode, status, stage, attempts, max_attempts,
|
|
|
+ params_json, result_json, last_error, email, plan_type, cpa_file_name,
|
|
|
+ created_at, updated_at, started_at, finished_at)
|
|
|
+ VALUES (?, ?, 'queued', '', 0, ?, ?, NULL, NULL, NULL, NULL, NULL, ?, ?, NULL, NULL)
|
|
|
+ """,
|
|
|
+ (task_id, mode, max_attempts, _dump(params or {}), now, now),
|
|
|
+ )
|
|
|
+ return _task_row(c.execute("SELECT * FROM tasks WHERE task_id = ?", (task_id,)).fetchone())
|
|
|
+
|
|
|
+
|
|
|
+def update_task(task_id: str, fields: dict) -> dict | None:
|
|
|
+ init_db()
|
|
|
+ if not fields:
|
|
|
+ return get_task(task_id)
|
|
|
+ sets = ["updated_at = ?"]
|
|
|
+ args: list[Any] = [_now_ms()]
|
|
|
+ for col in (
|
|
|
+ "status", "stage", "attempts", "max_attempts", "last_error",
|
|
|
+ "email", "plan_type", "cpa_file_name", "started_at", "finished_at",
|
|
|
+ ):
|
|
|
+ if col in fields and fields[col] is not None:
|
|
|
+ sets.append(f"{col} = ?")
|
|
|
+ args.append(fields[col])
|
|
|
+ if "result" in fields:
|
|
|
+ sets.append("result_json = ?")
|
|
|
+ args.append(_dump(fields["result"]))
|
|
|
+ elif "result_json" in fields and fields["result_json"] is not None:
|
|
|
+ sets.append("result_json = ?")
|
|
|
+ args.append(fields["result_json"])
|
|
|
+ if "params" in fields:
|
|
|
+ sets.append("params_json = ?")
|
|
|
+ args.append(_dump(fields["params"]))
|
|
|
+ args.append(task_id)
|
|
|
+ with _conn() as c:
|
|
|
+ c.execute(f"UPDATE tasks SET {', '.join(sets)} WHERE task_id = ?", args)
|
|
|
+ row = c.execute("SELECT * FROM tasks WHERE task_id = ?", (task_id,)).fetchone()
|
|
|
+ return _task_row(row) if row else None
|
|
|
+
|
|
|
+
|
|
|
+def get_task(task_id: str) -> dict | None:
|
|
|
+ init_db()
|
|
|
+ with _conn() as c:
|
|
|
+ r = c.execute("SELECT * FROM tasks WHERE task_id = ?", (task_id,)).fetchone()
|
|
|
+ return _task_row(r) if r else None
|
|
|
+
|
|
|
+
|
|
|
+def list_tasks(limit: int = 100, status: str | None = None) -> list[dict]:
|
|
|
+ init_db()
|
|
|
+ with _conn() as c:
|
|
|
+ if status:
|
|
|
+ rows = c.execute(
|
|
|
+ "SELECT * FROM tasks WHERE status = ? ORDER BY created_at DESC LIMIT ?",
|
|
|
+ (status, limit),
|
|
|
+ ).fetchall()
|
|
|
+ else:
|
|
|
+ rows = c.execute(
|
|
|
+ "SELECT * FROM tasks ORDER BY created_at DESC LIMIT ?",
|
|
|
+ (limit,),
|
|
|
+ ).fetchall()
|
|
|
+ return [_task_row(r) for r in rows]
|
|
|
+
|
|
|
+
|
|
|
+def get_next_queued_task() -> dict | None:
|
|
|
+ """取一个最早的 queued 任务(FIFO)。"""
|
|
|
+ init_db()
|
|
|
+ with _conn() as c:
|
|
|
+ r = c.execute(
|
|
|
+ "SELECT * FROM tasks WHERE status = 'queued' ORDER BY created_at ASC LIMIT 1"
|
|
|
+ ).fetchone()
|
|
|
+ return _task_row(r) if r else None
|
|
|
+
|
|
|
+
|
|
|
+def _task_row(row: sqlite3.Row | None) -> dict | None:
|
|
|
+ if row is None:
|
|
|
+ return None
|
|
|
+ d = {k: row[k] for k in row.keys()}
|
|
|
+ for k in ("params_json", "result_json"):
|
|
|
+ raw = d.get(k)
|
|
|
+ if raw:
|
|
|
+ try:
|
|
|
+ d[k.replace("_json", "")] = json.loads(raw)
|
|
|
+ except Exception:
|
|
|
+ d[k.replace("_json", "")] = None
|
|
|
+ return d
|