Просмотр исходного кода

feat(roi): schedule guarded internal automation report

刘立冬 16 часов назад
Родитель
Сommit
83480c76b9

+ 10 - 2
examples/auto_put_ad_mini/.env.example

@@ -30,9 +30,15 @@ DAILY_SYNC_DELIVERY_TEMPLATE=0
 
 # 日级 ROI 北极星指标与表格逐行审批。默认关闭,分阶段启用。
 DAILY_ROI_ENABLED=0
-DAILY_ROI_HOUR=11
+DAILY_ROI_HOUR=9
 DAILY_ROI_MINUTE=0
 DAILY_ROI_RUN_ON_STARTUP=0
+# 每天14:30运行独立内部测试,只将“自动化投放”分表发送到“内部”Webhook,不进入腾讯审批执行。
+ROI_INTERNAL_TEST_ENABLED=0
+ROI_INTERNAL_TEST_HOUR=14
+ROI_INTERNAL_TEST_MINUTE=30
+# 生产默认去重;测试环境显式设0时每次触发使用独立测试批次。
+ROI_INTERNAL_TEST_DEDUP_ENABLED=1
 ROI_APPLY_ENABLED=0
 ROI_SHEET_APPROVAL_ENABLED=0
 ROI_SHEET_APPROVAL_POLL_SECONDS=60
@@ -56,8 +62,10 @@ ROI_FISSION_PARAMETER_VERSION=20260712_A0-A15_v2
 # 审批表链接为获得链接者可编辑,黄色审批列选择“批准”即自动执行。
 # 为空时依次复用 FEISHU_AD_PROJECT_CHAT_ID / RTC_COMMAND_CHAT_ID
 ROI_FEISHU_CHAT_ID=
+# 日级ROI定时任务失败告警群;为空时回退FEISHU_OPERATOR_CHAT_ID。
+ROI_FAILURE_FEISHU_CHAT_ID=
 # 代理表由飞书应用上传后,通过各代理自定义机器人发送在线表卡片。
-# 默认关闭;JSON 的 key 使用代理简称(如“棱镜”),完整 webhook 只写入真实 .env/密钥系统。
+# 默认关闭;JSON 的 key 使用代理简称(如“棱镜”),“内部”为内部测试专用路由
 ROI_AGENCY_WEBHOOK_ENABLED=0
 ROI_AGENCY_WEBHOOKS_JSON={}
 

+ 3 - 0
examples/auto_put_ad_mini/DEPLOYMENT.md

@@ -170,6 +170,7 @@ curl -X POST http://localhost:8080/trigger | jq .
 
 ### 日级 ROI 三日口径
 
+- 正式任务每天北京时间09:00执行;在创建数据库批次和发送消息前,必须确认 `loghubods.opengid_base_data` 已存在精确T-1分区,且小程序投流和公众号投流相关数据均非空。未就绪时禁止退回旧分区。
 - 日级 ROI 读取 T-1 至 T-3 三个连续完整数据日,数据深度过滤使用最新业务 SQL 的 `usersharedepth<=1`。
 - 正式阈值样本为连续三天每天首层 UV>200、成本>0 且 ROI 有效的小程序创意级和公众号实体;两类实体合并后按实体等权计算整体 P20 关停线。
 - 小程序广告级从 ODPS 原始数据按广告直接 `COUNT(DISTINCT mid)`,复用同一 P20,但不重复进入样本池;低于关停线且广告 age>3 天时生成广告级关停建议。
@@ -183,6 +184,8 @@ curl -X POST http://localhost:8080/trigger | jq .
 - 汇总表将三日总 T0 裂变人数 / 三日总首层 UV 的加权比例展示为“日均T0裂变率”。汇总顺序为关停、扩量、中间观察、条件不足观察;第一条扩量行和扩量后第一条观察行顶部均使用粗线分隔。“当日效率ROI”和“预测总效率ROI”均使用深红—黄—绿色阶,绿色代表表现好。普通数值显示两位小数,UV、人数、数量和广告年龄显示整数;两个“裂变系数-总裂变UV/…”比率列固定显示两位小数。“传播裂变系数匹配”保留审计数据但默认隐藏。
 - “审批选择”默认隐藏,需取消隐藏后审批;创意级“当前创意状态”只对关停建议只读腾讯状态,显示正常、已停止或读取失败,其他行留空。
 - 完整主报表保存后,流程会额外按“小程序投流”渠道的“代理名称”生成一代理一份调控建议工作簿;文件名为 `YYYYMMDD_代理名称_调控建议.xlsx`。代理版只包含小程序创意级和广告级三日汇总,广告级 Sheet 继续隐藏,不包含每日明细。包名、广告age、日均首层UV、建议说明、收入、预测总效率ROI、P20/P80/P30、排名、两个裂变系数、日均T0裂变人数/率、审批、执行状态和幂等键均不会写入代理文件;内部“当日效率ROI”仅改名为两位小数的“评分”展示,建议动作仍按预测总效率ROI计算并保留“关停创意/关停广告/扩量/观察”。代理表默认只落本地;开启 `ROI_AGENCY_WEBHOOK_ENABLED` 后,飞书应用先上传在线表,再由 `ROI_AGENCY_WEBHOOKS_JSON` 中精确匹配的代理机器人发送卡片。未配置代理直接跳过,不回退总群;完整 webhook 只能放在真实 `.env` 或密钥系统,数据库和日志只保存哈希指纹。
