refactor(status): derive site ready/business_date from PG instead of ingest_state
_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 <noreply@anthropic.com>
This commit is contained in:
@@ -16,7 +16,7 @@ import socket
|
|||||||
import subprocess
|
import subprocess
|
||||||
import time
|
import time
|
||||||
import urllib.request
|
import urllib.request
|
||||||
from datetime import datetime, timedelta
|
from datetime import date, datetime, timedelta
|
||||||
|
|
||||||
import yaml
|
import yaml
|
||||||
from playwright.sync_api import sync_playwright
|
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"))
|
_write("undelivered", date if date else off("expected"))
|
||||||
|
|
||||||
|
|
||||||
def _fresh_ingest(ingest, site, kind):
|
def _ready_flags(site):
|
||||||
"""ingest_state[site][kind] 是否「今天入库成功」——ready 派生的 DB 真相来源。"""
|
"""从 PG 业务表派生单站三就绪态(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")
|
|
||||||
|
|
||||||
|
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):
|
返回 (flags: {kind: bool}, dates: {kind: target_date_str})。
|
||||||
"""从 ingest_state 派生单站三就绪态(ready = DB 入库真相)。
|
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 == "百世":
|
if site == "百世":
|
||||||
return {
|
has_und, _ = store.has_data(site, "undelivered", today_str)
|
||||||
"expected": False,
|
return (
|
||||||
"actual": False,
|
{"expected": False, "actual": False, "undelivered": has_und},
|
||||||
"undelivered": _fresh_ingest(ingest, site, "undelivered"),
|
{"undelivered": today_str},
|
||||||
}
|
)
|
||||||
exp = _fresh_ingest(ingest, site, "expected")
|
|
||||||
act = _fresh_ingest(ingest, site, "actual")
|
exp_off = state_store.get_offset(site, "expected")
|
||||||
return {"expected": exp, "actual": act, "undelivered": exp and act}
|
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):
|
def _apply_ready(site, flags, dates=None):
|
||||||
"""写入单站就绪态(仅 ready,不碰 business_date/generated_at)。失败仅告警。"""
|
"""写入单站就绪态 + 业务日期(同源:ready 与 business_date 均据 PG + offset 派生)。
|
||||||
|
ready=True 时同步写入 target_date 作为 business_date,消除不同源导致的日期标签漂移。
|
||||||
|
失败仅告警。"""
|
||||||
for k, rdy in flags.items():
|
for k, rdy in flags.items():
|
||||||
try:
|
try:
|
||||||
state_store.set_ready(site, k, rdy)
|
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:
|
except Exception as e:
|
||||||
print(f">> [状态] 置就绪态失败 {site}/{k}: {e}")
|
print(f">> [状态] 置就绪态失败 {site}/{k}: {e}")
|
||||||
|
|
||||||
|
|
||||||
def _refresh_ready(site):
|
def _refresh_ready(site):
|
||||||
"""入库后立即据最新 ingest_state 派生并写入该站就绪态(省 30s 心跳等待,与心跳同源)。"""
|
"""入库后立即据 PG 派生并写入该站就绪态(省 30s 心跳等待,与心跳同源)。"""
|
||||||
_apply_ready(site, _ready_flags(state_store.get_all_ingest_state(), site))
|
flags, dates = _ready_flags(site)
|
||||||
|
_apply_ready(site, flags, dates)
|
||||||
|
|
||||||
|
|
||||||
def _persist_to_db(site, kind):
|
def _persist_to_db(site, kind):
|
||||||
@@ -692,16 +708,14 @@ def dispatch_task(ctx, task_spec):
|
|||||||
|
|
||||||
|
|
||||||
def run_heartbeat(ctx):
|
def run_heartbeat(ctx):
|
||||||
"""一轮心跳:探测各站登录态 + 据 ingest_state 派生数据就绪态;登录态变化时提示。
|
"""一轮心跳:探测各站登录态 + 据 PG 业务表派生数据就绪态;登录态变化时提示。
|
||||||
|
|
||||||
ready 不再读 Excel、也不在启动时重置,而是每轮从持久的 ingest_state 派生(重启不丢):
|
ready 直接查询 PG 业务表(expected_record / actual_record / baishi_daily_stats),
|
||||||
expected/actual_ready = 该类今天入库成功;
|
以「目标业务日期是否有数据」为唯一依据,彻底消除 ingest_state 日期比对带来的每日零点重置。
|
||||||
undelivered_ready = 百世(原生未到)今天入库成功 / 4 站(派生未到)= expected ∧ actual。
|
_refresh_ready 在入库瞬间即据 PG 派生(省 30s 等待),心跳同源复核。
|
||||||
_refresh_ready 在入库瞬间即据 ingest_state 派生(省 30s 等待),心跳同源复核。
|
|
||||||
只在 Playwright 所属线程调用。
|
只在 Playwright 所属线程调用。
|
||||||
"""
|
"""
|
||||||
prev = state_store.get_all_status()
|
prev = state_store.get_all_status()
|
||||||
ingest = state_store.get_all_ingest_state()
|
|
||||||
for site_name in ctx.sites_to_watch:
|
for site_name in ctx.sites_to_watch:
|
||||||
logged_in = probe_site_login(site_name, ctx.pages_map)
|
logged_in = probe_site_login(site_name, ctx.pages_map)
|
||||||
prev_login = prev.get(site_name, {}).get("login_state")
|
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
|
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):
|
if prev_login and prev_login not in (now_login, state_store.LOGIN_UNKNOWN):
|
||||||
print(f"\n ⚠️【{site_name}】登录态变化: {prev_login} → {now_login}")
|
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)
|
||||||
|
|||||||
@@ -482,6 +482,52 @@ def get_existing_handover_nos(site):
|
|||||||
return set()
|
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)
|
||||||
|
|
||||||
|
|
||||||
# ============================== 命令行 ==============================
|
# ============================== 命令行 ==============================
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user