From 53c71aeeac1e34b77b162ff55c843eda0137a132 Mon Sep 17 00:00:00 2001 From: Misaka Date: Sun, 2 Aug 2026 09:22:54 +0800 Subject: [PATCH] refactor(status): derive site ready/business_date from PG instead of ingest_state MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit _ready_flags 改从 PostgreSQL 直接查询(expected_record / actual_record / baishi_daily_stats),target_date = today − offset。 消除因 ingested_at 日期比对导致的每日零点全站 ready 集体重置。 store.py: 新增 has_data(site, kind, target_date) 查 PG 数据存在性 runtime.py: _ready_flags 返回 (flags, dates) 同源元组,_apply_ready 同步写入 business_date,修正 ready 与 business_date 不同源导致的 前端日期标签漂移(如实到就绪却显示'前天') Co-Authored-By: Claude --- inbound_verify/runtime.py | 77 +++++++++++++++++++++++---------------- inbound_verify/store.py | 46 +++++++++++++++++++++++ 2 files changed, 92 insertions(+), 31 deletions(-) diff --git a/inbound_verify/runtime.py b/inbound_verify/runtime.py index ace1f5d..191f32b 100644 --- a/inbound_verify/runtime.py +++ b/inbound_verify/runtime.py @@ -16,7 +16,7 @@ import socket import subprocess import time import urllib.request -from datetime import datetime, timedelta +from datetime import date, datetime, timedelta import yaml from playwright.sync_api import sync_playwright @@ -581,43 +581,59 @@ def _record_business_date(site, kind, date=None): _write("undelivered", date if date else off("expected")) -def _fresh_ingest(ingest, site, kind): - """ingest_state[site][kind] 是否「今天入库成功」——ready 派生的 DB 真相来源。""" - v = ingest.get(site, {}).get(kind) - if not v or not v.get("ok"): - return False - return (v.get("ingested_at") or "")[:10] == datetime.now().strftime("%Y-%m-%d") +def _ready_flags(site): + """从 PG 业务表派生单站三就绪态(ready = DB 数据真相)。 + expected/actual = PG 中存在对应 target_date(today − offset)的数据; + 百世 undelivered = baishi_daily_stats 中存在 target_date 的数据; + 4 站 undelivered = expected_ready ∧ actual_ready(派生)。 + PG 不可达时返回全 False(降级安全,不阻塞心跳)。 -def _ready_flags(ingest, site): - """从 ingest_state 派生单站三就绪态(ready = DB 入库真相)。 + 返回 (flags: {kind: bool}, dates: {kind: target_date_str})。 + dates 与 flags 同源——ready=True 时 business_date 即该 target_date, + 彻底消除 ready 与 business_date 不同源导致的日期标签漂移。""" + from inbound_verify import store # 懒导入:避免模块级循环 + + today = date.today() + today_str = today.isoformat() - expected/actual = 该类今天入库成功;undelivered = 百世(原生未到)今天入库成功 - / 4 站(派生未到)= expected ∧ actual。 - """ if site == "百世": - return { - "expected": False, - "actual": False, - "undelivered": _fresh_ingest(ingest, site, "undelivered"), - } - exp = _fresh_ingest(ingest, site, "expected") - act = _fresh_ingest(ingest, site, "actual") - return {"expected": exp, "actual": act, "undelivered": exp and act} + has_und, _ = store.has_data(site, "undelivered", today_str) + return ( + {"expected": False, "actual": False, "undelivered": has_und}, + {"undelivered": today_str}, + ) + + exp_off = state_store.get_offset(site, "expected") + act_off = state_store.get_offset(site, "actual") + exp_date = (today - timedelta(days=exp_off)).isoformat() + act_date = (today - timedelta(days=act_off)).isoformat() + + has_exp, _ = store.has_data(site, "expected", exp_date) + has_act, _ = store.has_data(site, "actual", act_date) + return ( + {"expected": has_exp, "actual": has_act, "undelivered": has_exp and has_act}, + {"expected": exp_date, "actual": act_date, "undelivered": exp_date}, + ) -def _apply_ready(site, flags): - """写入单站就绪态(仅 ready,不碰 business_date/generated_at)。失败仅告警。""" +def _apply_ready(site, flags, dates=None): + """写入单站就绪态 + 业务日期(同源:ready 与 business_date 均据 PG + offset 派生)。 + ready=True 时同步写入 target_date 作为 business_date,消除不同源导致的日期标签漂移。 + 失败仅告警。""" for k, rdy in flags.items(): try: state_store.set_ready(site, k, rdy) + if rdy and dates and dates.get(k): + state_store.set_business_date(site, k, dates[k]) except Exception as e: print(f">> [状态] 置就绪态失败 {site}/{k}: {e}") def _refresh_ready(site): - """入库后立即据最新 ingest_state 派生并写入该站就绪态(省 30s 心跳等待,与心跳同源)。""" - _apply_ready(site, _ready_flags(state_store.get_all_ingest_state(), site)) + """入库后立即据 PG 派生并写入该站就绪态(省 30s 心跳等待,与心跳同源)。""" + flags, dates = _ready_flags(site) + _apply_ready(site, flags, dates) def _persist_to_db(site, kind): @@ -692,16 +708,14 @@ def dispatch_task(ctx, task_spec): def run_heartbeat(ctx): - """一轮心跳:探测各站登录态 + 据 ingest_state 派生数据就绪态;登录态变化时提示。 + """一轮心跳:探测各站登录态 + 据 PG 业务表派生数据就绪态;登录态变化时提示。 - ready 不再读 Excel、也不在启动时重置,而是每轮从持久的 ingest_state 派生(重启不丢): - expected/actual_ready = 该类今天入库成功; - undelivered_ready = 百世(原生未到)今天入库成功 / 4 站(派生未到)= expected ∧ actual。 - _refresh_ready 在入库瞬间即据 ingest_state 派生(省 30s 等待),心跳同源复核。 + ready 直接查询 PG 业务表(expected_record / actual_record / baishi_daily_stats), + 以「目标业务日期是否有数据」为唯一依据,彻底消除 ingest_state 日期比对带来的每日零点重置。 + _refresh_ready 在入库瞬间即据 PG 派生(省 30s 等待),心跳同源复核。 只在 Playwright 所属线程调用。 """ prev = state_store.get_all_status() - ingest = state_store.get_all_ingest_state() for site_name in ctx.sites_to_watch: logged_in = probe_site_login(site_name, ctx.pages_map) prev_login = prev.get(site_name, {}).get("login_state") @@ -709,4 +723,5 @@ def run_heartbeat(ctx): now_login = state_store.LOGIN_IN if logged_in else state_store.LOGIN_OUT if prev_login and prev_login not in (now_login, state_store.LOGIN_UNKNOWN): print(f"\n ⚠️【{site_name}】登录态变化: {prev_login} → {now_login}") - _apply_ready(site_name, _ready_flags(ingest, site_name)) + flags, dates = _ready_flags(site_name) + _apply_ready(site_name, flags, dates) diff --git a/inbound_verify/store.py b/inbound_verify/store.py index c9b73af..1f71fa1 100644 --- a/inbound_verify/store.py +++ b/inbound_verify/store.py @@ -482,6 +482,52 @@ def get_existing_handover_nos(site): return set() +# ============================== PG 数据存在性查询 ============================== + + +def has_data(site, kind, target_date): + """查询 PG:指定站点在 target_date 是否有业务数据。 + target_date: str 'YYYY-MM-DD' 或 date 对象。 + 返回 (has_rows: bool, count: int)。 + PG 不可达时返回 (False, 0),不抛异常——调用方按「未确认存在」处理。 + + kind 路由: + expected → expected_record (business_date) + actual → actual_record (scan_time::date) + undelivered → 百世: baishi_daily_stats;4 站: 不单独查(由调用方 expected∧actual 派生) + """ + if site == "百世" and kind == "undelivered": + sql = ( + "SELECT COUNT(*) FROM baishi_daily_stats" + " WHERE site = %s AND business_date = %s" + ) + params = (site, target_date) + elif kind == "expected": + sql = ( + "SELECT COUNT(*) FROM expected_record" + " WHERE site = %s AND business_date = %s" + ) + params = (site, target_date) + elif kind == "actual": + sql = ( + "SELECT COUNT(*) FROM actual_record" + " WHERE site = %s AND scan_time::date = %s" + ) + params = (site, target_date) + else: + return (False, 0) + try: + with _connect(_load_pg_config()["dbname"]) as conn: + with conn.cursor() as cur: + cur.execute(sql, params) + row = cur.fetchone() + cnt = int(row[0]) if row else 0 + return (cnt > 0, cnt) + except Exception as e: + print(f">> [状态] PG 查询 {site}/{kind}/{target_date} 失败: {e}") + return (False, 0) + + # ============================== 命令行 ==============================