Compare commits

..

6 Commits

Author SHA1 Message Date
zxr
38a86e1fd8 fix 2026-08-19 19:25:46 +08:00
zxr
2fd95c9d75 fix 2026-08-10 20:56:12 +08:00
zxr
3d95965c95 docs: 清理旧文档 2026-08-07 17:36:25 +08:00
zxr
db3cb24b4e chore: 使用环境变量管理内部密钥 2026-08-04 15:48:59 +08:00
zxr
8495fdf039 fix: correlate log alert recovery lifecycle 2026-07-21 22:15:59 +08:00
zxr
eaca3d0ac8 fix: preserve log alert recovery events 2026-07-21 21:21:47 +08:00
14 changed files with 354 additions and 737 deletions

View File

@@ -1,87 +0,0 @@
# Syslog-Trap 接入与重放
## 接入目标
`logs` 服务负责接收 Syslog 与 SNMP Trap按字典和规则解析后写入 `logs_events`,并通过 `logs_alert_outbox` 异步转发到 `alert` 的原始事件池:
```text
Syslog / Trap -> logs_events -> logs_alert_outbox -> Alert/v1/raw-events/ingest
```
转发使用 `X-Internal-Key`,配置来自 `AlertForward.internal_key`。解析成功的事件 `parse_status=parsed`,未命中字典或规则的事件仍保存原始报文,并以 `parse_status=unparsed` 入队,便于规则调整后重放。
## 部署配置
`logs` 当前内置 UDP 接收器:
```yaml
Ingest:
syslog_listen_addr: "0.0.0.0:5140"
trap_listen_addr: "0.0.0.0:1620"
rule_refresh_secs: 30
AlertForward:
enabled: true
base_url: "http://127.0.0.1:18080"
internal_key: "change-me"
default_policy_id: 1
```
生产环境如需标准端口 `514/162`,建议由 systemd socket、firewalld rich rule、iptables REDIRECT 或外层采集网关转发到非特权端口。TCP Syslog 接入建议在网关层启用 TCP listener再转发到 UDP 或调用后续 HTTP ingest 入口;开启 TCP 时必须保留原始来源 IP 和 trace ID。
## 字典与规则
Trap 字典字段:
- `vendor`:厂商,例如 `H3C`
- `oid`:精确 Trap OID。
- `oid_prefix`OID 前缀,兼容旧字典。
- `name` / `title`:展示名称。
- `severity_mapping_json`:级别映射 JSON。
- `parse_expression`:解析 varbind 的表达式或正则。
Syslog 规则字段:
- `source_match`:来源 IP、主机名或原始行子串。
- `message_regex`:消息正文正则。
- `severity_mapping_json`:按正则映射平台级别。
- `resource_uid_extract_regex`:提取 `resource_uid`,优先使用命名分组 `resource_uid`
示例 Syslog
```text
<189>Jun 24 10:00:01 h3c-core-01 IFNET/4/LINK_DOWN: Interface GigabitEthernet1/0/1 is down, resource_uid=network:h3c-core-01
```
示例 H3C Trap OID
```text
1.3.6.1.6.3.1.1.5.3
```
## 未解析队列与重放
未解析事件仍写入 `logs_events`,并创建 outbox payload
- `source_type=syslog``trap`
- `parse_status=unparsed`
- `raw_payload` 保存原始报文或 varbind 摘要
重放接口:
```http
POST /Logs/v1/entries/{id}/replay
Authorization: Bearer <jwt>
```
成功响应会返回新的 `outbox_id`。重放 payload 使用 `parse_status=replayed`,并带上 `labels.replay_of_log_event_id`,前端可在“日志查询 -> 重放结果”中查看发送结果,失败任务可人工重试。
## Smoke 样例
输出 H3C Syslog 与 Trap 示例载荷:
```powershell
C:\Users\27105\.cache\codex-runtimes\codex-primary-runtime\dependencies\python\python.exe scripts\test_alert_receive_smoke.py --print-log-samples
```
这些样例用于准备 UDP/TCP 接收器 smoke 数据,也可作为联调 alert 原始事件池时的期望字段参考。

View File