+- 可选内部测试任务每天北京时间14:30执行独立复算,只将“自动化投放”分表发送到“内部”Webhook,不发送完整主表、正式总群或代理群,也不生成可审批腾讯动作。缺少该分表时任务失败并告警。正式任务始终去重;内部测试使用 `ROI_INTERNAL_TEST_DEDUP_ENABLED` 控制,设为`0`时每次显式触发都发送新的测试批次。内部任务无startup入口,重启服务不会自动补发。
+- 09:00正式任务或14:30内部测试执行失败时,调度服务通过飞书应用向 `ROI_FAILURE_FEISHU_CHAT_ID` 发送红色告警卡片;未配置时回退 `FEISHU_OPERATOR_CHAT_ID`。告警失败只记日志,不会掩盖原任务失败。
 
 离线复算不发送飞书:
 

+ 14 - 3
examples/auto_put_ad_mini/docs/unified_services_deployment.md

@@ -70,9 +70,14 @@ DAILY_REVIEW_RUN_ON_STARTUP=0
 DAILY_SYNC_DELIVERY_TEMPLATE=0
 
 DAILY_ROI_ENABLED=0
-DAILY_ROI_HOUR=11
+DAILY_ROI_HOUR=9
 DAILY_ROI_MINUTE=0
 DAILY_ROI_RUN_ON_STARTUP=0
+ROI_INTERNAL_TEST_ENABLED=0
+ROI_INTERNAL_TEST_HOUR=14
+ROI_INTERNAL_TEST_MINUTE=30
+ROI_INTERNAL_TEST_DEDUP_ENABLED=1
+ROI_FAILURE_FEISHU_CHAT_ID=
 ROI_APPLY_ENABLED=0
 ROI_SHEET_APPROVAL_ENABLED=0
 ROI_SHEET_APPROVAL_POLL_SECONDS=60
@@ -189,10 +194,16 @@ ROI 表格使用获得链接者可编辑权限。默认隐藏的黄色【审批
 创意级“当前创意状态”只对关停建议只读腾讯状态并显示正常、已停止或读取失败,其他行留空;“审批选择”默认隐藏。
 完整主报表写入成功后,服务额外按“小程序投流”渠道的“代理名称”生成一代理一份 `roi_agency_advice_v7` 工作簿,文件名为 `YYYYMMDD_代理名称_调控建议.xlsx`。代理版可见 Sheet 名为“小程序创意调控建议”,并保留默认隐藏的“小程序广告级三日汇总”,不包含每日明细;包名、广告age、日均首层UV、建议说明、收入、预测总效率ROI、阈值、排名、两个裂变系数、日均T0裂变人数/率、审批和执行审计列均物理删除。内部“当日效率ROI”仅以两位小数的“评分”列对外展示,不参与代理动作计算;建议动作仍按预测总效率ROI生成,并明确显示“关停创意”或“关停广告”。代理表生成不改变主报表、数据库快照或主审批链接。
 
+正式任务每天北京时间09:00执行。任务固定以T-1为结束日,并在创建MySQL运行批次、生成报表和发送飞书前校验 `loghubods.opengid_base_data` 的精确T-1分区;小程序投流或公众号投流相关数据任一为空即失败关闭,禁止静默回退到更早分区。
+
 总表继续由飞书应用按 `ROI_FEISHU_CHAT_ID` / `FEISHU_OPERATOR_CHAT_ID` 发布。代理表发布由 `ROI_AGENCY_WEBHOOK_ENABLED` 独立控制,完整机器人地址仅允许保存在密钥系统或真实 `.env` 的 `ROI_AGENCY_WEBHOOKS_JSON`,JSON key 使用代理简称。飞书应用先上传并开放在线表链接,再由对应代理机器人发送可点击卡片;未配置代理记录为跳过,不回退到总群。同一 `run_id + 代理 + 代理报表版本` 幂等审计,已上传表格在通知重试时复用,避免重复建表。机器人地址只保存哈希指纹,不能写入数据库或日志。
 
+`ROI_INTERNAL_TEST_ENABLED=1` 时,服务每天北京时间14:30注册独立内部任务。该任务复算后使用空管理账户集合,不生成可审批腾讯动作,只把“自动化投放”分表发送到 `ROI_AGENCY_WEBHOOKS_JSON` 的“内部”路由,不发送完整主表、正式总群或代理群;缺少该分表时任务失败并告警。它不注册startup任务,服务重启不会主动补发。正式任务始终去重;内部测试由 `ROI_INTERNAL_TEST_DEDUP_ENABLED` 控制,生产默认`1`,测试环境显式设为`0`时每次触发生成独立测试run并重新发送。
+
+09:00正式任务或14:30内部测试的子进程启动失败或非零退出时,调度服务通过飞书应用向 `ROI_FAILURE_FEISHU_CHAT_ID` 发送红色失败告警;该变量为空时回退 `FEISHU_OPERATOR_CHAT_ID`。失败告警本身发送失败只记录错误日志,不会吞掉或替换原任务异常。
+
 1. 首次部署先执行下方参数发布命令并完成回读校验。
-2. 设置 `DAILY_ROI_ENABLED=1`、`ROI_APPLY_ENABLED=0`、`ROI_SHEET_APPROVAL_ENABLED=0`,观察 11:00 的 ODPS 计算、数据库快照和可编辑飞书表。
+2. 设置 `DAILY_ROI_ENABLED=1`、`ROI_APPLY_ENABLED=0`、`ROI_SHEET_APPROVAL_ENABLED=0`,观察09:00的T-1来源校验、ODPS计算、数据库快照和可编辑飞书表。
 3. 核对阈值、账户可执行范围、黄色审批列和动作数量。
 4. 再设置 `ROI_APPLY_ENABLED=1`、`ROI_SHEET_APPROVAL_ENABLED=1`,重建 `ad-control-service`;批准一条低风险测试行,核对腾讯回读、表格结果和群通知。
 5. 生产审批有效期默认 120 分钟;过期行不能执行。表格逐行选择“拒绝”不会产生腾讯写操作。
