Merge remote-tracking branch 'origin/main' into fix/moneypagefix

This commit is contained in:
左辰勇
2026-07-21 14:56:09 +08:00
63 changed files with 7106 additions and 87 deletions
+2
View File
@@ -115,6 +115,7 @@
| 40a | `GET /api/v1/platform/flags` | 无 | [详情](./platform/platform-flags.md)(客户端运营 feature flag,比价/领券期广告开关等,拉取后缓存) |
| 40b | `GET /api/v1/platform/ad-config` | 无 | [详情](./platform/platform-ad-config.md)(客户端拉广告配置:穿山甲 app_id+各位ID+各场景开关;不含验签密钥) |
| 40c | `GET /api/v1/platform/app-version` | 无 | [详情](./platform/platform-app-version.md)(最新 App 版本,OTA 检查更新;与本机 versionCode 比) |
| 40d | `GET /api/v1/platform/huawei-review` | 无 | [详情](./platform/platform-huawei-review.md)(华为审核开关:快速设置权限步能否被用户关闭;仅华为 ROM 客户端拉) |
| **微信支付回调**(前缀 `/api/v1/wxpay` |||
| W1 | `POST /api/v1/wxpay/transfer-auth-notify` | 无 | 免确认收款授权结果通知(一期 stub:仅应答 200 不验签不改账,授权状态靠主动查询兜底)(无单独文档) |
| **CPS 群发短链落地**(**无前缀**,挂域名根;公网不鉴权) |||
@@ -153,6 +154,7 @@
| A12 | `GET /admin/api/ad-revenue-report` | admin | [详情](./admin/ad/admin-ad-revenue-report.md)(广告收益报表:分页/场景/`app_env` 筛 + **DAU/ARPU** #120;真实收益侧接穿山甲日表 #92) |
| A13 | `GET / PATCH /admin/api/ad-config` | operator/finance | 广告配置(穿山甲 ID/验签密钥/各场景开关;C 端只读版见 40b)(无单独文档,见 `app/admin/routers/ad_config.py`) |
| A14 | `GET /admin/api/config``PATCH /config/{key}` | operator/finance | 运营可配置项([app_config](../database/app_config.md):奖励常量/提现地板价等;#117 修系统配置下发)(无单独文档,见 `app/admin/routers/config.py`) |
| A16 | `GET / PATCH /admin/api/huawei-review` | operator/tech | 华为审核开关(快速设置权限步能否被用户关闭,落 `app_config.huawei_review`;C 端只读版见 40d)(无单独文档,见 `app/admin/routers/huawei_review.py`) |
| **A·管理员与角色**(super_admin):`GET`/`POST` `/admins``PATCH`/`DELETE` `/admins/{id}`(#126 删除+`pages_override`)、`GET`/`POST` `/roles``GET /roles/catalog``PATCH`/`DELETE` `/roles/{id}`(#117/#126 自定义角色) ||| [列表](./admin/admins/admin-admins-list.md) / [](./admin/admins/admin-admin-create.md) / [改+删](./admin/admins/admin-admin-update.md) / [角色](./admin/admin-roles.md) |
| A15 | `GET /admin/api/audit-logs` | admin | [详情](./admin/admin-audit-logs.md) |
| **A·CPS 运营台**:群/活动 CRUD、`POST /referral-links``POST /orders/reconcile`(美团+京东 #90)、`GET /orders``/stats`、群 `timeseries`/`daily`/`wx-users`/`day-users`(#79) ||| [详情](./admin/admin-cps.md) |
@@ -0,0 +1,36 @@
# GET /api/v1/platform/huawei-review — 华为审核开关
> 所属:Platform 组(前缀 `/api/v1/platform` | 鉴权:无 | [← 返回 API 索引](../README.md)
华为应用市场审核要求:新手引导的**「快速设置」权限步必须可被用户关闭**(引导视频页不在要求内)。本端点把运营后台配的开关下发给客户端,决定该步左上角是否出现退出按钮。
**不鉴权**:引导页在登录之前就展示,此时客户端必然没有 token。
值来自 `app_config` 表的 `huawei_review` 行(admin 页 `GET / PATCH /admin/api/huawei-review` 可改),空库回退 `default`
## 入参
无。
## 出参
响应 `200`:`HuaweiReviewOut`
| 字段 | 类型 | 说明 |
|---|---|---|
| `mode` | string | `default`(强制展示,不可关闭 = 上线至今的现状)/ `review`(可关闭,过审用)。仅供排查时看后台切成了哪态 |
| `onboarding_closable` | bool | 快速设置权限步是否允许用户退出(= `mode == "review"`)。**客户端只读这一个字段决策** |
Mock 出参:
```json
{
"mode": "review",
"onboarding_closable": true
}
```
## 说明
- **只有华为 ROM 客户端会来拉**HarmonyOS / EMUI;荣耀 MagicOS 不拉)。机型判断在客户端做(`OemDetector`),服务端不看 UA,也就不用维护机型名单。
- 客户端在进新手引导前拉一次并本地缓存;请求失败/超时用上次缓存值,从未拉到过则按 `onboarding_closable=false`(宁可不给退出按钮,也不误放开)。
- 脏值兜底:DB 里 `mode` 不在枚举内时服务端一律回退 `default`
- 切回 `default` 即可一键收回退出按钮(审核通过后无需发版)。
+2
View File
@@ -41,6 +41,8 @@
| `ad_ecpm_record` | 广告展示 eCPM 上报(收益对账) | `models/ad_ecpm.py` | [详情](./ad_ecpm_record.md) |
| `ad_feed_reward_record` | 信息流/Draw 广告结算记录(10 秒一份,client_event_id 幂等;`ad_type`+`feed_scene` 分形态/场景) | `models/ad_feed_reward.py` | [详情](./ad_feed_reward_record.md) |
| `ad_pangle_daily_revenue` | 穿山甲 GroMore 后台收益日表(定时拉取,收益报表/大盘真实收益源,#92) | `models/ad_pangle_revenue.py` | [详情](./ad_pangle_daily_revenue.md) |
| `inactivity_reset_log` | 15 天不活跃清零审计(每次清零一行;清零前三桶余额快照+原因+不活跃天数;只清金币+现金,邀请金仅快照) | `models/inactivity.py` | [详情](./inactivity_reset_log.md) |
| `inactivity_notification_log` | 不活跃清零前预警记录(余额快照+档位+通道+状态;streak 去重依据 + 占位 outbox) | `models/inactivity.py` | [详情](./inactivity_notification_log.md) |
### 比价 / 省钱
| 表 | 用途 | 模型 | 文档 |
+9
View File
@@ -29,3 +29,12 @@
## 注意
- 不缓存:配置读频率低(每次福利操作读一次,主键查极快),admin 改了立即生效、跨进程一致(多 worker 也对)。
- 新增可配项 = 在 `CONFIG_DEFS` 加一条 + 业务处改用 `app_config.get_value(db, key)` 读;不需要建迁移(行是动态插的,表结构不变)。
## 专用 key(借表不进 CONFIG_DEFS)
有自己的语义与专用管理页的配置,复用本表但**不注册进 `CONFIG_DEFS`**——混进通用「系统配置」页只会显示成一个没头没尾的 on/off。它们各有一对 `get_*` / `set_*` 函数(仍在 `repositories/app_config.py`),`value` 存 dict,空行回退各自的模块内默认值。
| key | 管理页 / admin 端点 | C 端读取 | 说明 |
|---|---|---|---|
| `ad_config` | `GET/PATCH /admin/api/ad-config` | `GET /api/v1/platform/ad-config`(去密钥) | 穿山甲 app_id / 各代码位 / 各场景开关 |
| `app_version` | 内部写入(`X-Internal-Secret`) | `GET /api/v1/platform/app-version` | OTA 最新版本信息 |
| `huawei_review` | `GET/PATCH /admin/api/huawei-review` | `GET /api/v1/platform/huawei-review` | 华为审核开关:`{"mode": "default""review"}`,决定新手引导「快速设置」权限步能否被用户关闭。脏值/空行一律回退 `default`(不给退出按钮) |
@@ -0,0 +1,36 @@
# inactivity_notification_log — 不活跃清零前预警记录
> 模型 `app/models/inactivity.py` · 仓库 `app/repositories/inactivity.py` · 通知器 `app/integrations/notifier.py` · [← 索引](./README.md) · [总览](./OVERVIEW.md)
清零前按可配置节奏(`INACTIVITY_WARN_DAYS_BEFORE`,默认清零前 7 天、2 天各一次)向用户预警"账户里的 xx 金币和 xx 现金将被清零"。每发一次预警写一行,记推送时的余额快照 + 提前天数档 + 通道 + 状态。兼作两用:**预警去重**依据(同 streak 内 `stage==k 且 created_at > last_active` 即已推过、不重推)与**占位 outbox**(v1 通道=`log`,只打日志不真推;后续接 JPush/短信同层扩展)。append-only,不更新。**预警只涉及会被清的金币 + 折算现金;邀请奖励金不清、不预警**(`invite_cash_balance_cents` 仅作账户状态快照)。
## 用在哪 / 增删改查
- **C(插入)**:`inactivity.run_warn_once` 命中预警档、且本 streak 未推过时,调 `notifier.warn` 后写一行(`status` = 通知器返回,占位实现为 `placeholder`)。
- **U / D**:无(append-only)。
- **R**:预警去重查询(`user_id + stage + created_at > last_active`);未来接真实推送时作待推送 outbox。
## 字段
| 列 | 类型 | 约束 / 默认 | 说明(取值 / join) |
|---|---|---|---|
| `id` | Integer | **PK**, autoincrement | 主键 |
| `user_id` | Integer | NOT NULL, index | 预警对象;只索引不设外键(同 `analytics_event`) |
| `stage` | Integer | NOT NULL | 提前天数档(如 `7` / `2`,即清零前第几天推) |
| `inactive_days` | Integer | NOT NULL | 推送时的不活跃天数(北京自然日) |
| `coin_balance` | Integer | NOT NULL | 推送时金币余额快照(将被清) |
| `cash_balance_cents` | Integer | NOT NULL | 推送时折算现金余额快照(分,将被清) |
| `invite_cash_balance_cents` | Integer | NOT NULL | 推送时**邀请奖励金**余额快照(分,**不清、不在预警额度内**) |
| `channel` | String(16) | NOT NULL | 通道:`log`(占位) / `jpush` / `sms` |
| `status` | String(16) | NOT NULL | 状态:`placeholder`(占位未真推) / `sent` / `failed` |
| `created_at` | DateTime(tz) | server_default now(), index | 推送时刻;去重比 `created_at > last_active`(用户回归后 `last_active` 前移 → 旧行自然失效、开启新 streak) |
## 关系 / Join Key
- `user_id``user.id`(无外键直连,靠 `user_id` 关联)。
-`inactivity_reset_log` 无直接外键;同一 streak 内先有若干预警行,到期后有一行清零。
## 索引与约束
- PK `id`;`ix_inactivity_notification_log_user_id``ix_inactivity_notification_log_created_at`
## 注意
- **预警去重按 streak**:判据是 `created_at > last_active`;用户一有活跃(`home_view`/比价/领券),`last_active` 前移,旧预警行"失效",回归后可重新进入预警。
- **占位实现**:v1 `LogNotifier``logger.warning("[inactivity-warn] ...")`、返回 `placeholder`,不真推(参照心跳告警"本期先不接推送"先例)。
- **漏跑补发**:worker 漏跑数天后某用户可能同时满足多档,只补发**最紧急的未推档**(最小提前天数),避免刷屏。
+35
View File
@@ -0,0 +1,35 @@
# inactivity_reset_log — 15 天不活跃清零审计
> 模型 `app/models/inactivity.py` · 仓库 `app/repositories/inactivity.py` · worker `app/core/inactivity_reset_worker.py` · [← 索引](./README.md) · [总览](./OVERVIEW.md)
连续 15 天不活跃(北京自然日,活跃口径见 `app/repositories/activity.py`:`home_view` + 发起比价 + 发起领券,**不含登录**)的用户,worker 每日自动清零其**金币 + 折算现金**。每清一个用户写一行,记清零前三桶余额快照 + 原因 + 判定时的活跃时间/不活跃天数,供纠纷排查。清零同时另写 2 条钱包流水(`coin_transaction` / `cash_transaction`,`biz_type=inactivity_reset`,`ref_id=` 本表 `id`),资金流可逐笔回溯、人工恢复。**邀请奖励金(`invite_cash_balance_cents`)是产品红线、不清零**,本表 `invite_cash_balance_cents_before` 仅为清零时仍保留的邀请金快照(非被清金额)。append-only,不更新。
## 用在哪 / 增删改查
- **C(插入)**:`inactivity.clear_user` 逐用户清零(独立事务、行锁)时写一行,`db.flush()``id` 作流水 `ref_id` 交叉链接。
- **U / D**:无(append-only 审计)。
- **R**:纠纷排查 / 对账(与 `coin_transaction` / `cash_transaction``ref_id` 交叉核对)。
## 字段
| 列 | 类型 | 约束 / 默认 | 说明(取值 / join) |
|---|---|---|---|
| `id` | Integer | **PK**, autoincrement | 主键;作 `ref_id` 写入两条清零流水 |
| `user_id` | Integer | NOT NULL, index | 被清零用户;只索引不设外键(同 `analytics_event`,避免删用户级联 / 留历史) |
| `coin_balance_before` | Integer | NOT NULL | 清零前金币余额(个数);= 对应 `coin_transaction.amount` 绝对值 |
| `cash_balance_cents_before` | Integer | NOT NULL | 清零前折算现金余额(分);= 对应 `cash_transaction.amount_cents` 绝对值 |
| `invite_cash_balance_cents_before` | Integer | NOT NULL | 清零时的**邀请奖励金**余额快照(分)——**不清、原封保留**,仅记录以证明"未动邀请金" |
| `last_active_at` | DateTime(tz) | nullable | 判定时的最近活跃时刻(UTC);无任何活跃信号时兜底为 `user.created_at` |
| `inactive_days` | Integer | NOT NULL | 判定时的不活跃天数(北京自然日) |
| `reason` | String(32) | NOT NULL | 清零原因,如 `inactive_15d` |
| `reset_at` | DateTime(tz) | server_default now(), index | 清零时刻 |
## 关系 / Join Key
- `user_id``user.id`(无外键直连,靠 `user_id` 关联)。
- `id``coin_transaction.ref_id` / `cash_transaction.ref_id`(`biz_type=inactivity_reset`):审计行 ↔ 资金流水交叉对账。
## 索引与约束
- PK `id`;`ix_inactivity_reset_log_user_id`(按用户查)、`ix_inactivity_reset_log_reset_at`(按时间查)。
## 注意
- **只清 2 桶**:金币 + 折算现金;**邀请现金不清**(两本账物理隔离,见 [`coin_account`](./coin_account.md) / `wallet.CoinAccount` 注释)。
- **天然幂等**:清完余额=0,次日不再匹配;worker 重启 / 多次唤醒 / 补跑都不会重复清零或重复流水。
- **总闸默认关**(`INACTIVITY_RESET_ENABLED=false`),灰度验证清零名单后再开。
@@ -0,0 +1,897 @@
# 接口 QPS + 耗时可观测(OpenObserve)实现计划
> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.
**Goal:** 给 app-server 每个接口采集 QPS + 耗时 + 错误率,经轻量 ASGI 中间件 + 后台 worker 批量直采到本地 Docker 的 OpenObserve。
**Architecture:** 纯 ASGI 中间件测每请求耗时/抓路由模板+状态码 → 非阻塞入有界队列(满则丢、绝不阻塞)→ 后台 asyncio worker 批量 POST 到 OpenObserve `_json` ingest 端点。请求路径零 I/O;未配置观测则整套 no-op;上报失败丢批不重试。
**Tech Stack:** FastAPI / Starlette ASGI 中间件、`asyncio.Queue``httpx.AsyncClient`(已有依赖)、pydantic-settings、OpenObserveDocker)。
参考 spec[docs/superpowers/specs/2026-07-06-openobserve-api-metrics-design.md](2026-07-06-openobserve-api-metrics-design.md)
---
## 文件结构
| 文件 | 职责 |
|---|---|
| `app/core/config.py`(改) | 新增 `OBSERVE_*` 配置 + `observe_configured` 门槛属性 |
| `app/core/observe.py`(新) | 有界事件队列 + `record_event` + 路由模板解析 + `RequestMetricsMiddleware` |
| `app/core/observe_worker.py`(新) | 后台批量上报 worker:`_collect_batch` / `_post_batch` / `start_*` / `stop_*` |
| `app/main.py`(改) | 挂中间件(最外层)+ lifespan 启停 worker |
| `.env.example`(改) | 新增 `OBSERVE_*` 注释段 |
| `deploy/openobserve/docker-compose.yml`(新) | 本地 OpenObserve 容器 |
| `deploy/openobserve/README.md`(新) | 部署步骤 + 查询/仪表盘 SQL |
| `tests/test_observe.py`(新) | 配置门槛 / 队列 / 中间件 / worker 单测 |
**关键接口契约(跨任务一致,勿改名):**
- `app.core.observe.get_queue() -> asyncio.Queue[dict]`
- `app.core.observe.record_event(event: dict) -> None`
- `app.core.observe.take_dropped() -> int`
- `app.core.observe.RequestMetricsMiddleware`ASGI class`__init__(self, app)`
- 事件字段:`_timestamp`(µs int) / `service` / `env` / `method` / `route` / `status` / `duration_ms`(float)
- `app.core.observe_worker.start_observe_worker() -> asyncio.Task | None`
- `app.core.observe_worker.stop_observe_worker(task) -> None`
- `settings.observe_configured -> bool`
---
## Task 1: 配置项 `OBSERVE_*` + `observe_configured`
**Files:**
- Modify: `app/core/config.py`(在 `cors_origins_list` property 之后、`is_prod` property 之前插入)
- Test: `tests/test_observe.py`(新建)
- [ ] **Step 1: 写失败测试**
新建 `tests/test_observe.py`
```python
"""接口指标可观测(observe)单测:配置门槛 / 队列 / 中间件 / worker。
沿用仓库约定:TestClient + monkeypatch,绝不打真网络。observe 默认关(conftest 未设
OBSERVE_*),需要开启的用例用 monkeypatch 改 settings 单例属性。
"""
from __future__ import annotations
from app.core.config import settings
def test_observe_configured_requires_switch_and_creds(monkeypatch):
# 开关开 + endpoint(默认 localhost)+ user + password 齐全 → True
monkeypatch.setattr(settings, "OBSERVE_ENABLED", True)
monkeypatch.setattr(settings, "OBSERVE_USER", "u")
monkeypatch.setattr(settings, "OBSERVE_PASSWORD", "p")
assert settings.observe_configured is True
# 缺密码 → False
monkeypatch.setattr(settings, "OBSERVE_PASSWORD", "")
assert settings.observe_configured is False
# 开关关 → False(即便凭证齐全)
monkeypatch.setattr(settings, "OBSERVE_PASSWORD", "p")
monkeypatch.setattr(settings, "OBSERVE_ENABLED", False)
assert settings.observe_configured is False
```
- [ ] **Step 2: 运行,确认失败**
Run: `pytest tests/test_observe.py::test_observe_configured_requires_switch_and_creds -q`
Expected: FAIL —— `AttributeError``settings``OBSERVE_ENABLED` / 无 `observe_configured`
- [ ] **Step 3: 实现配置**
`app/core/config.py``cors_origins_list` property 之后、`is_prod` property 之前插入:
```python
# ===== 可观测(OpenObserve 接口指标)=====
# 采集每个接口的 QPS + 耗时 + 错误率,批量直采到 OpenObserve(本地 Docker)。
# 默认关(prod 安全):未开启 → 中间件透传、worker 不启动,整套 no-op。
# 开启需 ENABLED=true 且 ENDPOINT/USER/PASSWORD 齐全(见 observe_configured)。
OBSERVE_ENABLED: bool = False
OBSERVE_ENDPOINT: str = "http://localhost:5080" # OpenObserve base URL
OBSERVE_ORG: str = "default" # 组织名
OBSERVE_STREAM: str = "app_requests" # stream 名(首次上报自动建)
OBSERVE_USER: str = "" # Basic auth 邮箱
OBSERVE_PASSWORD: str = "" # Basic auth 密码/token
OBSERVE_FLUSH_INTERVAL_SEC: float = 5.0 # worker 最长攒批间隔
OBSERVE_BATCH_MAX: int = 200 # 单批最大事件数
OBSERVE_QUEUE_MAX: int = 10000 # 有界队列上限,满则丢
OBSERVE_TIMEOUT_SEC: float = 5.0 # 上报 HTTP 超时
@property
def observe_configured(self) -> bool:
"""观测上报可用 = 总开关开 且 endpoint/账号/密码齐全(缺则整套 no-op)。"""
return bool(
self.OBSERVE_ENABLED
and self.OBSERVE_ENDPOINT
and self.OBSERVE_USER
and self.OBSERVE_PASSWORD
)
```
- [ ] **Step 4: 运行,确认通过**
Run: `pytest tests/test_observe.py::test_observe_configured_requires_switch_and_creds -q`
Expected: PASS
- [ ] **Step 5: 提交**
```bash
git add app/core/config.py tests/test_observe.py
git commit -m "feat(observe): 加 OBSERVE_* 配置与 observe_configured 门槛"
```
---
## Task 2: 事件队列 + `record_event` + `take_dropped`
**Files:**
- Create: `app/core/observe.py`
- Test: `tests/test_observe.py`(追加)
- [ ] **Step 1: 写失败测试**
`tests/test_observe.py` 顶部 import 区补 `import asyncio``from app.core import observe`,并追加:
```python
def test_record_event_enqueues(monkeypatch):
q = asyncio.Queue(maxsize=10)
monkeypatch.setattr(observe, "_queue", q)
observe.record_event({"route": "/x"})
assert q.get_nowait() == {"route": "/x"}
def test_record_event_drops_when_full(monkeypatch):
q = asyncio.Queue(maxsize=1)
monkeypatch.setattr(observe, "_queue", q)
monkeypatch.setattr(observe, "_dropped", 0)
observe.record_event({"n": 1}) # 占满
observe.record_event({"n": 2}) # 满 → 丢弃当前,不抛异常
assert observe.take_dropped() == 1
assert observe.take_dropped() == 0 # 取出后清零
assert q.get_nowait() == {"n": 1} # 保留的是先到的
```
- [ ] **Step 2: 运行,确认失败**
Run: `pytest tests/test_observe.py -q -k record_event`
Expected: FAIL —— `ModuleNotFoundError: app.core.observe` 或无 `record_event`
- [ ] **Step 3: 实现 `app/core/observe.py`(先只放队列部分)**
> 注意:本步只放队列相关代码。中间件用到的 `os`/`time`/`Match` 及 `_SKIP_PATHS`/`_UNMATCHED`/`_SERVICE` 常量放到 Task 3 一并加入——否则本步提交时 ruff 会报 F401 未用导入。
新建 `app/core/observe.py`
```python
"""接口指标埋点:有界事件队列 + 纯 ASGI 中间件。
每个 HTTP 请求测总耗时、抓路由模板 + 状态码,非阻塞塞进有界队列;由 observe_worker
后台批量上报到 OpenObserve。请求路径上无任何 I/O。未配置观测时中间件直接透传。
"""
from __future__ import annotations
import asyncio
from app.core.config import settings
# 有界事件队列(懒创建,见 get_queue):首次取用时在运行中的 loop 里建,避免 import 期
# 无 loop 的边角问题;put_nowait/get_nowait 不需运行中的 loop → 可在无 loop 下测试。
_queue: asyncio.Queue[dict] | None = None
# 队列满时的丢弃计数,worker 定期取出打日志。
_dropped = 0
def get_queue() -> asyncio.Queue[dict]:
"""返回全局有界事件队列(懒创建)。测试可 monkeypatch 模块级 _queue 换成小队列。"""
global _queue
if _queue is None:
_queue = asyncio.Queue(maxsize=settings.OBSERVE_QUEUE_MAX)
return _queue
def take_dropped() -> int:
"""取出并清零累计丢弃数(供 worker 打点)。"""
global _dropped
n, _dropped = _dropped, 0
return n
def record_event(event: dict) -> None:
"""非阻塞入队;队列满则丢弃当前事件并计数。永不抛异常、永不阻塞请求。"""
global _dropped
try:
get_queue().put_nowait(event)
except asyncio.QueueFull:
_dropped += 1
```
- [ ] **Step 4: 运行,确认通过**
Run: `pytest tests/test_observe.py -q -k record_event`
Expected: PASS
- [ ] **Step 5: 提交**
```bash
git add app/core/observe.py tests/test_observe.py
git commit -m "feat(observe): 加有界事件队列与 record_event(满则丢)"
```
---
## Task 3: `RequestMetricsMiddleware`(路由模板 + 状态码 + 耗时)
**Files:**
- Modify: `app/core/observe.py`(追加 `_resolve_route` 和中间件 class
- Test: `tests/test_observe.py`(追加)
- [ ] **Step 1: 写失败测试**
`tests/test_observe.py` 顶部 import 区补:
```python
import pytest
from fastapi import FastAPI
from fastapi.testclient import TestClient
```
并追加:
```python
def _make_probe_app() -> FastAPI:
"""独立最小 app:只挂中间件 + 两个无鉴权路由,不碰真业务 DB/auth。"""
app = FastAPI()
app.add_middleware(observe.RequestMetricsMiddleware)
@app.get("/things/{tid}")
def get_thing(tid: str):
return {"tid": tid}
@app.get("/health")
def health():
return {"ok": True}
return app
@pytest.fixture
def observe_on(monkeypatch):
"""开启观测 + 换一个干净小队列,返回该队列供断言。"""
monkeypatch.setattr(settings, "OBSERVE_ENABLED", True)
monkeypatch.setattr(settings, "OBSERVE_USER", "u")
monkeypatch.setattr(settings, "OBSERVE_PASSWORD", "p")
q = asyncio.Queue(maxsize=100)
monkeypatch.setattr(observe, "_queue", q)
return q
def test_middleware_records_route_template(observe_on):
client = TestClient(_make_probe_app())
r = client.get("/things/42")
assert r.status_code == 200
evt = observe_on.get_nowait()
assert evt["route"] == "/things/{tid}" # 模板,不是 /things/42
assert evt["method"] == "GET"
assert evt["status"] == 200
assert evt["duration_ms"] >= 0
assert evt["service"] and "env" in evt and isinstance(evt["_timestamp"], int)
def test_middleware_skips_health(observe_on):
client = TestClient(_make_probe_app())
client.get("/health")
assert observe_on.empty()
def test_middleware_unmatched_route_is_normalized(observe_on):
client = TestClient(_make_probe_app())
r = client.get("/definitely-not-a-route")
assert r.status_code == 404
evt = observe_on.get_nowait()
assert evt["route"] == "__unmatched__"
assert evt["status"] == 404
def test_middleware_noop_when_disabled(monkeypatch):
monkeypatch.setattr(settings, "OBSERVE_ENABLED", False)
q = asyncio.Queue(maxsize=100)
monkeypatch.setattr(observe, "_queue", q)
client = TestClient(_make_probe_app())
client.get("/things/1")
assert q.empty() # 未配置观测 → 零入队
```
- [ ] **Step 2: 运行,确认失败**
Run: `pytest tests/test_observe.py -q -k middleware`
Expected: FAIL —— `AttributeError: module 'app.core.observe' has no attribute 'RequestMetricsMiddleware'`
- [ ] **Step 3a: 给 `app/core/observe.py` 补中间件用的导入与常量**
把顶部 import 段从
```python
from __future__ import annotations
import asyncio
from app.core.config import settings
```
改成
```python
from __future__ import annotations
import asyncio
import os
import time
from starlette.routing import Match
from app.core.config import settings
# 不采集的路径(纯噪音):健康检查。
_SKIP_PATHS = frozenset({"/health"})
# 未匹配路由(404/扫描器)归一到此,防维度爆炸。
_UNMATCHED = "__unmatched__"
# service 字段:与 logging.py 同源(LOG_SERVICE_NAME),默认 app-server。
_SERVICE = os.getenv("LOG_SERVICE_NAME", "app-server")
```
`_queue` / `_dropped` / `get_queue` / `take_dropped` / `record_event` 保持不动。)
- [ ] **Step 3b: 实现中间件(追加到 `app/core/observe.py` 末尾)**
```python
def _resolve_route(scope) -> str:
"""从 scope 取路由模板(如 /things/{tid})。优先 scope['route'](现代 Starlette
路由后写入);取不到则手动匹配一次(老版本兜底);仍无 → __unmatched__(404/扫描器)。"""
route = scope.get("route")
path = getattr(route, "path", None)
if path:
return path
app_ = scope.get("app")
router = getattr(app_, "router", None)
for candidate in getattr(router, "routes", []):
try:
match, _ = candidate.matches(scope)
except Exception: # noqa: BLE001 - 匹配兜底,任一路由异常不影响整体
continue
if match == Match.FULL and getattr(candidate, "path", None):
return candidate.path
return _UNMATCHED
class RequestMetricsMiddleware:
"""纯 ASGI 中间件:测每个 http 请求耗时,记 method/route/status/duration。
放在最外层(main.py 里 CORS 之后 add),测到含 CORS 的完整耗时。未配置观测 → 透传。
"""
def __init__(self, app) -> None:
self.app = app
async def __call__(self, scope, receive, send) -> None:
if scope["type"] != "http" or not settings.observe_configured:
await self.app(scope, receive, send)
return
if scope.get("path") in _SKIP_PATHS:
await self.app(scope, receive, send)
return
start = time.perf_counter()
status_holder = {"status": 500} # 下游异常未产出 response 时兜底 500
async def send_wrapper(message) -> None:
if message["type"] == "http.response.start":
status_holder["status"] = message["status"]
await send(message)
try:
await self.app(scope, receive, send_wrapper)
finally:
duration_ms = (time.perf_counter() - start) * 1000.0
record_event({
"_timestamp": int(time.time() * 1_000_000), # µs,OpenObserve 时间列
"service": _SERVICE,
"env": settings.APP_ENV,
"method": scope.get("method", ""),
"route": _resolve_route(scope),
"status": status_holder["status"],
"duration_ms": round(duration_ms, 3),
})
```
- [ ] **Step 4: 运行,确认通过**
Run: `pytest tests/test_observe.py -q -k middleware`
Expected: PASS4 个中间件用例全过)
> 若 `test_middleware_records_route_template` 拿到的是 `/things/42` 而非模板,说明该 Starlette 版本未在 `scope["route"]` 写模板——此时 `_resolve_route` 的手动匹配兜底应已生效并返回模板;若仍不对,检查兜底分支是否被 import 顺序影响。
- [ ] **Step 5: 提交**
```bash
git add app/core/observe.py tests/test_observe.py
git commit -m "feat(observe): 加 RequestMetricsMiddleware(路由模板+状态码+耗时)"
```
---
## Task 4: 后台上报 worker
**Files:**
- Create: `app/core/observe_worker.py`
- Test: `tests/test_observe.py`(追加)
- [ ] **Step 1: 写失败测试**
`tests/test_observe.py` 顶部 import 区补:
```python
import httpx
from app.core import observe_worker
```
并追加:
```python
async def test_collect_batch_drains_up_to_batch_max(monkeypatch):
q = asyncio.Queue(maxsize=100)
monkeypatch.setattr(observe, "_queue", q)
monkeypatch.setattr(settings, "OBSERVE_FLUSH_INTERVAL_SEC", 0.1)
monkeypatch.setattr(settings, "OBSERVE_BATCH_MAX", 200)
for i in range(3):
q.put_nowait({"n": i})
batch = await observe_worker._collect_batch()
assert [e["n"] for e in batch] == [0, 1, 2]
async def test_collect_batch_timeout_returns_empty(monkeypatch):
q = asyncio.Queue(maxsize=100)
monkeypatch.setattr(observe, "_queue", q)
monkeypatch.setattr(settings, "OBSERVE_FLUSH_INTERVAL_SEC", 0.05)
batch = await observe_worker._collect_batch()
assert batch == []
async def test_post_batch_hits_json_ingest_url(monkeypatch):
monkeypatch.setattr(settings, "OBSERVE_ORG", "default")
monkeypatch.setattr(settings, "OBSERVE_STREAM", "app_requests")
captured = {}
def handler(request: httpx.Request) -> httpx.Response:
captured["url"] = str(request.url)
captured["json"] = request.content
return httpx.Response(200, json={"code": 200})
client = httpx.AsyncClient(
base_url="http://oo", transport=httpx.MockTransport(handler)
)
await observe_worker._post_batch(client, [{"route": "/x", "status": 200}])
await client.aclose()
assert captured["url"] == "http://oo/api/default/app_requests/_json"
assert b"/x" in captured["json"]
def test_start_observe_worker_noop_when_not_configured(monkeypatch):
monkeypatch.setattr(settings, "OBSERVE_ENABLED", False)
assert observe_worker.start_observe_worker() is None
```
- [ ] **Step 2: 运行,确认失败**
Run: `pytest tests/test_observe.py -q -k "collect_batch or post_batch or start_observe"`
Expected: FAIL —— `ModuleNotFoundError: app.core.observe_worker`
- [ ] **Step 3: 实现 `app/core/observe_worker.py`**
新建 `app/core/observe_worker.py`
```python
"""接口指标后台上报 worker:批量 drain 事件队列 → POST 到 OpenObserve。
对齐 heartbeat_monitor_worker 等的 start_*/stop_* 形态。best-effort 遥测:catch 全部
异常,上报失败直接丢批不重试。未配置观测 → start 返回 None(不启动),整套 no-op。
"""
from __future__ import annotations
import asyncio
import contextlib
import logging
import httpx
from app.core.config import settings
from app.core.observe import get_queue, take_dropped
logger = logging.getLogger("shagua.observe")
# 上报用的 httpx client,start 时建、stop 时关。
_client: httpx.AsyncClient | None = None
async def _collect_batch() -> list[dict]:
"""等到 ≥1 条(或到 flush 间隔)后,连抽到 BATCH_MAX 条或抽空。超时且空 → 返回 []。"""
queue = get_queue()
batch: list[dict] = []
try:
first = await asyncio.wait_for(
queue.get(), timeout=settings.OBSERVE_FLUSH_INTERVAL_SEC
)
except asyncio.TimeoutError:
return batch
batch.append(first)
while len(batch) < settings.OBSERVE_BATCH_MAX:
try:
batch.append(queue.get_nowait())
except asyncio.QueueEmpty:
break
return batch
async def _post_batch(client: httpx.AsyncClient, batch: list[dict]) -> None:
"""POST 一批事件到 OpenObserve 的 _json ingest 端点。非 2xx 仅告警。"""
url = f"/api/{settings.OBSERVE_ORG}/{settings.OBSERVE_STREAM}/_json"
resp = await client.post(url, json=batch)
if resp.status_code >= 300:
logger.warning(
"observe ingest failed status=%s body=%s",
resp.status_code,
resp.text[:200],
)
async def _run_loop(client: httpx.AsyncClient) -> None:
try:
while True:
batch = await _collect_batch()
dropped = take_dropped()
if dropped:
logger.warning("observe dropped %d events (queue full)", dropped)
if not batch:
continue
try:
await _post_batch(client, batch)
except Exception: # noqa: BLE001 - best-effort 遥测,失败丢批不重试、不退出
logger.warning(
"observe post batch failed, dropped %d events",
len(batch),
exc_info=True,
)
except asyncio.CancelledError:
logger.info("observe worker stopped")
raise
def start_observe_worker() -> asyncio.Task | None:
"""启动上报 worker。未配置观测 → 返回 None(no-op)。"""
global _client
if not settings.observe_configured:
return None
_client = httpx.AsyncClient(
base_url=settings.OBSERVE_ENDPOINT,
auth=(settings.OBSERVE_USER, settings.OBSERVE_PASSWORD),
timeout=settings.OBSERVE_TIMEOUT_SEC,
)
logger.info(
"observe worker started endpoint=%s org=%s stream=%s",
settings.OBSERVE_ENDPOINT,
settings.OBSERVE_ORG,
settings.OBSERVE_STREAM,
)
return asyncio.create_task(_run_loop(_client), name="observe-worker")
async def stop_observe_worker(task: asyncio.Task | None) -> None:
"""收尾:cancel worker → best-effort 发最后一批 → 关 client。"""
global _client
if task is None:
return
task.cancel()
with contextlib.suppress(asyncio.CancelledError):
await task
if _client is not None:
# worker 已停,安全 drain 剩余并 best-effort 发最后一批(短超时,不卡关停)。
try:
queue = get_queue()
final: list[dict] = []
while len(final) < settings.OBSERVE_BATCH_MAX:
try:
final.append(queue.get_nowait())
except asyncio.QueueEmpty:
break
if final:
await asyncio.wait_for(
_post_batch(_client, final), timeout=settings.OBSERVE_TIMEOUT_SEC
)
except Exception: # noqa: BLE001 - 关停期尽力而为,失败忽略
pass
await _client.aclose()
_client = None
```
- [ ] **Step 4: 运行,确认通过**
Run: `pytest tests/test_observe.py -q -k "collect_batch or post_batch or start_observe"`
Expected: PASS
- [ ] **Step 5: 提交**
```bash
git add app/core/observe_worker.py tests/test_observe.py
git commit -m "feat(observe): 加后台批量上报 worker(失败丢批不重试)"
```
---
## Task 5: 接线到 `app/main.py`(挂中间件 + lifespan 启停)
**Files:**
- Modify: `app/main.py`import 区、lifespan、CORS 之后)
- Test: `tests/test_observe.py`(追加)
- [ ] **Step 1: 写失败测试**
`tests/test_observe.py` 追加:
```python
def test_app_has_metrics_middleware():
from app.main import app
names = [m.cls.__name__ for m in app.user_middleware]
assert "RequestMetricsMiddleware" in names
```
- [ ] **Step 2: 运行,确认失败**
Run: `pytest tests/test_observe.py::test_app_has_metrics_middleware -q`
Expected: FAIL —— 断言失败(中间件尚未挂载)
- [ ] **Step 3: 实现接线**
3a. 在 `app/main.py` import 区(`withdraw_reconcile_worker` import 块之后)加:
```python
from app.core.observe import RequestMetricsMiddleware
from app.core.observe_worker import (
start_observe_worker,
stop_observe_worker,
)
```
3b. lifespan 里加启停(现有 `daily_exchange_task = start_daily_exchange_worker()` 之后、`try:` 之前加一行;`finally` 里在 `stop_daily_exchange_worker` 之后加一行):
```python
daily_exchange_task = start_daily_exchange_worker()
observe_task = start_observe_worker()
try:
yield
finally:
await stop_heartbeat_monitor(heartbeat_task)
await stop_withdraw_reconcile_worker(reconcile_task)
await stop_daily_exchange_worker(daily_exchange_task)
await stop_observe_worker(observe_task)
await aclose_pricebot_client()
logger.info("shutting down")
```
3c. 挂中间件——在 CORS 的 `if settings.cors_origins_list:` 整块之后加(使其成为最外层,测到含 CORS 的完整耗时):
```python
# 接口指标埋点(最外层:测含 CORS 的完整耗时)。未配置观测时中间件自 no-op。
app.add_middleware(RequestMetricsMiddleware)
```
- [ ] **Step 4: 运行,确认通过**
Run: `pytest tests/test_observe.py::test_app_has_metrics_middleware -q`
Expected: PASS
- [ ] **Step 5: 跑整套 observe 测试 + 全量回归,确认无破坏**
Run: `pytest tests/test_observe.py -q && pytest -q`
Expected: 全 PASS(现有用例不受影响:conftest 未设 `OBSERVE_*` → 观测关 → worker no-op、中间件透传)
- [ ] **Step 6: 提交**
```bash
git add app/main.py tests/test_observe.py
git commit -m "feat(observe): main.py 挂中间件 + lifespan 启停上报 worker"
```
---
## Task 6: OpenObserve 本地部署(compose + README + .env.example
**Files:**
- Create: `deploy/openobserve/docker-compose.yml`
- Create: `deploy/openobserve/README.md`
- Modify: `.env.example`(追加 `OBSERVE_*` 段)
- [ ] **Step 1: 写 docker-compose**
新建 `deploy/openobserve/docker-compose.yml`
```yaml
# 本地开发用 OpenObserve(单容器 = local 模式)。用于接收 app-server 的接口指标。
# 启动: cd deploy/openobserve && docker compose up -d
# Web UI: http://localhost:5080 (账号见下方 env)
services:
openobserve:
image: public.ecr.aws/zinclabs/openobserve:latest
container_name: openobserve
ports:
- "5080:5080"
environment:
ZO_ROOT_USER_EMAIL: "admin@shaguabijia.local"
ZO_ROOT_USER_PASSWORD: "Complexpass#123"
ZO_DATA_DIR: "/data"
volumes:
- ./data:/data
restart: unless-stopped
```
- [ ] **Step 2: 写 README**
新建 `deploy/openobserve/README.md`
````markdown
# OpenObserve 本地部署(接口 QPS / 耗时可观测)
app-server 通过中间件采集每个接口的 QPS + 耗时 + 错误率,批量上报到这里。
## 启动
```bash
cd deploy/openobserve
docker compose up -d
```
- Web UIhttp://localhost:5080
- 登录:`admin@shaguabijia.local` / `Complexpass#123`(见 `docker-compose.yml`
- 数据落 `deploy/openobserve/data/`(已挂卷持久化;`data/` 建议 gitignore
## 让 app-server 上报
在项目根的 `.env` 打开观测(`OBSERVE_*`,账号密码与 compose 里 root 一致):
```dotenv
OBSERVE_ENABLED=true
OBSERVE_ENDPOINT=http://localhost:5080
OBSERVE_ORG=default
OBSERVE_STREAM=app_requests
OBSERVE_USER=admin@shaguabijia.local
OBSERVE_PASSWORD=Complexpass#123
```
重启 app-server,随便打几个接口。stream `app_requests` **首次上报自动创建**
在 UI 的 Logs → 选 `app_requests` 就能看到逐条请求事件。
## 查询(Logs 页 SQL,或建 Dashboard 面板)
各接口 QPS(1 分钟分桶,面板里再除 60 得每秒):
```sql
SELECT route, histogram(_timestamp, '1 minute') AS ts, count(*) AS cnt
FROM app_requests GROUP BY route, ts ORDER BY ts
```
各接口 P95 耗时(毫秒):
```sql
SELECT route, approx_percentile_cont(duration_ms, 0.95) AS p95_ms
FROM app_requests GROUP BY route ORDER BY p95_ms DESC
```
各接口错误率(5xx 占比):
```sql
SELECT route,
count(*) FILTER (WHERE status >= 500) * 100.0 / count(*) AS err_pct
FROM app_requests GROUP BY route ORDER BY err_pct DESC
```
## 停止 / 清数据
```bash
docker compose down # 停止(保留数据)
docker compose down -v && rm -rf data # 停止并清空数据
```
> 生产部署(持久化规格、独立 ingest 账号、鉴权收紧)见 spec 第 9 节,本期不做。
````
- [ ] **Step 3: 追加 `.env.example`**
在 `.env.example` 末尾追加:
```dotenv
# ===== 可观测(OpenObserve 接口指标)=====
# 采集每个接口 QPS + 耗时 + 错误率,批量直采到 OpenObserve(本地 Docker,见 deploy/openobserve/)。
# 默认关;开启需 ENABLED=true 且填 USER/PASSWORD(与 docker-compose 里 root 账号一致)。
OBSERVE_ENABLED=false
OBSERVE_ENDPOINT=http://localhost:5080
OBSERVE_ORG=default
OBSERVE_STREAM=app_requests
OBSERVE_USER=admin@shaguabijia.local
OBSERVE_PASSWORD=Complexpass#123
# 进阶(一般不用改):攒批间隔秒 / 单批最大条数 / 有界队列上限(满则丢) / 上报超时秒
OBSERVE_FLUSH_INTERVAL_SEC=5
OBSERVE_BATCH_MAX=200
OBSERVE_QUEUE_MAX=10000
OBSERVE_TIMEOUT_SEC=5
```
- [ ] **Step 4: 校验 compose 语法(不需真拉镜像)**
Run: `docker compose -f deploy/openobserve/docker-compose.yml config`
Expected: 打印规整后的配置、无报错(若本机无 docker,可跳过,标注为手动验证项)
- [ ] **Step 5: 提交**
```bash
git add deploy/openobserve/docker-compose.yml deploy/openobserve/README.md .env.example
git commit -m "feat(observe): 加 OpenObserve 本地 compose + README + .env.example"
```
---
## Task 7: 端到端手动验证 + 全量 lint/test 收尾
**Files:** 无(验证 + 收尾)
- [ ] **Step 1: 起 OpenObserve**
Run: `cd deploy/openobserve && docker compose up -d`
Expected: 容器起来,浏览器打开 http://localhost:5080 能登录
- [ ] **Step 2: 本地开观测起 app-server**
在根 `.env` 设 `OBSERVE_ENABLED=true` + `OBSERVE_USER/PASSWORD`(同 compose),然后:
Run: `./run.sh`Windows 用 `python -m uvicorn app.main:app --port 8770 --reload --reload-dir app`
Expected: 启动日志出现 `observe worker started endpoint=http://localhost:5080 ...`
- [ ] **Step 3: 打几个接口产生数据**
Run: `curl http://localhost:8770/health && curl http://localhost:8770/things-does-not-exist -i`(或正常业务接口若干)
Expected: 稍等 ≤5sflush 间隔),OpenObserve UI 的 Logs → `app_requests` 出现事件;`/health` 不应出现;不存在的路径 route 为 `__unmatched__`
- [ ] **Step 4: 验证三条查询**
在 OpenObserve UI 分别粘贴 README 里的 QPS / P95 / 错误率 SQL,确认能出数。
- [ ] **Step 5: lint(仅本改动涉及文件)+ 全量测试**
> 说明:仓库基线有 ~558 个既有 ruff 错误、且未强制 ruff 通过。不要去清历史欠债(范围蔓延)。只要求**本次新增/改动的文件**零 ruff 错误。
Run: `ruff check app/core/observe.py app/core/observe_worker.py tests/test_observe.py && python -m pytest -q`
Expected: 上述三个新文件 ruff 无错;测试里 `tests/test_observe.py` 全 PASS,且**全量失败数不超过基线的 4 个**(test_compare_proxy ×2 / test_coupon_proxy ×1 / test_invite ×1,均与本功能无关)。
额外确认我对既有文件的改动没有引入**新的** ruff 错误:`ruff check app/core/config.py app/main.py`(数量应与基线一致,不因本改动增加)。
- [ ] **Step 6: 关观测复跑一次,确认降级**
把 `.env` 的 `OBSERVE_ENABLED` 改回 `false``ruff check .` 不涉及,直接 `pytest -q`
Expected: 全 PASS(验证 observe 关闭时零副作用)
- [ ] **Step 7: 收尾提交(如有 .env 之外的改动)**
```bash
git add -A
git commit -m "chore(observe): 端到端验证与收尾" --allow-empty
```
> `.env` 不入 git(已 gitignore);本任务只验证,不提交 `.env`。
---
## Self-Review(写完计划后自查)
- **Spec 覆盖**Docker 部署→Task 6/7;事件 schema→Task 3`record_event` 事件字段);中间件→Task 3worker→Task 4;配置→Task 1main 接线→Task 5;查询/仪表盘→Task 6 README;测试→Task 1-5;决策(a)队列满丢→Task 2(b)失败不重试→Task 4(c)跳过 /health→Task 3。全覆盖。
- **占位符**:无 TBD/TODO;每个代码步骤含完整代码。
- **类型/命名一致**`get_queue` / `record_event` / `take_dropped` / `RequestMetricsMiddleware` / `start_observe_worker` / `stop_observe_worker` / `observe_configured` / 事件字段名,跨 Task 1-5 与文件结构表一致。
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,236 @@
# 接口 QPS + 耗时可观测(OpenObserve)设计
- **日期**2026-07-06
- **状态**:已评审通过,待写实现计划
- **范围**:仅 app-server8770);admin8771)暂不接入
- **方案**:A —— 轻量自研 ASGI 中间件 + 后台 worker 批量直采到 OpenObserve
## 1. 背景与目标
app-server 目前除 CORS 外无任何中间件,也无接口级可观测。需要按**每个接口**采集:
- **QPS**(每秒请求数,可按接口/时间分桶)
- **耗时**P50/P95/P99 等分位)
顺带低成本拿到**错误率**`status >= 500` 占比)。落地目标是:本地 Docker 跑一个 OpenObserve 实例接收数据,服务侧加埋点上报,在 OpenObserve 仪表盘上看各接口 QPS + 耗时。
### 非目标(YAGNI
- 不做分布式 trace / span 关联(只要接口聚合指标)。
- 不引入 OpenTelemetry / Prometheus 客户端等重依赖。
- 不采集请求体 / query / 用户身份等,任何 PII 都不进上报。
- admin(8771)本期不接(中间件写成可复用,未来一行挂载即可)。
- 上报失败不做持久化重试 / 落盘补偿(best-effort)。
## 2. 方案选型
对比过三条路(详见评审记录):
- **A 轻量自研中间件 + JSON 直采**(选中):零新依赖(`httpx` 已在依赖里),完全贴合本仓库「后台 worker + JSON 事件 + `*_configured` 优雅降级」的既有习惯,恰好满足「每接口 QPS + 耗时 + 错误率」并保留原始事件下钻能力。
- B OpenTelemetry 自动埋点 + OTLP:行业标准、顺带 trace,但多 5–6 个依赖、概念多、数据量/成本高于需求,与精简代码库风格相悖。
- C Prometheus 进程内聚合 + remote_write/抓取:数据量最小,但 remote_write 编码复杂或需额外抓取进程,丢失单请求下钻,最不贴合 OpenObserve 的 log-first 强项。
**结论:A。**
## 3. 架构与数据流
```
每个 HTTP 请求
→ RequestMetricsMiddleware(最外层:测总耗时 / 抓路由模板 + 状态码)
→ record_event() 非阻塞入队(有界队列,满则丢最旧,绝不阻塞、绝不 OOM)
→ observe_worker(后台 asyncio.Task,随 lifespan 启停)批量 drain
→ httpx POST {ENDPOINT}/api/{ORG}/{STREAM}/_json → OpenObserve
→ 仪表盘 SQL 聚合出 QPS / 分位耗时 / 错误率
```
**核心不变量**
1. 请求路径上只做「测时 + 构建一个小 dict + `put_nowait`」,**无任何网络/磁盘 I/O**。
2. 所有上报 I/O 在后台 worker;worker 捕获全部异常,绝不让埋点影响请求。
3. 未配置观测(`observe_configured=False`)→ 中间件透传、worker 不启动,整套 no-op。
4. OpenObserve 不可用 → 队列填满后丢弃事件 + 限流告警,业务零影响。
## 4. 组件设计
### 4.1 OpenObserve 本地部署 —— `deploy/openobserve/docker-compose.yml`(新增)
```yaml
services:
openobserve:
image: public.ecr.aws/zinclabs/openobserve:latest
container_name: openobserve
ports: ["5080:5080"]
environment:
ZO_ROOT_USER_EMAIL: "admin@shaguabijia.local"
ZO_ROOT_USER_PASSWORD: "Complexpass#123"
ZO_DATA_DIR: "/data"
volumes: ["./data:/data"]
restart: unless-stopped
```
- `docker compose up -d` 启动;Web UI `http://localhost:5080`,用上面邮箱/密码登录。
- 单容器 = local 模式,数据落 `./data`(已挂卷持久化)。
- **stream 首次上报自动创建**,无需预建 `app_requests`
- 上报鉴权:HTTP Basic auth`email:password`),本地直接用 root 账号;生产应另建仅具 ingest 权限的用户/服务账号(本期不涉及)。
### 4.2 事件 schema(一请求一行 JSON
```json
{
"_timestamp": 1720000000000000, // 微秒(µs)整数,请求完成时刻。OpenObserve 默认时间列 _timestamp 以微秒计
"service": "app-server", // 取 LOG_SERVICE_NAME / 固定值
"env": "dev", // settings.APP_ENV
"method": "POST",
"route": "/api/v1/coupon/step", // 路由模板(非实际 path
"status": 200,
"duration_ms": 42.7 // float 毫秒
}
```
- **只存路由模板**(如 `/c/{code}``/media` 静态归一),避免 path 参数把维度打爆。
- 未匹配路由(404 / 扫描器)归一到常量 `__unmatched__`
- 只采 method / route / status / duration —— 无 body、无 query、无 PII。
### 4.3 埋点中间件 —— `app/core/observe.py`(新增)
**纯 ASGI 中间件**(比 `BaseHTTPMiddleware` 开销低;能可靠读到路由与最终状态码;scope 按引用透传,内层 router 的 `scope["route"]` 外层可见)。
职责:
1.`http` 请求、或 `not settings.observe_configured` → 直接透传,不测。
2. `perf_counter()` 记起点;包一层 `send``http.response.start``status`(默认兜底 500,覆盖下游抛异常未产出 response 的情况)。
3. `finally` 里算 `duration_ms`,从 `scope` 取路由模板(见下),构建事件,调 `record_event()`
4. 跳过路径集合 `_SKIP_PATHS = {"/health"}`(纯噪音)。
**路由模板解析(跨 Starlette 版本稳健)**
```python
route = scope.get("route")
template = getattr(route, "path", None)
if template is None: # 未匹配 / 老版本未写 scope["route"]
template = "__unmatched__"
```
(若实测某 Starlette 版本不写 `scope["route"]`,回退用 `request.app.router.routes` 逐个 `route.matches(scope)==Match.FULL` 找模板;实现时以实际版本为准,优先 `scope["route"]`。)
**入队(`record_event`**:模块级 `asyncio.Queue(maxsize=OBSERVE_QUEUE_MAX)`。用 `put_nowait``QueueFull` 则丢弃并累加一个 `_dropped` 计数(每累计 N 条限流打一条 WARNING)。**永不 `await put()`、永不阻塞请求**。
> 决策(a):队列满 → **丢弃**(不阻塞请求)。
### 4.4 上报 worker —— `app/core/observe_worker.py`(新增)
对齐现有 `heartbeat_monitor_worker.py` / `daily_exchange_worker.py` / `withdraw_reconcile_worker.py``start_*` / `stop_*` 形态。
- `start_observe_worker() -> asyncio.Task | None`
- `not settings.observe_configured` → 返回 `None`no-op)。
- 否则建专用 `httpx.AsyncClient``base_url=ENDPOINT``auth=(USER, PASSWORD)``timeout=OBSERVE_TIMEOUT_SEC`),起 `_run_loop` task。
- `_run_loop()`:循环
1. `_collect_batch()``await asyncio.wait_for(queue.get(), timeout=FLUSH_INTERVAL)` 拿到首条(超时且空 → 返回空,continue);再 `get_nowait()` 连抽到 `BATCH_MAX` 条或抽空。
2. `POST /api/{ORG}/{STREAM}/_json`body 为事件数组。
3. **catch 所有异常**:失败限流打 WARNING,**直接丢弃该批,不重试**。
- `stop_observe_worker(task)`best-effort 收尾 flush(短超时)→ `task.cancel()``await`(吞 `CancelledError`)→ 关 client。
> 决策(b):上报失败 → **直接丢弃,不重试**(best-effort 遥测)。
### 4.5 配置 —— `app/core/config.py`(改)
新增一段 `# ===== 可观测(OpenObserve 接口指标)=====`,默认全关(prod 安全):
| 配置 | 默认 | 说明 |
|---|---|---|
| `OBSERVE_ENABLED` | `False` | 总开关;默认关,opt-in |
| `OBSERVE_ENDPOINT` | `http://localhost:5080` | OpenObserve base URL |
| `OBSERVE_ORG` | `default` | 组织名 |
| `OBSERVE_STREAM` | `app_requests` | stream 名 |
| `OBSERVE_USER` | `""` | Basic auth 邮箱 |
| `OBSERVE_PASSWORD` | `""` | Basic auth 密码/token |
| `OBSERVE_FLUSH_INTERVAL_SEC` | `5.0` | worker 最长攒批间隔 |
| `OBSERVE_BATCH_MAX` | `200` | 单批最大事件数 |
| `OBSERVE_QUEUE_MAX` | `10000` | 有界队列上限,满则丢 |
| `OBSERVE_TIMEOUT_SEC` | `5.0` | 上报 HTTP 超时 |
```python
@property
def observe_configured(self) -> bool:
return bool(self.OBSERVE_ENABLED and self.OBSERVE_ENDPOINT
and self.OBSERVE_USER and self.OBSERVE_PASSWORD)
```
`.env.example` 同步补一段带注释的 `OBSERVE_*`(沿用该文件重注释风格),`OBSERVE_ENABLED=false`
### 4.6 接线 —— `app/main.py`(改)
- import `RequestMetricsMiddleware``start_observe_worker` / `stop_observe_worker`
- `app.add_middleware(RequestMetricsMiddleware)`:放在 CORS `add_middleware` **之后** → 成为最外层,测到含 CORS 的完整耗时。无条件挂载(内部自 no-op)。
- `lifespan`:启动 `observe_task = start_observe_worker()``finally``await stop_observe_worker(observe_task)`,与现有 worker 并列。
### 4.7 OpenObserve 查询 / 仪表盘 —— `deploy/openobserve/README.md`(新增)
含:compose 启停、登录、stream 自动创建说明、`.env` 接线,以及可直接粘的示例 SQL
- **各接口 QPS**(1 分钟分桶):
```sql
SELECT route, histogram(_timestamp, '1 minute') AS ts, count(*) AS cnt
FROM app_requests GROUP BY route, ts ORDER BY ts
```
(面板按 `cnt/60` 展示每秒;或用 OpenObserve 图表的 rate 能力。)
- **各接口 P95 耗时**
```sql
SELECT route, approx_percentile_cont(duration_ms, 0.95) AS p95_ms
FROM app_requests GROUP BY route ORDER BY p95_ms DESC
```
- **各接口错误率**
```sql
SELECT route,
count(*) FILTER (WHERE status >= 500) * 100.0 / count(*) AS err_pct
FROM app_requests GROUP BY route ORDER BY err_pct DESC
```
## 5. 关键设计决策汇总
- **(a) 队列满 → 丢弃**(不阻塞请求):遥测让路于业务可用性。
- **(b) 上报失败 → 不重试**best-effort;避免 poison batch 堆积与队列无限增长。
- **(c) 跳过 `/health`**:健康检查是纯噪音,硬编码在 `_SKIP_PATHS`。
- **只存路由模板 + `__unmatched__`**:防维度爆炸。
- **默认 OFF、opt-in**prod 安全默认;开启后仍全异步 + 有界。
- **纯 ASGI 中间件 + `perf_counter`**:请求路径开销微秒级,无 I/O。
## 6. 安全 / 性能保证
- 请求路径新增开销 ≈ 一次 `perf_counter` 差 + 一个小 dict + 一次 `put_nowait`(微秒级),无锁竞争的显著热点。
- 失败隔离:入队丢弃 + worker 全异常捕获;OpenObserve 宕机不影响任何请求。
- 有界内存:队列 `maxsize` 封顶,最坏丢事件不涨内存。
- 无 PII:仅 method / route / status / duration。
## 7. 测试策略 —— `tests/test_observe.py`(新增)
沿用仓库约定(`TestClient` + `monkeypatch`,绝不打真网络;`conftest` 在 import 前设 env):
1. 埋点入队字段正确:模板路由、`status`、`duration_ms > 0`。
2. 参数化路由 → 取到**模板**而非实际 path。
3. 未匹配路径(404)→ `route == "__unmatched__"`。
4. `OBSERVE_ENABLED=false` → 零入队、零 HTTP(现有测试不受影响)。
5. 队列满 → `record_event` 不抛异常(走丢弃分支)。
6. worker 批量 POST 的 URL / payload 正确(monkeypatch httpx client / `_post`,不打网络)。
7. `/health` 被跳过 → 不入队。
> `settings` 是 `lru_cache` 单例;需要开启观测的用例通过 monkeypatch `settings` 属性或直接调 `record_event` / 中间件并 patch `observe_configured` 实现,避免全局 env 改动波及他用例。
## 8. 文件清单
| 文件 | 动作 |
|---|---|
| `deploy/openobserve/docker-compose.yml` | 新增(OpenObserve 容器)|
| `deploy/openobserve/README.md` | 新增(部署步骤 + 查询/仪表盘)|
| `app/core/observe.py` | 新增(中间件 + 有界队列 + `record_event` + 路由解析)|
| `app/core/observe_worker.py` | 新增(后台批量上报 worker)|
| `app/core/config.py` | 改(`OBSERVE_*` + `observe_configured`|
| `app/main.py` | 改(挂中间件 + lifespan 启停 worker|
| `.env.example` | 改(新增 `OBSERVE_*` 注释段)|
| `tests/test_observe.py` | 新增 |
## 9. 未来工作(本期不做)
- admin(8771)接入同一套中间件(`service` 字段区分)。
- 生产部署 OpenObserve(持久化、独立 ingest 账号、资源规格、鉴权收紧)。
- 上报字段扩展(如按 user/设备维度、上游 pricebot 透传耗时拆分)。
@@ -0,0 +1,296 @@
# 15 天不活跃自动清零(金币 + 现金)设计
- **日期**2026-07-16
- **状态**Draft — 待评审
- **所属**app-server`app/`),含一处 admin 侧重构 + 一项 Android 端埋点依赖
- **一句话**:连续 15 天不活跃的用户,自动清零其金币与现金;清零前按可配置节奏预警;全过程留审计以备纠纷排查。
---
## 1. 背景与目标
运营需要对**长期不活跃**用户的钱包余额做清理。两条硬性要求:
1. **可审计**:记录清零原因与**清零前的三桶余额**,便于后续排查与处理客户纠纷。
2. **临清预警**:在临近清零前推送信息告知用户"因账号不活跃,账户里的 xx 金币和 xx 现金将被清零"。
### 非目标(本期不做)
- 不做真实推送通道(极光 JPush / 短信)的对接 —— 仅做**可插拔通知器 + 日志占位**,接口预留、后续无缝替换。
- 不改动提现(`WithdrawOrder`)流程。
- 不新增 `User.last_active_at` 列、不改鉴权热路径。
---
## 2. 需求
| # | 需求 | 落地 |
|---|---|---|
| R1 | 连续 15 天不活跃 → 清零金币 + 现金 | 每日 worker 扫描 + 逐用户事务清零(§6) |
| R2 | 记录清零原因 + 清零前余额 | `inactivity_reset_log` 审计表 + 3 条钱包流水(§5、§7) |
| R3 | 临清前预警"xx 金币 xx 现金将清零" | 阶段 A 预警 + `inactivity_notification_log`(§6、§7 |
| R4 | 活跃口径与"用户管理"一致 | 抽共享模块 `activity.py`admin 与 worker 共用(§4、§12 |
| R5 | 预警时机完全可配置 | `INACTIVITY_*` 配置项(§8 |
---
## 3. 决策记录(来自评审问答)
| 决策点 | 结论 | 理由 |
|---|---|---|
| **活跃口径** | 与"用户管理"一致:`max(首页可见 home_visible, 比价, 领券)`**不含 last_login_at**;无任何信号时以 `created_at` 为非空基线 | 比价可从**浮窗**触发、不进首页;`last_login_at` 只在登录/换绑动作更新(re-login 也算),代表不了"在用 App",故彻底排除 |
| **"进首页"信号落地** | **方案 A:前端上报 `home_visible` 埋点**(复用 `/analytics/events`),非新接口 | 三个活跃信号统一为同类埋点事件;零新接口零新列;与 admin 口径天然一致。B(鉴权接口 + 列)"更权威"的优势是假的——比价/领券仍是端上报事件,最弱环决定整体可信度 |
| **清零范围** | **金币 + 折算现金**(**邀请现金不清**——产品红线,仅快照入审计) | 对应"账户里的金币和现金";邀请奖励金与金币现金物理隔离、不可累加,见 `wallet.CoinAccount` 注释 |
| **预警推送** | **可插拔通知器 + 日志占位**v1),后续接 JPush/短信 | 现状无真实推送能力;先把清零主流程 + 审计做扎实,不阻塞 |
| **预警时机** | **完全可配置**(提前天数列表 + 次数 + 执行点 + 通道) | R5 |
| **触发方式** | **进程内每日 worker**,仿 `daily_exchange_worker` | 与项目最新模式一致,无需外部 cron |
| **admin 共享口径** | 共享模块 + **重构 admin 改用它** | 单一真源,永不漂移(R4 |
### 已知取舍(可接受)
- analytics 的 `user_id` 是**端上报、未鉴权**(可伪造)。但伪造只能"保自己活跃、避免被清",无收益,且正是本功能要防的行为,风险良性。活跃时间的非空基线由服务端权威的 `User.created_at` 提供(见 §4),不再依赖 `last_login_at`。与"用户管理"口径一致。
---
## 4. 活跃口径与共享模块 `app/repositories/activity.py`(新建)
活跃口径的**唯一真源**。app 侧模块,admin 可 import`app.main` 不 import `app.admin`,反向允许)。
### 口径
```
last_active = max(
User.created_at, # 注册基线(恒非空;re-login 不推进,只有真实使用才推进)
max AnalyticsEvent.created_at WHERE event IN ACTIVE_EVENTS,
max CouponPromptEngagement.created_at WHERE engage_type == "claim_started",
)
不活跃判定:按北京自然日、0 点对齐(非从末次活跃时刻滚动 15×24h)
last_active_date = 北京(last_active).date() # 末次活跃的北京日,记为「第 1 日」
清零边界 = 北京 00:00 of (last_active_date + RESET_DAYS 天) =「第 (RESET_DAYS+1) 日 0 点」 # 15 → 第16日0点
应清零 ⟺ (cn_today() last_active_date).days ≥ RESET_DAYS
⟺ last_active < cutoff cutoff = 北京 00:00 of (cn_today() (RESET_DAYS1)) # 供 SQL 比较
inactive_days = (cn_today() last_active_date).days # 清零当日恰 = RESET_DAYS
例:末次活跃 1/1 → 1/16 00:00(第16日0点)清零,当日 inactive_days=151/15 及之前不清
```
### 模块内容
- 常量:
- **首页可见活跃信号已定名:`event=home_visible`**(前端最终确认;曾用过渡期 `show`+`page=home` 组合,已废弃)。活跃行为过滤见 `activity.active_event_condition()`:首页可见 `home_visible` 比价 `real_compare_start` 领券 `real_coupon_start`——三者均为纯 event 名,全部收进 `ACTIVE_EVENTS`
- `ACTIVE_ENGAGE_TYPE = "claim_started"`
- `last_active_subqueries(db)` —— 复刻现 admin `queries._last_active_parts()`:两个按 `user_id``GROUP BY max(created_at)` 聚合子查询。
- `last_active_expr(base_col, ev_sub, eng_sub, dialect)` —— 生成 `greatest`/`max`PG `func.greatest`SQLite `func.max`);子聚合缺失时 `coalesce(子聚合, User.created_at)` 兜底(注册基线恒非空,**替代原 last_login_at**)。
- `_norm_utc()` —— 沿用现 admin 的 naive→UTC 归一(SQLite naive / PG aware 混算保护)。
- `reset_cutoff(reset_days)` / `warn_cutoff(reset_days, k)` —— 生成**北京 0 点对齐**的边界 datetime(见口径):`reset_cutoff = 北京 00:00 of (cn_today() (reset_days1))`,供下面查询按 `last_active < cutoff` 比较。
- `select_inactive_users(db, *, cutoff, with_balance=True)` —— **worker 专用**join `CoinAccount`,筛 `last_active < cutoff`(cutoff = 北京 0 点对齐边界,见口径)且(`coin_balance>0 OR cash_balance_cents>0`;**邀请现金不清、不计入候选**),返回 `(user, account, last_active, inactive_days)`
- `select_warn_targets(db, *, reset_days, warn_days_before)` —— **worker 专用**:返回 `(user, account, last_active, inactive_days, stage)` 元组——各"提前天数"窗口内、有余额、本 streak 未推过档 `stage` 的用户(去重结合 `notification_log`,逻辑见 §9)。
> **参考现状**:现口径散落在 `app/admin/repositories/queries.py:38,91-124,199-204``_ACTIVE_EVENTS`/`_last_active_parts`/`greatest`)与 `app/admin/repositories/stats.py:51-52,138-146``COMPARE_START_EVENT`/`COUPON_START_EVENT`/活跃用户集)。这些改为从 `activity.py` 导入(§12)。
---
## 5. 数据模型(2 张新表,不动 `User`)
两表均登记进 `app/models/__init__.py`;一个 Alembic 迁移建两表(`render_as_batch`SQLite 兼容)。
### ① `inactivity_reset_log` —— 清零审计(R2
仿 `app/models/phone_rebind_log.py` 的简单审计表风格。
| 字段 | 类型 | 说明 |
|---|---|---|
| `id` | int PK autoincrement | |
| `user_id` | int, index, not null | |
| `coin_balance_before` | int, not null | 清零前金币 |
| `cash_balance_cents_before` | int, not null | 清零前折算现金(分) |
| `invite_cash_balance_cents_before` | int, not null | 清零前邀请现金(分) |
| `last_active_at` | DateTime(tz), nullable | 判定时的最近活跃时间 |
| `inactive_days` | int, not null | 判定时不活跃天数 |
| `reason` | String(32), not null | 如 `"inactive_15d"` |
| `reset_at` | DateTime(tz), server_default now(), index, not null | 清零时刻 |
### ② `inactivity_notification_log` —— 预警记录 + 去重 + 占位 outboxR3
| 字段 | 类型 | 说明 |
|---|---|---|
| `id` | int PK autoincrement | |
| `user_id` | int, index, not null | |
| `stage` | int, not null | 提前天数档(如 7 / 2 |
| `inactive_days` | int, not null | 推送时不活跃天数 |
| `coin_balance` | int, not null | 推送快照:告知用户的金币数 |
| `cash_balance_cents` | int, not null | 推送快照:折算现金 |
| `invite_cash_balance_cents` | int, not null | 推送快照:邀请现金 |
| `channel` | String(16), not null | `"log"` / `"jpush"` / `"sms"` |
| `status` | String(16), not null | `"placeholder"` / `"sent"` / `"failed"` |
| `created_at` | DateTime(tz), server_default now(), index, not null | 去重锚点(见 §9 |
> 备注:不新增 `User.last_active_at` 列,不改 `get_current_user`。活跃时间由 §4 口径**实时计算**。
---
## 6. 清零 worker `app/core/inactivity_reset_worker.py`(新建)
**完全仿 [`app/core/daily_exchange_worker.py`](../../../app/core/daily_exchange_worker.py)**App 启动自带 asyncio 任务,文件锁(`data/inactivity_reset.lock`)防同机多进程并发。**worker 常驻**;`INACTIVITY_RESET_ENABLED` 只决定是否**真清**:false(默认)= 只记审计名单、不动钱(dry-run),true = 真清。
### 调度
-`INACTIVITY_RESET_CHECK_INTERVAL_SEC` 秒醒一次;`last_run: date` 守卫**北京日**,保证每日只跑一轮。
- 仅当 `cn_today() != last_run` 且当前北京小时 `>= INACTIVITY_RESET_RUN_HOUR` 时执行(启动补跑同 daily_exchange 语义)。
- **清零资格边界 = 第 16 日 0 点(北京,见 §4),与 worker 执行点解耦**worker 于当日 `RUN_HOUR`(默认 3 点)跑,把已过边界者一并清;若要严格 0 点触发可置 `RUN_HOUR=0`,但注意与 `daily_auto_exchange` 的 0 点任务错峰。
- lifespan 里 `start_inactivity_reset_worker()` / `stop_...`(仿 `start_daily_exchange_worker``app/main.py` 的接线)。
### 一轮 `run_once(db)` 两阶段(同一次运行、各自逐用户独立 commit)
**阶段 A — 预警**
```
for user, acc, last_active, inactive_days, stage in activity.select_warn_targets(...):
notifier.send_inactivity_warning(user, balances=snapshot(acc), stage=stage, days_until_reset=RESET_DAYS-inactive_days)
db.add(InactivityNotificationLog(..., channel=notifier.channel, status=notifier.last_status))
db.commit() # 逐条独立
```
**阶段 B — 清零**`biz_type="inactivity_reset"`
```
for user, acc, last_active, inactive_days in activity.select_inactive_users(db, cutoff=activity.reset_cutoff(RESET_DAYS)): # 北京 00:00 of (今天−(RESET_DAYS1))
try:
acc = wallet.get_or_create_account(db, user.id, commit=False, lock=True) # 行锁
before = (acc.coin_balance, acc.cash_balance_cents, acc.invite_cash_balance_cents)
if acc.coin_balance == 0 and acc.cash_balance_cents == 0: continue # 邀请现金不清,不算可清余额
log = InactivityResetLog(user_id=user.id, coin_balance_before=before[0],
cash_balance_cents_before=before[1], invite_cash_balance_cents_before=before[2], # 邀请现金仅快照
last_active_at=last_active, inactive_days=inactive_days, reason=f"inactive_{RESET_DAYS}d")
db.add(log); db.flush() # 拿 log.id 作 ref_id 交叉链接
if acc.coin_balance: wallet.grant_coins(db, user.id, -acc.coin_balance, biz_type="inactivity_reset", ref_id=str(log.id), remark="15天不活跃清零")
if acc.cash_balance_cents: wallet.grant_cash(db, user.id, -acc.cash_balance_cents, biz_type="inactivity_reset", ref_id=str(log.id), remark="15天不活跃清零")
# 邀请现金(invite_cash_balance_cents)不清:产品红线、两本账物理隔离,仅快照记入审计。
db.commit()
except SQLAlchemyError:
db.rollback(); stats["failed"] += 1
```
- `grant_*` 负数出账、`balance_after=0`、写**两条**流水(金币 + 折算现金;**邀请现金不清**);`grant_coins` 负数**不**动 `total_coin_earned`(历史累计保留)。
- 逐用户独立 commit:一个失败不影响其余。返回 `stats = {warned, warn_skipped, warn_failed, scanned, cleared, failed}``logger.info`。**预警逐用户 try/except 隔离、且预警整段异常也绝不阻塞清零**(清零是不可逆资金操作,不能被通知故障拖住)。
---
## 7. 预警与可插拔通知器
`app/integrations/notifier.py` 定义协议(外部投递属 integrations 层):
```python
class InactivityNotifier(Protocol):
channel: str # "log" / "jpush" / "sms"
last_status: str # "placeholder" / "sent" / "failed"
def send_inactivity_warning(self, user, *, balances, stage, days_until_reset) -> None: ...
```
- **v1 `LogNotifier`**`channel="log"`):`logger.warning("[inactivity-warn] user=%s coin=%s cash=%s invite=%s T-%s", ...)``last_status="placeholder"`。参照 `heartbeat_monitor_worker` 先例("本期先不接推送,用终端打印代替")。
- 未来 `JPushNotifier` / `SmsNotifier`:实现同协议即可替换,worker 不改。
- 选择:`INACTIVITY_NOTIFY_CHANNEL` → 工厂返回对应实现(未配到真实实现时回退 `LogNotifier`)。
- 预警文案数据来自快照 `balances`,满足 R3"告知 xx 金币 xx 现金"。
---
## 8. 配置项(`app/core/config.py`
```
INACTIVITY_RESET_ENABLED = False # false(默认)=只记审计名单(dry-run,不动钱);true=真清
INACTIVITY_RESET_DAYS = 15 # 不活跃阈值(天)
INACTIVITY_WARN_DAYS_BEFORE = "7,2" # 清零前几天各推一次;空串=不推。逗号分隔,降序解析
INACTIVITY_RESET_RUN_HOUR = 3 # 北京时间每日执行点(0-23)
INACTIVITY_NOTIFY_CHANNEL = "log" # log(占位) / jpush / sms
INACTIVITY_RESET_CHECK_INTERVAL_SEC = 1800 # worker 唤醒间隔(可复用现有间隔常量)
```
- 清零范围(三桶)固定为常量,不做配置。
- `INACTIVITY_WARN_DAYS_BEFORE` 语义(`inactive_days` 为北京自然日,见 §4):档位 `k` ⟹ 当 `inactive_days >= RESET_DAYS-k``< RESET_DAYS` 且本 streak 未推过档 `k` 时预警,即在北京日 `last_active_date + (RESET_DAYSk)` 触发(漏跑某天时补发最紧急未推档,§9)。
- `INACTIVITY_RESET_RUN_HOUR` 只决定 worker 每日执行点,**不改变**"第 16 日 0 点"这一资格边界(§4/§6)。
---
## 9. 幂等与重新活跃
- **重新活跃自动退出**`inactive_days` 由 §4 口径**实时算**。用户一有 `home_visible`/比价/领券(**登录本身不算**),`last_active` 前移,自动移出预警与清零队列。**无需**显式"重置标记"。
- **预警去重**`inactivity_notification_log` 中存在 `stage==k 且 created_at > last_active` 的行 ⟹ 本 streak 已推过档 `k`,不重推。用户回归后 `last_active` 前移,旧预警行自然"失效",开启新 streak。
- **清零幂等**:阶段 B 只处理三桶非全 0 者;清完 = 0,次日不再匹配。worker 重启 / 多次唤醒 / 补跑均安全,不产生重复清零或重复流水。
- **稳健补发**:worker 漏跑数天后,某用户可能同时满足多档;只补发**最紧急的未推档**(最小 `k`),避免一次刷屏。
---
## 10. 边界与安全
| 场景 | 处理 |
|---|---|
| 新用户 | `created_at` 作活跃基线(恒非空)→ 注册即"第 1 日活跃";注册后连续 15 天无 home_visible/比价/领券 才清 |
| 在途提现 | 提现申请时现金已扣入 `WithdrawOrder`,当前余额已不含在途;只清当前余额、不动提现单。提现失败退款到已清账户 = 用户的钱,正常 |
| 与 `daily_auto_exchange` 并存 | 各自逐用户幂等;金币多已日结折现金,三桶全清正好覆盖 |
| 时区/日界 | 统一北京(`rewards.cn_today()`/`CN_TZ`);**清零/预警按北京自然日 0 点对齐**(末次活跃记为第 1 日 → 第 16 日 0 点清零,见 §4),非滚动 24h;流水 `created_at` 沿用北京 wall-clock naive |
| 误清防护 | worker 常驻默认 **dry-run**`ENABLED=false` 只记审计名单、不动钱);看准名单再置 `true` 真清(§13 |
---
## 11. 前端依赖:`home_visible` 埋点(跨仓 — Android
- **Android 端**`shaguabijia-app-android`)需在**首页可见**`onResume`/Tab 切入)时,向现有 `POST /api/v1/analytics/events` 批量上报里加一条 `event=home_visible`(前端已定名)的事件,**携带登录后的 `user_id`**。
- 客户端按会话/前台去重即可(服务端只取 `max(created_at)`,多报无害)。
- **上线顺序依赖**`home_visible` 全量覆盖前,"进首页"信号缺失,只有比价/领券能推进活跃、其余落到 `created_at` 基线("只开首页不操作"且注册满 15 天的用户会被误清)—— 故**开真清(`ENABLED=true`)必须待 `home_visible` 铺满后再开**(§13);dry-run 只记名单不动钱、可先开着看。
---
## 12. admin 重构范围与影响(R4
- `app/admin/repositories/queries.py`:删本地 `_ACTIVE_EVENTS`/`_last_active_parts()`,改用 `activity.py` 的常量与子查询构造;`list_users``greatest(...)` 排序/筛选、`_attach_last_active` 均改走共享构造器。
- `app/admin/repositories/stats.py``COMPARE_START_EVENT`/`COUPON_START_EVENT`/活跃用户集(`:138-146`)改用共享常量与口径。
- **行为变化(预期内、需产品知会)**:admin 的"最近活跃 / DAU"口径变化——**移除 `last_login_at`(登录不再计为活跃)、以 `created_at` 为基线、纳入 `home_visible`**。net`home_visible` 铺满后更准(真正把"开首页"算进活跃);铺满前"只登录不操作"的用户活跃度会下降。
- **回归底线**:现有 admin 用户列表 / stats 测试按新口径**更新预期**last_login_at 移除 + created_at 基线 + home_visible 纳入);非活跃口径部分行为不变。
---
## 13. 灰度与上线顺序(安全优先)
1. **后端先行**:合入共享模块 + 两表 + worker + 通知器,`INACTIVITY_RESET_ENABLED=False`;活跃口径以 `created_at` 为非空基线、**不含 last_login_at**。
2. **Android 发版**:上报 `home_visible`;观察 analytics 覆盖率。
3. **dry-run 灰度(默认即是)**`INACTIVITY_RESET_ENABLED=False` 时 worker 常驻只写审计名单(`reason=inactive_Nd_dryrun`)、不动钱、不预警;核对名单准确。
4. **开真清**:确认无误后置 `INACTIVITY_RESET_ENABLED=True`(转为真清 + 预警)。
5. **收尾/监控**:持续观察 `home_visible` 覆盖率与预警/清零名单;发现"活跃却被判不活跃"的漏报即回查埋点覆盖(口径已不含 last_login_at,登录不再兜底)。
---
## 14. 测试计划
- **活跃口径(共享模块)**`home_visible`/比价/领券 各单独命中都算活跃;**纯登录不算**;无信号用户以 `created_at` 计;`max` 取最新;naive/aware 混算不崩。
- **admin 回归**:用户列表 / stats 按新口径更新预期(移除 last_login_at + created_at 基线 + home_visible)。
- **不活跃判定**`last_active` 分别 `<15d / =15d / >15d` × 有/无余额 的命中矩阵。
- **清零**:三桶归零;`inactivity_reset_log` 清前值正确;三条流水 `biz_type=inactivity_reset``balance_after=0``ref_id=log.id``total_coin_earned` 不变。
- **预警**:命中窗口调 notifier + 写 `notification_log`;同 streak 不重推;回归后 `last_active` 前移可再次预警;漏跑补发最紧急档。
- **worker**:常驻;`ENABLED=false` 走 dry-run(只记审计名单、不清、不预警);文件锁互斥;逐用户失败隔离(一个抛错不影响其余,`failed` 计数);重复跑幂等。
- **配置**`INACTIVITY_WARN_DAYS_BEFORE` 解析(含空串=不推);`RESET_DAYS`/`RUN_HOUR` 生效。
- 沿用 `tests/conftest.py`(临时 SQLite、`RATE_LIMIT_ENABLED=false`);外部通知 monkeypatch。
---
## 15. 未来工作
- 接真实 `JPushNotifier`(需用户级 `registration_id` 覆盖 + JPush push API/ `SmsNotifier`
- 如需 admin 后台可视化:不活跃/预警/清零名单与历史查询接口。
- 如量级增长导致每日 join 扫描变慢:再考虑物化 `last_active_at`(当前每日一次可接受)。
---
## 附:涉及文件清单
**新增**
- `app/repositories/activity.py` — 活跃口径唯一真源
- `app/models/inactivity_reset_log.py` — 审计表
- `app/models/inactivity_notification_log.py` — 预警/占位表
- `app/core/inactivity_reset_worker.py` — 每日 worker(仿 daily_exchange_worker
- `app/integrations/notifier.py` — 通知器协议 + `LogNotifier`(真实 JPush/短信后续同层扩展)
- `alembic/versions/<...>_add_inactivity_tables.py` — 建两表迁移
- `docs/database/inactivity_reset_log.md` / `inactivity_notification_log.md` — 表字典(随实现补)
- 对应 `tests/test_inactivity_reset.py`
**改动**
- `app/models/__init__.py` — 注册两模型
- `app/core/config.py``INACTIVITY_*` 配置
- `app/main.py` — lifespan 接线 start/stop worker
- `app/admin/repositories/queries.py``stats.py` — 改用 `activity.py`(§12