@@ -1,498 +0,0 @@
# Ops Logs 前端页面设计文档Log Mgmt
## 1. 背景与目标
`Logs` 服务负责采集并归一化设备侧日志Syslog / SNMP Trap并提供规则与字典等配置能力。前端需要在统一的后台界面中完成
1. 日志查询(查看归一化后的日志事件及详情)
2. Syslog 规则配置
3. Trap 规则配置
4. Trap 字典配置
5. Trap 屏蔽/抑制规则配置
本设计以当前代码库的后端模型与前端实现为准:后端路由在 `internal/routers/register.go`,前端页面在 `front/src/views/ops/pages/log-mgmt/**/index.vue`
---
## 2. 范围(页面数量与路由)
本模块共 5 个页面,对应后端 5 组资源(列表+CRUD 或列表+详情抽屉)。
| 页面 | 菜单/路由路径 | 前端组件 |
|---|---|---|
| 日志查询 | `/log-mgmt/entries` | `front/src/views/ops/pages/log-mgmt/entries/index.vue` |
| Syslog 匹配规则 | `/log-mgmt/syslog-rules` | `front/src/views/ops/pages/log-mgmt/syslog-rules/index.vue` |
| SNMP Trap 匹配规则 | `/log-mgmt/trap-rules` | `front/src/views/ops/pages/log-mgmt/trap-rules/index.vue` |
| Trap 字典 | `/log-mgmt/trap-dictionary` | `front/src/views/ops/pages/log-mgmt/trap-dictionary/index.vue` |
| Trap 屏蔽/抑制 | `/log-mgmt/trap-suppressions` | `front/src/views/ops/pages/log-mgmt/trap-suppressions/index.vue` |
路由与菜单配置参考:
- `front/src/router/local-menu-flat.ts` / `front/src/router/local-menu-items.ts`
- `front/src/views/ops/pages/system-settings/system-logs/index.vue`(页面入口按钮)
- `front/src/views/ops/pages/monitor/log/index.vue`(嵌入 `LogMgmtEntries`
---
## 3. 数据对象与接口映射
后端认证API 路由组启用 `middleware.JwtAuth(true)`
前端请求的 API Base`front/src/api/ops/logs.ts` 中为 `/Logs/v1`
### 3.1 日志事件entries
- 接口:`GET /Logs/v1/entries`
- 返回结构(前端类型):`LogEntriesResult``total``page``page_size``items`
- 日志事件字段(前端类型 `LogEvent`
- `id`
- `created_at`
- `source_kind``syslog` / `snmp_trap`
- `remote_addr`
- `raw_payload`
- `normalized_summary`
- `normalized_detail`
- `device_name`
- `severity_code`
- `trap_oid`
- `alert_sent`
后端实现:`internal/models/log_event.go``internal/logic/controllers/crud.go``ListLogEvents`)。
### 3.2 Syslog 规则syslog-rules
- 接口:
- `GET /Logs/v1/syslog-rules`
- `POST /Logs/v1/syslog-rules`
- `PUT /Logs/v1/syslog-rules/:id`
- `DELETE /Logs/v1/syslog-rules/:id`
- 规则字段(前端类型 `SyslogRule` / 后端 `SyslogRule`
- `id``created_at``updated_at`
- `name`
- `enabled`
- `priority`
- `device_name_contains`
- `keyword_regex`
- `alert_name`
- `severity_code`
- `policy_id`
后端实现:`internal/models/syslog_rule.go``internal/logic/controllers/crud.go`
### 3.3 Trap 规则trap-rules
- 接口:
- `GET /Logs/v1/trap-rules`
- `POST /Logs/v1/trap-rules`
- `PUT /Logs/v1/trap-rules/:id`
- `DELETE /Logs/v1/trap-rules/:id`
- 规则字段(前端类型 `TrapRule` / 后端 `TrapRule`
- `name`
- `enabled`
- `priority`
- `oid_prefix`
- `varbind_match_regex`
- `alert_name`
- `severity_code`
- `policy_id`
后端实现:`internal/models/trap_rule.go``internal/logic/controllers/crud.go`
### 3.4 Trap 字典trap-dictionary
- 接口:
- `GET /Logs/v1/trap-dictionary`
- `POST /Logs/v1/trap-dictionary`
- `PUT /Logs/v1/trap-dictionary/:id`
- `DELETE /Logs/v1/trap-dictionary/:id`
- 字典条目字段(前端类型 `TrapDictionaryEntry` / 后端 `TrapDictionaryEntry`
- `oid_prefix`后端约束uniqueIndex
- `title`
- `description`
- `severity_code`
- `recovery_message`
- `enabled`
后端实现:`internal/models/trap_dictionary.go``internal/logic/controllers/crud.go`
### 3.5 Trap 屏蔽/抑制trap-suppressions
- 接口:
- `GET /Logs/v1/trap-suppressions`
- `POST /Logs/v1/trap-suppressions`
- `PUT /Logs/v1/trap-suppressions/:id`
- `DELETE /Logs/v1/trap-suppressions/:id`
- 屏蔽规则字段(前端类型 `TrapShield` / 后端 `TrapShield`
- `name`
- `enabled`
- `source_ip_cidr`
- `oid_prefix`
- `interface_hint`
- `time_windows_json`JSON 字符串)
后端实现:`internal/models/trap_shield.go``internal/logic/controllers/crud.go`
---
## 4. 页面设计详情(逐页)
### 4.1 日志查询页(`/log-mgmt/entries`
目标:以“可筛选的列表 + 详情抽屉”方式查看归一化日志事件。
#### 1顶部筛选区
- 使用 `search-table` 组件
- 筛选项:`source_kind`(下拉)
- `全部`value=''
- `Syslog`value='syslog'
- `SNMP Trap`value='snmp_trap'
筛选触发:`@search` 调用 `handleSearch`,重置则 `@reset` 调用 `handleReset`
#### 2列表表格列Columns
表格由 `columns` 定义,主要列:
- `ID`
- `来源``source_kind`,通过 `sourceKindLabel()` 显示(`syslog`->`Syslog``snmp_trap`->`SNMP Trap`
- `时间``created_at`
- `来源地址``remote_addr`
- `设备``device_name`
- `级别``severity_code`
- `OID``trap_oid`
- `原始报文``raw_payload`
- 使用 slot `raw_payload`:省略显示,保留 `tooltip`
- `已告警``alert_sent`
- 使用 slot `alert_sent``a-tag`(已转发/否)
- `操作`slot `operations`
- `详情`:打开右侧抽屉
#### 3详情抽屉a-drawer
- 打开逻辑:点击表格行操作中的 `详情`,调用 `openDetail(record)`
- 抽屉展示:`a-descriptions`1 列bordered
- 展示字段:
- 来源类型(`source_kind`
- 采集时间(`created_at`
- 来源地址(`remote_addr`,空则 `-`
- 设备名(`device_name`
- 严重级别(`severity_code`
- Trap OID`trap_oid`
- 已转发告警(`alert_sent`
- 摘要(`normalized_summary`
- 详情(`normalized_detail``pre-block` 预格式化展示)
- 原始报文(`raw_payload``pre-block` 预格式化展示)
#### 4分页策略
- 分页参数由前端 `pagination.current/pageSize` 控制,并随筛选条件一起请求后端:
- 调用 `fetchLogEntries({ page, page_size, source_kind })`
### 4.2 Syslog 规则页(`/log-mgmt/syslog-rules`
目标:规则的“列表 + 新建/编辑弹窗 + 删除确认”。
#### 1通用列表与本地过滤
- 使用 `search-table`,并在前端进行“关键词本地过滤”,过滤字段:
- `name`
- `alert_name`
- `keyword_regex`
- 搜索输入字段:
- `keyword`label`关键词`placeholder`规则名 / 告警名`
说明:该页(以及 trap-*、dictionary、suppressions 三类列表页)采用“先拉取全量 -> 本地过滤 -> 切片分页”的方式。
#### 2表格列
- `ID`
- `名称``name`
- `优先级``priority`
- `启用``enabled`slot `enabled`tag启用/禁用)
- `设备名包含``device_name_contains`
- `关键字正则``keyword_regex`
- `告警名``alert_name`
- `级别``severity_code`
- `策略ID``policy_id`
- `操作`:编辑/删除
#### 3新建/编辑弹窗a-modal
- 弹窗标题:
- 新建:`新建 Syslog 规则`
- 编辑:`编辑规则 #${editingId}`
- 表单 `a-form`(布局 `vertical`
- 表单字段:
- `name``a-input`(必填)
- `enabled``a-switch`
- `priority``a-input-number`
- `device_name_contains``a-input`
- `keyword_regex``a-input`
- `alert_name``a-input`
- `severity_code``a-input`
- `policy_id``a-input-number`min=0
提交逻辑:
- 编辑:`updateSyslogRule(editingId, { ...formData })`
- 新建:`createSyslogRule({ ...formData })`
- 成功后关闭弹窗并刷新列表 `fetchList()`
#### 4删除确认
- `Modal.confirm` 二次确认
- 删除接口:`deleteSyslogRule(id)`
### 4.3 Trap 规则页(`/log-mgmt/trap-rules`
目标TrapRule 的列表+弹窗 CRUD与 Syslog 规则页同构。
#### 1本地过滤关键词
- 字段:`keyword`
- 匹配来源:
- `name`
- `oid_prefix`
- `alert_name`
#### 2表格列
- `ID``名称``优先级``启用`
- `OID 前缀``oid_prefix`
- `Varbind 正则``varbind_match_regex`
- `告警名``alert_name`
- `级别``severity_code`
- `策略ID``policy_id`
- 操作:编辑/删除
#### 3弹窗表单字段
- `name`(必填)
- `enabled`
- `priority`
- `oid_prefix`
- `varbind_match_regex`
- `alert_name`
- `severity_code`
- `policy_id`min=0
### 4.4 Trap 字典页(`/log-mgmt/trap-dictionary`
目标TrapDictionaryEntry 的列表+弹窗 CRUD。
#### 1本地过滤关键词
- 匹配字段:
- `oid_prefix`
- `title`
- `description`
#### 2表格列
- `ID`
- `OID 前缀``oid_prefix`
- `标题``title`
- `级别``severity_code`
- `启用``enabled`
- `描述``description`
- 操作:编辑/删除
#### 3弹窗表单字段
- `oid_prefix`(必填,建议提示“唯一前缀”)
- `title`(必填)
- `description``a-textarea`rows=3
- `severity_code`
- `enabled`
- `recovery_message``a-textarea`rows=2
### 4.5 Trap 屏蔽/抑制页(`/log-mgmt/trap-suppressions`
目标TrapShield 的列表+弹窗 CRUD并对 `time_windows_json` 做前端校验。
#### 1本地过滤关键词
- 匹配字段:
- `name`
- `oid_prefix`
- `source_ip_cidr`
#### 2表格列
- `ID`
- `名称``name`
- `启用``enabled`
- `源 IP/CIDR``source_ip_cidr`
- `OID 前缀``oid_prefix`
- `接口提示``interface_hint`
- 操作:编辑/删除
#### 3弹窗表单字段
- `name`(必填)
- `enabled`
- `source_ip_cidr`
- `oid_prefix`
- `interface_hint`
- `time_windows_json``a-textarea`rows=4placeholder=`{}`
#### 4time_windows_json JSON 校验
-`time_windows_json` 非空时:
-`trim` 后尝试 `JSON.parse(tw)`
- 校验失败:`Message.warning('时间窗 JSON 格式无效')` 并阻止提交
---
## 5. 页面交互一致性要求(实现要点)
为了保证各列表页体验一致,本模块约定:
1. 列表页使用统一的 `search-table` 布局(顶部搜索、表格、分页、刷新)
2. 规则类/字典/屏蔽页采用“拉取全量 -> 本地过滤 -> 切片分页”的方式
3. 创建/编辑统一使用 `a-modal`,提交按钮触发 `formRef.validate()`
4. 删除统一使用 `Modal.confirm`,成功后刷新列表并给出 `Message.success`
5. `trap-suppressions``time_windows_json` 进行 JSON 字符串合法性校验
---
## 6. 数据流(简图)
```mermaid
flowchart LR
UI[前端页面search-table + 表格/弹窗/抽屉)] --> API[front/src/api/ops/logs.ts]
API --> BE[后端路由 internal/routers/register.go]
BE --> DB[(Postgres)]
BE --> Refresh[ingest.Global.Refresh()(规则/字典/屏蔽变更后触发)]
```
---
## 7. 中优先级待办(已立项,未完成)
本节用于记录当前版本可用但尚未产品化完善的中优先级项,作为后续迭代输入。
### 7.1 Outbox 可观测性增强
当前状态:
- 已支持 `alert_outbox` 入队、重试、死信、手动重试;
- 已有基础列表查询接口和前端入口。
待完善内容:
- 增加 outbox 指标接口或埋点:
- `pending_count`
- `retrying_count`
- `dead_count`
- `dispatch_success_rate`
- `dispatch_latency_p95`
- 增加失败原因聚合视图(按 `last_error` 分类统计)。
- 增加任务生命周期字段(首次入队时间、最后发送时间)用于问题排查。
建议落地文件:
- 后端:`internal/logic/controllers/outbox.go``internal/ingest/alert_outbox.go`
- 前端:`front/src/views/ops/pages/log-mgmt/entries/index.vue`
### 7.2 分发状态模型统一(替代 bool
当前状态:
- `logs_events` 已新增 `dispatch_status`,并在 outbox 流程中维护状态。
- 历史字段 `alert_sent` 仍保留,用于兼容旧页面展示。
待完善内容:
- 明确状态枚举为:`not_applicable/pending/retrying/sent/dead`
- 前后端统一以 `dispatch_status` 作为主状态字段,`alert_sent` 逐步降级为派生字段或移除。
- 页面文案由“已告警”升级为“分发状态”主展示,避免语义歧义。
建议落地文件:
- 后端:`internal/models/log_event.go``internal/logic/controllers/crud.go`
- 前端:`front/src/api/ops/logs.ts``front/src/views/ops/pages/log-mgmt/entries/index.vue`
### 7.3 关键路径测试补齐
当前状态:
- 已有基础单测覆盖核心函数。
待完善内容:
- 增加资源事件安全链路测试:
- 验签失败/成功
- 超时事件拒绝
- 幂等事件重复提交
- 增加 outbox 重试链路测试:
- 发送成功更新状态
- 重试次数递增
- 超过阈值转 `dead`
- 增加资源冲突优先级测试:
- `server > collector > device`
建议落地文件:
- `internal/logic/controllers/resource_event_test.go`
- `internal/ingest/alert_outbox_test.go`
- `internal/ingest/resource_resolver_test.go`
---
## 8. 后续产品化规划Phase 3
本节对应“可运维与产品化”阶段,优先级低于中优先级修复项,但会显著提升系统可管理性。
### 8.1 规则发布流draft / publish / rollback
目标:
- 规则配置与生效状态解耦,降低误操作风险。
范围:
- 引入规则草稿态与发布态;
- 支持发布记录、回滚到历史版本;
- 变更需记录操作人、时间、变更说明。
接口建议:
- `POST /Logs/v1/rule-sets/:id/publish`
- `POST /Logs/v1/rule-sets/:id/rollback`
- `GET /Logs/v1/rule-sets/:id/history`
### 8.2 规则仿真/回放能力
目标:
- 上线前可验证规则命中结果,减少误报漏报。
范围:
- 输入样本报文syslog/trap执行仿真
- 返回命中链路(命中/未命中原因);
- 支持历史事件回放。
接口建议:
- `POST /Logs/v1/rule-sets/:id/simulate`
- `POST /Logs/v1/rule-sets/:id/replay`
### 8.3 指标与审计面板
目标:
- 建立“采集-匹配-分发”全链路可观测性。
范围:
- 采集侧:接收速率、解析失败率;
- 匹配侧:命中率、规则耗时;
- 分发侧:成功率、重试率、死信量;
- 安全侧:验签失败次数、重放拦截次数。
前端建议:
- 在日志管理模块增加“运行指标”页签;
- 对死信和验签失败提供快捷定位入口。
---
## 9. 未完成项执行顺序(建议)
为降低风险,建议按以下顺序推进:
1. **中优先级先完成**
- outbox 指标与失败聚合
- `dispatch_status` 主状态化
- 关键路径测试补齐
2. **再做产品化**
- 规则发布流
- 规则仿真/回放
- 指标与审计面板
验收建议:
- 每项功能完成后执行“单项验证 + 回归验证”,最后统一做端到端联调。

View File

@@ -23,9 +23,9 @@ Ingest:
AlertForward: AlertForward:
enabled: true enabled: true
base_url: https://ops-api.apinb.com base_url: https://ops-api.apinb.com
internal_key: "ops-alert" internal_key: ${LOGS_ALERT_SECRET}
default_policy_id: 0 default_policy_id: 0
ResourceEvent: ResourceEvent:
hmac_secret: "replace-with-dc-control-shared-secret" hmac_secret: ${DC_CONTROL_LOGS_EVENT_SECRET}
max_skew_secs: 300 max_skew_secs: 300

View File

@@ -23,9 +23,9 @@ Ingest:
AlertForward: AlertForward:
enabled: true enabled: true
base_url: https://ops-api.apinb.com base_url: https://ops-api.apinb.com
internal_key: "ops-alert" internal_key: ${LOGS_ALERT_SECRET}
default_policy_id: 0 default_policy_id: 0
ResourceEvent: ResourceEvent:
hmac_secret: "replace-with-dc-control-shared-secret" hmac_secret: ${DC_CONTROL_LOGS_EVENT_SECRET}
max_skew_secs: 300 max_skew_secs: 300

28
etc/ops-logs.service Normal file
View File

@@ -0,0 +1,28 @@
[Unit]
Description=OPS Logs Service
Wants=network-online.target
Requires=ops-mgt.service
After=network-online.target ops-mgt.service
PartOf=ops-stack.target
[Service]
Type=simple
WorkingDirectory=/data/app
EnvironmentFile=/data/app/etc/ops.env
Environment=BSM_RuntimeMode=prod
Environment=RUN_MODE=prod
Environment=BSM_Prefix=/data/app
ExecStart=/data/app/ops-logs
ExecStartPost=/data/app/systemd/wait-http.sh ops-logs http://127.0.0.1:12440/Logs/v1/ping/hello 60
Restart=on-failure
RestartSec=5s
TimeoutStartSec=75s
TimeoutStopSec=30s
KillSignal=SIGTERM
StandardOutput=append:/data/app/logs/logs.log
StandardError=inherit
SyslogIdentifier=ops-logs
LimitNOFILE=1048576
[Install]
WantedBy=ops-stack.target

View File

@@ -2,6 +2,7 @@ package config
import ( import (
"net" "net"
"strings"
"git.apinb.com/bsm-sdk/core/conf" "git.apinb.com/bsm-sdk/core/conf"
) )
@@ -29,16 +30,16 @@ type ResourceEventConf struct {
} }
type SrvConfig struct { type SrvConfig struct {
conf.Base `yaml:",inline"` conf.Base `yaml:",inline"`
Databases *conf.DBConf `yaml:"Databases"` Databases *conf.DBConf `yaml:"Databases"`
MicroService *conf.MicroServiceConf `yaml:"MicroService"` MicroService *conf.MicroServiceConf `yaml:"MicroService"`
Rpc map[string]conf.RpcConf `yaml:"Rpc"` Rpc map[string]conf.RpcConf `yaml:"Rpc"`
Gateway *conf.GatewayConf `yaml:"Gateway"` Gateway *conf.GatewayConf `yaml:"Gateway"`
Apm *conf.ApmConf `yaml:"APM"` Apm *conf.ApmConf `yaml:"APM"`
Etcd *conf.EtcdConf `yaml:"Etcd"` Etcd *conf.EtcdConf `yaml:"Etcd"`
AlertForward *AlertForwardConf `yaml:"AlertForward"` AlertForward *AlertForwardConf `yaml:"AlertForward"`
Ingest IngestConf `yaml:"Ingest"` Ingest IngestConf `yaml:"Ingest"`
ResourceEvent ResourceEventConf `yaml:"ResourceEvent"` ResourceEvent ResourceEventConf `yaml:"ResourceEvent"`
} }
func New(srvKey string) { func New(srvKey string) {
@@ -46,6 +47,12 @@ func New(srvKey string) {
Spec.Port = conf.CheckPort(Spec.Port) Spec.Port = conf.CheckPort(Spec.Port)
Spec.BindIP = conf.CheckIP(Spec.BindIP) Spec.BindIP = conf.CheckIP(Spec.BindIP)
Spec.Addr = net.JoinHostPort(Spec.BindIP, Spec.Port) Spec.Addr = net.JoinHostPort(Spec.BindIP, Spec.Port)
conf.NotNil(Spec.Service, Spec.Cache) Spec.ResourceEvent.HMACSecret = strings.TrimSpace(Spec.ResourceEvent.HMACSecret)
conf.NotNil(Spec.Service, Spec.Cache, Spec.ResourceEvent.HMACSecret)
if Spec.AlertForward != nil && Spec.AlertForward.Enabled {
Spec.AlertForward.BaseURL = strings.TrimSpace(Spec.AlertForward.BaseURL)
Spec.AlertForward.InternalKey = strings.TrimSpace(Spec.AlertForward.InternalKey)
conf.NotNil(Spec.AlertForward.BaseURL, Spec.AlertForward.InternalKey)
}
conf.PrintInfo(Spec.Addr) conf.PrintInfo(Spec.Addr)
} }

View File

@@ -6,6 +6,7 @@ import (
"git.apinb.com/bsm-sdk/core/cache/redis" "git.apinb.com/bsm-sdk/core/cache/redis"
"git.apinb.com/bsm-sdk/core/conf" "git.apinb.com/bsm-sdk/core/conf"
"git.apinb.com/bsm-sdk/core/database" "git.apinb.com/bsm-sdk/core/database"
"git.apinb.com/bsm-sdk/core/types"
"git.apinb.com/bsm-sdk/core/vars" "git.apinb.com/bsm-sdk/core/vars"
"gorm.io/gorm" "gorm.io/gorm"
) )
@@ -14,7 +15,13 @@ func newDatabase(cfg *conf.DBConf) *gorm.DB {
if cfg == nil || len(cfg.Source) == 0 { if cfg == nil || len(cfg.Source) == 0 {
panic("database source is required") panic("database source is required")
} }
db, err := database.NewDatabase(cfg.Driver, cfg.Source, nil) db, err := database.NewDatabase(cfg.Driver, cfg.Source, &types.SqlOptions{
MaxIdleConns: vars.SqlOptionMaxIdleConns,
MaxOpenConns: vars.SqlOptionMaxOpenConns,
ConnMaxLifetime: vars.SqlOptionConnMaxLifetime,
LogStdout: false,
Debug: false,
})
if err != nil { if err != nil {
panic(fmt.Sprintf("database init failed: %v", err)) panic(fmt.Sprintf("database init failed: %v", err))
} }

View File

@@ -38,6 +38,7 @@ type AlertReceiveBody struct {
Labels map[string]string `json:"labels"` Labels map[string]string `json:"labels"`
Agent string `json:"agent"` Agent string `json:"agent"`
PolicyID uint `json:"policy_id"` PolicyID uint `json:"policy_id"`
Fingerprint string `json:"fingerprint,omitempty"`
State string `json:"state,omitempty"` State string `json:"state,omitempty"`
SourceEventKey string `json:"source_event_key"` SourceEventKey string `json:"source_event_key"`
TraceID string `json:"trace_id"` TraceID string `json:"trace_id"`
@@ -104,11 +105,14 @@ func forwardAlert(body AlertReceiveBody) error {
(details.Status != "firing" && details.Status != "resolved") || !validFingerprint(details.Fingerprint) { (details.Status != "firing" && details.Status != "resolved") || !validFingerprint(details.Fingerprint) {
return fmt.Errorf("Alert 响应缺少完整写入结果,无法确认转发成功;请稍后重试") return fmt.Errorf("Alert 响应缺少完整写入结果,无法确认转发成功;请稍后重试")
} }
if body.Fingerprint != "" && details.Fingerprint != body.Fingerprint {
return fmt.Errorf("Alert 返回的告警指纹与请求不一致;请稍后重试")
}
return nil return nil
} }
func postAlertPayload(cfg *config.AlertForwardConf, path string, payload []byte, traceID string) (*alertForwardResponse, error) { func postAlertPayload(cfg *config.AlertForwardConf, path string, payload []byte, traceID string) (*alertForwardResponse, error) {
req, err := http.NewRequest(http.MethodPost, cfg.BaseURL+path, bytes.NewReader(payload)) req, err := http.NewRequest(http.MethodPost, strings.TrimRight(strings.TrimSpace(cfg.BaseURL), "/")+path, bytes.NewReader(payload))
if err != nil { if err != nil {
return nil, fmt.Errorf("创建 Alert 转发请求失败:%w", err) return nil, fmt.Errorf("创建 Alert 转发请求失败:%w", err)
} }

View File

@@ -156,14 +156,19 @@ func processOneOutbox(db *gorm.DB, row models.AlertOutbox) error {
} }
return markOutboxDead(db, row, row.RetryCount, "invalid_payload: "+err.Error(), now) return markOutboxDead(db, row, row.RetryCount, "invalid_payload: "+err.Error(), now)
} }
forwardErr := forwardOutboxPayload(row.PayloadJSON, body, fmt.Sprintf("logs-outbox:%d", row.ID)) fallbackKey := fmt.Sprintf("logs-outbox:%d", row.ID)
occurredAt := row.CreatedAt.UTC()
if occurredAt.IsZero() {
occurredAt = time.Now().UTC()
}
forwardErr := forwardOutboxPayload(row.PayloadJSON, body, fallbackKey, occurredAt)
now, err := databaseNow(db) now, err := databaseNow(db)
if err != nil { if err != nil {
return err return err
} }
if forwardErr != nil { if forwardErr != nil {
if errors.Is(forwardErr, errAlertForwardDisabled) { if errors.Is(forwardErr, errAlertForwardDisabled) {
return markOutboxDead(db, row, row.RetryCount, forwardErr.Error(), now) return markOutboxWaiting(db, row, forwardErr.Error(), now)
} }
return markOutboxRetry(db, row, forwardErr.Error(), now) return markOutboxRetry(db, row, forwardErr.Error(), now)
} }
@@ -178,16 +183,35 @@ func databaseNow(db *gorm.DB) (time.Time, error) {
return now.UTC(), nil return now.UTC(), nil
} }
func forwardOutboxPayload(payloadJSON string, legacyBody AlertReceiveBody, fallbackSourceEventKey string) error { func forwardOutboxPayload(payloadJSON string, legacyBody AlertReceiveBody, fallbackSourceEventKey string, occurredAt time.Time) error {
var rawEvent RawEventIngestBody var rawEvent RawEventIngestBody
if err := json.Unmarshal([]byte(payloadJSON), &rawEvent); err == nil && rawEvent.SourceType != "" && len(rawEvent.RawPayload) > 0 { if err := json.Unmarshal([]byte(payloadJSON), &rawEvent); err == nil && rawEvent.SourceType != "" && len(rawEvent.RawPayload) > 0 {
if strings.TrimSpace(rawEvent.SourceEventKey) == "" {
rawEvent.SourceEventKey = fallbackSourceEventKey
}
if rawEvent.EventTime.IsZero() {
rawEvent.EventTime = occurredAt
}
rawEvent.TraceID = ensureAlertTraceID(rawEvent.TraceID, firstNonEmpty(rawEvent.SourceEventKey, fallbackSourceEventKey)) rawEvent.TraceID = ensureAlertTraceID(rawEvent.TraceID, firstNonEmpty(rawEvent.SourceEventKey, fallbackSourceEventKey))
return forwardRawEvent(rawEvent) return forwardRawEvent(rawEvent)
} }
if strings.TrimSpace(legacyBody.SourceEventKey) == "" {
legacyBody.SourceEventKey = fallbackSourceEventKey
}
if legacyBody.OccurredAt.IsZero() {
legacyBody.OccurredAt = occurredAt
}
legacyBody.TraceID = ensureAlertTraceID(legacyBody.TraceID, firstNonEmpty(legacyBody.SourceEventKey, fallbackSourceEventKey)) legacyBody.TraceID = ensureAlertTraceID(legacyBody.TraceID, firstNonEmpty(legacyBody.SourceEventKey, fallbackSourceEventKey))
return forwardAlert(legacyBody) return forwardAlert(legacyBody)
} }
func markOutboxWaiting(db *gorm.DB, row models.AlertOutbox, msg string, now time.Time) error {
return updateClaimedOutbox(db, row, map[string]interface{}{
"status": outboxStatusRetrying, "next_retry_at": now.Add(30 * time.Second),
"last_error": truncateError(msg, 1024), "lease_until": nil, "lease_owner": "",
}, "retrying")
}
func forwardRawEvent(body RawEventIngestBody) error { func forwardRawEvent(body RawEventIngestBody) error {
cfg := config.Spec.AlertForward cfg := config.Spec.AlertForward
if cfg == nil || !cfg.Enabled || cfg.BaseURL == "" { if cfg == nil || !cfg.Enabled || cfg.BaseURL == "" {
@@ -246,9 +270,16 @@ func markOutboxSent(db *gorm.DB, row models.AlertOutbox, now time.Time) error {
if result.RowsAffected == 0 { if result.RowsAffected == 0 {
return nil return nil
} }
return tx.Model(&models.LogEvent{}).Where("id = ? AND dispatch_outbox_id = ?", row.LogEventID, row.ID).Updates(map[string]interface{}{ eventResult := tx.Model(&models.LogEvent{}).Where("id = ? AND dispatch_outbox_id = ?", row.LogEventID, row.ID).Updates(map[string]interface{}{
"alert_sent": true, "dispatch_status": "sent", "alert_sent": true, "dispatch_status": "sent",
}).Error })
if eventResult.Error != nil {
return eventResult.Error
}
if eventResult.RowsAffected != 1 {
return fmt.Errorf("log event %d does not belong to outbox %d", row.LogEventID, row.ID)
}
return nil
}) })
} }
@@ -286,7 +317,14 @@ func updateClaimedOutbox(db *gorm.DB, row models.AlertOutbox, updates map[string
if result.RowsAffected == 0 { if result.RowsAffected == 0 {
return nil return nil
} }
return tx.Model(&models.LogEvent{}).Where("id = ? AND dispatch_outbox_id = ?", row.LogEventID, row.ID).Update("dispatch_status", eventStatus).Error eventResult := tx.Model(&models.LogEvent{}).Where("id = ? AND dispatch_outbox_id = ?", row.LogEventID, row.ID).Update("dispatch_status", eventStatus)
if eventResult.Error != nil {
return eventResult.Error
}
if eventResult.RowsAffected != 1 {
return fmt.Errorf("log event %d does not belong to outbox %d", row.LogEventID, row.ID)
}
return nil
}) })
} }

View File

@@ -1,6 +1,8 @@
package ingest package ingest
import ( import (
"crypto/sha256"
"encoding/hex"
"encoding/json" "encoding/json"
"fmt" "fmt"
"log" "log"
@@ -271,8 +273,9 @@ func (e *Engine) HandleSyslog(addr *net.UDPAddr, payload []byte) {
AlertName: matched.AlertName, Summary: summary, Description: summary, AlertName: matched.AlertName, Summary: summary, Description: summary,
SeverityCode: firstNonEmpty(matchDetails.SeverityCode, firstNonEmpty(matched.SeverityCode, sev)), SeverityCode: firstNonEmpty(matchDetails.SeverityCode, firstNonEmpty(matched.SeverityCode, sev)),
Value: parsed.Message, Labels: labels, Agent: "logs-syslog", PolicyID: matched.PolicyID, Value: parsed.Message, Labels: labels, Agent: "logs-syslog", PolicyID: matched.PolicyID,
State: "firing", SourceEventKey: logSourceEventKey(stored, "alert"), OccurredAt: occurredAt, RawData: rawBytes, State: matchDetails.State, SourceEventKey: logSourceEventKey(stored, "alert"), OccurredAt: occurredAt, RawData: rawBytes,
} }
body.Fingerprint = alertLifecycleFingerprint(body, matched.LifecycleKey, "")
return enqueueAlertWithDB(tx, stored.ID, body) return enqueueAlertWithDB(tx, stored.ID, body)
}); err != nil { }); err != nil {
log.Printf("logs: persist matched syslog event: %v", err) log.Printf("logs: persist matched syslog event: %v", err)
@@ -281,6 +284,7 @@ func (e *Engine) HandleSyslog(addr *net.UDPAddr, payload []byte) {
type syslogRuleMatch struct { type syslogRuleMatch struct {
Matched bool Matched bool
State string
ResourceUID string ResourceUID string
SeverityCode string SeverityCode string
Captures map[string]string Captures map[string]string
@@ -291,11 +295,12 @@ func syslogRuleMatches(rule *models.SyslogRule, device, message, rawLine string)
} }
func syslogRuleMatchDetails(rule *models.SyslogRule, device, message, rawLine string) syslogRuleMatch { func syslogRuleMatchDetails(rule *models.SyslogRule, device, message, rawLine string) syslogRuleMatch {
result := syslogRuleMatch{Captures: map[string]string{}} result := syslogRuleMatch{State: "firing", Captures: map[string]string{}}
deviceContains := strings.TrimSpace(rule.DeviceNameContains) deviceContains := strings.TrimSpace(rule.DeviceNameContains)
sourceMatch := strings.TrimSpace(rule.SourceMatch) sourceMatch := strings.TrimSpace(rule.SourceMatch)
keywordRegex := strings.TrimSpace(rule.KeywordRegex) keywordRegex := strings.TrimSpace(rule.KeywordRegex)
messageRegex := strings.TrimSpace(rule.MessageRegex) messageRegex := strings.TrimSpace(rule.MessageRegex)
recoveryRegex := strings.TrimSpace(rule.RecoveryMatchRegex)
if deviceContains == "" && sourceMatch == "" && keywordRegex == "" && messageRegex == "" { if deviceContains == "" && sourceMatch == "" && keywordRegex == "" && messageRegex == "" {
return result return result
} }
@@ -312,24 +317,35 @@ func syslogRuleMatchDetails(rule *models.SyslogRule, device, message, rawLine st
return result return result
} }
} }
for _, pattern := range []string{keywordRegex, messageRegex} { if recoveryRegex != "" {
if pattern == "" { re, err := regexp.Compile(recoveryRegex)
continue
}
re, err := regexp.Compile(pattern)
if err != nil { if err != nil {
return result return result
} }
matches := re.FindStringSubmatch(message) matches := firstRegexMatch(re, message, rawLine)
if matches == nil { if matches != nil {
matches = re.FindStringSubmatch(rawLine) mergeNamedCaptures(result.Captures, re, matches)
result.State = "resolved"
result.Matched = true
} }
if matches == nil {
return result
}
mergeNamedCaptures(result.Captures, re, matches)
} }
result.Matched = true if !result.Matched {
for _, pattern := range []string{keywordRegex, messageRegex} {
if pattern == "" {
continue
}
re, err := regexp.Compile(pattern)
if err != nil {
return result
}
matches := firstRegexMatch(re, message, rawLine)
if matches == nil {
return result
}
mergeNamedCaptures(result.Captures, re, matches)
}
result.Matched = true
}
if uid := extractWithNamedRegex(rule.ResourceUIDExtractRegex, "resource_uid", message, rawLine); uid != "" { if uid := extractWithNamedRegex(rule.ResourceUIDExtractRegex, "resource_uid", message, rawLine); uid != "" {
result.ResourceUID = normalizeExtractedResourceUID(uid) result.ResourceUID = normalizeExtractedResourceUID(uid)
} else if uid := result.Captures["resource_uid"]; uid != "" { } else if uid := result.Captures["resource_uid"]; uid != "" {
@@ -339,6 +355,15 @@ func syslogRuleMatchDetails(rule *models.SyslogRule, device, message, rawLine st
return result return result
} }
func firstRegexMatch(re *regexp.Regexp, values ...string) []string {
for _, value := range values {
if matches := re.FindStringSubmatch(value); matches != nil {
return matches
}
}
return nil
}
func mergeNamedCaptures(dst map[string]string, re *regexp.Regexp, matches []string) { func mergeNamedCaptures(dst map[string]string, re *regexp.Regexp, matches []string) {
names := re.SubexpNames() names := re.SubexpNames()
for i, name := range names { for i, name := range names {
@@ -366,7 +391,9 @@ func extractWithNamedRegex(pattern, groupName, message, rawLine string) string {
names := re.SubexpNames() names := re.SubexpNames()
for i, name := range names { for i, name := range names {
if i > 0 && name == groupName && i < len(matches) { if i > 0 && name == groupName && i < len(matches) {
return strings.TrimSpace(matches[i]) if value := strings.TrimSpace(matches[i]); value != "" {
return value
}
} }
} }
for i := 1; i < len(matches); i++ { for i := 1; i < len(matches); i++ {
@@ -506,8 +533,8 @@ func (e *Engine) HandleTrap(addr *net.UDPAddr, pkt *gosnmp.SnmpPacket) {
rules := e.trapRules rules := e.trapRules
e.mu.RUnlock() e.mu.RUnlock()
matched := firstMatchingTrapRule(rules, trapOID, fp) match := firstMatchingTrapRule(rules, trapOID, fp)
if matched == nil { if match.Rule == nil {
rawBytes, mErr := json.Marshal(fp) rawBytes, mErr := json.Marshal(fp)
if mErr != nil { if mErr != nil {
return return
@@ -537,6 +564,7 @@ func (e *Engine) HandleTrap(addr *net.UDPAddr, pkt *gosnmp.SnmpPacket) {
} }
return return
} }
matched := match.Rule
desc := readable desc := readable
if dict != nil && dict.RecoveryMessage != "" { if dict != nil && dict.RecoveryMessage != "" {
@@ -551,6 +579,10 @@ func (e *Engine) HandleTrap(addr *net.UDPAddr, pkt *gosnmp.SnmpPacket) {
"instance": addr.IP.String(), "instance": addr.IP.String(),
"job": "logs-trap", "job": "logs-trap",
} }
trapInstance := trapInstanceKey(pkt)
if trapInstance != "" {
labels["trap_instance"] = trapInstance
}
if matched.ID != 0 { if matched.ID != 0 {
labels["resource_type"] = "trap_rule" labels["resource_type"] = "trap_rule"
labels["resource_id"] = strconv.FormatUint(uint64(matched.ID), 10) labels["resource_id"] = strconv.FormatUint(uint64(matched.ID), 10)
@@ -587,9 +619,10 @@ func (e *Engine) HandleTrap(addr *net.UDPAddr, pkt *gosnmp.SnmpPacket) {
body := AlertReceiveBody{ body := AlertReceiveBody{
AlertName: firstNonEmpty(matched.AlertName, "SNMP Trap"), Summary: readable, Description: desc, AlertName: firstNonEmpty(matched.AlertName, "SNMP Trap"), Summary: readable, Description: desc,
SeverityCode: firstNonEmpty(matched.SeverityCode, sev), Value: string(vbJSON), Labels: labels, SeverityCode: firstNonEmpty(matched.SeverityCode, sev), Value: string(vbJSON), Labels: labels,
Agent: "logs-trap", PolicyID: matched.PolicyID, State: "firing", Agent: "logs-trap", PolicyID: matched.PolicyID, State: match.State,
SourceEventKey: logSourceEventKey(stored, "alert"), OccurredAt: occurredAt, RawData: rawBytes, SourceEventKey: logSourceEventKey(stored, "alert"), OccurredAt: occurredAt, RawData: rawBytes,
} }
body.Fingerprint = alertLifecycleFingerprint(body, matched.LifecycleKey, trapInstance)
return enqueueAlertWithDB(tx, stored.ID, body) return enqueueAlertWithDB(tx, stored.ID, body)
}); err != nil { }); err != nil {
log.Printf("logs: persist matched trap event: %v", err) log.Printf("logs: persist matched trap event: %v", err)
@@ -634,6 +667,27 @@ func trapVarbinds(pkt *gosnmp.SnmpPacket) []map[string]string {
return out return out
} }
func trapInstanceKey(pkt *gosnmp.SnmpPacket) string {
if pkt == nil {
return ""
}
for _, prefix := range []string{
"1.3.6.1.2.1.31.1.1.1.1.",
"1.3.6.1.2.1.2.2.1.2.",
"1.3.6.1.2.1.2.2.1.1.",
} {
for _, variable := range pkt.Variables {
name := normOID(variable.Name)
if strings.HasPrefix(name, prefix) {
if index := strings.TrimSpace(strings.TrimPrefix(name, prefix)); index != "" {
return index
}
}
}
}
return ""
}
func buildTrapReadable(trapOID string, dict *models.TrapDictionaryEntry, varbindSummary string) string { func buildTrapReadable(trapOID string, dict *models.TrapDictionaryEntry, varbindSummary string) string {
if dict != nil && firstNonEmpty(dict.Name, dict.Title) != "" { if dict != nil && firstNonEmpty(dict.Name, dict.Title) != "" {
return firstNonEmpty(dict.Name, dict.Title) + " (" + trapOID + ")" return firstNonEmpty(dict.Name, dict.Title) + " (" + trapOID + ")"
@@ -644,34 +698,49 @@ func buildTrapReadable(trapOID string, dict *models.TrapDictionaryEntry, varbind
return truncate(varbindSummary, 256) return truncate(varbindSummary, 256)
} }
func trapRuleMatches(rule *models.TrapRule, trapOID, varbindFP string) bool { type trapRuleMatch struct {
Rule *models.TrapRule
State string
}
func trapRuleState(rule *models.TrapRule, trapOID, varbindFP string) (string, bool) {
hasOID := strings.TrimSpace(rule.OIDPrefix) != "" hasOID := strings.TrimSpace(rule.OIDPrefix) != ""
hasRE := strings.TrimSpace(rule.VarbindMatchRegex) != "" hasRE := strings.TrimSpace(rule.VarbindMatchRegex) != ""
hasRecoveryRE := strings.TrimSpace(rule.RecoveryMatchRegex) != ""
if !hasOID && !hasRE { if !hasOID && !hasRE {
return false return "", false
} }
if hasOID && !strings.HasPrefix(normOID(trapOID), normOID(rule.OIDPrefix)) { if hasOID && !strings.HasPrefix(normOID(trapOID), normOID(rule.OIDPrefix)) {
return false return "", false
}
if hasRecoveryRE {
re, err := regexp.Compile(rule.RecoveryMatchRegex)
if err != nil {
return "", false
}
if re.MatchString(trapOID) || re.MatchString(varbindFP) {
return "resolved", true
}
} }
if hasRE { if hasRE {
re, err := regexp.Compile(rule.VarbindMatchRegex) re, err := regexp.Compile(rule.VarbindMatchRegex)
if err != nil { if err != nil {
return false return "", false
} }
if !re.MatchString(varbindFP) { if !re.MatchString(varbindFP) {
return false return "", false
} }
} }
return true return "firing", true
} }
func firstMatchingTrapRule(rules []models.TrapRule, trapOID, varbindFP string) *models.TrapRule { func firstMatchingTrapRule(rules []models.TrapRule, trapOID, varbindFP string) trapRuleMatch {
for i := range rules { for i := range rules {
if trapRuleMatches(&rules[i], trapOID, varbindFP) { if state, matched := trapRuleState(&rules[i], trapOID, varbindFP); matched {
return &rules[i] return trapRuleMatch{Rule: &rules[i], State: state}
} }
} }
return nil return trapRuleMatch{}
} }
func firstNonEmpty(a, b string) string { func firstNonEmpty(a, b string) string {
@@ -681,6 +750,29 @@ func firstNonEmpty(a, b string) string {
return b return b
} }
func alertLifecycleFingerprint(body AlertReceiveBody, lifecycleKey, instance string) string {
lifecycleKey = strings.TrimSpace(lifecycleKey)
if lifecycleKey == "" {
return ""
}
resourceUID := ""
sourceIP := ""
if body.Labels != nil {
resourceUID = strings.TrimSpace(body.Labels["resource_uid"])
sourceIP = strings.TrimSpace(body.Labels["ip"])
}
identity := strings.Join([]string{
strings.TrimSpace(body.Agent),
resourceUID,
sourceIP,
strconv.FormatUint(uint64(body.PolicyID), 10),
lifecycleKey,
strings.TrimSpace(instance),
}, "\x00")
sum := sha256.Sum256([]byte(identity))
return hex.EncodeToString(sum[:])
}
func (e *Engine) resolveResource(sourceIP, hostname string) (resourceRef, string) { func (e *Engine) resolveResource(sourceIP, hostname string) (resourceRef, string) {
e.mu.RLock() e.mu.RLock()
ipMap := e.resourceByIP ipMap := e.resourceByIP

View File

@@ -16,6 +16,7 @@ func validateSyslogRule(rule *models.SyslogRule) error {
}{ }{
{name: "keyword_regex", pattern: rule.KeywordRegex}, {name: "keyword_regex", pattern: rule.KeywordRegex},
{name: "message_regex", pattern: rule.MessageRegex}, {name: "message_regex", pattern: rule.MessageRegex},
{name: "recovery_match_regex", pattern: rule.RecoveryMatchRegex},
{name: "resource_uid_extract_regex", pattern: rule.ResourceUIDExtractRegex}, {name: "resource_uid_extract_regex", pattern: rule.ResourceUIDExtractRegex},
} }
for _, field := range regexFields { for _, field := range regexFields {
@@ -48,18 +49,25 @@ func validateSyslogRule(rule *models.SyslogRule) error {
strings.TrimSpace(rule.MessageRegex) == "" { strings.TrimSpace(rule.MessageRegex) == "" {
return fmt.Errorf("Syslog 规则的匹配条件全部为空,运行时永远不会命中;请至少填写 device_name_contains、source_match、keyword_regex、message_regex 中的一项") return fmt.Errorf("Syslog 规则的匹配条件全部为空,运行时永远不会命中;请至少填写 device_name_contains、source_match、keyword_regex、message_regex 中的一项")
} }
if strings.TrimSpace(rule.RecoveryMatchRegex) != "" && strings.TrimSpace(rule.LifecycleKey) == "" {
return fmt.Errorf("配置 recovery_match_regex 时 lifecycle_key 不能为空")
}
return nil return nil
} }
func validateTrapRule(rule *models.TrapRule) error { func validateTrapRule(rule *models.TrapRule) error {
if strings.TrimSpace(rule.VarbindMatchRegex) != "" { if err := validateOptionalRegex("varbind_match_regex", rule.VarbindMatchRegex); err != nil {
if _, err := regexp.Compile(rule.VarbindMatchRegex); err != nil { return err
return fmt.Errorf("varbind_match_regex 不是有效正则表达式:%v请修正 varbind_match_regex 后重试", err) }
} if err := validateOptionalRegex("recovery_match_regex", rule.RecoveryMatchRegex); err != nil {
return err
} }
if strings.TrimSpace(rule.OIDPrefix) == "" && strings.TrimSpace(rule.VarbindMatchRegex) == "" { if strings.TrimSpace(rule.OIDPrefix) == "" && strings.TrimSpace(rule.VarbindMatchRegex) == "" {
return fmt.Errorf("Trap 规则的匹配条件全部为空,运行时永远不会命中;请至少填写 oid_prefix、varbind_match_regex 中的一项") return fmt.Errorf("Trap 规则的匹配条件全部为空,运行时永远不会命中;请至少填写 oid_prefix、varbind_match_regex 中的一项")
} }
if strings.TrimSpace(rule.RecoveryMatchRegex) != "" && strings.TrimSpace(rule.LifecycleKey) == "" {
return fmt.Errorf("配置 recovery_match_regex 时 lifecycle_key 不能为空")
}
return nil return nil
} }

View File

@@ -108,6 +108,8 @@ func seedDefaultSyslogRules(db *gorm.DB) error {
KeywordRegex: "(?i)(link down|interface .* down|port .* down)", KeywordRegex: "(?i)(link down|interface .* down|port .* down)",
SourceMatch: "", SourceMatch: "",
MessageRegex: "(?i)(link down|interface .* down|port .* down|LINK_DOWN)", MessageRegex: "(?i)(link down|interface .* down|port .* down|LINK_DOWN)",
RecoveryMatchRegex: `(?i)(link[ _-]?up|interface .* up|port .* up|ifup)`,
LifecycleKey: "syslog-link-state",
AlertName: "Syslog链路中断", AlertName: "Syslog链路中断",
SeverityCode: "major", SeverityCode: "major",
SeverityMappingJSON: `{"(?i)(critical|fatal|emergency)":"critical","(?i)(error|LINK_DOWN|down)":"major","(?i)(warning|warn)":"warning"}`, SeverityMappingJSON: `{"(?i)(critical|fatal|emergency)":"critical","(?i)(error|LINK_DOWN|down)":"major","(?i)(warning|warn)":"warning"}`,
@@ -118,8 +120,10 @@ func seedDefaultSyslogRules(db *gorm.DB) error {
Name: "H3C-Syslog-接口中断", Name: "H3C-Syslog-接口中断",
Enabled: true, Enabled: true,
Priority: 120, Priority: 120,
SourceMatch: "h3c", DeviceNameContains: "h3c",
MessageRegex: `(?i)(LINK_DOWN|Interface .* down|port .* down)`, MessageRegex: `(?i)(LINK_DOWN|Interface .* down|port .* down)`,
RecoveryMatchRegex: `(?i)(link[ _-]?up|interface .* up|port .* up|ifup)`,
LifecycleKey: "h3c-syslog-interface-state",
AlertName: "H3C Syslog接口中断", AlertName: "H3C Syslog接口中断",
SeverityCode: "major", SeverityCode: "major",
SeverityMappingJSON: `{"(?i)(LINK_DOWN|down)":"major","(?i)(LINK_UP|up)":"info"}`, SeverityMappingJSON: `{"(?i)(LINK_DOWN|down)":"major","(?i)(LINK_UP|up)":"info"}`,
@@ -147,6 +151,8 @@ func seedDefaultSyslogRules(db *gorm.DB) error {
"source_match", "source_match",
"keyword_regex", "keyword_regex",
"message_regex", "message_regex",
"recovery_match_regex",
"lifecycle_key",
"alert_name", "alert_name",
"severity_code", "severity_code",
"severity_mapping_json", "severity_mapping_json",
@@ -162,14 +168,16 @@ func seedDefaultSyslogRules(db *gorm.DB) error {
func seedDefaultTrapRules(db *gorm.DB) error { func seedDefaultTrapRules(db *gorm.DB) error {
rows := []TrapRule{ rows := []TrapRule{
{ {
Name: "默认-Trap链路中断", Name: "默认-Trap链路中断",
Enabled: true, Enabled: true,
Priority: 100, Priority: 100,
OIDPrefix: "1.3.6.1.6.3.1.1.5", OIDPrefix: "1.3.6.1.6.3.1.1.5",
VarbindMatchRegex: "(?i)(linkdown|ifdown|down)", VarbindMatchRegex: `(?i)(1\.3\.6\.1\.6\.3\.1\.1\.5\.3([^0-9]|$)|\b(linkdown|ifdown|down)\b)`,
AlertName: "SNMP Trap链路中断", RecoveryMatchRegex: `(?i)(1\.3\.6\.1\.6\.3\.1\.1\.5\.4([^0-9]|$)|\b(linkup|ifup)\b)`,
SeverityCode: "major", LifecycleKey: "snmp-interface-link-state",
PolicyID: 0, AlertName: "SNMP Trap链路中断",
SeverityCode: "major",
PolicyID: 0,
}, },
} }
for _, row := range rows { for _, row := range rows {
@@ -190,6 +198,8 @@ func seedDefaultTrapRules(db *gorm.DB) error {
"priority", "priority",
"o_id_prefix", "o_id_prefix",
"varbind_match_regex", "varbind_match_regex",
"recovery_match_regex",
"lifecycle_key",
"alert_name", "alert_name",
"severity_code", "severity_code",
"policy_id", "policy_id",

View File

@@ -24,6 +24,10 @@ type SyslogRule struct {
KeywordRegex string `gorm:"size:512" json:"keyword_regex"` KeywordRegex string `gorm:"size:512" json:"keyword_regex"`
// MessageRegex 表示消息正文匹配的正则表达式。 // MessageRegex 表示消息正文匹配的正则表达式。
MessageRegex string `gorm:"size:1024" json:"message_regex"` MessageRegex string `gorm:"size:1024" json:"message_regex"`
// RecoveryMatchRegex 匹配同一生命周期的恢复消息。
RecoveryMatchRegex string `gorm:"size:1024" json:"recovery_match_regex"`
// LifecycleKey 将故障和恢复事件绑定到同一告警生命周期。
LifecycleKey string `gorm:"size:256" json:"lifecycle_key"`
// AlertName 表示告警名称。 // AlertName 表示告警名称。
AlertName string `gorm:"size:256" json:"alert_name"` AlertName string `gorm:"size:256" json:"alert_name"`
// SeverityCode 表示严重级别编码。 // SeverityCode 表示严重级别编码。

View File

@@ -20,6 +20,10 @@ type TrapRule struct {
OIDPrefix string `gorm:"size:512" json:"oid_prefix"` OIDPrefix string `gorm:"size:512" json:"oid_prefix"`
// VarbindMatchRegex 表示对 varbind 内容的正则匹配条件。 // VarbindMatchRegex 表示对 varbind 内容的正则匹配条件。
VarbindMatchRegex string `gorm:"size:512" json:"varbind_match_regex"` VarbindMatchRegex string `gorm:"size:512" json:"varbind_match_regex"`
// RecoveryMatchRegex 匹配同一生命周期的恢复 Trap OID 或 varbind。
RecoveryMatchRegex string `gorm:"size:1024" json:"recovery_match_regex"`
// LifecycleKey 将故障和恢复事件绑定到同一告警生命周期。
LifecycleKey string `gorm:"size:256" json:"lifecycle_key"`
// AlertName 表示告警名称。 // AlertName 表示告警名称。
AlertName string `gorm:"size:256" json:"alert_name"` AlertName string `gorm:"size:256" json:"alert_name"`
// SeverityCode 表示严重级别编码。 // SeverityCode 表示严重级别编码。