@@ -265,7 +276,7 @@ docker compose --env-file /dev/null logs --since=30m ad_daily_service
 - 飞书只有一个 WebSocket 消费者。
 - 10:30 只出现一次创建任务。
 - 审核扫描每 2 小时最多一个实例。
-- 11:00 只生成一个相同指标/策略版本的 ROI 批次,同日重跑不会重复发布
+- 09:00正式任务只生成一个相同指标/策略版本的ROI批次,同日重跑不会重复发布;14:30内部测试是否去重由独立变量控制
 - ROI 审批执行与实时 CPM 调控不并发写腾讯,基础出价调整后两边状态一致。
 
 回滚时将 `VERSION` 改为上一个镜像版本并重新部署。数据库新增表和列可以保留,

+ 4 - 1
examples/auto_put_ad_mini/roi_control/agency_delivery.py

@@ -149,6 +149,8 @@ def publish_agency_reports(
                         {
                             "agency_name": agency_name,
                             "status": "SENT",
+                            "sheet_url": str(delivery.get("sheet_url") or ""),
+                            "sheet_token": str(delivery.get("sheet_token") or ""),
                             "reused": True,
                         }
                     )
@@ -181,7 +183,7 @@ def publish_agency_reports(
                     )
                 response_code = sender.send(
                     webhook_url,
-                    title=path.stem,
+                    title=str(report.get("title") or path.stem),
                     sheet_url=sheet_url,
                     creative_rows=int(report.get("creative_rows") or 0),
                     ad_rows=int(report.get("ad_rows") or 0),
@@ -198,6 +200,7 @@ def publish_agency_reports(
                         "agency_name": agency_name,
                         "status": "SENT",
                         "sheet_url": sheet_url,
+                        "sheet_token": sheet_token,
                         "reused": False,
                     }
                 )

+ 15 - 2
examples/auto_put_ad_mini/roi_control/config.py

@@ -75,8 +75,12 @@ class RoiConfig:
     apply_enabled: bool = False
     sheet_approval_enabled: bool = False
     sheet_approval_poll_seconds: int = 60
-    report_hour: int = 11
+    report_hour: int = 9
     report_minute: int = 0
+    internal_test_enabled: bool = False
+    internal_test_hour: int = 14
+    internal_test_minute: int = 30
+    internal_test_dedup_enabled: bool = True
     approval_ttl_minutes: int = 120
     scale_ratio: Decimal = Decimal("1.10")
     scale_cooldown_days: int = 3
@@ -103,8 +107,15 @@ class RoiConfig:
             sheet_approval_poll_seconds=int(
                 os.getenv("ROI_SHEET_APPROVAL_POLL_SECONDS", "60")
             ),
-            report_hour=int(os.getenv("DAILY_ROI_HOUR", "11")),
+            report_hour=int(os.getenv("DAILY_ROI_HOUR", "9")),
             report_minute=int(os.getenv("DAILY_ROI_MINUTE", "0")),
+            internal_test_enabled=env_flag("ROI_INTERNAL_TEST_ENABLED"),
+            internal_test_hour=int(os.getenv("ROI_INTERNAL_TEST_HOUR", "14")),
+            internal_test_minute=int(os.getenv("ROI_INTERNAL_TEST_MINUTE", "30")),
+            internal_test_dedup_enabled=env_flag(
+                "ROI_INTERNAL_TEST_DEDUP_ENABLED",
+                True,
+            ),
             approval_ttl_minutes=int(
                 os.getenv("ROI_APPROVAL_TTL_MINUTES", "120")
             ),
@@ -160,6 +171,8 @@ class RoiConfig:
             raise ValueError("ROI_SHEET_APPROVAL_POLL_SECONDS must be at least 30")
         if not 0 <= self.report_hour <= 23 or not 0 <= self.report_minute <= 59:
             raise ValueError("DAILY_ROI_HOUR/MINUTE is invalid")
+        if not 0 <= self.internal_test_hour <= 23 or not 0 <= self.internal_test_minute <= 59:
+            raise ValueError("ROI_INTERNAL_TEST_HOUR/MINUTE is invalid")
         if self.approval_ttl_minutes < 1:
             raise ValueError("ROI_APPROVAL_TTL_MINUTES must be positive")
         if self.scale_ratio <= 1:

+ 71 - 12
examples/auto_put_ad_mini/roi_control/data_source.py

@@ -16,6 +16,10 @@ TABLE_NAME = "loghubods.opengid_base_data"
 SHANGHAI = ZoneInfo("Asia/Shanghai")
 
 
+class SourceDataNotReadyError(RuntimeError):
+    """The exact ROI source partition required for this run is not ready."""
+
+
 def parse_yyyymmdd(value: str) -> datetime:
     try:
         return datetime.strptime(value, "%Y%m%d")
@@ -28,21 +32,76 @@ def date_window(end_date: str) -> Tuple[str, str]:
     return (end - timedelta(days=2)).strftime("%Y%m%d"), end.strftime("%Y%m%d")
 
 
-def resolve_end_date(client: ODPSClient, requested: str | None = None) -> str:
-    """使用显式日期,或选择不晚于昨日的最新可用分区。"""
+def build_source_readiness_sql(end_date: str) -> str:
+    parse_yyyymmdd(end_date)
+    return f"""
+SELECT
+  COUNT(1) AS row_count,
+  SUM(CASE WHEN channel = '{SELF_CHANNEL}' THEN 1 ELSE 0 END) AS self_rows,
+  SUM(CASE WHEN channel = '{GZH_CHANNEL}' THEN 1 ELSE 0 END) AS gzh_rows
+FROM {TABLE_NAME}
+WHERE dt = '{end_date}'
+  AND usersharedepth <= 1
+  AND videoid IS NOT NULL
+  AND NVL(hotsencetype, '') <> '1167'
+""".strip()
 
