mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-10 17:26:38 +08:00
feat(push): add telegram bot push notification channel
Implement TelegramPusher in pkg/push, register config schema in channels_definition.go, add validation in model/push_channel.go, update task routing, and update settings-tab.tsx UI validation.
This commit is contained in:
@@ -107,70 +107,55 @@ import (
|
|||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## 模板模板渲染与支持的系统变量 (Template Rendering & Variables)
|
## 模板渲染与支持的系统变量 (Template Rendering & Variables)
|
||||||
|
|
||||||
消息的 `title`、`content` 以及 `ext` 字段中的字符串值都支持变量占位符替换,采用双花括号形式 `{{variable}}`。
|
消息的 `title`、`content` 以及 `ext` 字段中的字符串值都支持变量占位符替换,采用双花括号形式 `{{variable}}`。
|
||||||
|
|
||||||
### 1. 通用事件参数 (Common Variables)
|
### 1. 通用事件参数 (Common Variables)
|
||||||
任何通知事件触发时,均支持以下通用参数的渲染:
|
在 Wavelet 系统中,`user` 是一个通用的、必传的事件参数。如果在触发通知事件时未提供 `user`(或为 `nil`),底层 `EventTrigger` 会自动注入一个系统的虚拟用户(ID 为 999,昵称为“系统”)。因此,以下变量是所有通知事件均支持的通用渲染参数:
|
||||||
|
|
||||||
- `{{time}}`:事件发生/触发的具体时间(格式:`2006-01-02 15:04:05`)
|
- `{{time}}`:事件发生/触发的具体时间(格式:`2006-01-02 15:04:05`)
|
||||||
|
- `{{user.id}}`:触发用户/系统用户的 ID
|
||||||
|
- `{{user.username}}`:触发用户/系统用户的用户名
|
||||||
|
- `{{user.nickname}}`:触发用户/系统用户的昵称
|
||||||
|
- `{{user.email}}`:触发用户/系统用户的电子邮箱
|
||||||
|
- `{{user.phone}}`:触发用户/系统用户的手机号
|
||||||
|
- `{{user.bio}}`:触发用户/系统用户的个人简介
|
||||||
|
- `{{user.gender}}`:触发用户/系统用户的性别
|
||||||
|
- `{{user.location}}`:触发用户/系统用户的所在地
|
||||||
|
- `{{user.website}}`:触发用户/系统用户的个人网站
|
||||||
|
|
||||||
|
*(注:系统中的任何自定义事件,若传入了对应的复杂结构体,其结构体 JSON 字段均可通过扁平化点路径方式直接在模板中进行引用。)*
|
||||||
|
|
||||||
### 2. 特定事件携带的业务变量 (Event Specific Variables)
|
### 2. 特定事件携带的业务变量 (Event Specific Variables)
|
||||||
特定事件在触发时会携带复杂的业务对象(例如 `user`),可以通过点(`.`)路径语法获取其属性:
|
除了通用的 `user` 和 `time` 外,特定事件在触发时还可以携带额外的上下文参数:
|
||||||
|
|
||||||
- **管理员登录提醒 (`admin_login`)**
|
- **管理员登录提醒 (`admin_login`)**
|
||||||
- `{{user.id}}`:管理员 ID
|
|
||||||
- `{{user.username}}`:管理员用户名
|
|
||||||
- `{{user.email}}`:管理员邮箱地址
|
|
||||||
- `{{ip}}`:管理员登录来源的客户端 IP
|
- `{{ip}}`:管理员登录来源的客户端 IP
|
||||||
- `{{time}}`:管理员登录成功时间
|
- `{{time}}`:管理员登录成功时间
|
||||||
|
|
||||||
- **新用户注册提醒 (`user_registered`)**
|
### 3. 自定义消息通道的请求体变量说明 (Custom Channel JSON Variables)
|
||||||
- `{{user.id}}`:新注册用户 ID
|
在配置“自定义消息通道”时,其请求体 (JSON Schema) 支持以 `$` 开头的变量替换。支持的替换变量如下:
|
||||||
- `{{user.username}}`:新注册用户名
|
|
||||||
- `{{user.email}}`:新注册用户邮箱地址
|
```json
|
||||||
- `{{time}}`:注册成功时间
|
{
|
||||||
|
"title": "$title",
|
||||||
|
"description": "$description",
|
||||||
|
"content": "$content",
|
||||||
|
"url": "$url",
|
||||||
|
"to": "$to"
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
- `$title`:通知的标题(如:“管理员登录提醒”)
|
||||||
|
- `$description`:当前通知事件的描述
|
||||||
|
- `$content`:通知的具体渲染后正文内容
|
||||||
|
- `$url`:附加的操作或详情链接(若有)
|
||||||
|
- `$to`:当前派发的推送目标(如邮箱、ID 或 Chat ID,即 resolved target)
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## 严格遵循事项与防线 (Guardrails)
|
## 严格遵循事项与防线 (Guardrails)
|
||||||
|
|
||||||
### 1. 严格禁止包循环引用 (No Circular Dependencies)
|
### 1. 禁止绕过统一触发器 (Always Use EventTrigger)
|
||||||
- `custom_events` 包在定义事件时会导入并依赖 `push` 包的 `EventMetadata` 与 `DefaultTrigger` 等底层逻辑。
|
|
||||||
- **因此,`push` 包本身绝对不能导入 `custom_events` 包**(否则编译会抛出 package dependency cycle 错误)。
|
|
||||||
- 对新事件的自动注册只能在 `custom_events` 内部使用 `init()` 调用 `push.RegisterBuiltInEvent` 来实现。
|
|
||||||
|
|
||||||
### 2. 单元测试防循环依赖隔离 (Unit Testing Isolation)
|
|
||||||
- 在对 `push` 包自身进行单元测试(如 `push_test.go`)时,由于 `custom_events` 依赖 `push`,测试文件 `push_test.go` 也无法直接导入 `custom_events`。
|
|
||||||
- **防线策略**:在 `push_test.go` 的 `init()` 中本地声明并使用 `RegisterBuiltInEvent` 注册测试专用的 `EventMetadata`,以此来完成 `push` 包的隔离自测。
|
|
||||||
|
|
||||||
### 3. 禁止绕过统一触发器 (Always Use EventTrigger)
|
|
||||||
- 所有推送请求必须经过 `EventTrigger.Trigger`,以确保进行“事件是否启用”、“目标渠道过滤”、“全局推送配置读取”及“发送日志审计”等流程。
|
- 所有推送请求必须经过 `EventTrigger.Trigger`,以确保进行“事件是否启用”、“目标渠道过滤”、“全局推送配置读取”及“发送日志审计”等流程。
|
||||||
|
|
||||||
### 4. 代码质量与零 Lint 警报
|
|
||||||
- **魔法值防范**:对于推送日志级别(如 `"INFO"`),不要在多个文件里写硬编码字符串,应统一在 `constants.go` 中引用 `defaultLevelInfo` 常量。
|
|
||||||
- **命名规范**:不要定义容易造成 Stuttering 的导出类型,例如在 `push` 包内不要使用 `PushSendPayload`,应重命名为 `SendPayload`。
|
|
||||||
|
|
||||||
---
|
|
||||||
|
|
||||||
## 验证计划 (Verification Guide)
|
|
||||||
|
|
||||||
1. **授权许可头部补全**:
|
|
||||||
新增/修改 Go 文件后,运行自动许可证生成:
|
|
||||||
```bash
|
|
||||||
make license
|
|
||||||
```
|
|
||||||
2. **Swagger 文档重新构建**:
|
|
||||||
如果修改了路由或 Swagger 注释:
|
|
||||||
```bash
|
|
||||||
make swagger
|
|
||||||
```
|
|
||||||
3. **静态代码质量门禁 (0 Issues)**:
|
|
||||||
运行静态检查,必须保证后端与前端均无任何警告:
|
|
||||||
```bash
|
|
||||||
make code-check
|
|
||||||
```
|
|
||||||
4. **单元测试通过**:
|
|
||||||
```bash
|
|
||||||
go test ./internal/apps/admin/push/...
|
|
||||||
```
|
|
||||||
|
|||||||
@@ -28,6 +28,7 @@
|
|||||||
- Go skills:使用针对性的 `go-*` skills 来获取 Go 实现细节,如测试、错误处理、包、Context、并发、日志、文档和审查。
|
- Go skills:使用针对性的 `go-*` skills 来获取 Go 实现细节,如测试、错误处理、包、Context、并发、日志、文档和审查。
|
||||||
- `shadcn`:在添加、修改或组合 shadcn/ui 组件时使用。
|
- `shadcn`:在添加、修改或组合 shadcn/ui 组件时使用。
|
||||||
- `code-review-skill`:在提交 PR 之前使用,检查代码质量、样式、潜在错误和最佳实践。
|
- `code-review-skill`:在提交 PR 之前使用,检查代码质量、样式、潜在错误和最佳实践。
|
||||||
|
- `push-notification`:在添加或修改系统通知推送事件、修改消息推送底层设计、调用统一触发器投递消息、或开发带消息推送功能的业务功能时使用。
|
||||||
|
|
||||||
## 严格遵循事项 (Guardrails)
|
## 严格遵循事项 (Guardrails)
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -1249,7 +1249,7 @@ const docTemplate = `{
|
|||||||
"SessionCookie": []
|
"SessionCookie": []
|
||||||
}
|
}
|
||||||
],
|
],
|
||||||
"description": "返回系统支持的所有消息通道类型(如飞书、邮件、自定义)的动态表单定义,需要管理员权限",
|
"description": "返回系统支持的所有消息通道类型(如飞书、邮件、自定义、Telegram)的动态表单定义,需要管理员权限",
|
||||||
"produces": [
|
"produces": [
|
||||||
"application/json"
|
"application/json"
|
||||||
],
|
],
|
||||||
|
|||||||
+1
-1
@@ -1242,7 +1242,7 @@
|
|||||||
"SessionCookie": []
|
"SessionCookie": []
|
||||||
}
|
}
|
||||||
],
|
],
|
||||||
"description": "返回系统支持的所有消息通道类型(如飞书、邮件、自定义)的动态表单定义,需要管理员权限",
|
"description": "返回系统支持的所有消息通道类型(如飞书、邮件、自定义、Telegram)的动态表单定义,需要管理员权限",
|
||||||
"produces": [
|
"produces": [
|
||||||
"application/json"
|
"application/json"
|
||||||
],
|
],
|
||||||
|
|||||||
+1
-1
@@ -2199,7 +2199,7 @@ paths:
|
|||||||
- admin-push
|
- admin-push
|
||||||
/api/v1/admin/push/channels/definitions:
|
/api/v1/admin/push/channels/definitions:
|
||||||
get:
|
get:
|
||||||
description: 返回系统支持的所有消息通道类型(如飞书、邮件、自定义)的动态表单定义,需要管理员权限
|
description: 返回系统支持的所有消息通道类型(如飞书、邮件、自定义、Telegram)的动态表单定义,需要管理员权限
|
||||||
produces:
|
produces:
|
||||||
- application/json
|
- application/json
|
||||||
responses:
|
responses:
|
||||||
|
|||||||
@@ -178,9 +178,9 @@ export function SettingsTab() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// 协议安全校验(非邮件服务强制 HTTPS 协议)
|
// 协议安全校验(非邮件服务且配置了地址时,强制 HTTPS 协议)
|
||||||
if (channelType !== "email") {
|
if (channelType !== "email") {
|
||||||
if (!channelUrl.startsWith("https://")) {
|
if (channelUrl && !channelUrl.startsWith("https://")) {
|
||||||
toast.error("地址必须以 https:// 开头以确保安全性")
|
toast.error("地址必须以 https:// 开头以确保安全性")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -20,7 +20,7 @@ import (
|
|||||||
|
|
||||||
// ListChannelDefinitions 获取各种消息通道的表单配置定义列表
|
// ListChannelDefinitions 获取各种消息通道的表单配置定义列表
|
||||||
// @Summary 获取所有消息通道配置字段定义
|
// @Summary 获取所有消息通道配置字段定义
|
||||||
// @Description 返回系统支持的所有消息通道类型(如飞书、邮件、自定义)的动态表单定义,需要管理员权限
|
// @Description 返回系统支持的所有消息通道类型(如飞书、邮件、自定义、Telegram)的动态表单定义,需要管理员权限
|
||||||
// @Tags admin-push
|
// @Tags admin-push
|
||||||
// @Produce json
|
// @Produce json
|
||||||
// @Security SessionCookie
|
// @Security SessionCookie
|
||||||
@@ -279,6 +279,7 @@ func TestChannel(c *gin.Context) {
|
|||||||
c.JSON(http.StatusBadRequest, util.Err(err.Error()))
|
c.JSON(http.StatusBadRequest, util.Err(err.Error()))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
url = tempChannel.URL
|
||||||
|
|
||||||
var config pkgpush.Config
|
var config pkgpush.Config
|
||||||
var renderedJSON string
|
var renderedJSON string
|
||||||
@@ -298,6 +299,13 @@ func TestChannel(c *gin.Context) {
|
|||||||
Key: token,
|
Key: token,
|
||||||
Secret: other,
|
Secret: other,
|
||||||
}
|
}
|
||||||
|
case channelTelegram:
|
||||||
|
config = pkgpush.Config{
|
||||||
|
Channel: channelTelegram,
|
||||||
|
URL: url,
|
||||||
|
Secret: token,
|
||||||
|
Key: other,
|
||||||
|
}
|
||||||
default:
|
default:
|
||||||
config = pkgpush.Config{
|
config = pkgpush.Config{
|
||||||
Channel: channelCustom,
|
Channel: channelCustom,
|
||||||
|
|||||||
@@ -56,8 +56,8 @@ func ListDefinitions() []Definition {
|
|||||||
defMu.RLock()
|
defMu.RLock()
|
||||||
defer defMu.RUnlock()
|
defer defMu.RUnlock()
|
||||||
|
|
||||||
// We want a stable order: custom, lark, email
|
// We want a stable order: custom, lark, telegram, email
|
||||||
order := []string{channelCustom, channelLark, channelEmail}
|
order := []string{channelCustom, channelLark, channelTelegram, channelEmail}
|
||||||
res := make([]Definition, 0, len(definitions))
|
res := make([]Definition, 0, len(definitions))
|
||||||
for _, t := range order {
|
for _, t := range order {
|
||||||
if d, ok := definitions[t]; ok {
|
if d, ok := definitions[t]; ok {
|
||||||
@@ -140,6 +140,39 @@ func init() {
|
|||||||
},
|
},
|
||||||
})
|
})
|
||||||
|
|
||||||
|
// Register Telegram channel
|
||||||
|
RegisterChannelDefinition(Definition{
|
||||||
|
Type: channelTelegram,
|
||||||
|
Name: "Telegram 机器人",
|
||||||
|
Description: "配置 Telegram 机器人推送消息。",
|
||||||
|
Fields: []Field{
|
||||||
|
{
|
||||||
|
Key: KeyURL,
|
||||||
|
Label: "API 基础地址 (可选)",
|
||||||
|
Type: TypeText,
|
||||||
|
Required: false,
|
||||||
|
Placeholder: "https://api.telegram.org",
|
||||||
|
Description: "接口请求的 HTTPS 基础地址,留空默认为 https://api.telegram.org",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
Key: KeyToken,
|
||||||
|
Label: "机器人 Token (Bot Token)",
|
||||||
|
Type: TypePassword,
|
||||||
|
Required: true,
|
||||||
|
Placeholder: "在此输入 Telegram 机器人的 Bot Token",
|
||||||
|
Description: "通过 BotFather 申请到的机器人 Access Token",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
Key: KeyOther,
|
||||||
|
Label: "默认会话 ID (Chat ID) (可选)",
|
||||||
|
Type: TypeText,
|
||||||
|
Required: false,
|
||||||
|
Placeholder: "例如 -100123456789 或 @channel_name",
|
||||||
|
Description: "默认的消息接收 Chat ID。如果通知事件中未配置 targets,将推送到此 ID",
|
||||||
|
},
|
||||||
|
},
|
||||||
|
})
|
||||||
|
|
||||||
// Register Email channel
|
// Register Email channel
|
||||||
RegisterChannelDefinition(Definition{
|
RegisterChannelDefinition(Definition{
|
||||||
Type: channelEmail,
|
Type: channelEmail,
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ const (
|
|||||||
channelCustom = "custom"
|
channelCustom = "custom"
|
||||||
channelEmail = "email"
|
channelEmail = "email"
|
||||||
channelLark = "lark"
|
channelLark = "lark"
|
||||||
|
channelTelegram = "telegram"
|
||||||
defaultLevelInfo = "INFO"
|
defaultLevelInfo = "INFO"
|
||||||
keyTitle = "title"
|
keyTitle = "title"
|
||||||
keyContent = "content"
|
keyContent = "content"
|
||||||
|
|||||||
@@ -362,6 +362,13 @@ func (t *EventTrigger) enqueueSingleCustomPushChannelTask(ctx context.Context, m
|
|||||||
Key: token, // SMTP Username
|
Key: token, // SMTP Username
|
||||||
Secret: other, // SMTP Password
|
Secret: other, // SMTP Password
|
||||||
}
|
}
|
||||||
|
case channelTelegram:
|
||||||
|
config = pkgpush.Config{
|
||||||
|
Channel: channelTelegram,
|
||||||
|
URL: channel.URL,
|
||||||
|
Secret: channel.Token, // Telegram Bot Token
|
||||||
|
Key: channel.Other, // Default Chat ID
|
||||||
|
}
|
||||||
default: // custom
|
default: // custom
|
||||||
config = pkgpush.Config{
|
config = pkgpush.Config{
|
||||||
Channel: channelCustom,
|
Channel: channelCustom,
|
||||||
|
|||||||
@@ -620,6 +620,17 @@ func TestPushChannelAPI(t *testing.T) {
|
|||||||
c6 := &model.PushChannel{Name: "lark_channel", Type: "lark", URL: "https://open.feishu.cn", Other: ""}
|
c6 := &model.PushChannel{Name: "lark_channel", Type: "lark", URL: "https://open.feishu.cn", Other: ""}
|
||||||
assert.NoError(t, c6.Validate())
|
assert.NoError(t, c6.Validate())
|
||||||
|
|
||||||
|
// Telegram 渠道校验
|
||||||
|
cTelegramErr := &model.PushChannel{Name: "tg_channel", Type: "telegram", URL: "https://api.telegram.org", Token: "", Other: ""}
|
||||||
|
assert.Error(t, cTelegramErr.Validate())
|
||||||
|
|
||||||
|
cTelegramErr2 := &model.PushChannel{Name: "tg_channel", Type: "telegram", URL: "http://api.telegram.org", Token: "123:abc", Other: ""}
|
||||||
|
assert.Error(t, cTelegramErr2.Validate())
|
||||||
|
|
||||||
|
cTelegramOk := &model.PushChannel{Name: "tg_channel", Type: "telegram", URL: "", Token: "123:abc", Other: "-100123"}
|
||||||
|
assert.NoError(t, cTelegramOk.Validate())
|
||||||
|
assert.Equal(t, "https://api.telegram.org", cTelegramOk.URL)
|
||||||
|
|
||||||
// 邮件配置校验:允许空配置以复用系统全局设置
|
// 邮件配置校验:允许空配置以复用系统全局设置
|
||||||
c7 := &model.PushChannel{Name: "email_channel", Type: "email", URL: "", Token: "", Other: ""}
|
c7 := &model.PushChannel{Name: "email_channel", Type: "email", URL: "", Token: "", Other: ""}
|
||||||
assert.NoError(t, c7.Validate())
|
assert.NoError(t, c7.Validate())
|
||||||
|
|||||||
@@ -12,6 +12,7 @@ import (
|
|||||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// StorageReadOnly checks if the storage system is in read-only maintenance mode.
|
||||||
func StorageReadOnly(ctx context.Context) bool {
|
func StorageReadOnly(ctx context.Context) bool {
|
||||||
execution, ok, err := latestStorageMigrationExecution(ctx)
|
execution, ok, err := latestStorageMigrationExecution(ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -19,6 +19,8 @@ const (
|
|||||||
TypeCustom = "custom"
|
TypeCustom = "custom"
|
||||||
// TypeEmail 邮件推送消息通道类型
|
// TypeEmail 邮件推送消息通道类型
|
||||||
TypeEmail = "email"
|
TypeEmail = "email"
|
||||||
|
// TypeTelegram 电报机器人推送消息通道类型
|
||||||
|
TypeTelegram = "telegram"
|
||||||
)
|
)
|
||||||
|
|
||||||
// PushChannel 消息通道模型
|
// PushChannel 消息通道模型
|
||||||
@@ -53,6 +55,10 @@ func (pc *PushChannel) Validate() error {
|
|||||||
pc.Type = TypeCustom
|
pc.Type = TypeCustom
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if pc.Type == TypeTelegram && pc.URL == "" {
|
||||||
|
pc.URL = "https://api.telegram.org"
|
||||||
|
}
|
||||||
|
|
||||||
if pc.Name == "" {
|
if pc.Name == "" {
|
||||||
return errors.New("channel name is required")
|
return errors.New("channel name is required")
|
||||||
}
|
}
|
||||||
@@ -77,6 +83,10 @@ func (pc *PushChannel) Validate() error {
|
|||||||
return validateJSON(pc.Other)
|
return validateJSON(pc.Other)
|
||||||
case TypeEmail:
|
case TypeEmail:
|
||||||
// Email channel SMTP configs fall back to global settings, so they are not required to be filled.
|
// Email channel SMTP configs fall back to global settings, so they are not required to be filled.
|
||||||
|
case TypeTelegram:
|
||||||
|
if pc.Token == "" {
|
||||||
|
return errors.New("telegram bot token is required")
|
||||||
|
}
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
// Copyright 2026 Arctel.net
|
// Copyright 2026 Arctel.net
|
||||||
// SPDX-License-Identifier: Apache-2.0
|
// SPDX-License-Identifier: Apache-2.0
|
||||||
|
|
||||||
|
// Package service implements Wavelet's core background service processes.
|
||||||
package service
|
package service
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
|||||||
+2
-2
@@ -115,7 +115,7 @@ func (p *LarkPusher) Send(ctx context.Context, cfg Config, _ string, body map[st
|
|||||||
content := msg.Content
|
content := msg.Content
|
||||||
level := strings.ToUpper(msg.Level)
|
level := strings.ToUpper(msg.Level)
|
||||||
if level == "" {
|
if level == "" {
|
||||||
level = "INFO"
|
level = levelInfo
|
||||||
}
|
}
|
||||||
|
|
||||||
headerColor := "blue"
|
headerColor := "blue"
|
||||||
@@ -170,7 +170,7 @@ func (p *LarkPusher) Send(ctx context.Context, cfg Config, _ string, body map[st
|
|||||||
content = strings.Join(parts, "\n")
|
content = strings.Join(parts, "\n")
|
||||||
}
|
}
|
||||||
|
|
||||||
level := "INFO"
|
level := levelInfo
|
||||||
if l, ok := body["level"].(string); ok && l != "" {
|
if l, ok := body["level"].(string); ok && l != "" {
|
||||||
level = strings.ToUpper(l)
|
level = strings.ToUpper(l)
|
||||||
}
|
}
|
||||||
|
|||||||
+4
-1
@@ -10,7 +10,10 @@ import (
|
|||||||
"sync"
|
"sync"
|
||||||
)
|
)
|
||||||
|
|
||||||
const defaultTitle = "系统通知"
|
const (
|
||||||
|
defaultTitle = "系统通知"
|
||||||
|
levelInfo = "INFO"
|
||||||
|
)
|
||||||
|
|
||||||
// Config 基础通知渠道配置
|
// Config 基础通知渠道配置
|
||||||
type Config struct {
|
type Config struct {
|
||||||
|
|||||||
@@ -0,0 +1,157 @@
|
|||||||
|
// Copyright 2026 Arctel.net
|
||||||
|
// SPDX-License-Identifier: Apache-2.0
|
||||||
|
|
||||||
|
package push
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
|
"fmt"
|
||||||
|
"net/http"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
func init() {
|
||||||
|
Register("telegram", &TelegramPusher{})
|
||||||
|
}
|
||||||
|
|
||||||
|
// TelegramPusher Telegram 机器人推送实现
|
||||||
|
type TelegramPusher struct{}
|
||||||
|
|
||||||
|
type telegramMessageRequest struct {
|
||||||
|
ChatID string `json:"chat_id"`
|
||||||
|
Text string `json:"text"`
|
||||||
|
ParseMode string `json:"parse_mode,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type telegramErrorResponse struct {
|
||||||
|
Ok bool `json:"ok"`
|
||||||
|
ErrorCode int `json:"error_code"`
|
||||||
|
Description string `json:"description"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// Send 执行 Telegram 消息发送
|
||||||
|
//
|
||||||
|
//nolint:cyclop
|
||||||
|
func (p *TelegramPusher) Send(ctx context.Context, cfg Config, target string, body map[string]any, template string, _ map[string]any) error {
|
||||||
|
if cfg.Secret == "" {
|
||||||
|
return errors.New("telegram: Bot Token (Secret) is required")
|
||||||
|
}
|
||||||
|
|
||||||
|
chatID := target
|
||||||
|
if chatID == "" {
|
||||||
|
chatID = cfg.Key // Use default chat ID (Key) if target is blank
|
||||||
|
}
|
||||||
|
if chatID == "" {
|
||||||
|
return errors.New("telegram: chat_id (target or default Key) is required")
|
||||||
|
}
|
||||||
|
|
||||||
|
baseURL := cfg.URL
|
||||||
|
if baseURL == "" {
|
||||||
|
baseURL = "https://api.telegram.org"
|
||||||
|
}
|
||||||
|
baseURL = strings.TrimSuffix(baseURL, "/")
|
||||||
|
|
||||||
|
title := defaultTitle
|
||||||
|
if t, ok := body["title"].(string); ok && t != "" {
|
||||||
|
title = t
|
||||||
|
}
|
||||||
|
content := ""
|
||||||
|
if c, ok := body["content"].(string); ok && c != "" {
|
||||||
|
content = c
|
||||||
|
} else {
|
||||||
|
var parts []string
|
||||||
|
for k, v := range body {
|
||||||
|
parts = append(parts, fmt.Sprintf("<b>%s</b>: %v", k, v))
|
||||||
|
}
|
||||||
|
content = strings.Join(parts, "\n")
|
||||||
|
}
|
||||||
|
level := levelInfo
|
||||||
|
if l, ok := body["level"].(string); ok && l != "" {
|
||||||
|
level = strings.ToUpper(l)
|
||||||
|
}
|
||||||
|
|
||||||
|
var text string
|
||||||
|
if template != "" {
|
||||||
|
text = ParseTemplate(template, body)
|
||||||
|
} else {
|
||||||
|
text = fmt.Sprintf("<b>[%s] %s</b>\n\n%s", escapeHTML(level), escapeHTML(title), escapeHTML(content))
|
||||||
|
}
|
||||||
|
|
||||||
|
// Try sending with HTML parse mode
|
||||||
|
err := p.sendMessage(ctx, baseURL, cfg.Secret, chatID, text, "HTML")
|
||||||
|
if err != nil {
|
||||||
|
// Fallback: send as plain text without parse mode
|
||||||
|
plainText := text
|
||||||
|
if template == "" {
|
||||||
|
plainText = fmt.Sprintf("[%s] %s\n\n%s", level, title, content)
|
||||||
|
}
|
||||||
|
fallbackErr := p.sendMessage(ctx, baseURL, cfg.Secret, chatID, plainText, "")
|
||||||
|
if fallbackErr != nil {
|
||||||
|
return fmt.Errorf("telegram: send message failed (fallback also failed): %w (original HTML error: %v)", fallbackErr, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// ValidateConfig 校验 Telegram 配置
|
||||||
|
func (p *TelegramPusher) ValidateConfig(cfg Config) error {
|
||||||
|
if cfg.Secret == "" {
|
||||||
|
return errors.New("bot Token (Secret) is required")
|
||||||
|
}
|
||||||
|
if cfg.URL != "" {
|
||||||
|
if !strings.HasPrefix(cfg.URL, "http://") && !strings.HasPrefix(cfg.URL, "https://") {
|
||||||
|
return errors.New("API base URL must start with http:// or https://")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (p *TelegramPusher) sendMessage(ctx context.Context, baseURL, token, chatID, text, parseMode string) error {
|
||||||
|
apiURL := fmt.Sprintf("%s/bot%s/sendMessage", baseURL, token)
|
||||||
|
|
||||||
|
reqPayload := telegramMessageRequest{
|
||||||
|
ChatID: chatID,
|
||||||
|
Text: text,
|
||||||
|
ParseMode: parseMode,
|
||||||
|
}
|
||||||
|
|
||||||
|
jsonData, err := json.Marshal(reqPayload)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("marshal request failed: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, apiURL, bytes.NewBuffer(jsonData))
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("create http request failed: %w", err)
|
||||||
|
}
|
||||||
|
httpReq.Header.Set("Content-Type", "application/json")
|
||||||
|
|
||||||
|
client := &http.Client{Timeout: 10 * time.Second}
|
||||||
|
resp, err := client.Do(httpReq)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("http request failed: %w", err)
|
||||||
|
}
|
||||||
|
defer func() { _ = resp.Body.Close() }()
|
||||||
|
|
||||||
|
if resp.StatusCode != http.StatusOK {
|
||||||
|
var errRes telegramErrorResponse
|
||||||
|
if decodeErr := json.NewDecoder(resp.Body).Decode(&errRes); decodeErr == nil {
|
||||||
|
return fmt.Errorf("http status %d: %s", resp.StatusCode, errRes.Description)
|
||||||
|
}
|
||||||
|
return fmt.Errorf("http status %s", resp.Status)
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func escapeHTML(s string) string {
|
||||||
|
s = strings.ReplaceAll(s, "&", "&")
|
||||||
|
s = strings.ReplaceAll(s, "<", "<")
|
||||||
|
s = strings.ReplaceAll(s, ">", ">")
|
||||||
|
return s
|
||||||
|
}
|
||||||
@@ -0,0 +1,116 @@
|
|||||||
|
// Copyright 2026 Arctel.net
|
||||||
|
// SPDX-License-Identifier: Apache-2.0
|
||||||
|
|
||||||
|
package push
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestTelegramPusher_Send(t *testing.T) {
|
||||||
|
t.Run("successful send with HTML parse mode", func(t *testing.T) {
|
||||||
|
var receivedReq telegramMessageRequest
|
||||||
|
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
assert.Equal(t, "/botmy-token/sendMessage", r.URL.Path)
|
||||||
|
assert.Equal(t, http.MethodPost, r.Method)
|
||||||
|
assert.Equal(t, "application/json", r.Header.Get("Content-Type"))
|
||||||
|
|
||||||
|
err := json.NewDecoder(r.Body).Decode(&receivedReq)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
w.WriteHeader(http.StatusOK)
|
||||||
|
_, _ = w.Write([]byte(`{"ok": true}`))
|
||||||
|
}))
|
||||||
|
defer server.Close()
|
||||||
|
|
||||||
|
pusher := &TelegramPusher{}
|
||||||
|
cfg := Config{
|
||||||
|
Channel: "telegram",
|
||||||
|
URL: server.URL,
|
||||||
|
Secret: "my-token",
|
||||||
|
}
|
||||||
|
body := map[string]any{
|
||||||
|
"title": "Alert",
|
||||||
|
"content": "Host down",
|
||||||
|
"level": "CRITICAL",
|
||||||
|
}
|
||||||
|
err := pusher.Send(context.Background(), cfg, "123456", body, "", nil)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
assert.Equal(t, "123456", receivedReq.ChatID)
|
||||||
|
assert.Contains(t, receivedReq.Text, "[CRITICAL] Alert")
|
||||||
|
assert.Contains(t, receivedReq.Text, "Host down")
|
||||||
|
assert.Equal(t, "HTML", receivedReq.ParseMode)
|
||||||
|
})
|
||||||
|
|
||||||
|
t.Run("fallback to plain text on HTML error", func(t *testing.T) {
|
||||||
|
var requests []*telegramMessageRequest
|
||||||
|
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
var req telegramMessageRequest
|
||||||
|
err := json.NewDecoder(r.Body).Decode(&req)
|
||||||
|
require.NoError(t, err)
|
||||||
|
requests = append(requests, &req)
|
||||||
|
|
||||||
|
if len(requests) == 1 {
|
||||||
|
w.WriteHeader(http.StatusBadRequest)
|
||||||
|
_, _ = w.Write([]byte(`{"ok": false, "error_code": 400, "description": "Bad Request: can't parse entities"}`))
|
||||||
|
} else {
|
||||||
|
w.WriteHeader(http.StatusOK)
|
||||||
|
_, _ = w.Write([]byte(`{"ok": true}`))
|
||||||
|
}
|
||||||
|
}))
|
||||||
|
defer server.Close()
|
||||||
|
|
||||||
|
pusher := &TelegramPusher{}
|
||||||
|
cfg := Config{
|
||||||
|
Channel: "telegram",
|
||||||
|
URL: server.URL,
|
||||||
|
Secret: "my-token",
|
||||||
|
}
|
||||||
|
body := map[string]any{
|
||||||
|
"title": "Alert & Info",
|
||||||
|
"content": "A < B comparison",
|
||||||
|
"level": "INFO",
|
||||||
|
}
|
||||||
|
err := pusher.Send(context.Background(), cfg, "123456", body, "", nil)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
require.Len(t, requests, 2)
|
||||||
|
assert.Equal(t, "HTML", requests[0].ParseMode)
|
||||||
|
assert.Equal(t, "", requests[1].ParseMode)
|
||||||
|
assert.Contains(t, requests[1].Text, "[INFO] Alert & Info")
|
||||||
|
assert.Contains(t, requests[1].Text, "A < B comparison")
|
||||||
|
})
|
||||||
|
|
||||||
|
t.Run("validation error", func(t *testing.T) {
|
||||||
|
pusher := &TelegramPusher{}
|
||||||
|
cfg := Config{
|
||||||
|
Channel: "telegram",
|
||||||
|
URL: "https://api.telegram.org",
|
||||||
|
}
|
||||||
|
err := pusher.ValidateConfig(cfg)
|
||||||
|
assert.Error(t, err)
|
||||||
|
|
||||||
|
cfg = Config{
|
||||||
|
Channel: "telegram",
|
||||||
|
URL: "ftp://api.telegram.org",
|
||||||
|
Secret: "token",
|
||||||
|
}
|
||||||
|
err = pusher.ValidateConfig(cfg)
|
||||||
|
assert.Error(t, err)
|
||||||
|
|
||||||
|
cfg = Config{
|
||||||
|
Channel: "telegram",
|
||||||
|
Secret: "token",
|
||||||
|
}
|
||||||
|
err = pusher.ValidateConfig(cfg)
|
||||||
|
assert.NoError(t, err)
|
||||||
|
})
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user