fix(growth): align sales freshness with report SLA
Some checks failed
CD Pipeline / deploy (push) Has been cancelled

This commit is contained in:
ogt
2026-07-15 01:38:57 +08:00
parent de6933c137
commit 11d3865dfc
17 changed files with 422 additions and 85 deletions

View File

@@ -30,6 +30,7 @@ DAILY_SALES_DETAIL_SHEET_HINTS = ("即時業績明細", "業績明細", "明細"
DAILY_SALES_HEADER_SCAN_ROWS = 15
from services.google_drive_service import drive_service
from services.pchome_sales_freshness_policy import evaluate_pchome_sales_freshness
from database.import_models import ImportJob, ImportConfig, Base
from database.manager import ensure_metadata_initialized
@@ -1213,17 +1214,19 @@ class ImportService:
if not files:
logger.info("沒有找到待匯入的檔案")
data_lag_days = None
latest_sales_date = None
freshness = evaluate_pchome_sales_freshness(
None,
now=datetime.now(TAIPEI_TZ),
)
# Staleness gate (critic-approved 2026-05-03)
# 'move-then-success' 反模式:成功 import 後 move_file 把 Excel 搬到
# 「已匯入」資料夾 → 後續排程 list 回空 → 走此分支 silent return success
# → 4/27~5/2 daily_sales_snapshot 停更 8 天無告警。補主動偵測:
# Drive 空 + DB ≥3 天無新資料時主動發催促告警(週末跨假期不誤觸)。
# Drive 空時依台北時間到檔 SLA 判斷,避免午夜假性過期,也避免
# cutoff 後仍把缺檔當成成功。
try:
from sqlalchemy import text
from datetime import date
from services.openclaw_strategist_service import _send_data_stale_alert
_stale_session = Session()
@@ -1235,19 +1238,15 @@ class ImportService:
_stale_session.close()
if last_date:
normalized_last_date = (
last_date
if isinstance(last_date, date)
else date.fromisoformat(str(last_date)[:10])
freshness = evaluate_pchome_sales_freshness(
last_date,
now=datetime.now(TAIPEI_TZ),
)
days_since = (datetime.now(TAIPEI_TZ).date() - normalized_last_date).days
data_lag_days = days_since
latest_sales_date = str(normalized_last_date)
if days_since >= 3 and send_stale_alert:
if not freshness["decision_ready"] and send_stale_alert:
_send_data_stale_alert(
report_type="upstream_drive",
last_date=str(normalized_last_date),
period=f"已停更 {days_since}",
last_date=str(freshness["latest_sales_date"]),
period=f"超過到檔 SLA {freshness['sla_lag_days']}",
)
except Exception:
logger.error(
@@ -1259,22 +1258,30 @@ class ImportService:
'success': True,
'status': (
'upstream_missing'
if data_lag_days is None
if freshness.get('latest_sales_date') is None
else 'upstream_stale'
if data_lag_days >= 3
if not freshness['decision_ready']
else 'no_pending_file'
),
'message': (
f'上游沒有新檔,最新業績已落後 {data_lag_days}'
if data_lag_days is not None and data_lag_days >= 3
f"上游沒有新檔,最新業績超過到檔 SLA {freshness['sla_lag_days']}"
if freshness.get('latest_sales_date') is not None
and not freshness['decision_ready']
else '沒有找到待匯入的檔案'
),
'file_count': 0,
'imported_count': 0,
'latest_sales_date': latest_sales_date,
'data_lag_days': data_lag_days,
'decision_ready': data_lag_days is not None and data_lag_days <= 1,
'requires_upstream_acquisition': data_lag_days is None or data_lag_days >= 3,
'latest_sales_date': freshness.get('latest_sales_date'),
'data_lag_days': freshness.get('data_lag_days'),
'sla_lag_days': freshness.get('sla_lag_days'),
'freshness_status': freshness.get('freshness_status'),
'expected_latest_sales_date': freshness.get('expected_latest_sales_date'),
'grace_period_active': freshness.get('grace_period_active'),
'report_cutoff_hour': freshness.get('report_cutoff_hour'),
'report_due_at': freshness.get('report_due_at'),
'freshness_contract': freshness.get('contract_version'),
'decision_ready': freshness.get('decision_ready'),
'requires_upstream_acquisition': freshness.get('requires_upstream_acquisition'),
}
# 處理每個檔案

View File

@@ -6,12 +6,16 @@ from __future__ import annotations
import json
import logging
from datetime import date, datetime
from datetime import date, datetime, time
from typing import Any
from sqlalchemy import bindparam, inspect, text
from services.external_market_offer_service import build_external_source_readiness
from services.pchome_sales_freshness_policy import (
TAIPEI_TZ,
evaluate_pchome_sales_freshness,
)
logger = logging.getLogger(__name__)
@@ -60,45 +64,15 @@ def _source_names_by_status(source_readiness: dict[str, Any], status_code: str)
]
def _sales_freshness(latest_sales_date: Any, *, today: date | None = None) -> dict[str, Any]:
raw = str(latest_sales_date or "").strip()
if not raw:
return {
"status": "missing",
"label": "尚無業績資料",
"age_days": None,
"decision_ready": False,
"next_action": "自動取得並匯入最新 PChome 業績檔",
}
try:
parsed = date.fromisoformat(raw[:10])
except ValueError:
return {
"status": "invalid",
"label": "業績日期格式異常",
"age_days": None,
"decision_ready": False,
"next_action": "修正業績資料日期後重新匯入",
}
age_days = max(0, ((today or datetime.now().date()) - parsed).days)
if age_days <= 1:
status, label, ready = "fresh", "資料新鮮", True
elif age_days == 2:
status, label, ready = "warning", "資料即將過期", False
else:
status, label, ready = "critical", "資料已過期", False
return {
"status": status,
"label": label,
"age_days": age_days,
"decision_ready": ready,
"next_action": (
"依最新業績執行今日作戰清單"
if ready
else "自動取得並匯入最新 PChome 業績檔"
),
}
def _sales_freshness(
latest_sales_date: Any,
*,
today: date | None = None,
now: datetime | None = None,
) -> dict[str, Any]:
if now is None and today is not None:
now = TAIPEI_TZ.localize(datetime.combine(today, time(hour=23, minute=59)))
return evaluate_pchome_sales_freshness(latest_sales_date, now=now)
def _table_exists(engine, table_name: str) -> bool:
@@ -158,7 +132,6 @@ def _fetch_sales_rows(conn, limit: int) -> tuple[list[dict[str, Any]], str | Non
dialect = conn.dialect.name
sku_col = _quote_identifier(cols["sku"])
name_col = _quote_identifier(cols["name"])
date_col = _quote_identifier(cols["date"])
revenue_expr = _numeric_expr(cols["revenue"], dialect)
qty_expr = _numeric_expr(cols["qty"], dialect) if cols.get("qty") else "0"

View File

@@ -6,7 +6,7 @@ import json
import logging
import os
import uuid
from datetime import date, datetime
from datetime import datetime
from pathlib import Path
from typing import Dict, List, Optional
@@ -24,6 +24,7 @@ from services.pchome_sales_acquisition_providers import (
env_int,
)
from services.pchome_sales_source_reconciliation import build_source_reconciliation
from services.pchome_sales_freshness_policy import evaluate_pchome_sales_freshness
logger = logging.getLogger(__name__)
@@ -50,23 +51,10 @@ class PChomeSalesAcquisitionService:
except Exception:
logger.warning("Unable to read PChome sales freshness", exc_info=True)
normalized = None
if latest:
try:
normalized = latest if isinstance(latest, date) else date.fromisoformat(str(latest)[:10])
except (TypeError, ValueError):
normalized = None
lag_days = (datetime.now(TAIPEI_TZ).date() - normalized).days if normalized else None
return {
"latest_sales_date": normalized.isoformat() if normalized else None,
"data_lag_days": lag_days,
"decision_ready": lag_days is not None and lag_days <= 1,
"freshness_status": (
"fresh" if lag_days is not None and lag_days <= 1
else "degraded" if lag_days is not None
else "missing"
),
}
return evaluate_pchome_sales_freshness(
latest,
now=datetime.now(TAIPEI_TZ),
)
def _latest_receipt(self) -> Optional[Dict[str, object]]:
session = Session()
@@ -590,6 +578,13 @@ class PChomeSalesAcquisitionService:
"date_range": date_range,
"latest_sales_date": verifier.get("latest_sales_date"),
"data_lag_days": verifier.get("data_lag_days"),
"sla_lag_days": verifier.get("sla_lag_days"),
"freshness_status": verifier.get("freshness_status"),
"expected_latest_sales_date": verifier.get("expected_latest_sales_date"),
"grace_period_active": verifier.get("grace_period_active"),
"report_cutoff_hour": verifier.get("report_cutoff_hour"),
"report_due_at": verifier.get("report_due_at"),
"freshness_contract": verifier.get("contract_version"),
"decision_ready": verifier.get("decision_ready"),
"requires_upstream_acquisition": not bool(verifier.get("decision_ready")),
"manual_review_required": False,

View File

@@ -0,0 +1,158 @@
"""Taipei-time availability SLA for PChome daily sales reports."""
from __future__ import annotations
import os
from datetime import date, datetime, time, timedelta
from typing import Any
import pytz
CONTRACT_VERSION = "pchome_sales_freshness_sla_v1"
TAIPEI_TZ = pytz.timezone("Asia/Taipei")
DEFAULT_REPORT_CUTOFF_HOUR = 20
def _bounded_cutoff_hour(value: Any = None) -> int:
raw = os.getenv("PCHOME_SALES_REPORT_CUTOFF_HOUR", "20") if value is None else value
try:
parsed = int(raw)
except (TypeError, ValueError):
parsed = DEFAULT_REPORT_CUTOFF_HOUR
return max(0, min(23, parsed))
def _taipei_now(value: datetime | None) -> datetime:
if value is None:
return datetime.now(TAIPEI_TZ)
if value.tzinfo is None:
return TAIPEI_TZ.localize(value)
return value.astimezone(TAIPEI_TZ)
def _parse_sales_date(value: Any) -> date | None:
if isinstance(value, datetime):
return value.date()
if isinstance(value, date):
return value
raw = str(value or "").strip()
if not raw:
return None
try:
return date.fromisoformat(raw[:10])
except ValueError:
return None
def evaluate_pchome_sales_freshness(
latest_sales_date: Any,
*,
now: datetime | None = None,
cutoff_hour: int | None = None,
) -> dict[str, Any]:
"""Evaluate freshness against the report arrival deadline, not midnight."""
observed_at = _taipei_now(now)
bounded_cutoff = _bounded_cutoff_hour(cutoff_hour)
report_due_at = TAIPEI_TZ.localize(
datetime.combine(observed_at.date(), time(hour=bounded_cutoff))
)
arrival_window_open = observed_at < report_due_at
expected_latest_date = observed_at.date() - timedelta(
days=2 if arrival_window_open else 1
)
raw = str(latest_sales_date or "").strip()
parsed = _parse_sales_date(latest_sales_date)
common = {
"contract_version": CONTRACT_VERSION,
"observed_at": observed_at.isoformat(),
"report_timezone": "Asia/Taipei",
"report_cutoff_hour": bounded_cutoff,
"report_due_at": report_due_at.isoformat(),
"arrival_window_open": arrival_window_open,
"expected_latest_sales_date": expected_latest_date.isoformat(),
}
if not raw:
return {
**common,
"latest_sales_date": None,
"status": "missing",
"freshness_status": "missing",
"label": "尚無業績資料",
"age_days": None,
"data_lag_days": None,
"sla_lag_days": None,
"grace_period_active": False,
"decision_ready": False,
"requires_upstream_acquisition": True,
"next_action": "自動取得並匯入最新 PChome 業績檔",
}
if parsed is None:
return {
**common,
"latest_sales_date": raw[:10],
"status": "invalid",
"freshness_status": "invalid",
"label": "業績日期格式異常",
"age_days": None,
"data_lag_days": None,
"sla_lag_days": None,
"grace_period_active": False,
"decision_ready": False,
"requires_upstream_acquisition": True,
"next_action": "修正業績資料日期後重新匯入",
}
age_days = (observed_at.date() - parsed).days
sla_lag_days = (expected_latest_date - parsed).days
if age_days < 0:
status = "future"
label = "業績日期超前"
ready = False
elif sla_lag_days <= 0:
ready = True
if arrival_window_open and age_days == 2:
status = "grace"
label = "到檔 SLA 內"
else:
status = "fresh"
label = "資料新鮮"
elif sla_lag_days == 1:
status = "warning"
label = "資料超過到檔 SLA"
ready = False
else:
status = "critical"
label = "資料已過期"
ready = False
grace_period_active = bool(ready and status == "grace")
return {
**common,
"latest_sales_date": parsed.isoformat(),
"status": status,
"freshness_status": status,
"label": label,
"age_days": max(0, age_days),
"data_lag_days": max(0, age_days),
"sla_lag_days": sla_lag_days,
"grace_period_active": grace_period_active,
"decision_ready": ready,
"requires_upstream_acquisition": not ready,
"next_action": (
"持續自動探測,於到檔 SLA 前使用最近完整日資料"
if grace_period_active
else "依最新業績執行今日作戰清單"
if ready
else "自動取得並匯入最新 PChome 業績檔"
),
}
__all__ = (
"CONTRACT_VERSION",
"DEFAULT_REPORT_CUTOFF_HOUR",
"TAIPEI_TZ",
"evaluate_pchome_sales_freshness",
)

View File

@@ -251,9 +251,14 @@ def build_source_reconciliation(
automation_decision = "reconcile_authorized_source_configuration"
safe_next_action = "reconcile_authorized_source_configuration"
lag_bucket = _lag_bucket(freshness.get("data_lag_days"))
effective_lag_days = freshness.get("sla_lag_days")
if effective_lag_days is None:
effective_lag_days = freshness.get("data_lag_days")
lag_bucket = _lag_bucket(effective_lag_days)
severity = (
"critical"
"info"
if decision_ready
else "critical"
if lag_bucket in {"critical", "missing"} or not ready_states
else "warning"
if lag_bucket in {"warning", "stale"}
@@ -261,6 +266,9 @@ def build_source_reconciliation(
)
fingerprint_payload = {
"lag_bucket": lag_bucket,
"decision_ready": decision_ready,
"expected_latest_sales_date": freshness.get("expected_latest_sales_date"),
"freshness_contract": freshness.get("contract_version"),
"source_states": [
[item["id"], item["runtime_state"], item.get("error_kind")]
for item in source_states