-    if requested:
-        parse_yyyymmdd(requested)
-        return requested
 
-    latest = (
-        client.odps.get_table("opengid_base_data")
-        .get_max_partition()
-        .partition_spec["dt"]
-    )
+def validate_source_ready(client: ODPSClient, end_date: str) -> dict[str, int]:
+    """Require an exact, non-empty T-1 partition for every reported channel."""
+
+    max_partition = client.odps.get_table("opengid_base_data").get_max_partition()
+    if max_partition is None:
+        raise SourceDataNotReadyError(
+            f"ROI来源表没有可用分区: table={TABLE_NAME}, required_dt={end_date}"
+        )
+    latest = max_partition.partition_spec["dt"]
     parse_yyyymmdd(latest)
-    yesterday = (datetime.now(SHANGHAI) - timedelta(days=1)).strftime("%Y%m%d")
-    return min(latest, yesterday)
+    if latest < end_date:
+        raise SourceDataNotReadyError(
+            f"ROI来源表未同步目标分区: table={TABLE_NAME}, "
+            f"required_dt={end_date}, latest_dt={latest}"
+        )
+
+    frame = client.execute_sql(build_source_readiness_sql(end_date))
+    if frame.empty:
+        raise SourceDataNotReadyError(
+            f"ROI来源表目标分区为空: table={TABLE_NAME}, dt={end_date}"
+        )
+    row = frame.iloc[0]
+    counts = {}
+    for name in ("row_count", "self_rows", "gzh_rows"):
+        value = row.get(name)
+        counts[name] = 0 if pd.isna(value) else int(value)
+    missing = [name for name, value in counts.items() if value <= 0]
+    if missing:
+        raise SourceDataNotReadyError(
+            f"ROI来源表目标分区业务数据未就绪: table={TABLE_NAME}, "
+            f"dt={end_date}, missing={','.join(missing)}"
+        )
+    return counts
+
+
+def resolve_end_date(
+    client: ODPSClient,
+    requested: str | None = None,
+    *,
+    now: datetime | None = None,
+) -> str:
+    """Resolve exactly the requested date or T-1, then require source readiness."""
+
+    if requested:
+        parse_yyyymmdd(requested)
+        target = requested
+    else:
+        current = now or datetime.now(SHANGHAI)
+        if current.tzinfo is None:
+            current = current.replace(tzinfo=SHANGHAI)
+        target = (current.astimezone(SHANGHAI) - timedelta(days=1)).strftime(
+            "%Y%m%d"
+        )
+    validate_source_ready(client, target)
+    return target
 
 
 def build_daily_sql(start_date: str, end_date: str) -> str:

+ 42 - 0
examples/auto_put_ad_mini/roi_control/feishu.py

@@ -254,6 +254,48 @@ class RoiFeishuPublisher:
                 )
                 self._json(response, "send ROI execution result")
 
