feat: 添加 P4 通知降噪机制 (Refs #1200) (#1260)

* feat: add notification noise controls

* docs: clarify notification timezone fallback

* fix: address notification noise review feedback
This commit is contained in:
Alfred
2026-05-11 07:45:41 +08:00
committed by GitHub
parent 41844a181d
commit 9f70705234
26 changed files with 1533 additions and 10 deletions
+9
View File
@@ -485,6 +485,15 @@ AGENT_SKILLS=
# NOTIFICATION_ALERT_CHANNELS=
# NOTIFICATION_SYSTEM_ERROR_CHANNELS=
#
# 【通知降噪机制】(Issue #1200 P4)
# 默认全部关闭;仅影响静态通知渠道,不影响机器人触发会话回执。
# NOTIFICATION_DEDUP_TTL_SECONDS=0 # 同一稳定去重 key 在 TTL 内只发送一次;0 关闭
# NOTIFICATION_COOLDOWN_SECONDS=0 # 同一冷却 key 在窗口内限频;0 关闭
# NOTIFICATION_QUIET_HOURS= # 静默时段,格式 HH:MM-HH:MM,支持跨午夜
# NOTIFICATION_TIMEZONE= # 静默时段时区,如 Asia/Shanghai;留空跟随 TZ/系统本地时区
# NOTIFICATION_MIN_SEVERITY= # info,warning,error,critical;留空保持现状
# NOTIFICATION_DAILY_DIGEST_ENABLED=false # 预留配置;当前不会发送每日摘要
#
# 【实时行情预取】(Issue #455)
# PREFETCH_REALTIME_QUOTES=true # 设为 false 可禁用,避免 efinance/akshare_em 全市场拉取;tushare 为单股接口,预取仅拉首股
+6
View File
@@ -324,6 +324,12 @@ jobs:
NOTIFICATION_REPORT_CHANNELS: ${{ vars.NOTIFICATION_REPORT_CHANNELS || secrets.NOTIFICATION_REPORT_CHANNELS }}
NOTIFICATION_ALERT_CHANNELS: ${{ vars.NOTIFICATION_ALERT_CHANNELS || secrets.NOTIFICATION_ALERT_CHANNELS }}
NOTIFICATION_SYSTEM_ERROR_CHANNELS: ${{ vars.NOTIFICATION_SYSTEM_ERROR_CHANNELS || secrets.NOTIFICATION_SYSTEM_ERROR_CHANNELS }}
NOTIFICATION_DEDUP_TTL_SECONDS: ${{ vars.NOTIFICATION_DEDUP_TTL_SECONDS || secrets.NOTIFICATION_DEDUP_TTL_SECONDS || '0' }}
NOTIFICATION_COOLDOWN_SECONDS: ${{ vars.NOTIFICATION_COOLDOWN_SECONDS || secrets.NOTIFICATION_COOLDOWN_SECONDS || '0' }}
NOTIFICATION_QUIET_HOURS: ${{ vars.NOTIFICATION_QUIET_HOURS || secrets.NOTIFICATION_QUIET_HOURS }}
NOTIFICATION_TIMEZONE: ${{ vars.NOTIFICATION_TIMEZONE || secrets.NOTIFICATION_TIMEZONE }}
NOTIFICATION_MIN_SEVERITY: ${{ vars.NOTIFICATION_MIN_SEVERITY || secrets.NOTIFICATION_MIN_SEVERITY }}
NOTIFICATION_DAILY_DIGEST_ENABLED: ${{ vars.NOTIFICATION_DAILY_DIGEST_ENABLED || secrets.NOTIFICATION_DAILY_DIGEST_ENABLED || 'false' }}
# ==========================================
# 自选股配置
@@ -35,6 +35,7 @@ export const Select: React.FC<SelectProps> = ({
}) => {
const selectId = useId();
const resolvedId = id ?? selectId;
const hasEmptyOption = options.some((option) => option.value === '');
return (
<div className={cn('flex flex-col', className)}>
@@ -51,7 +52,7 @@ export const Select: React.FC<SelectProps> = ({
disabled ? 'cursor-not-allowed opacity-50' : 'cursor-pointer',
)}
>
{placeholder && (
{placeholder && !hasEmptyOption && (
<option value="" disabled>
{placeholder}
</option>
@@ -83,6 +83,50 @@ describe('SettingsField', () => {
expect(screen.getAllByRole('button', { name: '删除' })).toHaveLength(2);
});
it('allows optional select fields to be cleared when schema provides an empty option', () => {
const onChange = vi.fn();
render(
<SettingsField
item={{
key: 'NOTIFICATION_MIN_SEVERITY',
value: 'warning',
rawValueExists: true,
isMasked: false,
schema: {
key: 'NOTIFICATION_MIN_SEVERITY',
title: 'Notification Minimum Severity',
category: 'notification',
dataType: 'string',
uiControl: 'select',
isSensitive: false,
isRequired: false,
isEditable: true,
options: [
{ label: 'Not set', value: '' },
{ label: 'info', value: 'info' },
{ label: 'warning', value: 'warning' },
{ label: 'error', value: 'error' },
{ label: 'critical', value: 'critical' },
],
validation: { enum: ['', 'info', 'warning', 'error', 'critical'] },
displayOrder: 69,
},
}}
value="warning"
onChange={onChange}
/>
);
const select = screen.getByLabelText('NOTIFICATION_MIN_SEVERITY');
expect(screen.getByRole('option', { name: 'Not set' })).not.toBeDisabled();
expect(screen.queryByRole('option', { name: '请选择' })).not.toBeInTheDocument();
fireEvent.change(select, { target: { value: '' } });
expect(onChange).toHaveBeenCalledWith('NOTIFICATION_MIN_SEVERITY', '');
});
it('renders localized custom webhook body template guidance', () => {
const onChange = vi.fn();
+1
View File
@@ -29,6 +29,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/).
- [新功能] Web 系统设置页开放 `.env` 配置备份导入/导出,复用键级覆盖、配置版本冲突保护和重载链路;Web 端在 `ADMIN_AUTH_ENABLED=false` 时该入口为禁用状态。
- [chore] 精简仓库根目录:将文档图片资源迁入 `docs/assets/`,将东方财富请求补丁迁入 `src/patches/`,并下移 CI 专用依赖文件与技能适配服务。
- [文档] 更新多语言 README 首页浅色工作台 GIF,并精简功能特性表,保留原有赞助商、快速开始和推送效果结构。
- [新功能] 通知网关新增默认关闭的进程内降噪配置,支持去重、冷却、静默时段和最低严重级别,并将每日摘要开关标记为预留能力。
- [文档] 恢复多语言 README 新闻源配置表中推荐项的加粗样式,统一相关项目章节层级,并精简顶部导航、联系文案和尾部展示。
## [3.16.0] - 2026-05-10
+13 -1
View File
@@ -102,7 +102,7 @@ daily_stock_analysis/
> *注:至少配置一个渠道,配置多个则同时推送
>
> 当前默认 `daily_analysis.yml` 只显式映射固定 Secret / Variable 名称,不会自动把 `STOCK_GROUP_1`、`EMAIL_GROUP_1` 这类任意编号变量导入运行环境。所以分组邮箱功能目前不适用于仓库自带默认 GitHub Actions workflow;它适用于本地 `.env`、Docker,或你自行显式扩展过 `env:` 映射的运行环境。Actions 已显式映射 `CUSTOM_WEBHOOK_BODY_TEMPLATE`、`WEBHOOK_VERIFY_SSL`、`FEISHU_WEBHOOK_SECRET`、`FEISHU_WEBHOOK_KEYWORD`、`PUSHPLUS_TOPIC` 以及 P3 通知路由键;`MARKDOWN_TO_IMAGE_CHANNELS` 和 `MERGE_EMAIL_NOTIFICATION` 仍作为行为开关不在默认 workflow 中自动映射。
> 当前默认 `daily_analysis.yml` 只显式映射固定 Secret / Variable 名称,不会自动把 `STOCK_GROUP_1`、`EMAIL_GROUP_1` 这类任意编号变量导入运行环境。所以分组邮箱功能目前不适用于仓库自带默认 GitHub Actions workflow;它适用于本地 `.env`、Docker,或你自行显式扩展过 `env:` 映射的运行环境。Actions 已显式映射 `CUSTOM_WEBHOOK_BODY_TEMPLATE`、`WEBHOOK_VERIFY_SSL`、`FEISHU_WEBHOOK_SECRET`、`FEISHU_WEBHOOK_KEYWORD`、`PUSHPLUS_TOPIC`、P3 通知路由键以及 P4 通知降噪键;`MARKDOWN_TO_IMAGE_CHANNELS` 和 `MERGE_EMAIL_NOTIFICATION` 仍作为行为开关不在默认 workflow 中自动映射。
#### 推送行为配置
@@ -123,6 +123,12 @@ daily_stock_analysis/
| `NOTIFICATION_REPORT_CHANNELS` | report 路由渠道(单股推送、聚合日报、大盘复盘、合并推送等);留空表示所有已配置渠道 | 可选 |
| `NOTIFICATION_ALERT_CHANNELS` | alert 路由渠道(EventMonitor 告警);留空表示所有已配置渠道 | 可选 |
| `NOTIFICATION_SYSTEM_ERROR_CHANNELS` | system_error 预留路由渠道;当前不新增自动系统错误生产者,留空表示所有已配置渠道 | 可选 |
| `NOTIFICATION_DEDUP_TTL_SECONDS` | 通知去重 TTL 秒数,`0` 关闭;同一稳定去重 key 在 TTL 内只发送一次 | 可选 |
| `NOTIFICATION_COOLDOWN_SECONDS` | 通知冷却秒数,`0` 关闭;同一冷却 key 在窗口内限频 | 可选 |
| `NOTIFICATION_QUIET_HOURS` | 通知静默时段,格式 `HH:MM-HH:MM`,支持跨午夜;留空关闭 | 可选 |
| `NOTIFICATION_TIMEZONE` | 静默时段使用的 IANA 时区,如 `Asia/Shanghai`;留空跟随 `TZ` 或系统本地时区 | 可选 |
| `NOTIFICATION_MIN_SEVERITY` | 最低通知级别:`info`、`warning`、`error`、`critical`;留空保持现状 | 可选 |
| `NOTIFICATION_DAILY_DIGEST_ENABLED` | 每日摘要预留开关;当前不会发送摘要或持久化摘要内容 | 可选 |
| `MARKDOWN_TO_IMAGE_MAX_CHARS` | 超过此长度不转图片,避免超大图片(默认 15000) | 可选 |
| `MD2IMG_ENGINE` | 转图引擎:`wkhtmltoimage`(默认,需 wkhtmltopdf)或 `markdown-to-file`(emoji 更好,需 `npm i -g markdown-to-file`) | 可选 |
| `PREFETCH_REALTIME_QUOTES` | 设为 `false` 可禁用实时行情预取,避免 efinance/akshare_em 全市场拉取(默认 true) | 可选 |
@@ -261,6 +267,12 @@ daily_stock_analysis/
| `NOTIFICATION_REPORT_CHANNELS` | report 路由渠道,逗号分隔;允许值:wechat,feishu,telegram,email,pushover,pushplus,serverchan3,custom,discord,slack,astrbot | 可选 |
| `NOTIFICATION_ALERT_CHANNELS` | alert 路由渠道,逗号分隔;留空保持全渠道 | 可选 |
| `NOTIFICATION_SYSTEM_ERROR_CHANNELS` | system_error 预留路由渠道,逗号分隔;留空保持全渠道 | 可选 |
| `NOTIFICATION_DEDUP_TTL_SECONDS` | 通知去重 TTL 秒数,`0` 关闭 | 可选 |
| `NOTIFICATION_COOLDOWN_SECONDS` | 通知冷却秒数,`0` 关闭 | 可选 |
| `NOTIFICATION_QUIET_HOURS` | 静默时段,格式 `HH:MM-HH:MM`,支持跨午夜 | 可选 |
| `NOTIFICATION_TIMEZONE` | 静默时段时区,如 `Asia/Shanghai`;留空跟随 `TZ` 或系统本地时区 | 可选 |
| `NOTIFICATION_MIN_SEVERITY` | 最低通知级别:info, warning, error, critical;留空保持现状 | 可选 |
| `NOTIFICATION_DAILY_DIGEST_ENABLED` | 每日摘要预留开关;当前不会发送摘要 | 可选 |
> 说明:默认 `daily_analysis` GitHub Actions workflow 只映射固定变量名,不会自动导入任意编号的 `STOCK_GROUP_N` / `EMAIL_GROUP_N`。因此分组邮箱目前仅在本地 `.env`、Docker 或其他已显式注入这些环境变量的运行环境中生效;若你要在自己的 GitHub Actions 中使用,需在 workflow 的 job `env:` 中逐组显式映射。
+13 -1
View File
@@ -103,7 +103,7 @@ Go to your forked repo → `Settings` → `Secrets and variables` → `Actions`
> *Note: Configure at least one channel; multiple channels will all receive notifications
>
> The default `daily_analysis.yml` in this repository only exports fixed Secret / Variable names. Arbitrary numbered env vars such as `STOCK_GROUP_1` and `EMAIL_GROUP_1` are not auto-injected into the job, so grouped email routing is not available in the stock workflow unless you explicitly extend the workflow's `env:` mapping in your own fork. Actions now maps `CUSTOM_WEBHOOK_BODY_TEMPLATE`, `WEBHOOK_VERIFY_SSL`, `FEISHU_WEBHOOK_SECRET`, `FEISHU_WEBHOOK_KEYWORD`, `PUSHPLUS_TOPIC`, and the P3 notification route keys; `MARKDOWN_TO_IMAGE_CHANNELS` and `MERGE_EMAIL_NOTIFICATION` remain behavior toggles outside the default workflow mapping.
> The default `daily_analysis.yml` in this repository only exports fixed Secret / Variable names. Arbitrary numbered env vars such as `STOCK_GROUP_1` and `EMAIL_GROUP_1` are not auto-injected into the job, so grouped email routing is not available in the stock workflow unless you explicitly extend the workflow's `env:` mapping in your own fork. Actions now maps `CUSTOM_WEBHOOK_BODY_TEMPLATE`, `WEBHOOK_VERIFY_SSL`, `FEISHU_WEBHOOK_SECRET`, `FEISHU_WEBHOOK_KEYWORD`, `PUSHPLUS_TOPIC`, the P3 notification route keys, and the P4 notification noise-control keys; `MARKDOWN_TO_IMAGE_CHANNELS` and `MERGE_EMAIL_NOTIFICATION` remain behavior toggles outside the default workflow mapping.
#### Push Behavior Configuration
@@ -121,6 +121,12 @@ Go to your forked repo → `Settings` → `Secrets and variables` → `Actions`
| `NOTIFICATION_REPORT_CHANNELS` | Report route channels for single-stock, aggregate daily, market review, merged push, and Feishu document success notifications. Empty means all configured channels | Optional |
| `NOTIFICATION_ALERT_CHANNELS` | Alert route channels for EventMonitor notifications. Empty means all configured channels | Optional |
| `NOTIFICATION_SYSTEM_ERROR_CHANNELS` | Reserved system_error route channels. No automatic system error producer is added in P3; empty means all configured channels | Optional |
| `NOTIFICATION_DEDUP_TTL_SECONDS` | Dedup TTL in seconds. `0` disables dedup; the same stable dedup key sends only once within the TTL | Optional |
| `NOTIFICATION_COOLDOWN_SECONDS` | Cooldown window in seconds. `0` disables cooldown; the same cooldown key is rate-limited within the window | Optional |
| `NOTIFICATION_QUIET_HOURS` | Quiet-hours window in `HH:MM-HH:MM` format, supports overnight ranges. Empty disables quiet hours | Optional |
| `NOTIFICATION_TIMEZONE` | IANA timezone for quiet hours, e.g. `Asia/Shanghai`. Empty follows `TZ` or the local system timezone | Optional |
| `NOTIFICATION_MIN_SEVERITY` | Minimum severity: `info`, `warning`, `error`, `critical`. Empty keeps current behavior | Optional |
| `NOTIFICATION_DAILY_DIGEST_ENABLED` | Reserved daily digest flag. The current implementation does not send or persist digests | Optional |
#### Other Configuration
@@ -233,6 +239,12 @@ For the P0 notification baseline and diagnostics, see [Notification Baseline](no
| `NOTIFICATION_REPORT_CHANNELS` | Report route channels, comma-separated. Allowed values: wechat,feishu,telegram,email,pushover,pushplus,serverchan3,custom,discord,slack,astrbot | Optional |
| `NOTIFICATION_ALERT_CHANNELS` | Alert route channels, comma-separated. Empty keeps all configured channels | Optional |
| `NOTIFICATION_SYSTEM_ERROR_CHANNELS` | Reserved system_error route channels, comma-separated. Empty keeps all configured channels | Optional |
| `NOTIFICATION_DEDUP_TTL_SECONDS` | Dedup TTL in seconds. `0` disables dedup | Optional |
| `NOTIFICATION_COOLDOWN_SECONDS` | Cooldown window in seconds. `0` disables cooldown | Optional |
| `NOTIFICATION_QUIET_HOURS` | Quiet-hours window in `HH:MM-HH:MM` format, supports overnight ranges | Optional |
| `NOTIFICATION_TIMEZONE` | Quiet-hours timezone, e.g. `Asia/Shanghai`; empty follows `TZ` or local system timezone | Optional |
| `NOTIFICATION_MIN_SEVERITY` | Minimum severity: info, warning, error, critical. Empty keeps current behavior | Optional |
| `NOTIFICATION_DAILY_DIGEST_ENABLED` | Reserved daily digest flag. It does not send digests yet | Optional |
> Note: the default `daily_analysis` GitHub Actions workflow only maps fixed variable names. It does not automatically import arbitrary numbered variables such as `STOCK_GROUP_N` / `EMAIL_GROUP_N`. This feature therefore works in local `.env`, Docker, or any runtime where you explicitly inject those variables.
+42 -2
View File
@@ -1,6 +1,6 @@
# 通知能力基线
本文档记录通知能力 P0-P3 基线:渠道、配置 key、GitHub Actions 映射、Web 设置元数据、CLI 诊断口径、Web 一键测试、自定义 Webhook Body 模板语义和通知路由策略。P0 只做基线与只读诊断;P1 增加 Web 单渠道真实测试;P2 产品化现有 Body 模板;P3 增加 report / alert / system_error 路由,不包含降噪、per-URL 模板或新增一等渠道。
本文档记录通知能力 P0-P4 基线:渠道、配置 key、GitHub Actions 映射、Web 设置元数据、CLI 诊断口径、Web 一键测试、自定义 Webhook Body 模板语义、通知路由策略和降噪机制。P0 只做基线与只读诊断;P1 增加 Web 单渠道真实测试;P2 产品化现有 Body 模板;P3 增加 report / alert / system_error 路由;P4 增加进程内降噪,不包含 per-URL 模板、跨进程持久化、真实每日摘要或新增一等渠道。
## 渠道基线
@@ -26,7 +26,8 @@
- Minimal key:足以启用一个通知渠道的最小配置。
- Advanced key:只影响认证、安全、格式、线程、群组、证书校验或展示行为,不能单独启用渠道。
- P3 的 `NOTIFICATION_*_CHANNELS` 属于 Advanced key:只收窄已启用渠道,不会单独启用渠道。
- 降噪、长尾渠道和更细粒度路由不在 P3 范围内;相关配置如未来引入,应先更新本文档、`.env.example`、Web 元数据与回归测试。
- P4 的 `NOTIFICATION_DEDUP_TTL_SECONDS`、`NOTIFICATION_COOLDOWN_SECONDS`、`NOTIFICATION_QUIET_HOURS`、`NOTIFICATION_TIMEZONE`、`NOTIFICATION_MIN_SEVERITY`、`NOTIFICATION_DAILY_DIGEST_ENABLED` 属于 Advanced key:只影响已启用静态渠道的发送策略,不会单独启用渠道。
- 长尾渠道、更细粒度路由、跨进程降噪和真实每日摘要不在 P4 范围内;相关配置如未来引入,应先更新本文档、`.env.example`、Web 元数据与回归测试。
## GitHub Actions 映射
@@ -44,6 +45,15 @@ P3 补齐以下通知路由映射:
- `NOTIFICATION_ALERT_CHANNELS`
- `NOTIFICATION_SYSTEM_ERROR_CHANNELS`
P4 补齐以下通知降噪映射:
- `NOTIFICATION_DEDUP_TTL_SECONDS`
- `NOTIFICATION_COOLDOWN_SECONDS`
- `NOTIFICATION_QUIET_HOURS`
- `NOTIFICATION_TIMEZONE`
- `NOTIFICATION_MIN_SEVERITY`
- `NOTIFICATION_DAILY_DIGEST_ENABLED`
默认 workflow 仍不映射 `MARKDOWN_TO_IMAGE_CHANNELS` 与 `MERGE_EMAIL_NOTIFICATION`。它们是发送形态或聚合行为开关,不是渠道凭证;在 Actions 中自动开始读取同名 Secret/Variable 会引入额外行为变化。
## CLI 诊断
@@ -126,6 +136,36 @@ P3 新增三类通知路由配置:
- `MERGE_EMAIL_NOTIFICATION` 不需要额外配置;只要 `email` 仍在 report 路由后的渠道中,现有合并邮件行为保持不变。
- `--check-notify` 会把未知渠道值报为 error,把合法但未启用的路由目标报为 warning。
## 通知降噪机制
P4 新增进程内降噪,只影响静态配置渠道,不影响 `send_to_context()` 的机器人触发会话回执。默认所有配置关闭,未设置时保持旧行为。
| 配置 key | 默认值 | 说明 |
| --- | --- | --- |
| `NOTIFICATION_DEDUP_TTL_SECONDS` | `0` | 同一稳定去重 key 在 TTL 内只发送一次;`0` 关闭 |
| `NOTIFICATION_COOLDOWN_SECONDS` | `0` | 同一冷却 key 在窗口内限频;`0` 关闭 |
| `NOTIFICATION_QUIET_HOURS` | 空 | 静默时段,格式 `HH:MM-HH:MM`,支持跨午夜 |
| `NOTIFICATION_TIMEZONE` | 空 | 静默时段时区,如 `Asia/Shanghai`;留空使用 Python 运行时本地时区(通常由进程 `TZ` 或系统时区决定) |
| `NOTIFICATION_MIN_SEVERITY` | 空 | `info`, `warning`, `error`, `critical`;留空不过滤 |
| `NOTIFICATION_DAILY_DIGEST_ENABLED` | `false` | 预留配置;当前不会发送每日摘要或持久化摘要内容 |
严重级别默认值:
- `report`:`info`
- `alert`:`warning`
- `system_error`:`error`
- 未知或未设置路由:`info`
实现边界:
- 去重 / 冷却状态是当前 Python 进程内 dict,适用于 `main.py` 单进程和 `--serve` 单 worker。
- `uvicorn --workers N`、多容器或多台机器场景下状态不共享,降噪为 per-worker 近似生效。
- pipeline 单股和聚合报告路径使用稳定 key,避免报告内生成时间变化击穿去重;其他未显式传入 `dedup_key` 的 report 通知按内容 hash 去重。
- 未显式传入 `cooldown_key` 的调用按路由和严重级别共享默认冷却槽位,例如 report / info 的普通通知会共用同一个槽位。
- 同一进程内相同 key 的并发发送会先占用短生命周期 in-flight 槽位,避免突发重复发送;静态渠道全部失败时释放该槽位,不写入正式去重 / 冷却状态。
- 降噪判断异常时 fail-open:记录日志并继续发送静态渠道。
- `NOTIFICATION_TIMEZONE` 留空时使用 `datetime.now().astimezone()` 解析到的运行时本地时区;Actions / Docker 场景建议显式配置 `NOTIFICATION_TIMEZONE` 以避免时区歧义。
## 场景占位
- Local:优先使用 `.env`,可用 `python main.py --check-notify` 做本地诊断。
+73
View File
@@ -25,6 +25,12 @@ from src.report_language import (
normalize_report_language,
)
from src.notification_routing import parse_notification_route_channels
from src.notification_noise import (
NOTIFICATION_SEVERITIES,
is_supported_notification_severity,
parse_notification_quiet_hours,
validate_notification_timezone,
)
logger = logging.getLogger(__name__)
@@ -753,6 +759,14 @@ class Config:
notification_alert_channels: List[str] = field(default_factory=list)
notification_system_error_channels: List[str] = field(default_factory=list)
# 通知降噪机制(Issue #1200 P4):默认全部关闭,仅对静态通知渠道生效
notification_dedup_ttl_seconds: int = 0
notification_cooldown_seconds: int = 0
notification_quiet_hours: str = ""
notification_timezone: str = ""
notification_min_severity: str = ""
notification_daily_digest_enabled: bool = False
# 单股推送模式:每分析完一只股票立即推送,而不是汇总后推送
single_stock_notify: bool = False
@@ -1484,6 +1498,25 @@ class Config:
notification_system_error_channels=parse_notification_route_channels(
os.getenv('NOTIFICATION_SYSTEM_ERROR_CHANNELS')
),
notification_dedup_ttl_seconds=parse_env_int(
os.getenv('NOTIFICATION_DEDUP_TTL_SECONDS'),
0,
field_name='NOTIFICATION_DEDUP_TTL_SECONDS',
minimum=0,
),
notification_cooldown_seconds=parse_env_int(
os.getenv('NOTIFICATION_COOLDOWN_SECONDS'),
0,
field_name='NOTIFICATION_COOLDOWN_SECONDS',
minimum=0,
),
notification_quiet_hours=(os.getenv('NOTIFICATION_QUIET_HOURS') or '').strip(),
notification_timezone=(os.getenv('NOTIFICATION_TIMEZONE') or '').strip(),
notification_min_severity=(os.getenv('NOTIFICATION_MIN_SEVERITY') or '').strip().lower(),
notification_daily_digest_enabled=parse_env_bool(
os.getenv('NOTIFICATION_DAILY_DIGEST_ENABLED'),
default=False,
),
single_stock_notify=os.getenv('SINGLE_STOCK_NOTIFY', 'false').lower() == 'true',
report_type=cls._parse_report_type(os.getenv('REPORT_TYPE', 'simple')),
report_language=cls._parse_report_language(report_language_raw),
@@ -2446,6 +2479,46 @@ class Config:
field="WECHAT_WEBHOOK_URL",
))
if self.notification_quiet_hours:
try:
parse_notification_quiet_hours(self.notification_quiet_hours)
except ValueError as exc:
issues.append(ConfigIssue(
severity="error",
message=f"通知静默时段配置无效:{exc}",
field="NOTIFICATION_QUIET_HOURS",
))
if self.notification_timezone:
try:
validate_notification_timezone(self.notification_timezone)
except ValueError as exc:
issues.append(ConfigIssue(
severity="error",
message=f"通知时区配置无效:{exc}",
field="NOTIFICATION_TIMEZONE",
))
if self.notification_min_severity and not is_supported_notification_severity(self.notification_min_severity):
issues.append(ConfigIssue(
severity="error",
message=(
"通知最低级别配置无效,允许值:"
f"{', '.join(NOTIFICATION_SEVERITIES)}"
),
field="NOTIFICATION_MIN_SEVERITY",
))
if self.notification_daily_digest_enabled:
issues.append(ConfigIssue(
severity="warning",
message=(
"NOTIFICATION_DAILY_DIGEST_ENABLED 当前为预留配置;"
"P4 不会发送每日摘要或持久化摘要内容。"
),
field="NOTIFICATION_DAILY_DIGEST_ENABLED",
))
has_feishu_app_id = bool((self.feishu_app_id or "").strip())
has_feishu_app_secret = bool((self.feishu_app_secret or "").strip())
has_feishu_app_credentials = has_feishu_app_id or has_feishu_app_secret
+88
View File
@@ -11,6 +11,7 @@ from copy import deepcopy
from typing import Any, Dict, List, Optional
from src.config import AGENT_MAX_STEPS_DEFAULT
from src.notification_noise import NOTIFICATION_SEVERITIES
from src.notification_routing import ROUTABLE_NOTIFICATION_CHANNELS
SCHEMA_VERSION = "2026-05-10"
@@ -1500,6 +1501,93 @@ _FIELD_DEFINITIONS: Dict[str, Dict[str, Any]] = {
"validation": {"allowed_values": list(ROUTABLE_NOTIFICATION_CHANNELS), "delimiter": ","},
"display_order": 64,
},
"NOTIFICATION_DEDUP_TTL_SECONDS": {
"title": "Notification Dedup TTL Seconds",
"description": "Suppress duplicate static notifications with the same dedup key within this TTL. 0 disables deduplication.",
"category": "notification",
"data_type": "integer",
"ui_control": "number",
"is_sensitive": False,
"is_required": False,
"is_editable": True,
"default_value": "0",
"options": [],
"validation": {"min": 0},
"display_order": 65,
},
"NOTIFICATION_COOLDOWN_SECONDS": {
"title": "Notification Cooldown Seconds",
"description": "Suppress repeated static notifications with the same cooldown key during this window. 0 disables cooldown.",
"category": "notification",
"data_type": "integer",
"ui_control": "number",
"is_sensitive": False,
"is_required": False,
"is_editable": True,
"default_value": "0",
"options": [],
"validation": {"min": 0},
"display_order": 66,
},
"NOTIFICATION_QUIET_HOURS": {
"title": "Notification Quiet Hours",
"description": "Quiet window in HH:MM-HH:MM format. Supports overnight ranges. Empty disables quiet hours.",
"category": "notification",
"data_type": "string",
"ui_control": "text",
"is_sensitive": False,
"is_required": False,
"is_editable": True,
"default_value": "",
"options": [],
"validation": {"pattern": r"^([01]\d|2[0-3]):[0-5]\d-([01]\d|2[0-3]):[0-5]\d$"},
"display_order": 67,
},
"NOTIFICATION_TIMEZONE": {
"title": "Notification Timezone",
"description": "IANA timezone for quiet hours, e.g. Asia/Shanghai. Empty follows TZ or the local system timezone.",
"category": "notification",
"data_type": "string",
"ui_control": "text",
"is_sensitive": False,
"is_required": False,
"is_editable": True,
"default_value": "",
"options": [],
"validation": {"timezone": True},
"display_order": 68,
},
"NOTIFICATION_MIN_SEVERITY": {
"title": "Notification Minimum Severity",
"description": "Suppress static notifications below this severity. Empty keeps current behavior.",
"category": "notification",
"data_type": "string",
"ui_control": "select",
"is_sensitive": False,
"is_required": False,
"is_editable": True,
"default_value": "",
"options": [
{"label": "Not set", "value": ""},
*({"label": severity, "value": severity} for severity in NOTIFICATION_SEVERITIES),
],
"validation": {"enum": ["", *NOTIFICATION_SEVERITIES]},
"display_order": 69,
},
"NOTIFICATION_DAILY_DIGEST_ENABLED": {
"title": "Notification Daily Digest Enabled (Reserved)",
"description": "Reserved P4 flag. It is visible for compatibility but does not send daily digests yet.",
"category": "notification",
"data_type": "boolean",
"ui_control": "switch",
"is_sensitive": False,
"is_required": False,
"is_editable": True,
"default_value": "false",
"options": [],
"validation": {},
"display_order": 70,
},
"SCHEDULE_TIME": {
"title": "Schedule Time",
"description": "Daily schedule time in HH:MM format.",
+40
View File
@@ -1883,6 +1883,9 @@ class StockAnalysisPipeline:
report_content,
email_stock_codes=[stock_code],
route_type="report",
severity="info",
dedup_key=f"report:single:{stock_code}:{report_type.value}",
cooldown_key=f"report:single:{stock_code}:{report_type.value}",
):
logger.info(f"[{stock_code}] 单股推送成功")
else:
@@ -1918,6 +1921,8 @@ class StockAnalysisPipeline:
results: 分析结果列表
skip_push: 是否跳过推送(仅保存到本地,用于单股推送模式)
"""
noise_decision = None
noise_finalized = False
try:
logger.info("生成决策仪表盘日报...")
report = self._generate_aggregate_report(results, report_type)
@@ -1931,6 +1936,22 @@ class StockAnalysisPipeline:
channels = self.notifier.get_available_channels()
channels = self.notifier.get_channels_for_route("report", channels=channels)
context_success = self.notifier.send_to_context(report)
if channels and hasattr(self.notifier, "evaluate_noise_control"):
report_type_key = report_type.value if isinstance(report_type, ReportType) else str(report_type)
codes_key = ",".join(
sorted(str(getattr(result, "code", "") or "") for result in results)
)
noise_key = f"report:aggregate:{report_type_key}:{codes_key}"
noise_decision = self.notifier.evaluate_noise_control(
report,
route_type="report",
severity="info",
dedup_key=noise_key,
cooldown_key=noise_key,
)
if not noise_decision.should_send:
logger.info(noise_decision.message)
return
# Issue #455: Markdown 转图片(与 notification.send 逻辑一致)
from src.md2img import markdown_to_image
@@ -2096,6 +2117,19 @@ class StockAnalysisPipeline:
logger.warning(f"未知通知渠道: {channel}")
success = wechat_success or non_wechat_success or context_success
if (
(wechat_success or non_wechat_success)
and noise_decision is not None
and hasattr(self.notifier, "record_noise_control")
):
self.notifier.record_noise_control(noise_decision)
noise_finalized = True
elif (
noise_decision is not None
and hasattr(self.notifier, "release_noise_control")
):
self.notifier.release_noise_control(noise_decision)
noise_finalized = True
if success:
logger.info("决策仪表盘推送成功")
else:
@@ -2104,6 +2138,12 @@ class StockAnalysisPipeline:
logger.info("通知渠道未配置,跳过推送")
except Exception as e:
if (
noise_decision is not None
and not noise_finalized
and hasattr(self.notifier, "release_noise_control")
):
self.notifier.release_noise_control(noise_decision)
import traceback
logger.error(f"发送通知失败: {e}\n{traceback.format_exc()}")
+56
View File
@@ -27,6 +27,12 @@ from src.notification_routing import (
get_notification_route_config,
split_notification_route_channels,
)
from src.notification_noise import (
NotificationNoiseDecision,
evaluate_notification_noise,
record_notification_noise,
release_notification_noise,
)
from src.report_language import (
get_localized_stock_name,
get_report_labels,
@@ -389,6 +395,35 @@ class NotificationService(
names.append("钉钉会话")
return ', '.join(names)
def evaluate_noise_control(
self,
content: str,
*,
route_type: Optional[str] = None,
severity: Optional[str] = None,
dedup_key: Optional[str] = None,
cooldown_key: Optional[str] = None,
) -> NotificationNoiseDecision:
"""Evaluate static-channel notification noise controls."""
return evaluate_notification_noise(
self._config,
content=content,
route_type=route_type,
severity=severity,
dedup_key=dedup_key,
cooldown_key=cooldown_key,
)
@staticmethod
def record_noise_control(decision: NotificationNoiseDecision) -> None:
"""Record static-channel notification noise state after a successful send."""
record_notification_noise(decision)
@staticmethod
def release_noise_control(decision: NotificationNoiseDecision) -> None:
"""Release static-channel in-flight noise reservation after send failure."""
release_notification_noise(decision)
# ===== Context channel =====
def _has_context_channel(self) -> bool:
"""判断是否存在基于消息上下文的临时渠道(如钉钉会话、飞书会话)"""
@@ -1627,6 +1662,9 @@ class NotificationService(
email_stock_codes: Optional[List[str]] = None,
email_send_to_all: bool = False,
route_type: Optional[str] = None,
severity: Optional[str] = None,
dedup_key: Optional[str] = None,
cooldown_key: Optional[str] = None,
) -> bool:
"""
统一发送接口 - 向所有已配置的渠道发送
@@ -1644,6 +1682,9 @@ class NotificationService(
email_stock_codes: 股票代码列表(可选,用于邮件渠道路由到对应分组邮箱,Issue #268)
email_send_to_all: 邮件是否发往所有配置邮箱(用于大盘复盘等无股票归属的内容)
route_type: 通知路由类型;None 保持旧行为,report/alert/system_error 按配置过滤静态渠道
severity: 通知严重级别;未设置时按路由类型推断
dedup_key: 可选稳定去重 key;未设置时使用内容 hash
cooldown_key: 可选冷却 key;未设置时使用路由/级别默认 key
Returns:
是否至少有一个渠道发送成功
@@ -1665,6 +1706,17 @@ class NotificationService(
logger.warning("通知路由 %s 未命中任何已配置渠道,跳过静态通知渠道", route_type)
return False
noise_decision = self.evaluate_noise_control(
content,
route_type=route_type,
severity=severity,
dedup_key=dedup_key,
cooldown_key=cooldown_key,
)
if not noise_decision.should_send:
logger.info(noise_decision.message)
return context_success
# Markdown to image (Issue #289): convert once if any channel needs it.
# Per-channel decision via _should_use_image_for_channel (see send() docstring for fallback rules).
image_bytes = None
@@ -1767,6 +1819,10 @@ class NotificationService(
fail_count += 1
logger.info(f"通知发送完成:成功 {success_count} 个,失败 {fail_count} 个")
if success_count > 0:
self.record_noise_control(noise_decision)
else:
self.release_noise_control(noise_decision)
return success_count > 0 or context_success
def save_report_to_file(
+406
View File
@@ -0,0 +1,406 @@
# -*- coding: utf-8 -*-
"""In-process notification noise-control helpers.
The state in this module is intentionally process-local. It provides a small
runtime guard for duplicate/cooldown/quiet-hour suppression without adding
persistent storage, file locks, or cross-worker coordination.
"""
from __future__ import annotations
import hashlib
import logging
import re
import threading
import uuid
from dataclasses import dataclass
from datetime import datetime
from typing import Dict, Optional, Tuple
try: # pragma: no cover - Python <3.9 fallback is not expected in CI.
from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
except Exception: # pragma: no cover
ZoneInfo = None # type: ignore
ZoneInfoNotFoundError = Exception # type: ignore
logger = logging.getLogger(__name__)
NOTIFICATION_SEVERITIES: Tuple[str, ...] = ("info", "warning", "error", "critical")
NOTIFICATION_SEVERITY_RANK = {severity: index for index, severity in enumerate(NOTIFICATION_SEVERITIES)}
DEFAULT_NOTIFICATION_SEVERITY_BY_ROUTE = {
"report": "info",
"alert": "warning",
"system_error": "error",
}
P4_NOISE_ENV_KEYS: Tuple[str, ...] = (
"NOTIFICATION_DEDUP_TTL_SECONDS",
"NOTIFICATION_COOLDOWN_SECONDS",
"NOTIFICATION_QUIET_HOURS",
"NOTIFICATION_TIMEZONE",
"NOTIFICATION_MIN_SEVERITY",
"NOTIFICATION_DAILY_DIGEST_ENABLED",
)
_QUIET_HOURS_RE = re.compile(r"^([01]\d|2[0-3]):([0-5]\d)-([01]\d|2[0-3]):([0-5]\d)$")
_INFLIGHT_RESERVATION_SECONDS = 300
@dataclass(frozen=True)
class NotificationNoiseDecision:
"""Decision returned by the notification noise-control gate."""
should_send: bool
reason_code: str = "allowed"
message: str = ""
route_type: str = "default"
severity: str = "info"
dedup_key: Optional[str] = None
cooldown_key: Optional[str] = None
dedup_ttl_seconds: int = 0
cooldown_seconds: int = 0
evaluated_at: Optional[datetime] = None
dedup_reserved: bool = False
cooldown_reserved: bool = False
reservation_token: Optional[str] = None
_dedup_expires_at: Dict[str, float] = {}
_cooldown_expires_at: Dict[str, float] = {}
_dedup_inflight_until: Dict[str, Tuple[float, str]] = {}
_cooldown_inflight_until: Dict[str, Tuple[float, str]] = {}
_state_lock = threading.Lock()
def reset_notification_noise_state() -> None:
"""Clear process-local notification noise state. Intended for tests."""
with _state_lock:
_dedup_expires_at.clear()
_cooldown_expires_at.clear()
_dedup_inflight_until.clear()
_cooldown_inflight_until.clear()
def is_supported_notification_severity(value: object) -> bool:
"""Return whether *value* is a supported severity string."""
return str(value or "").strip().lower() in NOTIFICATION_SEVERITY_RANK
def normalize_notification_severity(route_type: Optional[str], severity: Optional[str] = None) -> str:
"""Normalize explicit severity, or derive a default from route type."""
explicit = str(severity or "").strip().lower()
if explicit in NOTIFICATION_SEVERITY_RANK:
return explicit
route = str(route_type or "").strip().lower()
return DEFAULT_NOTIFICATION_SEVERITY_BY_ROUTE.get(route, "info")
def parse_notification_quiet_hours(value: Optional[str]) -> Optional[Tuple[int, int]]:
"""Parse ``HH:MM-HH:MM`` into start/end minute-of-day values."""
raw = str(value or "").strip()
if not raw:
return None
match = _QUIET_HOURS_RE.match(raw)
if not match:
raise ValueError("NOTIFICATION_QUIET_HOURS must be in HH:MM-HH:MM format")
start_hour, start_minute, end_hour, end_minute = [int(group) for group in match.groups()]
return start_hour * 60 + start_minute, end_hour * 60 + end_minute
def validate_notification_timezone(value: Optional[str]) -> None:
"""Validate an optional IANA timezone name."""
raw = str(value or "").strip()
if not raw:
return
if ZoneInfo is None:
raise ValueError("zoneinfo is unavailable")
try:
ZoneInfo(raw)
except ZoneInfoNotFoundError as exc:
raise ValueError(f"unknown timezone: {raw}") from exc
def is_time_in_quiet_hours(now: datetime, quiet_hours: Tuple[int, int]) -> bool:
"""Return whether *now* falls inside a quiet-hours interval."""
start_minute, end_minute = quiet_hours
minute_of_day = now.hour * 60 + now.minute
if start_minute == end_minute:
return False
if start_minute < end_minute:
return start_minute <= minute_of_day < end_minute
return minute_of_day >= start_minute or minute_of_day < end_minute
def _resolve_now(timezone_name: Optional[str], now: Optional[datetime]) -> datetime:
raw_timezone = str(timezone_name or "").strip()
if raw_timezone:
if ZoneInfo is None:
raise ValueError("zoneinfo is unavailable")
tz = ZoneInfo(raw_timezone)
if now is None:
return datetime.now(tz)
if now.tzinfo is None:
return now.replace(tzinfo=tz)
return now.astimezone(tz)
if now is None:
return datetime.now().astimezone()
if now.tzinfo is not None:
return now.astimezone()
return now
def _timestamp(now: datetime) -> float:
return now.timestamp()
def _cleanup_expired(now_ts: float) -> None:
expired_dedup = [key for key, expires_at in _dedup_expires_at.items() if expires_at <= now_ts]
for key in expired_dedup:
_dedup_expires_at.pop(key, None)
expired_cooldown = [key for key, expires_at in _cooldown_expires_at.items() if expires_at <= now_ts]
for key in expired_cooldown:
_cooldown_expires_at.pop(key, None)
expired_dedup_inflight = [
key
for key, (expires_at, _token) in _dedup_inflight_until.items()
if expires_at <= now_ts
]
for key in expired_dedup_inflight:
_dedup_inflight_until.pop(key, None)
expired_cooldown_inflight = [
key
for key, (expires_at, _token) in _cooldown_inflight_until.items()
if expires_at <= now_ts
]
for key in expired_cooldown_inflight:
_cooldown_inflight_until.pop(key, None)
def _stable_content_hash(content: str) -> str:
return hashlib.sha256((content or "").encode("utf-8")).hexdigest()
def _state_key(prefix: str, route_type: str, severity: str, key: str) -> str:
return f"{prefix}:{route_type}:{severity}:{key}"
def _build_keys(
*,
content: str,
route_type: str,
severity: str,
dedup_key: Optional[str],
cooldown_key: Optional[str],
) -> Tuple[str, str]:
dedup_part = str(dedup_key).strip() if dedup_key else _stable_content_hash(content)
cooldown_part = str(cooldown_key).strip() if cooldown_key else "default"
return (
_state_key("dedup", route_type, severity, dedup_part),
_state_key("cooldown", route_type, severity, cooldown_part),
)
def evaluate_notification_noise(
config: object,
*,
content: str,
route_type: Optional[str],
severity: Optional[str] = None,
dedup_key: Optional[str] = None,
cooldown_key: Optional[str] = None,
now: Optional[datetime] = None,
) -> NotificationNoiseDecision:
"""Evaluate whether static notification channels should be sent.
This function is fail-open: invalid runtime state or unexpected exceptions
produce an allow decision and a warning log rather than blocking notification.
"""
try:
return _evaluate_notification_noise(
config,
content=content,
route_type=route_type,
severity=severity,
dedup_key=dedup_key,
cooldown_key=cooldown_key,
now=now,
)
except Exception as exc: # pragma: no cover - defensive behavior is tested via monkeypatch.
logger.warning("通知降噪判断失败,将继续发送静态通知渠道: %s", exc)
return NotificationNoiseDecision(
should_send=True,
reason_code="noise_check_failed_open",
message="Noise-control check failed open.",
route_type=str(route_type or "default").strip().lower() or "default",
severity=normalize_notification_severity(route_type, severity),
)
def _evaluate_notification_noise(
config: object,
*,
content: str,
route_type: Optional[str],
severity: Optional[str],
dedup_key: Optional[str],
cooldown_key: Optional[str],
now: Optional[datetime],
) -> NotificationNoiseDecision:
route = str(route_type or "default").strip().lower() or "default"
resolved_severity = normalize_notification_severity(route, severity)
dedup_ttl = max(0, int(getattr(config, "notification_dedup_ttl_seconds", 0) or 0))
cooldown = max(0, int(getattr(config, "notification_cooldown_seconds", 0) or 0))
quiet_hours_raw = getattr(config, "notification_quiet_hours", "") or ""
timezone_name = getattr(config, "notification_timezone", "") or ""
min_severity_raw = str(getattr(config, "notification_min_severity", "") or "").strip().lower()
effective_now = _resolve_now(timezone_name, now)
now_ts = _timestamp(effective_now)
decision_base = {
"route_type": route,
"severity": resolved_severity,
"dedup_ttl_seconds": dedup_ttl,
"cooldown_seconds": cooldown,
"evaluated_at": effective_now,
}
if min_severity_raw:
if min_severity_raw not in NOTIFICATION_SEVERITY_RANK:
logger.warning("NOTIFICATION_MIN_SEVERITY=%s 无效,将忽略最低级别过滤", min_severity_raw)
elif NOTIFICATION_SEVERITY_RANK[resolved_severity] < NOTIFICATION_SEVERITY_RANK[min_severity_raw]:
return NotificationNoiseDecision(
should_send=False,
reason_code="min_severity",
message=(
f"通知级别 {resolved_severity} 低于最低级别 {min_severity_raw},"
"已跳过静态通知渠道。"
),
**decision_base,
)
quiet_hours = parse_notification_quiet_hours(quiet_hours_raw)
if quiet_hours and is_time_in_quiet_hours(effective_now, quiet_hours):
return NotificationNoiseDecision(
should_send=False,
reason_code="quiet_hours",
message=f"当前时间处于静默时段 {quiet_hours_raw},已跳过静态通知渠道。",
**decision_base,
)
dedup_state_key, cooldown_state_key = _build_keys(
content=content,
route_type=route,
severity=resolved_severity,
dedup_key=dedup_key,
cooldown_key=cooldown_key,
)
with _state_lock:
_cleanup_expired(now_ts)
if dedup_ttl > 0 and _dedup_expires_at.get(dedup_state_key, 0) > now_ts:
return NotificationNoiseDecision(
should_send=False,
reason_code="dedup",
message="通知内容在去重 TTL 内已发送,已跳过静态通知渠道。",
dedup_key=dedup_state_key,
cooldown_key=cooldown_state_key,
**decision_base,
)
dedup_inflight = _dedup_inflight_until.get(dedup_state_key)
if dedup_ttl > 0 and dedup_inflight and dedup_inflight[0] > now_ts:
return NotificationNoiseDecision(
should_send=False,
reason_code="dedup_inflight",
message="同一通知正在发送中,已跳过静态通知渠道。",
dedup_key=dedup_state_key,
cooldown_key=cooldown_state_key,
**decision_base,
)
if cooldown > 0 and _cooldown_expires_at.get(cooldown_state_key, 0) > now_ts:
return NotificationNoiseDecision(
should_send=False,
reason_code="cooldown",
message="通知冷却时间尚未结束,已跳过静态通知渠道。",
dedup_key=dedup_state_key,
cooldown_key=cooldown_state_key,
**decision_base,
)
cooldown_inflight = _cooldown_inflight_until.get(cooldown_state_key)
if cooldown > 0 and cooldown_inflight and cooldown_inflight[0] > now_ts:
return NotificationNoiseDecision(
should_send=False,
reason_code="cooldown_inflight",
message="同一通知正在发送中,已跳过静态通知渠道。",
dedup_key=dedup_state_key,
cooldown_key=cooldown_state_key,
**decision_base,
)
reservation_until = now_ts + _INFLIGHT_RESERVATION_SECONDS
dedup_reserved = dedup_ttl > 0
cooldown_reserved = cooldown > 0
reservation_token = uuid.uuid4().hex if dedup_reserved or cooldown_reserved else None
if dedup_reserved:
_dedup_inflight_until[dedup_state_key] = (reservation_until, reservation_token)
if cooldown_reserved:
_cooldown_inflight_until[cooldown_state_key] = (reservation_until, reservation_token)
return NotificationNoiseDecision(
should_send=True,
dedup_key=dedup_state_key,
cooldown_key=cooldown_state_key,
dedup_reserved=dedup_reserved,
cooldown_reserved=cooldown_reserved,
reservation_token=reservation_token,
**decision_base,
)
def _release_reserved_locked(decision: NotificationNoiseDecision) -> None:
if decision.dedup_reserved and decision.dedup_key:
dedup_inflight = _dedup_inflight_until.get(decision.dedup_key)
if dedup_inflight and dedup_inflight[1] == decision.reservation_token:
_dedup_inflight_until.pop(decision.dedup_key, None)
if decision.cooldown_reserved and decision.cooldown_key:
cooldown_inflight = _cooldown_inflight_until.get(decision.cooldown_key)
if cooldown_inflight and cooldown_inflight[1] == decision.reservation_token:
_cooldown_inflight_until.pop(decision.cooldown_key, None)
def release_notification_noise(decision: NotificationNoiseDecision) -> None:
"""Release in-flight reservation without recording dedup/cooldown state."""
if not decision.should_send:
return
try:
with _state_lock:
_release_reserved_locked(decision)
except Exception as exc: # pragma: no cover - defensive branch.
logger.warning("通知降噪发送中状态释放失败,忽略该错误: %s", exc)
def record_notification_noise(decision: NotificationNoiseDecision, now: Optional[datetime] = None) -> None:
"""Record dedup/cooldown state after a static notification send succeeds."""
if not decision.should_send or decision.evaluated_at is None:
return
try:
record_at = now
if record_at is None:
record_at = datetime.now(decision.evaluated_at.tzinfo)
now_ts = _timestamp(record_at)
with _state_lock:
_cleanup_expired(now_ts)
_release_reserved_locked(decision)
if decision.dedup_ttl_seconds > 0 and decision.dedup_key:
_dedup_expires_at[decision.dedup_key] = now_ts + decision.dedup_ttl_seconds
if decision.cooldown_seconds > 0 and decision.cooldown_key:
_cooldown_expires_at[decision.cooldown_key] = now_ts + decision.cooldown_seconds
except Exception as exc: # pragma: no cover - defensive branch.
logger.warning("通知降噪状态记录失败,忽略该错误: %s", exc)
+71 -1
View File
@@ -8,6 +8,13 @@ from typing import List, Literal, Optional, Sequence, Tuple
from src.config import Config
from src.notification import ChannelDetector, NotificationChannel, NotificationService
from src.notification_noise import (
NOTIFICATION_SEVERITIES,
P4_NOISE_ENV_KEYS,
is_supported_notification_severity,
parse_notification_quiet_hours,
validate_notification_timezone,
)
from src.notification_routing import (
NOTIFICATION_ROUTE_CONFIGS,
ROUTABLE_NOTIFICATION_CHANNELS,
@@ -187,6 +194,14 @@ KEY_SPECS: Tuple[NotificationKeySpec, ...] = tuple(
channel="routing",
)
for route in NOTIFICATION_ROUTE_CONFIGS.values()
) + tuple(
NotificationKeySpec(
key=key,
tier="advanced",
description="Optional notification noise-control setting.",
channel="noise",
)
for key in P4_NOISE_ENV_KEYS
)
P0_ACTIONS_ENV_KEYS: Tuple[str, ...] = (
@@ -201,6 +216,8 @@ P3_ROUTE_ENV_KEYS: Tuple[str, ...] = tuple(
route["env_key"] for route in NOTIFICATION_ROUTE_CONFIGS.values()
)
P4_NOISE_ACTIONS_ENV_KEYS: Tuple[str, ...] = P4_NOISE_ENV_KEYS
def _value(config: Config, attr: str):
return getattr(config, attr, None)
@@ -274,7 +291,7 @@ def run_notification_diagnostics(config: Config) -> NotificationDiagnosticResult
_issue(
"info",
"phase_scope",
"通知诊断会检查渠道基线、只读诊断、Web 测试和 P3 路由配置;降噪和长尾渠道留给后续 Phase。",
"通知诊断会检查渠道基线、只读诊断、Web 测试、P3 路由配置和 P4 降噪配置;长尾渠道留给后续 Phase。",
),
]
@@ -411,6 +428,59 @@ def run_notification_diagnostics(config: Config) -> NotificationDiagnosticResult
)
)
if getattr(config, "notification_quiet_hours", ""):
try:
parse_notification_quiet_hours(config.notification_quiet_hours)
except ValueError as exc:
errors.append(
_issue(
"error",
"invalid_quiet_hours",
f"NOTIFICATION_QUIET_HOURS 配置无效: {exc}",
key="NOTIFICATION_QUIET_HOURS",
)
)
if getattr(config, "notification_timezone", ""):
try:
validate_notification_timezone(config.notification_timezone)
except ValueError as exc:
errors.append(
_issue(
"error",
"invalid_notification_timezone",
f"NOTIFICATION_TIMEZONE 配置无效: {exc}",
key="NOTIFICATION_TIMEZONE",
)
)
min_severity = getattr(config, "notification_min_severity", "") or ""
if min_severity and not is_supported_notification_severity(min_severity):
errors.append(
_issue(
"error",
"invalid_notification_min_severity",
(
"NOTIFICATION_MIN_SEVERITY 配置无效;"
f"允许值: {', '.join(NOTIFICATION_SEVERITIES)}。"
),
key="NOTIFICATION_MIN_SEVERITY",
)
)
if getattr(config, "notification_daily_digest_enabled", False):
warnings.append(
_issue(
"warning",
"reserved_daily_digest",
(
"NOTIFICATION_DAILY_DIGEST_ENABLED 当前为预留配置;"
"P4 不会发送每日摘要或持久化摘要内容。"
),
key="NOTIFICATION_DAILY_DIGEST_ENABLED",
)
)
return NotificationDiagnosticResult(
configured_channels=configured,
errors=tuple(errors),
+45
View File
@@ -43,6 +43,7 @@ from src.core.config_registry import (
get_field_definition,
get_registered_field_keys,
)
from src.notification_noise import validate_notification_timezone
logger = logging.getLogger(__name__)
@@ -1626,6 +1627,35 @@ class SystemConfigService:
}
)
elif validation.get("pattern"):
pattern = validation["pattern"]
if not re.match(pattern, value.strip()):
issues.append(
{
"key": key,
"code": "invalid_format",
"message": "Value does not match the required format",
"severity": "error",
"expected": pattern,
"actual": value,
}
)
if validation.get("timezone") and value:
try:
validate_notification_timezone(value)
except ValueError as exc:
issues.append(
{
"key": key,
"code": "invalid_timezone",
"message": str(exc),
"severity": "error",
"expected": "valid IANA timezone or empty",
"actual": value,
}
)
if "enum" in validation and value and value not in validation["enum"]:
issues.append(
{
@@ -3019,6 +3049,21 @@ class SystemConfigService:
)
issues.extend(SystemConfigService._validate_llm_runtime_selection(effective_map=effective_map))
if parse_env_bool(effective_map.get("NOTIFICATION_DAILY_DIGEST_ENABLED"), default=False):
issues.append(
{
"key": "NOTIFICATION_DAILY_DIGEST_ENABLED",
"code": "reserved_notification_daily_digest",
"message": (
"NOTIFICATION_DAILY_DIGEST_ENABLED is reserved; "
"the current P4 implementation does not send daily digests."
),
"severity": "warning",
"expected": "reserved flag only",
"actual": effective_map.get("NOTIFICATION_DAILY_DIGEST_ENABLED", ""),
}
)
return issues
@staticmethod
+39
View File
@@ -242,5 +242,44 @@ class TestNotificationRouteFieldsRegistered(unittest.TestCase):
self.assertIn(key, field_keys, f"{key} missing from schema response")
class TestNotificationNoiseFieldsRegistered(unittest.TestCase):
"""P4 notification noise-control keys must be visible in settings schema."""
_NOISE_KEYS = (
"NOTIFICATION_DEDUP_TTL_SECONDS",
"NOTIFICATION_COOLDOWN_SECONDS",
"NOTIFICATION_QUIET_HOURS",
"NOTIFICATION_TIMEZONE",
"NOTIFICATION_MIN_SEVERITY",
"NOTIFICATION_DAILY_DIGEST_ENABLED",
)
def test_field_definitions_exist(self):
for key in self._NOISE_KEYS:
field = get_field_definition(key)
self.assertEqual(field["category"], "notification", f"{key} category")
self.assertFalse(field["is_sensitive"], f"{key} should not be sensitive")
self.assertFalse(field["is_required"], f"{key} should not be required")
self.assertEqual(get_field_definition("NOTIFICATION_DEDUP_TTL_SECONDS")["data_type"], "integer")
self.assertEqual(get_field_definition("NOTIFICATION_COOLDOWN_SECONDS")["data_type"], "integer")
self.assertEqual(get_field_definition("NOTIFICATION_DAILY_DIGEST_ENABLED")["data_type"], "boolean")
min_severity = get_field_definition("NOTIFICATION_MIN_SEVERITY")
self.assertEqual(min_severity["options"][0]["value"], "")
self.assertIn("", min_severity["validation"]["enum"])
self.assertIn("warning", min_severity["validation"]["enum"])
def test_schema_response_includes_noise_fields(self):
schema = build_schema_response()
notification_cat = next(
(c for c in schema["categories"] if c["category"] == "notification"),
None,
)
self.assertIsNotNone(notification_cat, "notification category missing")
field_keys = {f["key"] for f in notification_cat["fields"]}
for key in self._NOISE_KEYS:
self.assertIn(key, field_keys, f"{key} missing from schema response")
if __name__ == "__main__":
unittest.main()
+23
View File
@@ -370,6 +370,29 @@ class TestValidateStructuredNotification:
warn = [i for i in issues if i.severity == "warning"]
assert not any("FEISHU_APP_ID / FEISHU_APP_SECRET" in i.message for i in warn)
def test_invalid_notification_noise_config_reports_errors(self):
cfg = _make_config(
notification_quiet_hours="9:00-18:00",
notification_timezone="Mars/Olympus",
notification_min_severity="notice",
)
issues = cfg.validate_structured()
errors = {(i.field, i.severity) for i in issues}
assert ("NOTIFICATION_QUIET_HOURS", "error") in errors
assert ("NOTIFICATION_TIMEZONE", "error") in errors
assert ("NOTIFICATION_MIN_SEVERITY", "error") in errors
def test_daily_digest_reserved_flag_warns_without_blocking(self):
cfg = _make_config(notification_daily_digest_enabled=True)
issues = cfg.validate_structured()
assert any(
issue.field == "NOTIFICATION_DAILY_DIGEST_ENABLED"
and issue.severity == "warning"
for issue in issues
)
def test_no_search_engine_is_info(self):
cfg = _make_config(searxng_public_instances_enabled=False)
issues = cfg.validate_structured()
@@ -5,7 +5,7 @@ from pathlib import Path
import yaml
from src.services.notification_diagnostics import P0_ACTIONS_ENV_KEYS, P3_ROUTE_ENV_KEYS
from src.services.notification_diagnostics import P0_ACTIONS_ENV_KEYS, P3_ROUTE_ENV_KEYS, P4_NOISE_ENV_KEYS
ROOT_DIR = Path(__file__).resolve().parent.parent
@@ -43,6 +43,13 @@ def test_daily_analysis_maps_p3_notification_route_env_keys() -> None:
assert key in env
def test_daily_analysis_maps_p4_notification_noise_env_keys() -> None:
env = _load_daily_analysis_env()
for key in P4_NOISE_ENV_KEYS:
assert key in env
def test_daily_analysis_keeps_deferred_behavior_switches_unmapped() -> None:
env = _load_daily_analysis_env()
+80
View File
@@ -31,6 +31,7 @@ for optional_module in ("litellm", "json_repair"):
from src.config import Config
from src.notification import NotificationService, NotificationChannel
from src.notification_noise import reset_notification_noise_state
from src.analyzer import AnalysisResult
import requests
@@ -77,6 +78,9 @@ class TestNotificationServiceSendToMethods(unittest.TestCase):
"""
def setUp(self):
reset_notification_noise_state()
@mock.patch("src.notification.get_config")
def test_no_channels_service_unavailable_and_send_returns_false(self, mock_get_config):
mock_get_config.return_value = _make_config()
@@ -229,6 +233,82 @@ class TestNotificationServiceSendToMethods(unittest.TestCase):
mock_context.assert_called_once_with("content")
mock_custom.assert_not_called()
@mock.patch("src.notification.get_config")
def test_send_dedup_suppresses_static_channels_after_success(self, mock_get_config: mock.MagicMock):
cfg = _make_config(
custom_webhook_urls=["https://example.com/webhook"],
notification_dedup_ttl_seconds=60,
)
mock_get_config.return_value = cfg
service = NotificationService()
with mock.patch.object(service, "send_to_custom", return_value=True) as mock_custom:
self.assertTrue(service.send("content at 12:00", route_type="report", dedup_key="report:aggregate:simple:600519"))
self.assertFalse(service.send("content at 12:01", route_type="report", dedup_key="report:aggregate:simple:600519"))
mock_custom.assert_called_once_with("content at 12:00")
@mock.patch("src.notification.get_config")
def test_send_releases_noise_reservation_when_static_channels_fail(self, mock_get_config: mock.MagicMock):
cfg = _make_config(
custom_webhook_urls=["https://example.com/webhook"],
notification_dedup_ttl_seconds=60,
)
mock_get_config.return_value = cfg
service = NotificationService()
with mock.patch.object(service, "send_to_custom", side_effect=[False, True]) as mock_custom:
self.assertFalse(
service.send(
"content at 12:00",
route_type="report",
dedup_key="report:aggregate:simple:600519",
)
)
self.assertTrue(
service.send(
"content at 12:01",
route_type="report",
dedup_key="report:aggregate:simple:600519",
)
)
self.assertEqual(mock_custom.call_count, 2)
@mock.patch("src.notification.get_config")
def test_send_to_context_is_not_limited_by_noise_controls(self, mock_get_config: mock.MagicMock):
cfg = _make_config(
custom_webhook_urls=["https://example.com/webhook"],
notification_dedup_ttl_seconds=60,
)
mock_get_config.return_value = cfg
service = NotificationService()
with mock.patch.object(service, "send_to_context", return_value=True) as mock_context, \
mock.patch.object(service, "send_to_custom", return_value=True) as mock_custom:
self.assertTrue(service.send("content at 12:00", route_type="report", dedup_key="report:aggregate:simple:600519"))
self.assertTrue(service.send("content at 12:01", route_type="report", dedup_key="report:aggregate:simple:600519"))
self.assertEqual(mock_context.call_count, 2)
mock_custom.assert_called_once_with("content at 12:00")
@mock.patch("src.notification.get_config")
def test_noise_check_failure_does_not_block_static_send(self, mock_get_config: mock.MagicMock):
cfg = _make_config(custom_webhook_urls=["https://example.com/webhook"])
mock_get_config.return_value = cfg
service = NotificationService()
with mock.patch("src.notification_noise._evaluate_notification_noise", side_effect=RuntimeError("boom")), \
mock.patch.object(service, "send_to_custom", return_value=True) as mock_custom:
ok = service.send("content", route_type="report")
self.assertTrue(ok)
mock_custom.assert_called_once_with("content")
@mock.patch("src.notification.get_config")
@mock.patch("requests.post")
def test_send_to_discord_via_notification_service_with_webhook(
+36
View File
@@ -10,6 +10,7 @@ from src.services.notification_diagnostics import (
KEY_SPECS,
NotificationDiagnosticResult,
P3_ROUTE_ENV_KEYS,
P4_NOISE_ENV_KEYS,
format_notification_diagnostics,
run_notification_diagnostics,
)
@@ -42,6 +43,8 @@ class NotificationDiagnosticsTestCase(unittest.TestCase):
self.assertIn(("WEBHOOK_VERIFY_SSL", "advanced"), key_tiers)
for key in P3_ROUTE_ENV_KEYS:
self.assertIn((key, "advanced"), key_tiers)
for key in P4_NOISE_ENV_KEYS:
self.assertIn((key, "advanced"), key_tiers)
self.assertIn(("DISCORD_BOT_TOKEN", "minimal"), key_tiers)
self.assertIn(("SLACK_BOT_TOKEN", "minimal"), key_tiers)
self.assertNotIn(("DISCORD_BOT_TOKEN", "advanced"), key_tiers)
@@ -122,6 +125,39 @@ class NotificationDiagnosticsTestCase(unittest.TestCase):
self.assertEqual(warnings[0].key, "NOTIFICATION_ALERT_CHANNELS")
self.assertIn("telegram", warnings[0].message)
def test_noise_invalid_quiet_hours_reports_error(self):
result = run_notification_diagnostics(
_config(
wechat_webhook_url="https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=1",
notification_quiet_hours="9:00-18:00",
)
)
self.assertFalse(result.ok)
self.assertIn("invalid_quiet_hours", {item.code for item in result.errors})
def test_noise_invalid_timezone_reports_error(self):
result = run_notification_diagnostics(
_config(
wechat_webhook_url="https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=1",
notification_timezone="Mars/Olympus",
)
)
self.assertFalse(result.ok)
self.assertIn("invalid_notification_timezone", {item.code for item in result.errors})
def test_noise_daily_digest_reserved_reports_warning(self):
result = run_notification_diagnostics(
_config(
wechat_webhook_url="https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=1",
notification_daily_digest_enabled=True,
)
)
self.assertTrue(result.ok)
self.assertIn("reserved_daily_digest", {item.code for item in result.warnings})
if __name__ == "__main__":
unittest.main()
+305
View File
@@ -0,0 +1,305 @@
from datetime import datetime, timedelta, timezone
from types import SimpleNamespace
from src.notification_noise import (
evaluate_notification_noise,
record_notification_noise,
release_notification_noise,
reset_notification_noise_state,
)
def _config(**overrides):
defaults = {
"notification_dedup_ttl_seconds": 0,
"notification_cooldown_seconds": 0,
"notification_quiet_hours": "",
"notification_timezone": "",
"notification_min_severity": "",
"notification_daily_digest_enabled": False,
}
defaults.update(overrides)
return SimpleNamespace(**defaults)
def setup_function():
reset_notification_noise_state()
def test_dedup_ttl_suppresses_until_expiry_with_explicit_key():
config = _config(notification_dedup_ttl_seconds=60)
now = datetime(2026, 5, 10, 12, 0, tzinfo=timezone.utc)
first = evaluate_notification_noise(
config,
content="content at 12:00",
route_type="report",
dedup_key="report:aggregate:simple:600519",
now=now,
)
assert first.should_send
record_notification_noise(first, now=now)
duplicate = evaluate_notification_noise(
config,
content="content at 12:01",
route_type="report",
dedup_key="report:aggregate:simple:600519",
now=now + timedelta(seconds=10),
)
assert not duplicate.should_send
assert duplicate.reason_code == "dedup"
expired = evaluate_notification_noise(
config,
content="content at 12:02",
route_type="report",
dedup_key="report:aggregate:simple:600519",
now=now + timedelta(seconds=61),
)
assert expired.should_send
def test_cooldown_keys_are_independent():
config = _config(notification_cooldown_seconds=60)
now = datetime(2026, 5, 10, 12, 0, tzinfo=timezone.utc)
first = evaluate_notification_noise(
config,
content="one",
route_type="report",
cooldown_key="report:single:600519:simple",
now=now,
)
assert first.should_send
record_notification_noise(first, now=now)
same_key = evaluate_notification_noise(
config,
content="two",
route_type="report",
cooldown_key="report:single:600519:simple",
now=now + timedelta(seconds=1),
)
other_key = evaluate_notification_noise(
config,
content="three",
route_type="report",
cooldown_key="report:single:000001:simple",
now=now + timedelta(seconds=1),
)
assert not same_key.should_send
assert same_key.reason_code == "cooldown"
assert other_key.should_send
def test_inflight_reservation_suppresses_same_key_until_released():
config = _config(notification_dedup_ttl_seconds=60)
now = datetime(2026, 5, 10, 12, 0, tzinfo=timezone.utc)
first = evaluate_notification_noise(
config,
content="first",
route_type="report",
dedup_key="report:aggregate:simple:600519",
now=now,
)
assert first.should_send
duplicate = evaluate_notification_noise(
config,
content="second",
route_type="report",
dedup_key="report:aggregate:simple:600519",
now=now + timedelta(seconds=1),
)
assert not duplicate.should_send
assert duplicate.reason_code == "dedup_inflight"
release_notification_noise(first)
retried = evaluate_notification_noise(
config,
content="retry",
route_type="report",
dedup_key="report:aggregate:simple:600519",
now=now + timedelta(seconds=2),
)
assert retried.should_send
def test_cooldown_inflight_reservation_suppresses_same_key_until_released():
config = _config(notification_cooldown_seconds=60)
now = datetime(2026, 5, 10, 12, 0, tzinfo=timezone.utc)
first = evaluate_notification_noise(
config,
content="first",
route_type="report",
cooldown_key="report:single:600519:simple",
now=now,
)
assert first.should_send
duplicate = evaluate_notification_noise(
config,
content="second",
route_type="report",
cooldown_key="report:single:600519:simple",
now=now + timedelta(seconds=1),
)
assert not duplicate.should_send
assert duplicate.reason_code == "cooldown_inflight"
release_notification_noise(first)
retried = evaluate_notification_noise(
config,
content="retry",
route_type="report",
cooldown_key="report:single:600519:simple",
now=now + timedelta(seconds=2),
)
assert retried.should_send
def test_record_uses_success_time_for_expiry_not_evaluate_time():
config = _config(notification_dedup_ttl_seconds=60)
evaluated_at = datetime(2026, 5, 10, 12, 0, tzinfo=timezone.utc)
success_at = evaluated_at + timedelta(minutes=5)
first = evaluate_notification_noise(
config,
content="content",
route_type="report",
dedup_key="report:aggregate:simple:600519",
now=evaluated_at,
)
assert first.should_send
record_notification_noise(first, now=success_at)
duplicate = evaluate_notification_noise(
config,
content="content",
route_type="report",
dedup_key="report:aggregate:simple:600519",
now=success_at + timedelta(seconds=59),
)
assert not duplicate.should_send
assert duplicate.reason_code == "dedup"
expired = evaluate_notification_noise(
config,
content="content",
route_type="report",
dedup_key="report:aggregate:simple:600519",
now=success_at + timedelta(seconds=61),
)
assert expired.should_send
def test_stale_release_does_not_clear_newer_inflight_reservation():
config = _config(notification_dedup_ttl_seconds=60)
now = datetime(2026, 5, 10, 12, 0, tzinfo=timezone.utc)
first = evaluate_notification_noise(
config,
content="content",
route_type="report",
dedup_key="report:aggregate:simple:600519",
now=now,
)
assert first.should_send
newer = evaluate_notification_noise(
config,
content="content",
route_type="report",
dedup_key="report:aggregate:simple:600519",
now=now + timedelta(seconds=301),
)
assert newer.should_send
release_notification_noise(first)
duplicate = evaluate_notification_noise(
config,
content="content",
route_type="report",
dedup_key="report:aggregate:simple:600519",
now=now + timedelta(seconds=302),
)
assert not duplicate.should_send
assert duplicate.reason_code == "dedup_inflight"
def test_quiet_hours_same_day_and_overnight():
same_day = _config(notification_quiet_hours="09:00-17:00", notification_timezone="UTC")
assert not evaluate_notification_noise(
same_day,
content="quiet",
route_type="report",
now=datetime(2026, 5, 10, 10, 0, tzinfo=timezone.utc),
).should_send
assert evaluate_notification_noise(
same_day,
content="loud",
route_type="report",
now=datetime(2026, 5, 10, 18, 0, tzinfo=timezone.utc),
).should_send
overnight = _config(notification_quiet_hours="22:00-06:00", notification_timezone="UTC")
assert not evaluate_notification_noise(
overnight,
content="late",
route_type="report",
now=datetime(2026, 5, 10, 23, 0, tzinfo=timezone.utc),
).should_send
assert not evaluate_notification_noise(
overnight,
content="early",
route_type="report",
now=datetime(2026, 5, 11, 5, 30, tzinfo=timezone.utc),
).should_send
assert evaluate_notification_noise(
overnight,
content="day",
route_type="report",
now=datetime(2026, 5, 11, 12, 0, tzinfo=timezone.utc),
).should_send
def test_invalid_timezone_fails_open():
config = _config(
notification_quiet_hours="00:00-23:59",
notification_timezone="Mars/Olympus",
)
decision = evaluate_notification_noise(
config,
content="content",
route_type="report",
now=datetime(2026, 5, 10, 12, 0, tzinfo=timezone.utc),
)
assert decision.should_send
assert decision.reason_code == "noise_check_failed_open"
def test_min_severity_filters_lower_severity_only():
config = _config(notification_min_severity="warning")
report = evaluate_notification_noise(config, content="report", route_type="report")
alert = evaluate_notification_noise(config, content="alert", route_type="alert")
system_error = evaluate_notification_noise(config, content="error", route_type="system_error")
assert not report.should_send
assert report.reason_code == "min_severity"
assert alert.should_send
assert system_error.should_send
def test_daily_digest_reserved_flag_does_not_change_runtime_decision():
config = _config(notification_daily_digest_enabled=True)
decision = evaluate_notification_noise(config, content="content", route_type="report")
assert decision.should_send
@@ -149,7 +149,7 @@ class TestPipelineWechatOnlyImageRouting(unittest.TestCase):
class _FakeRoutedNotifier:
def __init__(self, routed_channels, image_channels=None):
def __init__(self, routed_channels, image_channels=None, noise_should_send=True):
self._markdown_to_image_channels = set(image_channels or [])
self._markdown_to_image_max_chars = 15000
self.generate_dashboard_report = MagicMock(side_effect=self._generate_dashboard_report)
@@ -164,6 +164,14 @@ class _FakeRoutedNotifier:
)
self.get_channels_for_route = MagicMock(return_value=list(routed_channels))
self.send_to_context = MagicMock(return_value=False)
self.evaluate_noise_control = MagicMock(
return_value=SimpleNamespace(
should_send=noise_should_send,
message="noise suppressed" if not noise_should_send else "",
)
)
self.record_noise_control = MagicMock()
self.release_noise_control = MagicMock()
self._should_use_image_for_channel = MagicMock(
side_effect=lambda channel, image_bytes: (
channel.value in self._markdown_to_image_channels and image_bytes is not None
@@ -202,6 +210,11 @@ class TestPipelineReportRouteFiltering(unittest.TestCase):
pipeline.notifier.send_to_telegram.assert_called_once_with("report:000001")
pipeline.notifier.send_to_wechat.assert_not_called()
pipeline.notifier.send_to_email.assert_not_called()
pipeline.notifier.evaluate_noise_control.assert_called_once()
noise_kwargs = pipeline.notifier.evaluate_noise_control.call_args.kwargs
self.assertEqual(noise_kwargs["dedup_key"], "report:aggregate:simple:000001")
self.assertEqual(noise_kwargs["cooldown_key"], "report:aggregate:simple:000001")
pipeline.notifier.record_noise_control.assert_called_once()
def test_markdown_to_image_uses_route_filtered_channels(self):
pipeline = StockAnalysisPipeline.__new__(StockAnalysisPipeline)
@@ -219,6 +232,35 @@ class TestPipelineReportRouteFiltering(unittest.TestCase):
pipeline.notifier.send_to_email.assert_called_once_with("report:000001")
pipeline.notifier.send_to_telegram.assert_not_called()
def test_noise_suppression_happens_before_markdown_to_image(self):
pipeline = StockAnalysisPipeline.__new__(StockAnalysisPipeline)
pipeline.notifier = _FakeRoutedNotifier(
[NotificationChannel.TELEGRAM],
image_channels={"telegram"},
noise_should_send=False,
)
pipeline.config = SimpleNamespace(stock_email_groups=[])
results = [SimpleNamespace(code="000001")]
with patch("src.md2img.markdown_to_image", return_value=b"png") as mock_md2img:
pipeline._send_notifications(results, ReportType.SIMPLE)
mock_md2img.assert_not_called()
pipeline.notifier.send_to_telegram.assert_not_called()
pipeline.notifier.record_noise_control.assert_not_called()
def test_noise_reservation_released_when_pipeline_static_send_raises(self):
pipeline = StockAnalysisPipeline.__new__(StockAnalysisPipeline)
pipeline.notifier = _FakeRoutedNotifier([NotificationChannel.TELEGRAM])
pipeline.notifier.send_to_telegram.side_effect = RuntimeError("send failed")
pipeline.config = SimpleNamespace(stock_email_groups=[])
results = [SimpleNamespace(code="000001")]
pipeline._send_notifications(results, ReportType.SIMPLE)
pipeline.notifier.record_noise_control.assert_not_called()
pipeline.notifier.release_noise_control.assert_called_once()
if __name__ == "__main__":
unittest.main()
@@ -58,7 +58,15 @@ class _CriticalSectionTrackingNotifier:
self._enter("generate", result.code)
return f"single:{result.code}"
def _send(self, content: str, email_stock_codes=None, route_type=None) -> bool:
def _send(
self,
content: str,
email_stock_codes=None,
route_type=None,
severity=None,
dedup_key=None,
cooldown_key=None,
) -> bool:
stock_code = (email_stock_codes or ["unknown"])[0]
self._enter("send", stock_code)
return True
+12 -1
View File
@@ -42,7 +42,15 @@ class _TrackingNotifier:
)
self.send = MagicMock(side_effect=self._send)
def _send(self, content, email_stock_codes=None, route_type=None):
def _send(
self,
content,
email_stock_codes=None,
route_type=None,
severity=None,
dedup_key=None,
cooldown_key=None,
):
with self._lock:
self._inflight += 1
self.max_inflight = max(self.max_inflight, self._inflight)
@@ -143,6 +151,9 @@ class TestPipelineSingleStockNotify(unittest.TestCase):
"brief:600519",
email_stock_codes=["600519"],
route_type="report",
severity="info",
dedup_key="report:single:600519:brief",
cooldown_key="report:single:600519:brief",
)
def test_process_single_stock_direct_path_does_not_notify_when_failed(self):
+12
View File
@@ -114,6 +114,18 @@ class SystemConfigApiTestCase(unittest.TestCase):
self.assertTrue(stock_schema["examples"])
self.assertTrue(stock_schema["docs"])
def test_get_config_schema_includes_notification_noise_fields(self) -> None:
payload = system_config.get_system_config(include_schema=True, service=self.service).model_dump(by_alias=True)
item_map = {item["key"]: item for item in payload["items"]}
self.assertEqual(item_map["NOTIFICATION_DEDUP_TTL_SECONDS"]["schema"]["data_type"], "integer")
self.assertEqual(item_map["NOTIFICATION_COOLDOWN_SECONDS"]["schema"]["data_type"], "integer")
self.assertEqual(item_map["NOTIFICATION_DAILY_DIGEST_ENABLED"]["schema"]["data_type"], "boolean")
min_severity_schema = item_map["NOTIFICATION_MIN_SEVERITY"]["schema"]
self.assertEqual(min_severity_schema["options"][0]["value"], "")
self.assertIn("", min_severity_schema["validation"]["enum"])
self.assertIn("warning", min_severity_schema["validation"]["enum"])
def test_get_setup_status_returns_readiness_payload(self) -> None:
self.env_path.write_text(
"\n".join(
+57
View File
@@ -362,6 +362,63 @@ class SystemConfigServiceTestCase(unittest.TestCase):
)
)
def test_validate_reports_invalid_notification_quiet_hours(self) -> None:
validation = self.service.validate(
items=[{"key": "NOTIFICATION_QUIET_HOURS", "value": "9:00-18:00"}]
)
self.assertFalse(validation["valid"])
self.assertTrue(
any(
issue["key"] == "NOTIFICATION_QUIET_HOURS"
and issue["code"] == "invalid_format"
for issue in validation["issues"]
)
)
def test_validate_reports_invalid_notification_timezone(self) -> None:
validation = self.service.validate(
items=[{"key": "NOTIFICATION_TIMEZONE", "value": "Mars/Olympus"}]
)
self.assertFalse(validation["valid"])
self.assertTrue(
any(
issue["key"] == "NOTIFICATION_TIMEZONE"
and issue["code"] == "invalid_timezone"
for issue in validation["issues"]
)
)
def test_validate_reports_invalid_notification_min_severity(self) -> None:
validation = self.service.validate(
items=[{"key": "NOTIFICATION_MIN_SEVERITY", "value": "notice"}]
)
self.assertFalse(validation["valid"])
self.assertTrue(
any(
issue["key"] == "NOTIFICATION_MIN_SEVERITY"
and issue["code"] == "invalid_enum"
for issue in validation["issues"]
)
)
def test_validate_warns_daily_digest_is_reserved(self) -> None:
validation = self.service.validate(
items=[{"key": "NOTIFICATION_DAILY_DIGEST_ENABLED", "value": "true"}]
)
self.assertTrue(validation["valid"])
self.assertTrue(
any(
issue["key"] == "NOTIFICATION_DAILY_DIGEST_ENABLED"
and issue["code"] == "reserved_notification_daily_digest"
and issue["severity"] == "warning"
for issue in validation["issues"]
)
)
def test_validate_warns_when_feishu_app_credentials_are_used_without_webhook(self) -> None:
validation = self.service.validate(
items=[