فهرست منبع

feat(roi): evaluate weighted daily p25 over two days

刘立冬 4 روز پیش
والد
کامیت
d75533b480

+ 3 - 3
AGENTS.md

@@ -41,17 +41,17 @@
 
 - ROI 是指导渠道、账户、广告和创意动作的核心北极星指标,必须作为独立领域模块维护,不能散落在飞书、腾讯 API 或调度代码中。
 - ROI 指标计算必须保持纯数据输入/输出,不能依赖飞书审批、腾讯写操作或具体动作执行器。动作策略只能消费一个明确的 `metric_version`。
-- 当前日级 ROI 指标版本为 `north_star_roi_t15_v5`,报表版本为 `roi_report_v9`,使用已发布参数 `20260712_A0-A15_v2`。首层效率收入取 `SUM(效率收入)`,T0 实际裂变收入取当天的 `SUM(裂变效率收入)`;只预测 T1-T15 增量:`T0实际裂变收入 * (传播裂变系数 - 1)`,预测总收入为 `首层效率收入 + T0实际裂变收入 + 预测T1-T15裂变收入`,禁止继续读取已废弃的 `多层t0裂变收入`。
+- 当前日级 ROI 指标版本为 `north_star_roi_t15_v7`,策略版本为 `roi_policy_v6`,报表版本为 `roi_report_v15`,使用已发布参数 `20260712_A0-A15_v2`。T0 裂变人数取 `SUM(t0_fission_uv_root)`,不得使用旧字段 `t0裂变人数`。首层效率收入取 `SUM(效率收入)`,T0 实际裂变收入取当天的 `SUM(裂变效率收入)`;只预测 T1-T15 增量:`T0实际裂变收入 * (传播裂变系数 - 1)`,预测总收入为 `首层效率收入 + T0实际裂变收入 + 预测T1-T15裂变收入`,禁止继续读取已废弃的 `多层t0裂变收入`。
 - 小程序传播裂变系数按 `人群包+转化目标精确 -> 人群包回退 -> 转化目标回退 -> 渠道回退` 匹配;公众号按 `合作方+公众号精确 -> 合作方回退 -> 渠道回退` 匹配。样本不足的精确实体不能使用自身系数。
 - 小程序日级 ODPS 数据必须保留 `广告优化目标`;缺失目标只能进入人群包回退,不得默认成关键页面访问。企微 ROI 口径未复核,优先使用已发布的合作方 SQL 计算结果,合作方未匹配、样本不足或旧版本无企微参数时默认参考系数为 2.5 并明确标注;企微不进入阈值样本池且不得生成调控动作。
-- ROI 审批表金额和 UV 默认展示三日窗口的单日均值,UV 显示为整数;预测 ROI、实际 ROI、裂变率和阈值必须继续使用三日汇总后的加权口径,不能对每日比例做算术平均。原始三日汇总和每日明细字段保留为隐藏审计列
+- ROI 报表使用长表结构,同一实体每天一行并用 `dt` 区分,每行展示当日首层 UV、T0 裂变人数、T0 裂变率、首层/T0/预测总效率收入、成本、效率 ROI、两个裂变系数、消耗加权分位和 P25 关停线;动作、审批和执行状态只能出现在最新日期行,历史日期行不得保留动作幂等键。两日汇总 ROI 只保留在隐藏审计列,不作为可见决策字段。每日 ROI 必须按各日收入/成本独立计算,不能对每日比例做算术平均
 - ROI 审批表应展示所有进入阈值样本池的小程序和公众号实体,以及全部企微参考实体;不能只展示有动作的数据。每个渠道 Sheet 默认按预测 ROI 升序排列。
 - 成熟参数是独立版本化只读资产。审计后的参数必须通过更新脚本发布到 MySQL 的 `roi_fission_parameter_release` 和 `roi_fission_parameter_value`;日级 ROI 只从数据库消费 `ROI_FISSION_PARAMETER_VERSION` 指定且通过内容哈希校验的版本,不能在日常任务中现场重算或静默回退到镜像文件。离线复算只能生成草稿,经恒等式、行数、匹配率和新旧结果对账后才能发布。
 - ROI 报表只上传一次,同一在线表链接发送到 ROI 通知群和 `FEISHU_OPERATOR_CHAT_ID` 投放审批群;群 ID 相同时必须去重。
 - ROI 公式、收入/成本归属、裂变口径或实体粒度发生语义变化时,必须升级指标版本并保存新快照;不得覆盖或重算成旧版本历史结果。
 - 策略阈值和动作语义使用独立 `policy_version`;指标版本与策略版本必须同时写入每个运行批次。
 - 每次运行必须保存完整实体快照、阈值配置、动作建议、可执行性原因和最终执行审计,不能只保存候选或飞书表格。
-- 日级 ROI 当前读取 T-1 至 T-3 的连续三日数据。所有渠道都进入指标和通知,只有当前自动化腾讯账户的可执行行允许在表格逐行批准后执行。
+- 日级 ROI 当前读取 T-1 至 T-2 的连续两日数据。小程序与公众号按渠道、按日期分别计算阈值:关停线使用成本权重按 P95 封顶后的消耗加权 P25,扩量线保持实体等权 P80。只有连续两天单日首层 UV 都大于 200 且连续处于同一极端方向时才生成正式建议。最新日首层 UV 大于 100 但不满足连续两日样本条件的实体只标记为“观察”,不得生成腾讯写操作。所有渠道都进入指标和通知,只有当前自动化腾讯账户的可执行行允许在表格逐行批准后执行。
 - 当前低 ROI 动作只暂停 `dynamic_creative_id`,不能暂停整个广告;高 ROI 动作只调整广告永久基础出价,默认上调 10%。
 - 同广告永久基础出价 3 天内最多上调一次,且不超过首次纳管基础价的 2 倍。调整后必须同步实时 CPM 模块的基础出价,避免恢复旧值。
 - ROI 逐行审批有效期默认 120 分钟。腾讯写操作必须与实时调控共用数据库 advisory lock,执行前回读映射和状态,执行后再次回读校验;执行结果必须回写表格并发送飞书通知,通知失败只能重试通知,不能重复腾讯写操作。

+ 10 - 6
examples/auto_put_ad_mini/docs/unified_services_deployment.md

@@ -165,7 +165,7 @@ docker compose --env-file /dev/null run --rm \
 
 日级 ROI 单独分两阶段启用:
 
-当前 `north_star_roi_t15_v5` 使用版本化传播裂变参数,报表版本为 `roi_report_v9`。小程序按人群包和转化目标、
+当前 `north_star_roi_t15_v7` 使用版本化传播裂变参数,策略版本为 `roi_policy_v6`,报表版本为 `roi_report_v15`。T0 裂变人数取 `SUM(t0_fission_uv_root)`,不再读取旧字段 `t0裂变人数`。小程序按人群包和转化目标、
 公众号按合作方和公众号匹配传播裂变系数;企微优先按合作方精确匹配,合作方未匹配或样本不足时使用默认参考系数 2.5,不做渠道聚合回退。企微仅保留展示,不进入阈值和调控。首层效率收入读取 `效率收入`,T0 实际裂变收入读取当天的 `裂变效率收入`,
 并以 T0 实际裂变收入乘传播裂变系数预测完整裂变收入。日级任务不会现场重算 cohort 参数。部署前应保持
 `ROI_FISSION_PARAMETER_VERSION=20260712_A0-A15_v2`。参数必须先发布到 MySQL,
@@ -177,11 +177,15 @@ ROI 在线表只上传一次,并发送到 `ROI_FEISHU_CHAT_ID` 和
 ROI 表格使用获得链接者可编辑权限。黄色【审批选择】列只接受“批准”或“拒绝”,
 数据库隐藏幂等键决定真实执行目标,表格中的账户、广告、创意、成本和 ROI 不作为写入参数。
 
-审批表默认展示三日窗口的日均首层 UV、日均首层效率收入、日均 LTV 预测效率收入和
-日均成本;预测 ROI 仍按三日预测收入总和除以三日成本总和计算。三日汇总和每日明细
-保留为隐藏审计列,不参与日均展示。
-小程序和公众号 Sheet 展示所有进入阈值样本池的实体,企微 Sheet 展示全部参考实体;
-各 Sheet 默认按预测 ROI 升序排列,动作列不参与排序。
+审批表使用长表结构,同一实体每天一行并用 `dt` 区分,每行展示当日首层 UV、T0 裂变人数、
+T0 裂变率、首层/T0/预测总效率收入、成本、效率 ROI、两个裂变系数、消耗加权分位和
+P25 关停线;动作、审批和执行状态仅出现在最新日期行,
+历史日期行不得保留动作幂等键。两日汇总 ROI 仅保留为隐藏审计列。小程序与公众号分别按日期计算
+阈值,关停线使用成本权重按 P95 封顶后的消耗
+加权 P25,扩量线保持实体等权 P80;连续两天单日首层 UV 都大于 200 且
+连续处于同一极端方向时才生成正式建议。最新日首层 UV 大于 100 但未满足连续两日
+条件的实体作为“观察”放在正式样本之后,不进入阈值和腾讯执行。
+企微 Sheet 展示全部参考实体;各 Sheet 内正式样本按预测 ROI 升序排列。
 
 1. 首次部署先执行下方参数发布命令并完成回读校验。
 2. 设置 `DAILY_ROI_ENABLED=1`、`ROI_APPLY_ENABLED=0`、`ROI_SHEET_APPROVAL_ENABLED=0`,观察 11:00 的 ODPS 计算、数据库快照和可编辑飞书表。

+ 32 - 16
examples/auto_put_ad_mini/roi_control/config.py

@@ -28,13 +28,12 @@ class RoiConfig:
     max_base_ratio: Decimal = Decimal("2.00")
     self_stop_min_age: int = 5
     self_up_min_age: int = 3
-    self_min_avg_uv: float = 200
-    partner_min_avg_uv: float = 200
-    min_daily_cost: float = 100
-    stop_quantile: float = 0.20
+    self_min_daily_uv: float = 200
+    partner_min_daily_uv: float = 200
+    observe_min_latest_uv: float = 100
+    stop_quantile: float = 0.25
     up_quantile: float = 0.80
-    gzh_adjust_rank_min: float = 0.10
-    gzh_adjust_rank_max: float = 0.30
+    stop_weight_cap_quantile: float = 0.95
 
     @classmethod
     def from_env(cls) -> "RoiConfig":
@@ -55,18 +54,25 @@ class RoiConfig:
             max_base_ratio=Decimal(os.getenv("ROI_MAX_BASE_RATIO", "2.00")),
             self_stop_min_age=int(os.getenv("ROI_SELF_STOP_MIN_AGE", "5")),
             self_up_min_age=int(os.getenv("ROI_SELF_UP_MIN_AGE", "3")),