+    def send_service_alert(
+        self,
+        *,
+        title: str,
+        content: str,
+        chat_id: str | None = None,
+    ) -> str:
+        target_chat_id = (
+            (chat_id or "").strip()
+            or os.getenv("ROI_FAILURE_FEISHU_CHAT_ID", "").strip()
+            or os.getenv("FEISHU_OPERATOR_CHAT_ID", "").strip()
+        )
+        if not target_chat_id:
+            raise RuntimeError("Missing ROI failure alert chat_id")
+        card = {
+            "config": {"wide_screen_mode": True},
+            "header": {
+                "template": "red",
+                "title": {"tag": "plain_text", "content": title},
+            },
+            "elements": [
+                {
+                    "tag": "div",
+                    "text": {"tag": "lark_md", "content": content},
+                }
+            ],
+        }
+        token = self._token()
+        response = self.client.post(
+            f"{BASE_URL}/im/v1/messages",
+            headers={**self._headers(token), "Content-Type": "application/json"},
+            params={"receive_id_type": "chat_id"},
+            json={
+                "receive_id": target_chat_id,
+                "msg_type": "interactive",
+                "content": json.dumps(card, ensure_ascii=False),
+            },
+        )
+        return self._json(response, "send ROI service failure alert")["data"][
+            "message_id"
+        ]
+
     def publish(
         self,
         path: Path,

+ 4 - 0
examples/auto_put_ad_mini/roi_control/reporting.py

@@ -627,6 +627,7 @@ def write_agency_workbooks(
     rows: pd.DataFrame,
     output_dir: Path,
     report_date: str,
+    agency_names: set[str] | None = None,
 ) -> list[dict[str, object]]:
     """Create one miniapp control-advice workbook per agency."""
 
@@ -636,6 +637,9 @@ def write_agency_workbooks(
     ].copy()
     miniapp["_代理规范名"] = miniapp["代理名称"].map(_canonical_agency_name)
     agencies = sorted(name for name in miniapp["_代理规范名"].unique() if name)
+    if agency_names is not None:
+        canonical_names = {_canonical_agency_name(name) for name in agency_names}
+        agencies = [name for name in agencies if name in canonical_names]
     output_dir.mkdir(parents=True, exist_ok=True)
     outputs: list[dict[str, object]] = []
 

+ 138 - 22
examples/auto_put_ad_mini/roi_control/service.py

@@ -14,7 +14,7 @@ from zoneinfo import ZoneInfo
 from storage import initialize_schema, load_managed_accounts
 from tencent_client import ACTIVE_STATUS, TencentClient
 
-from .agency_delivery import publish_agency_reports
+from .agency_delivery import agency_route_key, publish_agency_reports
 from .config import AgencyWebhookConfig, RoiConfig
 from .data_source import (
     ODPSClient,
@@ -68,6 +68,44 @@ def _redact_agency_webhooks(
     return message
 
 
+def _internal_source_revision(config: RoiConfig, now: datetime) -> str:
+    if config.internal_test_dedup_enabled:
+        return "internal_test"
+    return f"it_{now.strftime('%H%M%S%f')}"
+
+
+def _build_internal_reports(
+    *,
+    agency_reports: list[dict[str, object]],
+) -> list[dict[str, object]]:
+    reports = [
+        {
+            **report,
+            "title": f"内部测试_{Path(str(report['report'])).stem}",
+        }
+        for report in agency_reports
+        if agency_route_key(str(report["agency_name"])) == "自动化投放"
+    ]
+    if len(reports) != 1:
+        raise RuntimeError(
+            "Internal ROI test requires exactly one 自动化投放 agency report"
+        )
+    return reports
+
+
+def _internal_webhook_config(
+    reports: list[dict[str, object]],
+    webhook_url: str,
+) -> AgencyWebhookConfig:
+    return AgencyWebhookConfig(
+        enabled=True,
+        webhooks={
+            agency_route_key(str(report["agency_name"])): webhook_url
+            for report in reports
+        },
+    )
+
+
 def _annotate_current_creative_status(
     rows,
     tencent: TencentClient | None = None,
@@ -179,11 +217,33 @@ def run_daily_roi(
     send_feishu: bool,
     now: datetime | None = None,
     source_revision: str | None = None,
+    internal_test: bool = False,
 ) -> dict[str, Any]:
     """Compute, snapshot, report, and optionally publish one ROI batch."""
 
     config = RoiConfig.from_env()
     agency_webhook_config = AgencyWebhookConfig.from_env()
+    if internal_test and send_feishu:
+        raise ValueError("internal_test cannot publish formal Feishu notifications")
+    if internal_test and not (agency_webhook_config.webhooks or {}).get("内部"):
+        raise ValueError("Internal ROI test requires agency webhook route: 内部")
+    effective_now = now or datetime.now(SHANGHAI)
+    if effective_now.tzinfo is None:
+        effective_now = effective_now.replace(tzinfo=SHANGHAI)
+    if internal_test:
+        if source_revision:
+            raise ValueError("internal_test manages source_revision automatically")
+        source_revision = _internal_source_revision(config, effective_now)
+
+    client = ODPSClient(project=os.getenv("ODPS_PROJECT", "loghubods"))
+    end_date = resolve_end_date(
+        client,
+        requested_end_date,
+        now=effective_now,
+    )
+    start_date, end_date = date_window(end_date)
+    expected_dates = _dates(start_date)
+
     initialize_schema()
     fission_version = os.getenv(
         "ROI_FISSION_PARAMETER_VERSION",
@@ -196,12 +256,9 @@ def run_daily_roi(
     run_config["fission_multiplier"] = fission_parameters.snapshot()
     run_config["report_version"] = REPORT_VERSION
     run_config["agency_webhook"] = agency_webhook_config.snapshot()
+    run_config["internal_test"] = internal_test
     if source_revision:
         run_config["source_revision"] = source_revision
-    client = ODPSClient(project=os.getenv("ODPS_PROJECT", "loghubods"))
-    end_date = resolve_end_date(client, requested_end_date)
-    start_date, end_date = date_window(end_date)
-    expected_dates = _dates(start_date)
     run_id, run_key = _run_identity(
         end_date,
         fission_parameters,
@@ -223,8 +280,9 @@ def run_daily_roi(
         }
     )
     reusable_statuses = FINAL_STATUSES | {"PENDING_APPROVAL", "EXECUTING"}
+    publish_requested = send_feishu or internal_test
     if run.get("status") in reusable_statuses or (
-        run.get("status") == "COMPUTED" and not send_feishu
+        run.get("status") == "COMPUTED" and not publish_requested
     ):
         logger.info("Reuse ROI run=%s status=%s", run_id, run.get("status"))
         reused_result: dict[str, Any] = {
@@ -284,10 +342,14 @@ def run_daily_roi(
             _rule_config(config),
             fission_parameters=fission_parameters,
         )
-        managed_ids = {
-            int(row["account_id"])
-            for row in load_managed_accounts()
-        }
+        managed_ids = (
+            set()
+            if internal_test
+            else {
+                int(row["account_id"])
+                for row in load_managed_accounts()
+            }
+        )
         annotated, snapshots, actions = annotate_execution(
             summary,
             managed_ids,
@@ -348,19 +410,31 @@ def run_daily_roi(
         if source_revision:
             batch_name = f"{batch_name}_{source_revision}"
         output_path = output_dir / f"{batch_name}.xlsx"
-        write_workbook(
-            report_rows,
-            thresholds,
-            expected_dates,
-            output_path,
-            run_config,
-            fission_match_summary,
-        )
-        agency_report_date = (now or datetime.now(SHANGHAI)).strftime("%Y%m%d")
+        if not internal_test:
+            write_workbook(
+                report_rows,
+                thresholds,
+                expected_dates,
+                output_path,
+                run_config,
+                fission_match_summary,
+            )
+        agency_report_date = effective_now.strftime("%Y%m%d")
         agency_reports = write_agency_workbooks(
             report_rows,
             output_dir / f"{agency_report_date}_调控建议",
             agency_report_date,
+            agency_names={"自动化投放"} if internal_test else None,
+        )
+        internal_reports = (
+            _build_internal_reports(agency_reports=agency_reports)
+            if internal_test
+            else []
+        )
+        result_report = (
+            str(internal_reports[0]["report"])
+            if internal_test
+            else str(output_path)
         )
 
         result: dict[str, Any] = {
@@ -378,10 +452,52 @@ def run_daily_roi(
             "fission_multiplier_matches": fission_match_summary.to_dict(
                 "records"
             ),
-            "report": str(output_path),
+            "report": result_report,
             "agency_reports": agency_reports,
         }
-        if not send_feishu:
+        if not publish_requested:
+            return result
+
+        if internal_test:
+            internal_url = (agency_webhook_config.webhooks or {})["内部"]
+            internal_config = _internal_webhook_config(
+                internal_reports,
+                internal_url,
+            )
+            publisher = RoiFeishuPublisher()
+            try:
+                deliveries = publish_agency_reports(
+                    run_id=run_id,
+                    reports=internal_reports,
+                    config=internal_config,
+                    publisher=publisher,
+                    now=effective_now,
+                )
+            finally:
+                publisher.close()
+            result["internal_deliveries"] = deliveries
+            failed = [row for row in deliveries if row.get("status") != "SENT"]
+            if failed:
+                raise RuntimeError(
+                    f"Internal ROI test delivery failed for {len(failed)} report(s)"
+                )
+            main_delivery = next(
+                row for row in deliveries if row["agency_name"] == "自动化投放"
+            )
+            mark_published(
+                run_id,
+                sheet_token=str(main_delivery["sheet_token"]),
+                sheet_url=str(main_delivery["sheet_url"]),
+                message_id="internal-webhook",
+                expires_at=None,
+                requires_approval=False,
+            )
+            result.update(
+                {
+                    "status": "COMPLETED",
+                    "sheet_url": main_delivery["sheet_url"],
+                }
+            )
             return result
 
         counts = recommendation_rows.groupby("动作").size().to_dict()
@@ -402,7 +518,7 @@ def run_daily_roi(
                 requires_approval=bool(actions),
             )
             expires_at = (
-                (now or datetime.now(SHANGHAI))
+                effective_now
                 + timedelta(minutes=config.approval_ttl_minutes)
                 if actions
                 else None

+ 6 - 0
examples/auto_put_ad_mini/run_daily_roi.py

@@ -39,6 +39,11 @@ def parse_args() -> argparse.Namespace:
         help="Stable revision label for corrected upstream data",
     )
     parser.add_argument("--send-feishu", action="store_true")
+    parser.add_argument(
+        "--internal-test",
+        action="store_true",
+        help="Run an isolated full calculation and notify only the internal webhook",
+    )
     return parser.parse_args()
 
 
@@ -49,6 +54,7 @@ def main() -> None:
         output_dir=args.output_dir,
         send_feishu=args.send_feishu,
         source_revision=args.source_revision,
+        internal_test=args.internal_test,
     )
     print(json.dumps(result, ensure_ascii=False, default=str))
 

+ 74 - 6
examples/auto_put_ad_mini/run_daily_service.py

@@ -29,7 +29,8 @@ load_dotenv(Path.cwd() / ".env", override=False)
 from config import TIME_SERIES_DEFAULT  # noqa: E402
 from db.connection import get_connection  # noqa: E402
 from storage import advisory_lock, initialize_schema  # noqa: E402
-from roi_control.config import RoiConfig  # noqa: E402
+from roi_control.config import AgencyWebhookConfig, RoiConfig  # noqa: E402
+from roi_control.feishu import RoiFeishuPublisher  # noqa: E402
 
 
 logger = logging.getLogger("auto_put_ad_mini.daily_service")
@@ -42,6 +43,30 @@ def _env_flag(name: str, default: bool) -> bool:
     return raw.strip().lower() in {"1", "true", "yes", "on"}
 
 
+def _notify_roi_failure(
+    *,
+    extra_args: list[str] | None,
+    error: str,
+) -> None:
+    mode = "内部14:30测试" if "--internal-test" in (extra_args or []) else "正式09:00任务"
+    publisher: RoiFeishuPublisher | None = None
+    try:
+        publisher = RoiFeishuPublisher()
+        publisher.send_service_alert(
+            title="日级ROI任务失败",
+            content=(
+                f"任务:**{mode}**\n"
+                f"错误:`{error[:500]}`\n"
+                "本次已停止,不会生成或发送不完整报表。"
+            ),
+        )
+    except Exception as exc:
+        logger.error("Failed to send ROI failure alert: %s", exc)
+    finally:
+        if publisher is not None:
+            publisher.close()
+
+
 def sync_enabled_delivery_templates() -> int:
     """Keep enabled DB templates aligned with the code's global delivery window."""
     serialized = json.dumps(TIME_SERIES_DEFAULT)
@@ -72,12 +97,22 @@ def _run_script(
             logger.warning("Skip %s: another instance holds %s", script_name, lock_name)
             return
         logger.info("Starting %s", script_name)
-        completed = subprocess.run(
-            [sys.executable, str(HERE / script_name), *(extra_args or [])],
-            cwd=HERE,
-            check=False,
-        )
+        try:
+            completed = subprocess.run(
+                [sys.executable, str(HERE / script_name), *(extra_args or [])],
+                cwd=HERE,
+                check=False,
+            )
+        except Exception as exc:
+            if script_name == "run_daily_roi.py":
+                _notify_roi_failure(extra_args=extra_args, error=str(exc))
+            raise
         if completed.returncode:
+            if script_name == "run_daily_roi.py":
+                _notify_roi_failure(
+                    extra_args=extra_args,
+                    error=f"process exited with code {completed.returncode}",
+                )
             raise RuntimeError(
                 f"{script_name} exited with code {completed.returncode}"
             )
@@ -106,8 +141,25 @@ def run_daily_roi() -> None:
     )
 
 
+def run_internal_roi_test() -> None:
+    _run_script(
+        "run_daily_roi.py",
+        os.getenv("DAILY_ROI_INTERNAL_TEST_LOCK_NAME", "ad_daily_roi"),
+        [
+            "--internal-test",
+            "--output-dir",
+            str(HERE / "outputs" / "roi_control" / "internal_test"),
+        ],
+    )
+
+
 def main() -> None:
     roi_config = RoiConfig.from_env()
+    agency_config = AgencyWebhookConfig.from_env()
+    if roi_config.internal_test_enabled and not (
+        agency_config.webhooks or {}
+    ).get("内部"):
+        raise ValueError("ROI internal test requires webhook route: 内部")
     initialize_schema()
     if _env_flag("DAILY_SYNC_DELIVERY_TEMPLATE", False):
         changed = sync_enabled_delivery_templates()
@@ -161,6 +213,22 @@ def main() -> None:
                 os.getenv("DAILY_ROI_MISFIRE_GRACE_SECONDS", "3600")
             ),
         )
+    if roi_config.internal_test_enabled:
+        scheduler.add_job(
+            run_internal_roi_test,
+            CronTrigger(
+                hour=roi_config.internal_test_hour,
+                minute=roi_config.internal_test_minute,
+                timezone="Asia/Shanghai",
+            ),
+            id="daily_roi_internal_test",
+            name="日级ROI内部端到端测试",
+            max_instances=1,
+            coalesce=True,
+            misfire_grace_time=int(
+                os.getenv("ROI_INTERNAL_TEST_MISFIRE_GRACE_SECONDS", "3600")
+            ),
+        )
     if _env_flag("DAILY_RUN_ON_STARTUP", False):
         scheduler.add_job(
             run_creation,

+ 153 - 2
examples/auto_put_ad_mini/test_roi_agency_delivery.py

@@ -1,7 +1,10 @@
 import json
+import logging
 import os
 import tempfile
 import unittest
+from contextlib import contextmanager
+from datetime import datetime
 from pathlib import Path
 from unittest.mock import Mock, patch
 
@@ -12,9 +15,14 @@ from roi_control.agency_delivery import (
     agency_route_key,
     publish_agency_reports,
 )
-from roi_control.config import AgencyWebhookConfig
+from roi_control.config import AgencyWebhookConfig, RoiConfig
 from roi_control.feishu import RoiFeishuPublisher
-from roi_control.service import _redact_agency_webhooks
+from roi_control.service import (
+    _build_internal_reports,
+    _internal_source_revision,
+    _internal_webhook_config,
+    _redact_agency_webhooks,
+)
 
 
 WEBHOOK = "https://open.feishu.cn/open-apis/bot/v2/hook/test-route"
@@ -39,6 +47,68 @@ class FakeNotifier:
 
 
 class AgencyDeliveryTest(unittest.TestCase):
+    def test_formal_and_internal_schedule_defaults(self):
+        with patch.dict(os.environ, {}, clear=True):
+            config = RoiConfig.from_env()
+        self.assertEqual((config.report_hour, config.report_minute), (9, 0))
+        self.assertEqual(
+            (config.internal_test_hour, config.internal_test_minute),
+            (14, 30),
+        )
+        self.assertTrue(config.internal_test_dedup_enabled)
+
+    def test_internal_test_dedup_switch_controls_run_revision(self):
+        current = datetime(2026, 8, 4, 14, 0, 1, 123456)
+        dedup = RoiConfig(internal_test_dedup_enabled=True)
+        no_dedup = RoiConfig(internal_test_dedup_enabled=False)
+        self.assertEqual(_internal_source_revision(dedup, current), "internal_test")
+        self.assertEqual(
+            _internal_source_revision(no_dedup, current),
+            "it_140001123456",
+        )
+
+    def test_internal_report_set_contains_only_automation_report(self):
+        reports = _build_internal_reports(
+            agency_reports=[
+                {
+                    "agency_name": "小程序-代投-棱镜",
+                    "report_version": "v7",
+                    "report": "mirror.xlsx",
+                    "creative_rows": 4,
+                    "ad_rows": 2,
+                },
+                {
+                    "agency_name": "自动化投放",
+                    "report_version": "v7",
+                    "report": "auto.xlsx",
+                    "creative_rows": 1,
+                    "ad_rows": 1,
+                },
+            ],
+        )
+        self.assertEqual(
+            [row["agency_name"] for row in reports],
+            ["自动化投放"],
+        )
+        self.assertTrue(all(str(row["title"]).startswith("内部测试_") for row in reports))
+        config = _internal_webhook_config(reports, WEBHOOK)
+        self.assertEqual(set(config.webhooks.values()), {WEBHOOK})
+        self.assertEqual(set(config.webhooks), {"自动化投放"})
+
+    def test_internal_report_set_rejects_missing_automation_report(self):
+        with self.assertRaisesRegex(RuntimeError, "自动化投放"):
+            _build_internal_reports(
+                agency_reports=[
+                    {
+                        "agency_name": "小程序-代投-棱镜",
+                        "report_version": "v7",
+                        "report": "mirror.xlsx",
+                        "creative_rows": 4,
+                        "ad_rows": 2,
+                    }
+                ]
+            )
+
     def test_config_snapshot_never_contains_webhook_secret(self):
         env = {
             "ROI_AGENCY_WEBHOOK_ENABLED": "1",
@@ -220,6 +290,87 @@ class AgencyDeliveryTest(unittest.TestCase):
         self.assertEqual(result["message_id"], "message-1")
         self.assertEqual(result["sheet_token"], "main-token")
 
+    def test_service_failure_alert_uses_explicit_chat_and_red_card(self):
+        publisher = RoiFeishuPublisher.__new__(RoiFeishuPublisher)
+        publisher.client = Mock()
+        publisher._token = Mock(return_value="tenant-token")
+        publisher._json = Mock(return_value={"data": {"message_id": "alert-1"}})
+
+        message_id = publisher.send_service_alert(
+            title="日级ROI任务失败",
+            content="任务执行失败",
+            chat_id="chat-internal",
+        )
+
+        self.assertEqual(message_id, "alert-1")
+        request = publisher.client.post.call_args
+        self.assertEqual(request.kwargs["json"]["receive_id"], "chat-internal")
+        self.assertEqual(request.kwargs["json"]["msg_type"], "interactive")
+        card = json.loads(request.kwargs["json"]["content"])
+        self.assertEqual(card["header"]["template"], "red")
+        self.assertEqual(card["header"]["title"]["content"], "日级ROI任务失败")
+
+    def test_service_failure_alert_falls_back_to_operator_chat(self):
+        publisher = RoiFeishuPublisher.__new__(RoiFeishuPublisher)
+        publisher.client = Mock()
+        publisher._token = Mock(return_value="tenant-token")
+        publisher._json = Mock(return_value={"data": {"message_id": "alert-2"}})
+
+        with patch.dict(
+            os.environ,
+            {
+                "ROI_FAILURE_FEISHU_CHAT_ID": "",
+                "FEISHU_OPERATOR_CHAT_ID": "chat-operator",
+            },
+            clear=False,
+        ):
+            message_id = publisher.send_service_alert(
+                title="日级ROI任务失败",
+                content="任务执行失败",
+            )
+
+        self.assertEqual(message_id, "alert-2")
+        request = publisher.client.post.call_args
+        self.assertEqual(request.kwargs["json"]["receive_id"], "chat-operator")
+
+    def test_daily_roi_nonzero_exit_triggers_failure_alert(self):
+        previous_logging_disable = logging.root.manager.disable
+        logging.disable(logging.CRITICAL)
+        try:
+            import run_daily_service as daily_service
+        finally:
+            logging.disable(previous_logging_disable)
+
+        @contextmanager
+        def acquired_lock():
+            yield True
+
+        internal_args = ["--internal-test"]
+        with (
+            patch.object(
+                daily_service,
+                "advisory_lock",
+                return_value=acquired_lock(),
+            ),
+            patch.object(
+                daily_service.subprocess,
+                "run",
+                return_value=Mock(returncode=7),
+            ),
+            patch.object(daily_service, "_notify_roi_failure") as notify,
+        ):
+            with self.assertRaisesRegex(RuntimeError, "exited with code 7"):
+                daily_service._run_script(
+                    "run_daily_roi.py",
+                    "test-roi-lock",
+                    internal_args,
+                )
+
+        notify.assert_called_once_with(
+            extra_args=internal_args,
+            error="process exited with code 7",
+        )
+
 
 if __name__ == "__main__":
     unittest.main()

+ 59 - 1
examples/auto_put_ad_mini/test_roi_control_metrics.py

@@ -1,11 +1,19 @@
 import tempfile
 import unittest
+from datetime import datetime
 from pathlib import Path
+from types import SimpleNamespace
 
 import pandas as pd
 from openpyxl import load_workbook
 
-from roi_control.data_source import build_daily_sql, date_window
+from roi_control.data_source import (
+    SourceDataNotReadyError,
+    build_daily_sql,
+    build_source_readiness_sql,
+    date_window,
+    resolve_end_date,
+)
 from roi_control.fission_multiplier import (
     DISPLAY_MULTIPLIER_COLUMN,
     load_fission_multiplier_parameters,
@@ -82,6 +90,46 @@ def row(entity_type, entity_id, dt, roi, *, uv=600, cost=200.0):
 
 
 class RoiThreeDayRulesTest(unittest.TestCase):
+    @staticmethod
+    def source_client(latest_dt, *, row_count=100, self_rows=60, gzh_rows=40):
+        partition = SimpleNamespace(partition_spec={"dt": latest_dt})
+        table = SimpleNamespace(get_max_partition=lambda: partition)
+        odps = SimpleNamespace(get_table=lambda _name: table)
+        return SimpleNamespace(
+            odps=odps,
+            execute_sql=lambda _sql: pd.DataFrame(
+                [
+                    {
+                        "row_count": row_count,
+                        "self_rows": self_rows,
+                        "gzh_rows": gzh_rows,
+                    }
+                ]
+            ),
+        )
+
+    def test_source_readiness_requires_exact_t_minus_one(self):
+        client = self.source_client("20260803")
+        end_date = resolve_end_date(
+            client,
+            now=datetime(2026, 8, 4, 9, 0),
+        )
+        self.assertEqual(end_date, "20260803")
+        sql = build_source_readiness_sql(end_date)
+        self.assertIn("dt = '20260803'", sql)
+        self.assertIn(SELF_CHANNEL, sql)
+        self.assertIn(GZH_CHANNEL, sql)
+
+    def test_source_readiness_never_falls_back_to_older_partition(self):
+        client = self.source_client("20260802")
+        with self.assertRaisesRegex(SourceDataNotReadyError, "required_dt=20260803"):
+            resolve_end_date(client, now=datetime(2026, 8, 4, 9, 0))
+
+    def test_source_readiness_rejects_missing_report_channel(self):
+        client = self.source_client("20260803", gzh_rows=0)
+        with self.assertRaisesRegex(SourceDataNotReadyError, "gzh_rows"):
+            resolve_end_date(client, now=datetime(2026, 8, 4, 9, 0))
+
     def build_daily(self):
         rows = []
         for dt in DATES:
@@ -306,6 +354,16 @@ class RoiThreeDayRulesTest(unittest.TestCase):
                 Path(bay["report"]).name,
                 "20260803_小程序-代投-贝湉_调控建议.xlsx",
             )
+            filtered = write_agency_workbooks(
+                agency_rows,
+                Path(directory) / "filtered",
+                "20260803",
+                agency_names={"代理B"},
+            )
+            self.assertEqual(
+                [row["agency_name"] for row in filtered],
+                ["代理B"],
+            )
 
             forbidden_fragments = (
                 "ROI",