-            self_min_avg_uv=float(os.getenv("ROI_SELF_MIN_AVG_UV", "200")),
-            partner_min_avg_uv=float(
-                os.getenv("ROI_PARTNER_MIN_AVG_UV", "200")
+            self_min_daily_uv=float(
+                os.getenv(
+                    "ROI_SELF_MIN_DAILY_UV",
+                    os.getenv("ROI_SELF_MIN_AVG_UV", "200"),
+                )
             ),
-            min_daily_cost=float(os.getenv("ROI_MIN_DAILY_COST", "100")),
-            stop_quantile=float(os.getenv("ROI_STOP_QUANTILE", "0.20")),
-            up_quantile=float(os.getenv("ROI_UP_QUANTILE", "0.80")),
-            gzh_adjust_rank_min=float(
-                os.getenv("ROI_GZH_ADJUST_RANK_MIN", "0.10")
+            partner_min_daily_uv=float(
+                os.getenv(
+                    "ROI_PARTNER_MIN_DAILY_UV",
+                    os.getenv("ROI_PARTNER_MIN_AVG_UV", "200"),
+                )
+            ),
+            observe_min_latest_uv=float(
+                os.getenv("ROI_OBSERVE_MIN_LATEST_UV", "100")
             ),
-            gzh_adjust_rank_max=float(
-                os.getenv("ROI_GZH_ADJUST_RANK_MAX", "0.30")
+            stop_quantile=float(os.getenv("ROI_STOP_QUANTILE", "0.25")),
+            up_quantile=float(os.getenv("ROI_UP_QUANTILE", "0.80")),
+            stop_weight_cap_quantile=float(
+                os.getenv("ROI_STOP_WEIGHT_CAP_QUANTILE", "0.95")
             ),
         )
         config.validate()
@@ -93,8 +99,18 @@ class RoiConfig:
             raise ValueError("ROI_MAX_BASE_RATIO must cover ROI_SCALE_RATIO")
         if self.scale_cooldown_days < 1:
             raise ValueError("ROI_SCALE_COOLDOWN_DAYS must be positive")
+        if min(
+            self.self_min_daily_uv,
+            self.partner_min_daily_uv,
+            self.observe_min_latest_uv,
+        ) <= 0:
+            raise ValueError("ROI UV thresholds must be positive")
         if not 0 < self.stop_quantile < self.up_quantile < 1:
             raise ValueError("ROI quantiles must satisfy 0 < stop < up < 1")
+        if not 0 < self.stop_weight_cap_quantile <= 1:
+            raise ValueError(
+                "ROI_STOP_WEIGHT_CAP_QUANTILE must satisfy 0 < value <= 1"
+            )
 
     def snapshot(self) -> dict[str, object]:
         values = asdict(self)

+ 4 - 4
examples/auto_put_ad_mini/roi_control/data_source.py

@@ -25,7 +25,7 @@ def parse_yyyymmdd(value: str) -> datetime:
 
 def date_window(end_date: str) -> Tuple[str, str]:
     end = parse_yyyymmdd(end_date)
-    return (end - timedelta(days=2)).strftime("%Y%m%d"), end.strftime("%Y%m%d")
+    return (end - timedelta(days=1)).strftime("%Y%m%d"), end.strftime("%Y%m%d")
 
 
 def resolve_end_date(client: ODPSClient, requested: str | None = None) -> str:
@@ -75,7 +75,7 @@ SELECT
   '' AS 合作方名,
   '' AS 公众号名,
   COUNT(DISTINCT mid) AS 首层UV,
-  SUM(NVL(t0裂变人数, 0)) AS T0裂变数,
+  SUM(NVL(t0_fission_uv_root, 0)) AS T0裂变数,
   SUM(NVL(成本, 0)) AS 成本,
   SUM(NVL(效率收入, 0)) AS 效率收入,
   SUM(NVL(裂变效率收入, 0)) AS 裂变效率收入
@@ -104,7 +104,7 @@ SELECT
   NVL(合作方名, '') AS 合作方名,
   NVL(公众号名, '') AS 公众号名,
   COUNT(DISTINCT mid) AS 首层UV,
-  SUM(NVL(t0裂变人数, 0)) AS T0裂变数,
+  SUM(NVL(t0_fission_uv_root, 0)) AS T0裂变数,
   SUM(NVL(成本, 0)) AS 成本,
   SUM(NVL(效率收入, 0)) AS 效率收入,
   SUM(NVL(裂变效率收入, 0)) AS 裂变效率收入
@@ -131,7 +131,7 @@ SELECT
   NVL(合作方名, '') AS 合作方名,
   MAX(NVL(公众号名, '')) AS 公众号名,
   COUNT(DISTINCT mid) AS 首层UV,
-  SUM(NVL(t0裂变人数, 0)) AS T0裂变数,
+  SUM(NVL(t0_fission_uv_root, 0)) AS T0裂变数,
   SUM(NVL(成本, 0)) AS 成本,
   SUM(NVL(效率收入, 0)) AS 效率收入,
   SUM(NVL(裂变效率收入, 0)) AS 裂变效率收入

+ 1 - 1
examples/auto_put_ad_mini/roi_control/feishu.py

@@ -208,7 +208,7 @@ class RoiFeishuPublisher:
                 details.append(
                     f"- 账户 `{row['account_id']}` / 广告 `{row['adgroup_id']}` / "
                     f"创意 `{row.get('dynamic_creative_id') or '-'}`\n"
-                    f"  {row.get('adgroup_name') or ''} | 三日成本 {cost:.2f} 元 | "
+                    f"  {row.get('adgroup_name') or ''} | 窗口成本 {cost:.2f} 元 | "
                     f"预测ROI {roi:.3f} | **{result}**"
                     + (f" | {error[:160]}" if error else "")
                 )

+ 34 - 20
examples/auto_put_ad_mini/roi_control/metrics.py

@@ -19,8 +19,8 @@ from .fission_multiplier import (
 )
 
 
-METRIC_VERSION = "north_star_roi_t15_v5"
-METRIC_RUN_SUFFIX = "m5"
+METRIC_VERSION = "north_star_roi_t15_v7"
+METRIC_RUN_SUFFIX = "m7"
 
 SELF_CHANNEL = "小程序投流-稳定"
 GZH_CHANNEL = "公众号合作-即转-稳定"
@@ -137,13 +137,13 @@ def _age_map(ad_age: pd.DataFrame | None) -> Dict[str, int]:
     return dict(zip(ages["广告id"], ages["广告age"]))
 
 
-def _expected_three_days(expected_dates: Iterable[str]) -> List[str]:
+def _expected_two_days(expected_dates: Iterable[str]) -> List[str]:
     dates = sorted({str(value) for value in expected_dates})
-    if len(dates) != 3:
-        raise ValueError(f"必须提供连续个日期,实际为: {dates}")
+    if len(dates) != 2:
+        raise ValueError(f"必须提供连续个日期,实际为: {dates}")
     parsed = pd.to_datetime(dates, format="%Y%m%d")
     gaps = parsed.to_series().diff().dropna().dt.days.tolist()
-    if gaps != [1, 1]:
+    if gaps != [1]:
         raise ValueError(f"日期必须连续,实际为: {dates}")
     return dates
 
@@ -183,8 +183,9 @@ def _summarize_entity(
         ]
         if inconsistent:
             raise ValueError(
-                f"同一实体日内传播裂变参数不一致: {key_values}; fields={inconsistent}"
+                f"同一实体日内传播裂变参数不一致: {key_values}; fields={inconsistent}"
             )
+        observed_dates = set(group["dt"].astype(str))
         by_date = group.groupby("dt", as_index=False)[
             [
                 "首层UV",
@@ -199,10 +200,10 @@ def _summarize_entity(
                 "预测全链路效率收入",
             ]
         ].sum()
-        if set(by_date["dt"]) != set(expected_dates):
-            continue
-
-        by_date = by_date.set_index("dt").loc[list(expected_dates)]
+        by_date = (
+            by_date.set_index("dt")
+            .reindex(list(expected_dates), fill_value=0)
+        )
         daily_actual_roi = np.where(
             by_date["成本"] > 0,
             by_date["实际全链路效率收入"] / by_date["成本"],
@@ -224,17 +225,19 @@ def _summarize_entity(
                 "",
             )
         record["entity_type"] = entity_type
-        record["覆盖天数"] = 3
-        record["三日首层UV"] = float(by_date["首层UV"].sum())
+        record["覆盖天数"] = len(observed_dates.intersection(expected_dates))
+        record["窗口首层UV"] = float(by_date["首层UV"].sum())
         record["日均首层UV"] = float(by_date["首层UV"].mean())
-        record["三日T0裂变数"] = float(by_date["T0裂变数"].sum())
+        record["窗口最小单日首层UV"] = float(by_date["首层UV"].min())
+        record["最新日首层UV"] = float(by_date.iloc[-1]["首层UV"])
+        record["窗口T0裂变数"] = float(by_date["T0裂变数"].sum())
         record["T0裂变率"] = (
-            record["三日T0裂变数"] / record["三日首层UV"]
-            if record["三日首层UV"] > 0
+            record["窗口T0裂变数"] / record["窗口首层UV"]
+            if record["窗口首层UV"] > 0
             else np.nan
         )
         record["成本"] = float(by_date["成本"].sum())
-        record["三日最小单日成本"] = float(by_date["成本"].min())
+        record["窗口最小单日成本"] = float(by_date["成本"].min())
         record["效率收入"] = float(by_date["效率收入"].sum())
         record["裂变效率收入"] = float(
             by_date["裂变效率收入"].sum()
@@ -274,6 +277,17 @@ def _summarize_entity(
             if record["成本"] > 0
             else np.nan
         )
+        record["前一日首层UV"] = float(by_date.iloc[0]["首层UV"])
+        record["前一日效率ROI"] = (
+            float(daily_predicted_roi[0])
+            if np.isfinite(daily_predicted_roi[0])
+            else np.nan
+        )
+        record["最新日效率ROI"] = (
+            float(daily_predicted_roi[-1])
+            if np.isfinite(daily_predicted_roi[-1])
+            else np.nan
+        )
         for index, dt in enumerate(expected_dates):
             actual_roi = (
                 float(daily_actual_roi[index])
@@ -309,9 +323,9 @@ def compute_roi_summary(
     *,
     fission_parameters: FissionMultiplierParameters,
 ) -> tuple[pd.DataFrame, list[str]]:
-    """Return complete three-day entity ROI snapshots and normalized dates."""
+    """Return two-day entity ROI snapshots and normalized dates."""
 
-    dates = _expected_three_days(expected_dates)
+    dates = _expected_two_days(expected_dates)
     daily = prepare_daily_metrics(raw_daily, fission_parameters)
     daily = daily[daily["dt"].isin(dates)].copy()
     ages = _age_map(ad_age)
@@ -330,7 +344,7 @@ def compute_roi_summary(
             summaries.append(entity_summary)
 
     if not summaries:
-        raise ValueError("三个目标渠道均没有可用的连续三日数据")
+        raise ValueError("三个目标渠道均没有可用的最近两日数据")
 
     summary = pd.concat(summaries, ignore_index=True, sort=False)
     return summary, dates

+ 5 - 3
examples/auto_put_ad_mini/roi_control/policy.py

@@ -78,7 +78,9 @@ def annotate_execution(
         action_type: str | None = None
 
         if recommendation:
-            if entity_type != ENTITY_SELF:
+            if recommendation == "观察":
+                execution_reason = "观察行仅展示,不执行腾讯写操作"
+            elif entity_type != ENTITY_SELF:
                 execution_reason = "合作渠道当前仅通知,不执行腾讯写操作"
             elif account_id not in managed_account_ids:
                 execution_reason = "账户不在自动化管理范围,仅通知"
@@ -190,8 +192,8 @@ def annotate_execution(
                 "fission_source": str(row.get("传播裂变系数来源") or ""),
                 "ad_age": safe_int(row.get("广告age")),
                 "avg_first_uv": _finite(row.get("日均首层UV")),
-                "first_uv": _finite(row.get("三日首层UV")),
-                "t0_fission_count": _finite(row.get("三日T0裂变数")),
+                "first_uv": _finite(row.get("窗口首层UV")),
+                "t0_fission_count": _finite(row.get("窗口T0裂变数")),
                 "t0_fission_rate": _finite(row.get("T0裂变率")),
                 "cost": _finite(row.get("成本")),
                 "efficiency_revenue": _finite(row.get("效率收入")),

+ 213 - 63
examples/auto_put_ad_mini/roi_control/reporting.py

@@ -24,10 +24,11 @@ HEADER_FONT = Font(color="FFFFFF", bold=True)
 STOP_FILL = PatternFill("solid", fgColor="F4CCCC")
 UP_FILL = PatternFill("solid", fgColor="D9EAD3")
 ADJUST_FILL = PatternFill("solid", fgColor="FFF2CC")
+OBSERVE_FILL = PatternFill("solid", fgColor="FFF2CC")
 APPROVAL_FILL = PatternFill("solid", fgColor="FFD966")
 APPROVAL_HEADER_FILL = PatternFill("solid", fgColor="BF9000")
-REPORT_VERSION = "roi_report_v9"
-REPORT_RUN_SUFFIX = "r9"
+REPORT_VERSION = "roi_report_v15"
+REPORT_RUN_SUFFIX = "r15"
 
 T0_FISSION_MULTIPLIER_COLUMN = "裂变系数-总裂变UV/T0裂变UV"
 TOTAL_FISSION_TO_FIRST_UV_COLUMN = DISPLAY_TOTAL_TO_FIRST_COLUMN
@@ -36,6 +37,7 @@ FINAL_ROI_COLUMN = "最终效率ROI"
 
 BASE_COLUMNS: Dict[str, Sequence[str]] = {
     "小程序投流": (
+        "dt",
         "渠道",
         "代理名称",
         "账号id",
@@ -45,53 +47,56 @@ BASE_COLUMNS: Dict[str, Sequence[str]] = {
         "广告优化目标",
         "创意id",
         "广告age",
-        "日均首层UV",
-        "首层效率收入",
-        "T0裂变效率收入",
-        "总预估效率收入",
-        "日均成本",
-        T0_FISSION_MULTIPLIER_COLUMN,
-        TOTAL_FISSION_TO_FIRST_UV_COLUMN,
-        FINAL_ROI_COLUMN,
-        "动作",
-        "关停线",
-        "审批选择",
-        "执行状态",
-        "执行结果",
     ),
     "公众号即转": (
+        "dt",
         "渠道",
         "合作方名",
         "公众号名",
-        "日均首层UV",
-        "首层效率收入",
-        "T0裂变效率收入",
-        "总预估效率收入",
-        "日均成本",
-        T0_FISSION_MULTIPLIER_COLUMN,
-        TOTAL_FISSION_TO_FIRST_UV_COLUMN,
-        FINAL_ROI_COLUMN,
-        "动作",
-        "关停线",
     ),
     "企微群合作": (
+        "dt",
         "渠道",
         "合作方名",
-        "日均首层UV",
-        "首层效率收入",
-        "T0裂变效率收入",
-        "总预估效率收入",
-        "日均成本",
-        T0_FISSION_MULTIPLIER_COLUMN,
-        TOTAL_FISSION_TO_FIRST_UV_COLUMN,
-        "传播裂变参数状态",
-        "调控参与状态",
-        FINAL_ROI_COLUMN,
+    ),
+}
+
+DAILY_COLUMNS = (
+    "首层UV",
+    "T0裂变人数",
+    "T0裂变率",
+    "首层效率收入",
+    "T0裂变效率收入",
+    "预测总效率收入",
+    "成本",
+    "效率ROI",
+    T0_FISSION_MULTIPLIER_COLUMN,
+    TOTAL_FISSION_TO_FIRST_UV_COLUMN,
+    "消耗加权分位",
+)
+
+DECISION_COLUMNS: Dict[str, Sequence[str]] = {
+    "小程序投流": (
         "动作",
-        "关停线",
+        "动作原因",
+        "阈值样本状态",
+        "审批选择",
+        "执行状态",
+        "执行结果",
+    ),
+    "公众号即转": (
+        "动作",
+        "动作原因",
+        "阈值样本状态",
     ),
+    "企微群合作": (),
 }
 
+QIWEI_STATUS_COLUMNS = (
+    "传播裂变参数状态",
+    "调控参与状态",
+)
+
 ENTITY_TO_SHEET = {
     ENTITY_SELF: "小程序投流",
     ENTITY_GZH: "公众号即转",
@@ -99,13 +104,42 @@ ENTITY_TO_SHEET = {
 }
 
 FREEZE_PANES = {
-    "小程序投流": "F2",
-    "公众号即转": "D2",
-    "企微群合作": "C2",
+    "小程序投流": "G2",
+    "公众号即转": "E2",
+    "企微群合作": "D2",
 }
 
 
-def _sheet_frame(candidates: pd.DataFrame, sheet_name: str) -> pd.DataFrame:
+def _visible_columns(
+    sheet_name: str,
+    expected_dates: Sequence[str],
+    stop_quantile: float,
+) -> list[str]:
+    stop_label = f"P{int(stop_quantile * 100)}"
+    threshold_columns = (
+        []
+        if sheet_name == "企微群合作"
+        else [f"消耗加权{stop_label}线"]
+    )
+    status_columns = (
+        list(QIWEI_STATUS_COLUMNS)
+        if sheet_name == "企微群合作"
+        else list(DECISION_COLUMNS[sheet_name])
+    )
+    return (
+        list(BASE_COLUMNS[sheet_name])
+        + list(DAILY_COLUMNS)
+        + threshold_columns
+        + status_columns
+    )
+
+
+def _sheet_frame(
+    candidates: pd.DataFrame,
+    sheet_name: str,
+    expected_dates: Sequence[str] | None = None,
+    stop_quantile: float = 0.25,
+) -> pd.DataFrame:
     entity_type = next(
         key for key, value in ENTITY_TO_SHEET.items() if value == sheet_name
     )
@@ -114,7 +148,7 @@ def _sheet_frame(candidates: pd.DataFrame, sheet_name: str) -> pd.DataFrame:
     coverage_source = (
         subset["覆盖天数"]
         if "覆盖天数" in subset
-        else pd.Series(3, index=subset.index)
+        else pd.Series(2, index=subset.index)
     )
     coverage_days = pd.to_numeric(
         coverage_source, errors="coerce"
@@ -148,7 +182,55 @@ def _sheet_frame(candidates: pd.DataFrame, sheet_name: str) -> pd.DataFrame:
     subset[FINAL_ROI_COLUMN] = subset["ROI"]
     subset["关停线"] = subset["t_stop"]
     subset["扩量线"] = subset["t_up"]
-    visible_columns = list(BASE_COLUMNS[sheet_name])
+    dates = list(expected_dates or sorted(
+        column.removeprefix("ROI_")
+        for column in subset.columns
+        if column.startswith("ROI_") and len(column) == 12
+    ))
+    stop_label = f"P{int(stop_quantile * 100)}"
+    daily_frames = []
+    latest_date = dates[-1] if dates else ""
+    for dt in dates:
+        daily = subset.copy()
+        daily["dt"] = dt
+        daily["首层UV"] = daily.get(f"首层UV_{dt}")
+        daily["T0裂变人数"] = daily.get(f"T0裂变数_{dt}")
+        daily["T0裂变率"] = np.where(
+            pd.to_numeric(daily["首层UV"], errors="coerce").gt(0),
+            pd.to_numeric(daily["T0裂变人数"], errors="coerce")
+            / pd.to_numeric(daily["首层UV"], errors="coerce"),
+            np.nan,
+        )
+        daily["首层效率收入"] = daily.get(f"效率收入_{dt}")
+        daily["T0裂变效率收入"] = daily.get(f"裂变效率收入_{dt}")
+        daily["成本"] = daily.get(f"成本_{dt}")
+        daily["效率ROI"] = daily.get(f"ROI_{dt}")
+        daily["预测总效率收入"] = (
+            pd.to_numeric(daily["效率ROI"], errors="coerce")
+            * pd.to_numeric(daily["成本"], errors="coerce")
+        )
+        daily["消耗加权分位"] = daily.get(f"消耗加权同日分位_{dt}")
+        if sheet_name != "企微群合作":
+            daily[f"消耗加权{stop_label}线"] = daily.get(f"t_stop_{dt}")
+        if dt != latest_date:
+            for column in (
+                "动作",
+                "动作原因",
+                "阈值样本状态",
+                "执行状态",
+                "执行结果",
+                "动作幂等键",
+                "执行模式",
+                "执行说明",
+            ):
+                daily[column] = ""
+            daily["审批选择"] = "历史日"
+        daily_frames.append(daily)
+    if daily_frames:
+        subset = pd.concat(daily_frames, ignore_index=True)
+    else:
+        subset["dt"] = ""
+    visible_columns = _visible_columns(sheet_name, dates, stop_quantile)
     for column in visible_columns:
         if column not in subset:
             subset[column] = ""
@@ -158,7 +240,32 @@ def _sheet_frame(candidates: pd.DataFrame, sheet_name: str) -> pd.DataFrame:
     ordered = visible_columns + hidden_columns
     subset = subset[ordered]
     if not subset.empty:
-        subset = subset.sort_values(FINAL_ROI_COLUMN, ascending=True)
+        subset["_日期排序"] = subset["dt"].map(
+            {dt: position for position, dt in enumerate(reversed(dates))}
+        ).fillna(len(dates))
+        subset["_补充观察排序"] = (
+            subset["阈值样本状态"].eq("观察_最新日UV达标").astype(int)
+        )
+        action_order = {"关停": 0, "扩量": 1, "观察": 2, "": 3}
+        subset["_动作排序"] = (
+            subset["动作"].fillna("").map(action_order).fillna(4)
+        )
+        subset["_当日成本排序"] = pd.to_numeric(
+            subset.get("成本"),
+            errors="coerce",
+        ).fillna(0)
+        subset = subset.sort_values(
+            ["_日期排序", "_补充观察排序", "_动作排序", "_当日成本排序"],
+            ascending=[True, True, True, False],
+            kind="stable",
+        ).drop(
+            columns=[
+                "_日期排序",
+                "_补充观察排序",
+                "_动作排序",
+                "_当日成本排序",
+            ]
+        )
     return subset
 
 
@@ -173,7 +280,11 @@ def _write_dataframe(ws, frame: pd.DataFrame) -> None:
         )
 
 
-def _format_sheet(ws, sheet_name: str) -> None:
+def _format_sheet(
+    ws,
+    sheet_name: str,
+    visible_columns: Sequence[str],
+) -> None:
     max_column = max(ws.max_column, 1)
     max_row = max(ws.max_row, 1)
     ws.freeze_panes = FREEZE_PANES[sheet_name]
@@ -202,7 +313,7 @@ def _format_sheet(ws, sheet_name: str) -> None:
         ws.add_data_validation(validation)
         for row in range(2, max_row + 1):
             cell = ws.cell(row, approval_column)
-            if cell.value != "不可执行":
+            if cell.value in (None, ""):
                 cell.fill = APPROVAL_FILL
                 validation.add(cell)
 
@@ -216,6 +327,8 @@ def _format_sheet(ws, sheet_name: str) -> None:
             "ROI" in str(header)
             or "裂变率" in str(header)
             or "百分位" in str(header)
+            or "分位" in str(header)
+            or "关停线" in str(header)
             or str(header).startswith("裂变系数-")
         ):
             for row in range(2, max_row + 1):
@@ -253,7 +366,7 @@ def _format_sheet(ws, sheet_name: str) -> None:
                 ),
             )
 
-        if column_index > len(BASE_COLUMNS[sheet_name]):
+        if header not in set(visible_columns):
             ws.column_dimensions[get_column_letter(column_index)].hidden = True
 
     action_column = headers.get("动作")
@@ -264,6 +377,7 @@ def _format_sheet(ws, sheet_name: str) -> None:
                 "关停": STOP_FILL,
                 "扩量": UP_FILL,
                 "调整封面&落地页视频": ADJUST_FILL,
+                "观察": OBSERVE_FILL,
             }.get(action)
             if fill:
                 ws.cell(row, action_column).fill = fill
@@ -276,9 +390,14 @@ def _write_summary(
     rule_config: Mapping[str, object] | None,
 ) -> None:
     ws = workbook.create_sheet("运行摘要")
-    self_min_uv = (rule_config or {}).get("self_min_avg_uv", 200)
-    partner_min_uv = (rule_config or {}).get("partner_min_avg_uv", 200)
-    min_daily_cost = (rule_config or {}).get("min_daily_cost", 100)
+    self_min_uv = (rule_config or {}).get("self_min_daily_uv", 200)
+    partner_min_uv = (rule_config or {}).get("partner_min_daily_uv", 200)
+    observe_min_uv = (rule_config or {}).get("observe_min_latest_uv", 100)
+    stop_quantile = float((rule_config or {}).get("stop_quantile", 0.25))
+    up_quantile = float((rule_config or {}).get("up_quantile", 0.80))
+    weight_cap_quantile = float(
+        (rule_config or {}).get("stop_weight_cap_quantile", 0.95)
+    )
     rows = [
         ("统计窗口", f"{expected_dates[0]} 至 {expected_dates[-1]}"),
         ("首层口径", "usersharedepth='0'"),
@@ -292,14 +411,31 @@ def _write_summary(
             "企微按合作方匹配已发布参数;未匹配或样本不足时默认取2.5。"
             "企微仅展示,不进入阈值样本池和调控。",
         ),
-        ("阈值口径", "小程序和公众号合格实体进入同一全局池:三日聚合ROI的P20=关停线、P80=扩量线;企微不参与"),
+        (
+            "阈值口径",
+            "小程序与公众号分渠道、分日期计算阈值;"
+            f"关停线为成本权重按P{int(weight_cap_quantile * 100)}封顶后的"
+            f"消耗加权P{int(stop_quantile * 100)},"
+            f"扩量线保持实体等权P{int(up_quantile * 100)};企微不参与。",
+        ),
         (
             "阈值样本条件",
-            f"连续3天数据完整;小程序日均首层UV>{self_min_uv:g};"
-            f"公众号日均首层UV>{partner_min_uv:g};连续3天每天成本>"
-            f"{min_daily_cost:g};三日成本>0且最终效率ROI有效;企微仅展示、不进入样本池",
+            f"连续2天分别满足:小程序单日首层UV>{self_min_uv:g},"
+            f"公众号单日首层UV>{partner_min_uv:g},且当日成本>0、当日ROI有效。",
+        ),
+        (
+            "动作口径",
+            f"连续2天均处于同渠道当日消耗加权后"
+            f"{int(stop_quantile * 100)}%才建议关停,"
+            f"连续2天均处于实体等权前{int((1 - up_quantile) * 100)}%"
+            "才建议扩量;"
+            "方向不一致或广告age不足时观察。",
+        ),
+        (
+            "观察补充",
+            f"未满足连续2天正式样本条件,但最新日首层UV>{observe_min_uv:g}的实体也进入报表末尾;"
+            "只展示最新日ROI,不进入阈值计算和腾讯执行。",
         ),
-        ("动作优先级", "关停 > 扩量 > 调整封面&落地页视频"),
         ("审批方式", "在小程序投流表黄色【审批选择】列逐行选择批准或拒绝;批准即为最终确认并自动执行"),
     ]
     for label, value in rows:
@@ -311,11 +447,12 @@ def _write_summary(
         for key in (
             "self_stop_min_age",
             "self_up_min_age",
-            "self_min_avg_uv",
-            "partner_min_avg_uv",
-            "min_daily_cost",
+            "self_min_daily_uv",
+            "partner_min_daily_uv",
+            "observe_min_latest_uv",
             "stop_quantile",
             "up_quantile",
+            "stop_weight_cap_quantile",
             "scale_ratio",
             "scale_cooldown_days",
             "max_base_ratio",
@@ -344,13 +481,15 @@ def _write_summary(
         columns={"t_stop": "关停线", "t_up": "扩量线"}
     )
     threshold_columns = [
-        "统计窗口",
+        "统计日期",
+        "渠道",
         "关停线",
         "扩量线",
         "阈值样本数",
-        "小程序样本数",
-        "公众号样本数",
-        "企微参考实体数",
+        "关停线口径",
+        "权重封顶成本",
+        "权重封顶分位",
+        "扩量线口径",
     ]
     ws.append(threshold_columns)
     for values in display_thresholds[threshold_columns].itertuples(
@@ -393,9 +532,20 @@ def write_workbook(
 
     for sheet_name in BASE_COLUMNS:
         ws = workbook.create_sheet(sheet_name)
-        frame = _sheet_frame(candidates, sheet_name)
+        stop_quantile = float((rule_config or {}).get("stop_quantile", 0.25))
+        visible_columns = _visible_columns(
+            sheet_name,
+            expected_dates,
+            stop_quantile,
+        )
+        frame = _sheet_frame(
+            candidates,
+            sheet_name,
+            expected_dates,
+            stop_quantile,
+        )
         _write_dataframe(ws, frame)
-        _format_sheet(ws, sheet_name)
+        _format_sheet(ws, sheet_name, visible_columns)
 
     if fission_match_summary is not None:
         ws = workbook.create_sheet("传播裂变系数匹配")

+ 271 - 103
examples/auto_put_ad_mini/roi_control/rules.py

@@ -1,4 +1,4 @@
-"""Versioned ROI thresholds and action recommendation policy."""
+"""Versioned daily ROI thresholds and two-day action policy."""
 
 from __future__ import annotations
 
@@ -8,159 +8,327 @@ from typing import Iterable
 import numpy as np
 import pandas as pd
 
+from .fission_multiplier import FissionMultiplierParameters
 from .metrics import (
     ENTITY_GZH,
     ENTITY_QIWEI,
     ENTITY_SELF,
     compute_roi_summary,
 )
-from .fission_multiplier import FissionMultiplierParameters
 
 
-POLICY_VERSION = "roi_policy_v2"
-POLICY_RUN_SUFFIX = "p2"
+POLICY_VERSION = "roi_policy_v6"
+POLICY_RUN_SUFFIX = "p6"
+
 
 @dataclass(frozen=True)
 class RuleConfig:
     self_stop_min_age: int = 5
     self_up_min_age: int = 3
-    self_min_avg_uv: float = 200
-    partner_min_avg_uv: float = 200
-    min_daily_cost: float = 100
-    stop_quantile: float = 0.20
+    self_min_daily_uv: float = 200
+    partner_min_daily_uv: float = 200
+    observe_min_latest_uv: float = 100
+    stop_quantile: float = 0.25
     up_quantile: float = 0.80
-    gzh_adjust_rank_min: float = 0.10
-    gzh_adjust_rank_max: float = 0.30
+    stop_weight_cap_quantile: float = 0.95
 
 
-def threshold_eligibility_mask(
+def _weighted_quantile(
+    values: pd.Series,
+    weights: pd.Series,
+    quantile: float,
+) -> float:
+    frame = pd.DataFrame(
+        {
+            "value": pd.to_numeric(values, errors="coerce"),
+            "weight": pd.to_numeric(weights, errors="coerce"),
+        }
+    ).dropna()
+    frame = frame[frame["weight"].gt(0)].sort_values("value", kind="stable")
+    if frame.empty:
+        raise ValueError("消耗加权分位没有有效ROI或成本")
+    cutoff = float(frame["weight"].sum()) * quantile
+    position = int(frame["weight"].cumsum().searchsorted(cutoff, side="left"))
+    return float(frame.iloc[min(position, len(frame) - 1)]["value"])
+
+
+def _weighted_percentiles(
+    values: pd.Series,
+    weights: pd.Series,
+) -> pd.Series:
+    frame = pd.DataFrame(
+        {
+            "value": pd.to_numeric(values, errors="coerce"),
+            "weight": pd.to_numeric(weights, errors="coerce"),
+        }
+    )
+    valid = frame["value"].notna() & frame["weight"].gt(0)
+    result = pd.Series(np.nan, index=frame.index)
+    if not valid.any():
+        return result
+    grouped = (
+        frame.loc[valid]
+        .groupby("value", sort=True)["weight"]
+        .sum()
+        .cumsum()
+    )
+    percentile_by_value = grouped / grouped.iloc[-1]
+    result.loc[valid] = frame.loc[valid, "value"].map(percentile_by_value)
+    return result
+
+
+def _daily_sample_mask(
     summary: pd.DataFrame,
+    entity_type: str,
+    dt: str,
     config: RuleConfig,
 ) -> pd.Series:
-    self_eligible = (
-        summary["entity_type"].eq(ENTITY_SELF)
-        & summary["日均首层UV"].gt(config.self_min_avg_uv)
+    min_uv = (
+        config.self_min_daily_uv
+        if entity_type == ENTITY_SELF
+        else config.partner_min_daily_uv
     )
-    gzh_eligible = (
-        summary["entity_type"].eq(ENTITY_GZH)
-        & summary["日均首层UV"].gt(config.partner_min_avg_uv)
+    return (
+        summary["entity_type"].eq(entity_type)
+        & pd.to_numeric(summary[f"首层UV_{dt}"], errors="coerce").gt(min_uv)
+        & pd.to_numeric(summary[f"成本_{dt}"], errors="coerce").gt(0)
+        & np.isfinite(pd.to_numeric(summary[f"ROI_{dt}"], errors="coerce"))
     )
+
+
+def threshold_eligibility_mask(
+    summary: pd.DataFrame,
+    config: RuleConfig,
+    expected_dates: list[str],
+) -> pd.Series:
+    eligible = pd.Series(False, index=summary.index)
+    for entity_type in (ENTITY_SELF, ENTITY_GZH):
+        entity_eligible = summary["entity_type"].eq(entity_type)
+        for dt in expected_dates:
+            entity_eligible &= _daily_sample_mask(
+                summary,
+                entity_type,
+                dt,
+                config,
+            )
+        eligible |= entity_eligible
+    return eligible
+
+
+def observation_mask(
+    summary: pd.DataFrame,
+    config: RuleConfig,
+    expected_dates: list[str],
+) -> pd.Series:
+    latest = expected_dates[-1]
+    eligible = threshold_eligibility_mask(summary, config, expected_dates)
     return (
-        (self_eligible | gzh_eligible)
-        & summary["成本"].gt(0)
-        & summary["三日最小单日成本"].gt(config.min_daily_cost)
-        & np.isfinite(summary["ROI"])
+        summary["entity_type"].isin([ENTITY_SELF, ENTITY_GZH])
+        & ~eligible
+        & pd.to_numeric(summary[f"首层UV_{latest}"], errors="coerce").gt(
+            config.observe_min_latest_uv
+        )
+        & pd.to_numeric(summary[f"成本_{latest}"], errors="coerce").gt(0)
+        & np.isfinite(pd.to_numeric(summary[f"ROI_{latest}"], errors="coerce"))
     )
 
 
-def compute_global_thresholds(
+def compute_daily_thresholds(
     summary: pd.DataFrame,
     expected_dates: list[str],
     config: RuleConfig,
 ) -> pd.DataFrame:
-    valid = summary.loc[threshold_eligibility_mask(summary, config)].copy()
-    if valid.empty:
-        raise ValueError("没有满足渠道UV门槛且三日聚合ROI有效的实体,无法计算全局阈值")
-
-    type_counts = valid["entity_type"].value_counts()
-    return pd.DataFrame(
-        [
-            {
-                "统计窗口": f"{expected_dates[0]} 至 {expected_dates[-1]}",
-                "t_stop": float(valid["ROI"].quantile(config.stop_quantile)),
-                "t_up": float(valid["ROI"].quantile(config.up_quantile)),
-                "阈值样本数": int(len(valid)),
-                "小程序样本数": int(type_counts.get(ENTITY_SELF, 0)),
-                "公众号样本数": int(type_counts.get(ENTITY_GZH, 0)),
-                "企微参考实体数": int(
-                    summary["entity_type"].eq(ENTITY_QIWEI).sum()
-                ),
-            }
-        ]
+    records: list[dict[str, object]] = []
+    labels = {ENTITY_SELF: "小程序投流", ENTITY_GZH: "公众号即转"}
+    formal_eligible = threshold_eligibility_mask(
+        summary,
+        config,
+        expected_dates,
     )
+    for dt in expected_dates:
+        for entity_type in (ENTITY_SELF, ENTITY_GZH):
+            mask = formal_eligible & summary["entity_type"].eq(entity_type)
+            valid = summary.loc[mask, f"ROI_{dt}"]
+            if valid.empty:
+                continue
+            raw_cost = pd.to_numeric(
+                summary.loc[mask, f"成本_{dt}"],
+                errors="coerce",
+            )
+            weight_cap = float(
+                raw_cost.quantile(config.stop_weight_cap_quantile)
+            )
+            capped_cost = raw_cost.clip(upper=weight_cap)
+            records.append(
+                {
+                    "统计日期": dt,
+                    "entity_type": entity_type,
+                    "渠道": labels[entity_type],
+                    "t_stop": _weighted_quantile(
+                        valid,
+                        capped_cost,
+                        config.stop_quantile,
+                    ),
+                    "t_up": float(valid.quantile(config.up_quantile)),
+                    "阈值样本数": int(len(valid)),
+                    "关停线口径": (
+                        f"消耗加权P{int(config.stop_quantile * 100)}"
+                    ),
+                    "权重封顶成本": weight_cap,
+                    "权重封顶分位": config.stop_weight_cap_quantile,
+                    "扩量线口径": (
+                        f"实体等权P{int(config.up_quantile * 100)}"
+                    ),
+                }
+            )
+    thresholds = pd.DataFrame(records)
+    if thresholds.empty:
+        raise ValueError("最近两日没有满足单日UV门槛的阈值样本")
+    return thresholds
 
 
 def apply_actions(
     summary: pd.DataFrame,
     thresholds: pd.DataFrame,
     config: RuleConfig,
+    expected_dates: list[str],
 ) -> pd.DataFrame:
     result = summary.copy()
     result["动作"] = ""
     result["动作原因"] = ""
-    result["渠道内ROI排名百分位"] = np.nan
-    t_stop = float(thresholds.iloc[0]["t_stop"])
-    t_up = float(thresholds.iloc[0]["t_up"])
-    result["t_stop"] = t_stop
-    result["t_up"] = t_up
-
-    gzh_mask = result["entity_type"].eq(ENTITY_GZH) & np.isfinite(result["ROI"])
-    if gzh_mask.any():
-        result.loc[gzh_mask, "渠道内ROI排名百分位"] = result.loc[gzh_mask, "ROI"].rank(
-            method="average", ascending=False, pct=True
+    eligible = threshold_eligibility_mask(result, config, expected_dates)
+    observe_only = observation_mask(result, config, expected_dates)
+    threshold_map = {
+        (str(row["统计日期"]), str(row["entity_type"])): (
+            float(row["t_stop"]),
+            float(row["t_up"]),
         )
+        for row in thresholds.to_dict("records")
+    }
+
+    for dt in expected_dates:
+        result[f"同日排名百分位_{dt}"] = np.nan
+        result[f"消耗加权同日分位_{dt}"] = np.nan
+        result[f"t_stop_{dt}"] = np.nan
+        result[f"t_up_{dt}"] = np.nan
+        for entity_type in (ENTITY_SELF, ENTITY_GZH):
+            key = (dt, entity_type)
+            if key not in threshold_map:
+                continue
+            mask = eligible & result["entity_type"].eq(entity_type)
+            stop, up = threshold_map[key]
+            result.loc[mask, f"同日排名百分位_{dt}"] = (
+                result.loc[mask, f"ROI_{dt}"].rank(
+                    method="average",
+                    ascending=True,
+                    pct=True,
+                )
+            )
+            raw_cost = pd.to_numeric(
+                result.loc[mask, f"成本_{dt}"],
+                errors="coerce",
+            )
+            weight_cap = float(
+                raw_cost.quantile(config.stop_weight_cap_quantile)
+            )
+            result.loc[mask, f"消耗加权同日分位_{dt}"] = (
+                _weighted_percentiles(
+                    result.loc[mask, f"ROI_{dt}"],
+                    raw_cost.clip(upper=weight_cap),
+                )
+            )
+            entity_mask = result["entity_type"].eq(entity_type)
+            result.loc[entity_mask, f"t_stop_{dt}"] = stop
+            result.loc[entity_mask, f"t_up_{dt}"] = up
+
+    first, latest = expected_dates
+    result["前一日同日排名百分位"] = result[f"同日排名百分位_{first}"]
+    result["最新日同日排名百分位"] = result[f"同日排名百分位_{latest}"]
+    result["前一日消耗加权分位"] = result[f"消耗加权同日分位_{first}"]
+    result["最新日消耗加权分位"] = result[f"消耗加权同日分位_{latest}"]
+    result["t_stop"] = result[f"t_stop_{latest}"]
+    result["t_up"] = result[f"t_up_{latest}"]
 
     for index, row in result.iterrows():
-        entity_type = row["entity_type"]
-        avg_uv = row["日均首层UV"]
-        roi = row["ROI"]
-        cost_confident = row["三日最小单日成本"] > config.min_daily_cost
-        stop_side = bool(np.isfinite(roi) and (roi <= t_stop or roi == 0))
-        up_side = bool(np.isfinite(roi) and roi > 0 and roi >= t_up)
+        entity_type = str(row["entity_type"])
+        if entity_type == ENTITY_QIWEI:
+            continue
+        if bool(observe_only.loc[index]):
+            result.at[index, "动作"] = "观察"
+            result.at[index, "动作原因"] = (
+                f"最新日首层UV>{config.observe_min_latest_uv:g},"
+                "但不满足连续两天首层UV均达到正式样本门槛;"
+                "已展示最新日ROI,不进入阈值和自动执行"
+            )
+            continue
+        if not bool(eligible.loc[index]):
+            continue
+
+        states: list[str] = []
+        for dt in expected_dates:
+            roi = float(row[f"ROI_{dt}"])
+            stop, up = threshold_map[(dt, entity_type)]
+            if roi <= stop:
+                states.append("低")
+            elif roi > 0 and roi >= up:
+                states.append("高")
+            else:
+                states.append("中")
+
+        both_low = states == ["低", "低"]
+        both_high = states == ["高", "高"]
+        inconsistent = len(set(states)) > 1 and any(
+            state in {"低", "高"} for state in states
+        )
 
         if entity_type == ENTITY_SELF:
-            if (
-                row["广告age"] >= config.self_stop_min_age
-                and stop_side
-                and avg_uv > config.self_min_avg_uv
-                and cost_confident
-            ):
+            raw_age = row.get("广告age")
+            age = int(raw_age) if pd.notna(raw_age) else 0
+            if both_low and age >= config.self_stop_min_age:
                 result.at[index, "动作"] = "关停"
                 result.at[index, "动作原因"] = (
-                    f"广告age≥{config.self_stop_min_age}天,三日聚合ROI不高于"
-                    f"全局t_stop(ROI=0明确关停),日均首层UV>"
-                    f"{config.self_min_avg_uv:g},连续3天每天成本>"
-                    f"{config.min_daily_cost:g}"
+                    f"连续两天分别处于同渠道当日消耗加权后"
+                    f"{config.stop_quantile:.0%},"
+                    f"广告age≥{config.self_stop_min_age}天"
                 )
-            elif (
-                row["广告age"] >= config.self_up_min_age
-                and up_side
-                and avg_uv > config.self_min_avg_uv
-                and cost_confident
-            ):
+            elif both_high and age >= config.self_up_min_age:
                 result.at[index, "动作"] = "扩量"
                 result.at[index, "动作原因"] = (
-                    f"广告age≥{config.self_up_min_age}天,三日聚合ROI不低于"
-                    f"全局t_up,日均首层UV>{config.self_min_avg_uv:g},"
-                    f"连续3天每天成本>{config.min_daily_cost:g}"
+                    f"连续两天分别处于同渠道当日前{1-config.up_quantile:.0%},"
+                    f"广告age≥{config.self_up_min_age}天"
+                )
+            elif both_low or both_high:
+                result.at[index, "动作"] = "观察"
+                required_age = (
+                    config.self_stop_min_age if both_low else config.self_up_min_age
+                )
+                result.at[index, "动作原因"] = (
+                    f"连续两天表现方向一致,但广告age<{required_age}天"
+                )
+            elif inconsistent:
+                result.at[index, "动作"] = "观察"
+                result.at[index, "动作原因"] = (
+                    f"连续两天表现不一致({first}:{states[0]},"
+                    f"{latest}:{states[1]})"
                 )
         elif entity_type == ENTITY_GZH:
-            if stop_side and avg_uv > config.partner_min_avg_uv and cost_confident:
+            if both_low:
                 result.at[index, "动作"] = "关停"
                 result.at[index, "动作原因"] = (
-                    "三日聚合ROI不高于全局t_stop(ROI=0明确关停),"
-                    f"日均首层UV>{config.partner_min_avg_uv:g},"
-                    f"连续3天每天成本>{config.min_daily_cost:g}"
+                    f"连续两天分别处于公众号当日消耗加权后"
+                    f"{config.stop_quantile:.0%}"
                 )
-            elif up_side and avg_uv > config.partner_min_avg_uv and cost_confident:
+            elif both_high:
                 result.at[index, "动作"] = "扩量"
                 result.at[index, "动作原因"] = (
-                    "三日聚合ROI不低于全局t_up,"
-                    f"日均首层UV>{config.partner_min_avg_uv:g},"
-                    f"连续3天每天成本>{config.min_daily_cost:g}"
+                    f"连续两天分别处于公众号当日前{1-config.up_quantile:.0%}"
+                )
+            elif inconsistent:
+                result.at[index, "动作"] = "观察"
+                result.at[index, "动作原因"] = (
+                    f"连续两天表现不一致({first}:{states[0]},"
+                    f"{latest}:{states[1]})"
                 )
-            else:
-                rank_pct = row["渠道内ROI排名百分位"]
-                if (
-                    np.isfinite(rank_pct)
-                    and config.gzh_adjust_rank_min <= rank_pct <= config.gzh_adjust_rank_max
-                ):
-                    result.at[index, "动作"] = "调整封面&落地页视频"
-                    result.at[index, "动作原因"] = (
-                        "三日综合ROI在公众号即转渠道排名"
-                        f"{config.gzh_adjust_rank_min:.0%}–"
-                        f"{config.gzh_adjust_rank_max:.0%}"
-                    )
     return result
 
 
@@ -178,13 +346,13 @@ def evaluate_rules(
         ad_age,
         fission_parameters=fission_parameters,
     )
-    thresholds = compute_global_thresholds(summary, dates, config)
-    evaluated = apply_actions(summary, thresholds, config)
-    evaluated["阈值样本状态"] = "未达到阈值样本门槛"
-    evaluated.loc[
-        threshold_eligibility_mask(evaluated, config),
-        "阈值样本状态",
-    ] = "进入阈值样本池"
+    thresholds = compute_daily_thresholds(summary, dates, config)
+    evaluated = apply_actions(summary, thresholds, config, dates)
+    eligible = threshold_eligibility_mask(evaluated, config, dates)
+    observe_only = observation_mask(evaluated, config, dates)
+    evaluated["阈值样本状态"] = "未达到连续两日样本门槛"
+    evaluated.loc[eligible, "阈值样本状态"] = "进入连续两日阈值样本池"
+    evaluated.loc[observe_only, "阈值样本状态"] = "观察_最新日UV达标"
     evaluated.loc[
         evaluated["entity_type"].eq(ENTITY_QIWEI),
         "阈值样本状态",

+ 19 - 12
examples/auto_put_ad_mini/roi_control/service.py

@@ -57,13 +57,12 @@ def _rule_config(config: RoiConfig) -> RuleConfig:
     return RuleConfig(
         self_stop_min_age=config.self_stop_min_age,
         self_up_min_age=config.self_up_min_age,
-        self_min_avg_uv=config.self_min_avg_uv,
-        partner_min_avg_uv=config.partner_min_avg_uv,
-        min_daily_cost=config.min_daily_cost,
+        self_min_daily_uv=config.self_min_daily_uv,
+        partner_min_daily_uv=config.partner_min_daily_uv,
+        observe_min_latest_uv=config.observe_min_latest_uv,
         stop_quantile=config.stop_quantile,
         up_quantile=config.up_quantile,
-        gzh_adjust_rank_min=config.gzh_adjust_rank_min,
-        gzh_adjust_rank_max=config.gzh_adjust_rank_max,
+        stop_weight_cap_quantile=config.stop_weight_cap_quantile,
     )
 
 
@@ -71,7 +70,7 @@ def _dates(start_date: str) -> list[str]:
     start = datetime.strptime(start_date, "%Y%m%d")
     return [
         (start + timedelta(days=offset)).strftime("%Y%m%d")
-        for offset in range(3)
+        for offset in range(2)
     ]
 
 
@@ -219,16 +218,24 @@ def run_daily_roi(
             fission_match_summary["实体数"] / channel_totals
         )
         fission_match_summary["参数版本"] = fission_parameters.release.version
-        actionable_candidates = annotated[annotated["动作"].ne("")].copy()
+        recommendation_rows = annotated[annotated["动作"].ne("")].copy()
         report_rows = annotated[
-            annotated["阈值样本状态"].eq("进入阈值样本池")
+            annotated["阈值样本状态"].isin(
+                [
+                    "进入连续两日阈值样本池",
+                    "观察_最新日UV达标",
+                ]
+            )
             | annotated["entity_type"].eq(ENTITY_QIWEI)
         ].copy()
         for row in snapshots:
             row["run_id"] = run_id
         for row in actions:
             row["run_id"] = run_id
-        threshold_record = thresholds.iloc[0].to_dict()
+        threshold_record = {
+            "统计窗口": f"{start_date} 至 {end_date}",
+            "每日渠道阈值": thresholds.to_dict("records"),
+        }
         replace_run_results(
             run_id,
             snapshots=snapshots,
@@ -259,7 +266,7 @@ def run_daily_roi(
             "end_date": end_date,
             "source_revision": source_revision,
             "entity_count": len(snapshots),
-            "candidate_count": len(actionable_candidates),
+            "candidate_count": len(recommendation_rows),
             "actionable_count": len(actions),
             "thresholds": threshold_record,
             "fission_multiplier_matches": fission_match_summary.to_dict(
@@ -270,12 +277,12 @@ def run_daily_roi(
         if not send_feishu:
             return result
 
-        counts = actionable_candidates.groupby("动作").size().to_dict()
+        counts = recommendation_rows.groupby("动作").size().to_dict()
         summary_text = (
             f"统计窗口:{start_date} - {end_date}\n"
             f"关停建议:{counts.get('关停', 0)}\n"
             f"扩量建议:{counts.get('扩量', 0)}\n"
-            f"素材调整建议:{counts.get('调整封面&落地页视频', 0)}\n"
+            f"观察:{counts.get('观察', 0)}\n"
             f"可执行动作:{len(actions)}\n"
             f"审批有效期:发送后 {config.approval_ttl_minutes} 分钟"
         )

+ 188 - 215
examples/auto_put_ad_mini/test_roi_control_metrics.py

@@ -3,6 +3,13 @@ from dataclasses import replace
 
 import pandas as pd
 
+from roi_control.fission_multiplier import (
+    DISPLAY_TOTAL_TO_FIRST_COLUMN,
+    MATCH_QIWEI_PARTNER,
+    QIWEI_REFERENCE_RUN_SUFFIX,
+    QIWEI_REFERENCE_VERSION,
+    load_fission_multiplier_parameters,
+)
 from roi_control.metrics import (
     ENTITY_GZH,
     ENTITY_QIWEI,
@@ -11,25 +18,19 @@ from roi_control.metrics import (
     QIWEI_CHANNEL,
     SELF_CHANNEL,
 )
-from roi_control.fission_multiplier import (
-    DISPLAY_TOTAL_TO_FIRST_COLUMN,
-    MATCH_QIWEI_PARTNER,
-    QIWEI_REFERENCE_RUN_SUFFIX,
-    QIWEI_REFERENCE_VERSION,
-    load_fission_multiplier_parameters,
-)
 from roi_control.reporting import (
-    BASE_COLUMNS,
     FINAL_ROI_COLUMN,
     T0_FISSION_MULTIPLIER_COLUMN,
     TOTAL_FISSION_TO_FIRST_UV_COLUMN,
     _sheet_frame,
+    _visible_columns,
 )
+from roi_control.data_source import build_daily_sql
 from roi_control.rules import evaluate_rules as _evaluate_rules
 from roi_control.service import _run_identity
 
 
-DATES = ["20260720", "20260721", "20260722"]
+DATES = ["20260721", "20260722"]
 FISSION_PARAMETERS = load_fission_multiplier_parameters()
 
 
@@ -38,7 +39,7 @@ def evaluate_rules(*args, **kwargs):
     return _evaluate_rules(*args, **kwargs)
 
 
-def row(entity_type, channel, entity_id, dt, roi, uv=600, age=10, cost=200.0):
+def row(entity_type, channel, entity_id, dt, roi, uv=600, cost=200.0):
     common = {
         "dt": dt,
         "entity_type": entity_type,
@@ -63,7 +64,7 @@ def row(entity_type, channel, entity_id, dt, roi, uv=600, age=10, cost=200.0):
         common.update(
             {
                 "代理名称": "代理",
-                "账号id": "account",
+                "账号id": "84502354",
                 "账号名称": "账户",
                 "广告id": entity_id,
                 "广告名称": f"广告{entity_id}",
@@ -80,10 +81,43 @@ def row(entity_type, channel, entity_id, dt, roi, uv=600, age=10, cost=200.0):
 
 
 class RoiRulesTest(unittest.TestCase):
-    def test_run_identity_tracks_qiwei_parameter_source(self):
-        temporary_id, temporary_key = _run_identity("20260727", FISSION_PARAMETERS)
-        self.assertIn(QIWEI_REFERENCE_RUN_SUFFIX, temporary_id)
-        self.assertIn(QIWEI_REFERENCE_VERSION, temporary_key)
+    def build_daily(self):
+        rows = []
+        for dt in DATES:
+            for index, roi in enumerate(
+                [0.1, 0.5, 0.8, 1.0, 1.2, 1.5, 2.0, 3.0, 4.0, 5.0]
+            ):
+                rows.append(
+                    row(ENTITY_SELF, SELF_CHANNEL, f"ad-{index}", dt, roi)
+                )
+            for index, roi in enumerate([0.1, 0.5, 1.0, 2.0, 4.0]):
+                rows.append(
+                    row(
+                        ENTITY_GZH,
+                        GZH_CHANNEL,
+                        f"公众号-{index}",
+                        dt,
+                        roi,
+                        uv=300,
+                    )
+                )
+            rows.append(
+                row(ENTITY_QIWEI, QIWEI_CHANNEL, "企微合作方", dt, 2.0, uv=300)
+            )
+        return pd.DataFrame(rows)
+
+    def test_run_identity_tracks_metric_policy_and_qiwei_versions(self):
+        run_id, run_key = _run_identity("20260727", FISSION_PARAMETERS)
+        self.assertIn("m7", run_id)
+        self.assertIn("p6", run_id)
+        self.assertIn("r15", run_id)
+        self.assertIn(QIWEI_REFERENCE_RUN_SUFFIX, run_id)
+        self.assertIn(QIWEI_REFERENCE_VERSION, run_key)
+
+    def test_daily_sql_uses_root_deduplicated_t0_fission_uv(self):
+        sql = build_daily_sql(DATES[0], DATES[1])
+        self.assertEqual(sql.count("SUM(NVL(t0_fission_uv_root, 0))"), 3)
+        self.assertNotIn("t0裂变人数", sql)
 
         formal = replace(
             FISSION_PARAMETERS,
@@ -95,240 +129,179 @@ class RoiRulesTest(unittest.TestCase):
             qiwei_exact_rows=1,
             qiwei_exact_available_rows=1,
         )
-        formal_id, formal_key = _run_identity("20260727", formal)
+        formal_id, _ = _run_identity("20260727", formal)
         self.assertIn(QIWEI_REFERENCE_RUN_SUFFIX, formal_id)
-        self.assertIn(QIWEI_REFERENCE_VERSION, formal_key)
-
-    def test_run_identity_isolates_corrected_source_revision(self):
-        default_id, default_key = _run_identity("20260729", FISSION_PARAMETERS)
-        revised_id, revised_key = _run_identity(
-            "20260729",
-            FISSION_PARAMETERS,
-            "datafix_20260730_1",
-        )
-
-        self.assertEqual(f"{default_id}_datafix_20260730_1", revised_id)
-        self.assertNotEqual(default_key, revised_key)
-        self.assertNotIn("datafix_20260730_1", revised_key)
-        self.assertLessEqual(len(revised_id), 64)
-        self.assertLessEqual(len(revised_key), 128)
-        self.assertEqual(
-            (revised_id, revised_key),
-            _run_identity(
-                "20260729",
-                FISSION_PARAMETERS,
-                "datafix_20260730_1",
-            ),
-        )
-        with self.assertRaisesRegex(ValueError, "source_revision"):
-            _run_identity("20260729", FISSION_PARAMETERS, "Invalid Revision")
-        with self.assertRaisesRegex(ValueError, "run_id"):
-            _run_identity("20260729", FISSION_PARAMETERS, "a" * 40)
 
-    def build_daily(self):
-        rows = []
-        # 每日形成稳定分布,低值实体连续低于后20%,高值实体连续高于前20%。
-        for dt in DATES:
-            for index, roi in enumerate([0.1, 0.5, 0.8, 1.0, 1.2, 1.5, 2.0, 3.0, 4.0, 5.0]):
-                rows.append(row(ENTITY_SELF, SELF_CHANNEL, f"ad-{index}", dt, roi))
-            rows.append(row(ENTITY_GZH, GZH_CHANNEL, "低ROI公众号", dt, 0.05, uv=300))
-            rows.append(row(ENTITY_GZH, GZH_CHANNEL, "高ROI公众号", dt, 6.0, uv=300))
-            rows.append(row(ENTITY_QIWEI, QIWEI_CHANNEL, "低ROI企微", dt, 0.02, uv=300))
-            rows.append(row(ENTITY_QIWEI, QIWEI_CHANNEL, "高ROI企微", dt, 7.0, uv=300))
-        return pd.DataFrame(rows)
-
-    def test_stop_and_up_actions(self):
-        daily = self.build_daily()
+    def test_two_daily_thresholds_drive_consecutive_actions(self):
         ages = pd.DataFrame(
             {
                 "广告id": [f"ad-{index}" for index in range(10)],
                 "广告age": [10] * 10,
             }
         )
-        candidates, thresholds, summary = evaluate_rules(daily, DATES, ages)
+        candidates, thresholds, summary = evaluate_rules(
+            self.build_daily(), DATES, ages
+        )
 
-        actions = set(candidates["动作"])
-        self.assertIn("关停", actions)
-        self.assertIn("扩量", actions)
-        self.assertEqual(len(thresholds), 1)
+        self.assertEqual(len(thresholds), 4)
         self.assertEqual(
-            thresholds.iloc[0]["阈值样本数"],
-            int(summary["entity_type"].ne(ENTITY_QIWEI).sum()),
+            set(zip(thresholds["统计日期"], thresholds["entity_type"])),
+            {
+                (DATES[0], ENTITY_SELF),
+                (DATES[1], ENTITY_SELF),
+                (DATES[0], ENTITY_GZH),
+                (DATES[1], ENTITY_GZH),
+            },
         )
-        self.assertTrue((summary["覆盖天数"] == 3).all())
-        self.assertTrue((summary["t_stop"] == thresholds.iloc[0]["t_stop"]).all())
-        self.assertTrue((summary["t_up"] == thresholds.iloc[0]["t_up"]).all())
-        qiwei = summary[summary["entity_type"].eq(ENTITY_QIWEI)]
-        self.assertEqual(len(qiwei), 2)
+        self.assertIn("关停", set(candidates["动作"]))
+        self.assertIn("扩量", set(candidates["动作"]))
+        self.assertTrue((summary["覆盖天数"] == 2).all())
         self.assertTrue(
-            qiwei["调控参与状态"].eq("仅展示_不进入阈值和调控").all()
+            summary[
+                summary["entity_type"].isin([ENTITY_SELF, ENTITY_GZH])
+            ]["阈值样本状态"].eq("进入连续两日阈值样本池").all()
         )
 
-    def test_qiwei_is_excluded_before_metric_and_action_evaluation(self):
-        daily = self.build_daily()
-        candidates, _, summary = evaluate_rules(daily, DATES)
-        self.assertTrue(summary["entity_type"].eq(ENTITY_QIWEI).any())
-        self.assertFalse(candidates["合作方名"].str.contains("企微").any())
-
-    def test_no_complete_three_day_entity_fails_clearly(self):
-        daily = self.build_daily()
-        incomplete = daily[daily["dt"].ne(DATES[-1])]
-        with self.assertRaisesRegex(ValueError, "连续三日数据"):
-            evaluate_rules(incomplete, DATES)
-
-    def test_daily_fission_income_is_scaled_once_without_duplicate_addition(self):
+    def test_inconsistent_daily_direction_is_observe(self):
         daily = self.build_daily()
-        mask = (
-            daily["entity_type"].eq(ENTITY_GZH)
-            & daily["公众号名"].eq("高ROI公众号")
-        )
-        daily.loc[mask, "效率收入"] = 100
-        daily.loc[mask, "裂变效率收入"] = 200
-        daily.loc[mask, "成本"] = 100
-
-        _, _, summary = evaluate_rules(daily, DATES)
-        target = summary[
-            summary["entity_type"].eq(ENTITY_GZH)
-            & summary["公众号名"].eq("高ROI公众号")
-        ].iloc[0]
-        expected = (
-            100 + 200 * FISSION_PARAMETERS.gzh_channel
-        ) / 100
-        self.assertEqual(target["裂变效率收入"], 600)
-        self.assertEqual(target["T0实际裂变收入"], 600)
-        self.assertAlmostEqual(target["ROI"], expected)
-        self.assertAlmostEqual(target["实际ROI"], 3.0)
+        mask = daily["广告id"].eq("ad-0")
+        daily.loc[mask & daily["dt"].eq(DATES[1]), "效率收入"] = 1000
+        ages = pd.DataFrame({"广告id": ["ad-0"], "广告age": [10]})
 
-    def test_self_age_blocks_actions(self):
-        daily = self.build_daily()
-        ages = pd.DataFrame({"广告id": ["ad-0", "ad-9"], "广告age": [2, 2]})
         candidates, _, _ = evaluate_rules(daily, DATES, ages)
-        blocked_ids = set(candidates.loc[candidates["entity_type"].eq(ENTITY_SELF), "广告id"])
-        self.assertNotIn("ad-0", blocked_ids)
-        self.assertNotIn("ad-9", blocked_ids)
-
-    def test_pyodps_lowercase_aliases_are_normalized(self):
-        daily = self.build_daily().rename(
-            columns={"首层UV": "首层uv", "T0裂变数": "t0裂变数"}
-        )
-        candidates, thresholds, _ = evaluate_rules(daily, DATES)
-        self.assertFalse(candidates.empty)
-        self.assertEqual(len(thresholds), 1)
+        target = candidates[candidates["广告id"].eq("ad-0")].iloc[0]
+        self.assertEqual(target["动作"], "观察")
+        self.assertIn("表现不一致", target["动作原因"])
 
-    def test_zero_roi_at_stop_boundary_is_stopped(self):
-        rows = []
-        for dt in DATES:
-            for index, roi in enumerate([0, 0, 0, 0, 0, 1, 2, 3, 4, 5]):
-                rows.append(row(ENTITY_SELF, SELF_CHANNEL, f"boundary-{index}", dt, roi))
-        ages = pd.DataFrame(
-            {
-                "广告id": [f"boundary-{index}" for index in range(10)],
-                "广告age": [10] * 10,
-            }
+    def test_latest_day_uv_over_100_is_reported_as_observe(self):
+        daily = self.build_daily()
+        extra = row(
+            ENTITY_SELF,
+            SELF_CHANNEL,
+            "latest-only",
+            DATES[1],
+            1.5,
+            uv=150,
         )
-        candidates, thresholds, _ = evaluate_rules(pd.DataFrame(rows), DATES, ages)
-        self.assertEqual(thresholds.iloc[0]["t_stop"], 0)
-        stopped = candidates[candidates["动作"].eq("关停")]
-        self.assertTrue(stopped["广告id"].eq("boundary-0").any())
+        daily = pd.concat([daily, pd.DataFrame([extra])], ignore_index=True)
 
-    def test_display_name_change_does_not_break_three_day_identity(self):
+        candidates, _, summary = evaluate_rules(daily, DATES)
+        target = summary[summary["广告id"].eq("latest-only")].iloc[0]
+        self.assertEqual(target["覆盖天数"], 1)
+        self.assertEqual(target["动作"], "观察")
+        self.assertEqual(target["阈值样本状态"], "观察_最新日UV达标")
+        self.assertAlmostEqual(target["最新日效率ROI"], 1.5)
+        self.assertTrue(candidates["广告id"].eq("latest-only").any())
+
+    def test_both_days_must_exceed_daily_uv_gate(self):
         daily = self.build_daily()
         mask = daily["广告id"].eq("ad-9")
-        daily.loc[mask & daily["dt"].eq(DATES[1]), "广告名称"] = "改名后的广告"
+        daily.loc[mask & daily["dt"].eq(DATES[0]), "首层UV"] = 180
         ages = pd.DataFrame({"广告id": ["ad-9"], "广告age": [10]})
-        candidates, _, summary = evaluate_rules(daily, DATES, ages)
-        target = summary[
-            summary["entity_type"].eq(ENTITY_SELF)
-            & summary["广告id"].eq("ad-9")
-        ]
-        self.assertEqual(len(target), 1)
-        self.assertTrue(candidates["广告id"].eq("ad-9").any())
 
-    def test_global_threshold_pool_only_uses_miniapp_and_gzh(self):
+        _, _, summary = evaluate_rules(daily, DATES, ages)
+        target = summary[summary["广告id"].eq("ad-9")].iloc[0]
+        self.assertEqual(target["动作"], "观察")
+        self.assertEqual(target["阈值样本状态"], "观察_最新日UV达标")
+
+    def test_stop_threshold_uses_cost_weighted_p30(self):
         rows = []
         specifications = [
-            (ENTITY_SELF, SELF_CHANNEL, "self-qualified", 1.0, 300),
-            (ENTITY_GZH, GZH_CHANNEL, "gzh-qualified", 2.0, 300),
-            (ENTITY_QIWEI, QIWEI_CHANNEL, "qiwei-qualified", 3.0, 300),
-            (ENTITY_GZH, GZH_CHANNEL, "gzh-low-uv-excluded", 100.0, 100),
+            ("weighted-0", 0.1, 10),
+            ("weighted-1", 0.2, 10),
+            ("weighted-2", 1.0, 1000),
+            ("weighted-3", 2.0, 1000),
         ]
-        for entity_type, channel, entity_id, roi, uv in specifications:
-            for dt in DATES:
-                rows.append(row(entity_type, channel, entity_id, dt, roi, uv=uv))
-
-        _, thresholds, summary = evaluate_rules(pd.DataFrame(rows), DATES)
-        threshold = thresholds.iloc[0]
-        expected = pd.Series([1.0, 2.0])
-        self.assertAlmostEqual(threshold["t_stop"], expected.quantile(0.20))
-        self.assertAlmostEqual(threshold["t_up"], expected.quantile(0.80))
-        self.assertEqual(threshold["阈值样本数"], 2)
-        self.assertEqual(threshold["小程序样本数"], 1)
-        self.assertEqual(threshold["公众号样本数"], 1)
-        self.assertEqual(threshold["企微参考实体数"], 1)
-        self.assertEqual(len(summary), 4)
-
-    def test_actions_use_three_day_aggregate_not_each_daily_roi(self):
-        rows = []
-        daily_rois = [0.0, 0.0, 6.0]
-        for dt, roi in zip(DATES, daily_rois):
-            rows.append(row(ENTITY_GZH, GZH_CHANNEL, "波动公众号", dt, roi, uv=300))
         for dt in DATES:
-            rows.append(row(ENTITY_SELF, SELF_CHANNEL, "基准创意", dt, 1.0, uv=600))
-            rows.append(row(ENTITY_GZH, GZH_CHANNEL, "基准公众号", dt, 3.0, uv=300))
-
-        candidates, thresholds, summary = evaluate_rules(pd.DataFrame(rows), DATES)
-        target = summary[summary["公众号名"].eq("波动公众号")].iloc[0]
-        self.assertAlmostEqual(target["ROI"], 2.0)
-        self.assertAlmostEqual(target["t_stop"], thresholds.iloc[0]["t_stop"])
-        self.assertAlmostEqual(target["t_up"], thresholds.iloc[0]["t_up"])
-        self.assertNotEqual(
-            candidates[candidates["公众号名"].eq("波动公众号")]["动作"].tolist(),
-            ["关停"],
+            for entity_id, roi, cost in specifications:
+                rows.append(
+                    row(
+                        ENTITY_SELF,
+                        SELF_CHANNEL,
+                        entity_id,
+                        dt,
+                        roi,
+                        cost=cost,
+                    )
+                )
+        _, thresholds, _ = evaluate_rules(pd.DataFrame(rows), DATES)
+        self_thresholds = thresholds[
+            thresholds["entity_type"].eq(ENTITY_SELF)
+        ]
+        self.assertTrue(self_thresholds["t_stop"].eq(1.0).all())
+        self.assertTrue(
+            self_thresholds["关停线口径"].eq("消耗加权P25").all()
         )
 
-    def test_uv_gate_uses_three_day_average_not_each_day(self):
-        daily = self.build_daily()
-        target_mask = daily["广告id"].eq("ad-9")
-        daily.loc[target_mask, "首层UV"] = [100, 250, 400]
-        ages = pd.DataFrame({"广告id": ["ad-9"], "广告age": [10]})
-
-        candidates, thresholds, summary = evaluate_rules(daily, DATES, ages)
-        target = summary[summary["广告id"].eq("ad-9")].iloc[0]
-        self.assertEqual(target["日均首层UV"], 250)
-        self.assertEqual(
-            candidates[candidates["广告id"].eq("ad-9")]["动作"].tolist(),
-            ["扩量"],
+    def test_qiwei_is_display_only(self):
+        candidates, _, summary = evaluate_rules(self.build_daily(), DATES)
+        self.assertFalse(candidates["entity_type"].eq(ENTITY_QIWEI).any())
+        qiwei = summary[summary["entity_type"].eq(ENTITY_QIWEI)]
+        self.assertTrue(
+            qiwei["调控参与状态"].eq("仅展示_不进入阈值和调控").all()
         )
-        self.assertEqual(thresholds.iloc[0]["小程序样本数"], 10)
 
-    def test_every_day_cost_must_be_above_100(self):
+    def test_report_uses_daily_rows_and_only_latest_day_is_actionable(self):
         daily = self.build_daily()
-        target_mask = daily["广告id"].eq("ad-9")
-        daily.loc[target_mask, "成本"] = [100, 200, 200]
-        daily.loc[target_mask, "效率收入"] = [500, 1000, 1000]
-        ages = pd.DataFrame({"广告id": ["ad-9"], "广告age": [10]})
-
-        candidates, thresholds, summary = evaluate_rules(daily, DATES, ages)
-        target = summary[summary["广告id"].eq("ad-9")].iloc[0]
-        self.assertEqual(target["三日最小单日成本"], 100)
-        self.assertFalse(candidates["广告id"].eq("ad-9").any())
-        self.assertEqual(thresholds.iloc[0]["小程序样本数"], 9)
-
-    def test_report_exposes_final_roi_and_both_fission_coefficients(self):
-        candidates, _, _ = evaluate_rules(self.build_daily(), DATES)
-        frame = _sheet_frame(candidates, "小程序投流")
-
-        visible_columns = list(BASE_COLUMNS["小程序投流"])
-        self.assertIn(FINAL_ROI_COLUMN, visible_columns)
-        self.assertIn(T0_FISSION_MULTIPLIER_COLUMN, visible_columns)
-        self.assertIn(TOTAL_FISSION_TO_FIRST_UV_COLUMN, visible_columns)
-        self.assertIn("T0裂变效率收入", visible_columns)
-        self.assertIn("总预估效率收入", visible_columns)
-        self.assertNotIn("LTV预测效率收入", visible_columns)
-        self.assertIn("裂变效率收入", frame.columns)
-        self.assertNotIn("裂变效率收入", visible_columns)
+        extra = row(
+            ENTITY_SELF,
+            SELF_CHANNEL,
+            "latest-only",
+            DATES[1],
+            0.01,
+            uv=150,
+        )
+        daily = pd.concat([daily, pd.DataFrame([extra])], ignore_index=True)
+        candidates, _, _ = evaluate_rules(daily, DATES)
+        frame = _sheet_frame(candidates, "小程序投流", DATES, 0.25)
+        visible = _visible_columns("小程序投流", DATES, 0.25)
+
+        for column in (
+            "dt",
+            "首层UV",
+            "T0裂变人数",
+            "T0裂变率",
+            "首层效率收入",
+            "T0裂变效率收入",
+            "预测总效率收入",
+            "成本",
+            "效率ROI",
+            T0_FISSION_MULTIPLIER_COLUMN,
+            TOTAL_FISSION_TO_FIRST_UV_COLUMN,
+            "消耗加权分位",
+            "消耗加权P25线",
+            "动作原因",
+            "阈值样本状态",
+        ):
+            self.assertIn(column, visible)
+        self.assertNotIn(FINAL_ROI_COLUMN, visible)
+        self.assertIn(FINAL_ROI_COLUMN, frame.columns)
+        latest_row = frame[
+            frame["广告id"].eq("ad-0") & frame["dt"].eq(DATES[1])
+        ].iloc[0]
+        self.assertEqual(latest_row["T0裂变人数"], 120)
+        self.assertAlmostEqual(latest_row["T0裂变率"], 0.2)
+        self.assertEqual(latest_row["首层效率收入"], 20)
+        self.assertEqual(latest_row["预测总效率收入"], 20)
+        self.assertEqual(set(frame["dt"]), set(DATES))
+        self.assertTrue((frame.groupby("广告id").size() == 2).all())
+
+        history = frame[frame["dt"].eq(DATES[0])]
+        latest = frame[frame["dt"].eq(DATES[1])]
+        self.assertTrue(history["动作"].fillna("").eq("").all())
+        self.assertTrue(history["审批选择"].eq("历史日").all())
+        self.assertTrue(history["动作幂等键"].fillna("").eq("").all())
+        self.assertTrue(latest["动作"].ne("").any())
         self.assertTrue(
-            frame[FINAL_ROI_COLUMN].equals(frame["ROI"])
+            latest["阈值样本状态"].eq("观察_最新日UV达标").any()
+        )
+        self.assertEqual(
+            frame.iloc[0]["dt"],
+            DATES[1],
+        )
+        self.assertEqual(
+            frame.iloc[-1]["dt"],
+            DATES[0],
         )
         pd.testing.assert_series_equal(
             frame[TOTAL_FISSION_TO_FIRST_UV_COLUMN],

+ 26 - 0
examples/auto_put_ad_mini/test_roi_control_policy.py

@@ -128,6 +128,32 @@ class RoiControlPolicyTest(unittest.TestCase):
         self.assertEqual(len(actions), 1)
         self.assertEqual(actions[0]["action_type"], ACTION_SCALE_BID)
 
+    def test_observe_recommendation_never_creates_tencent_action(self):
+        summary = pd.DataFrame(
+            [
+                {
+                    "entity_type": ENTITY_SELF,
+                    "channel": "小程序投流-稳定",
+                    "账号id": "84502354",
+                    "广告id": "1001",
+                    "创意id": "2001",
+                    "动作": "观察",
+                    "动作原因": "连续两天表现不一致",
+                }
+            ]
+        )
+        annotated, _, actions = annotate_execution(
+            summary,
+            {84502354},
+            "roi_20260730_m6_p3",
+        )
+        self.assertEqual(actions, [])
+        self.assertEqual(annotated.iloc[0]["执行模式"], MODE_NOTIFY_ONLY)
+        self.assertEqual(
+            annotated.iloc[0]["执行说明"],
+            "观察行仅展示,不执行腾讯写操作",
+        )
+
     def test_invalid_apply_without_daily_job_is_rejected(self):
         environment = {
             "DAILY_ROI_ENABLED": "0",