diff --git a/backend/go.mod b/backend/go.mod
index 3074113e..182fcc1e 100644
--- a/backend/go.mod
+++ b/backend/go.mod
@@ -77,13 +77,16 @@ require (
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.40.1 // indirect
github.com/aws/aws-sdk-go-v2/service/sts v1.47.1 // indirect
github.com/aws/smithy-go v1.28.1 // indirect
+ github.com/blinkbean/dingtalk v1.1.3 // indirect
github.com/boj/redistore v1.4.1 // indirect
+ github.com/bwmarrin/discordgo v0.29.0 // indirect
github.com/bytedance/gopkg v0.1.4 // indirect
github.com/bytedance/sonic v1.15.2 // indirect
github.com/bytedance/sonic/loader v0.5.2 // indirect
github.com/cenkalti/backoff/v5 v5.0.3 // indirect
github.com/cespare/xxhash/v2 v2.3.0 // indirect
github.com/cloudwego/base64x v0.1.7 // indirect
+ github.com/cschomburg/go-pushbullet v0.0.0-20171206132031-67759df45fbb // indirect
github.com/davecgh/go-spew v1.1.1 // indirect
github.com/dustin/go-humanize v1.0.1 // indirect
github.com/felixge/httpsnoop v1.1.0 // indirect
@@ -94,6 +97,7 @@ require (
github.com/go-faster/city v1.0.1 // indirect
github.com/go-faster/errors v0.7.1 // indirect
github.com/go-jose/go-jose/v4 v4.1.4 // indirect
+ github.com/go-lark/lark v1.16.0 // indirect
github.com/go-logr/logr v1.4.4 // indirect
github.com/go-logr/stdr v1.2.2 // indirect
github.com/go-openapi/jsonpointer v0.22.1 // indirect
@@ -112,6 +116,7 @@ require (
github.com/go-redis/redis_rate/v10 v10.0.1 // indirect
github.com/go-resty/resty/v2 v2.6.0 // indirect
github.com/go-sql-driver/mysql v1.10.0 // indirect
+ github.com/go-telegram-bot-api/telegram-bot-api v4.6.4+incompatible // indirect
github.com/go-viper/mapstructure/v2 v2.4.0 // indirect
github.com/goccy/go-json v0.10.6 // indirect
github.com/goccy/go-yaml v1.19.2 // indirect
@@ -119,6 +124,7 @@ require (
github.com/google/btree v1.0.0 // indirect
github.com/gorilla/context v1.1.2 // indirect
github.com/gorilla/securecookie v1.1.2 // indirect
+ github.com/gregdel/pushover v1.4.0 // indirect
github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 // indirect
github.com/hashicorp/go-version v1.9.0 // indirect
github.com/inconshreveable/mousetrap v1.1.0 // indirect
@@ -149,10 +155,13 @@ require (
github.com/sagikazarmark/locafero v0.12.0 // indirect
github.com/segmentio/asm v1.2.1 // indirect
github.com/sethvargo/go-retry v0.4.0 // indirect
+ github.com/slack-go/slack v0.29.0 // indirect
github.com/spf13/afero v1.15.0 // indirect
github.com/spf13/cast v1.10.0 // indirect
github.com/spf13/pflag v1.0.10 // indirect
+ github.com/stretchr/objx v0.5.3 // indirect
github.com/subosito/gotenv v1.6.0 // indirect
+ github.com/technoweenie/multipartstreamer v1.0.1 // indirect
github.com/tidwall/gjson v1.19.0 // indirect
github.com/tidwall/match v1.2.0 // indirect
github.com/tidwall/pretty v1.2.1 // indirect
diff --git a/backend/go.sum b/backend/go.sum
index 2c3a10df..4d17ad07 100644
--- a/backend/go.sum
+++ b/backend/go.sum
@@ -152,12 +152,16 @@ github.com/beorn7/perks v0.0.0-20180321164747-3a771d992973/go.mod h1:Dwedo/Wpr24
github.com/beorn7/perks v1.0.0/go.mod h1:KWe93zE9D1o94FZ5RNwFwVgaQK1VOXiVxmqh+CedLV8=
github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw=
github.com/bgentry/speakeasy v0.1.0/go.mod h1:+zsyZBPWlz7T6j88CTgSN5bM796AkVf0kBD4zp0CCIs=
+github.com/blinkbean/dingtalk v1.1.3 h1:MbidFZYom7DTFHD/YIs+eaI7kRy52kmWE/sy0xjo6E4=
+github.com/blinkbean/dingtalk v1.1.3/go.mod h1:9BaLuGSBqY3vT5hstValh48DbsKO7vaHaJnG9pXwbto=
github.com/boj/redistore v1.4.1 h1:lP9ZZWqKMq2RIqexlZX1w1ODSnegL+puxGIujkU5tIw=
github.com/boj/redistore v1.4.1/go.mod h1:c0Tvw6aMjslog4jHIAcNv6EtJM849YoOAhMY7JBbWpI=
github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs=
github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c=
github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA=
github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0=
+github.com/bwmarrin/discordgo v0.29.0 h1:FmWeXFaKUwrcL3Cx65c20bTRW+vOb6k8AnaP+EgjDno=
+github.com/bwmarrin/discordgo v0.29.0/go.mod h1:NJZpH+1AfhIcyQsPeuBKsUtYrRnjkyu0kIVMCHkZtRY=
github.com/bwmarrin/snowflake v0.3.0 h1:xm67bEhkKh6ij1790JB83OujPR5CzNe8QuQqAgISZN0=
github.com/bwmarrin/snowflake v0.3.0/go.mod h1:NdZxfVWX+oR6y2K0o6qAYv6gIOP9rjG0/E9WsDpxqwE=
github.com/bytedance/gopkg v0.1.4 h1:oZnQwnX82KAIWb7033bEwtxvTqXcYMxDBaQxo5JJHWM=
@@ -198,6 +202,8 @@ github.com/coreos/go-semver v0.3.0/go.mod h1:nnelYz7RCh+5ahJtPPxZlU+153eP4D4r3Ee
github.com/coreos/go-systemd/v22 v22.3.2/go.mod h1:Y58oyj3AT4RCenI/lSvhwexgC+NSVTIJ3seZv2GcEnc=
github.com/cpuguy83/go-md2man/v2 v2.0.6/go.mod h1:oOW0eioCTA6cOiMLiUPZOpcVxMig6NIQQ7OS05n1F4g=
github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E=
+github.com/cschomburg/go-pushbullet v0.0.0-20171206132031-67759df45fbb h1:7X9nrm+LNWdxzQOiCjy0G51rNUxbH35IDHCjAMvogyM=
+github.com/cschomburg/go-pushbullet v0.0.0-20171206132031-67759df45fbb/go.mod h1:RfQ9wji3fjcSEsQ+uFCtIh3+BXgcZum8Kt3JxvzYzlk=
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
@@ -258,6 +264,8 @@ github.com/go-jose/go-jose/v4 v4.1.4/go.mod h1:x4oUasVrzR7071A4TnHLGSPpNOm2a21K9
github.com/go-kit/kit v0.8.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as=
github.com/go-kit/kit v0.9.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as=
github.com/go-kit/log v0.1.0/go.mod h1:zbhenjAZHb184qTLMA9ZjW7ThYL0H2mk7Q6pNt4vbaY=
+github.com/go-lark/lark v1.16.0 h1:U6BwkLM9wrZedSM7cIiMofganr8PCvJN+M75w2lf2Gg=
+github.com/go-lark/lark v1.16.0/go.mod h1:6ltbSztPZRT6IaO9ZIQyVaY5pVp/KeMizDYtfZkU+vM=
github.com/go-logfmt/logfmt v0.3.0/go.mod h1:Qt1PoO58o5twSAckw1HlFXLmHsOX5/0LbT9GBnD5lWE=
github.com/go-logfmt/logfmt v0.4.0/go.mod h1:3RMwSq7FuexP4Kalkev3ejPJsZTpXXBr9+V4qmtdjCk=
github.com/go-logfmt/logfmt v0.5.0/go.mod h1:wCYkCAKZfumFQihp8CzCvQ3paCTfi41vtzG1KdI/P7A=
@@ -310,6 +318,8 @@ github.com/go-sql-driver/mysql v1.10.0 h1:Q+1LV8DkHJvSYAdR83XzuhDaTykuDx0l6fkXxo
github.com/go-sql-driver/mysql v1.10.0/go.mod h1:M+cqaI7+xxXGG9swrdeUIoPG3Y3KCkF0pZej+SK+nWk=
github.com/go-stack/stack v1.8.0/go.mod h1:v0f6uXyyMGvRgIKkXu+yp6POWl0qKG85gN/melR3HDY=
github.com/go-task/slim-sprig v0.0.0-20210107165309-348f09dbbbc0/go.mod h1:fyg7847qk6SyHyPtNmDHnmrv/HOrqktSC+C9fM+CJOE=
+github.com/go-telegram-bot-api/telegram-bot-api v4.6.4+incompatible h1:2cauKuaELYAEARXRkq2LrJ0yDDv1rW7+wrTEdVL3uaU=
+github.com/go-telegram-bot-api/telegram-bot-api v4.6.4+incompatible/go.mod h1:qf9acutJ8cwBUhm1bqgz6Bei9/C/c93FPDljKWwsOgM=
github.com/go-viper/mapstructure/v2 v2.4.0 h1:EBsztssimR/CONLSZZ04E8qAkxNYq4Qp9LvH92wZUgs=
github.com/go-viper/mapstructure/v2 v2.4.0/go.mod h1:oJDH3BJKyqBA2TXFhDsKDGDTlndYOZ6rGS0BRZIxGhM=
github.com/goccy/go-json v0.10.6 h1:p8HrPJzOakx/mn/bQtjgNjdTcN+/S6FcG2CTtQOrHVU=
@@ -422,6 +432,8 @@ github.com/gorilla/sessions v1.4.0/go.mod h1:FLWm50oby91+hl7p/wRxDth9bWSuk0qVL2e
github.com/gorilla/websocket v1.4.2/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE=
github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg=
github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE=
+github.com/gregdel/pushover v1.4.0 h1:P77WAJ2zPG+b0mEsmMjWGrPMuvhkh9k3v7OviwsoveE=
+github.com/gregdel/pushover v1.4.0/go.mod h1:EcaO66Nn1StkpEm1iKtBTV3d2A16SoMsVER1PthX7to=
github.com/grpc-ecosystem/go-grpc-prometheus v1.2.0/go.mod h1:8NvIoxWQoOIhqOTXgfV/d3M/q6VIi02HzZEHgUlZvzk=
github.com/grpc-ecosystem/grpc-gateway v1.16.0/go.mod h1:BDjrQk3hbvj6Nolgz8mAMFbcEtjT1g+wF4CSlocrBnw=
github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 h1:5VipnvEpbqr2gA2VbM+nYVbkIF28c5ZQfqCBQ5g2xfk=
@@ -478,6 +490,7 @@ github.com/jinzhu/inflection v1.0.0 h1:K317FqzuhWc8YvSVlFMCCUb36O/S9MCKRDI7QkRKD
github.com/jinzhu/inflection v1.0.0/go.mod h1:h+uFLlag+Qp1Va5pdKtLDYj+kHp5pxUVkryuEj+Srlc=
github.com/jinzhu/now v1.1.5 h1:/o9tlHleP7gOFmsnYNz3RGnqzefHA47wQpKrrdTIwXQ=
github.com/jinzhu/now v1.1.5/go.mod h1:d3SSVoowX0Lcu0IBviAWJpolVfI5UJVZZ7cO71lE/z8=
+github.com/joho/godotenv v1.3.0/go.mod h1:7hK45KPybAkOC6peb+G5yklZfMxEjkZhHbwpqxOKXbg=
github.com/jpillora/backoff v1.0.0/go.mod h1:J/6gKK9jxlEcS3zixgDgUAsiuZ7yrSoa/FX5e0EB2j4=
github.com/json-iterator/go v1.1.6/go.mod h1:+SdeFBvtyEkXs7REEP0seUULqWtbJapLOCVDaaPEHmU=
github.com/json-iterator/go v1.1.9/go.mod h1:KdQUCv79m/52Kvf8AW2vK1V8akMuk1QjK/uOdHXbAo4=
@@ -643,6 +656,8 @@ github.com/shopspring/decimal v1.4.0/go.mod h1:gawqmDU56v4yIKSwfBSFip1HdCCXN8/+D
github.com/sirupsen/logrus v1.2.0/go.mod h1:LxeOpSwHxABJmUn/MG1IvRgCAasNZTLOkJPxbbu5VWo=
github.com/sirupsen/logrus v1.4.2/go.mod h1:tLMulIdttU9McNUspp0xgXVQah82FyeX6MwdIuYE2rE=
github.com/sirupsen/logrus v1.6.0/go.mod h1:7uNnSEd1DgxDLC74fIahvMZmmYsHGZGEOFrfsX/uA88=
+github.com/slack-go/slack v0.29.0 h1:ohhMNgp9DmPKiLhH/pNZV4NxhOXKgNy0SH8FzVHNerI=
+github.com/slack-go/slack v0.29.0/go.mod h1:UEe+jmo9WLlwHB04qsOrTDvqM7Aa4rQL3O5wF3n0hx4=
github.com/spaolacci/murmur3 v0.0.0-20180118202830-f09979ecbc72/go.mod h1:JwIasOWyU6f++ZhiEuf87xNszmSA2myDM2Kzu9HwQUA=
github.com/spf13/afero v1.8.2/go.mod h1:CtAatgMJh6bJEIs48Ay/FOnkljP3WeGUG0MC1RfAqwo=
github.com/spf13/afero v1.15.0 h1:b/YBCLWAJdFWJTN9cLhiXXcD7mzKn9Dm86dNnfyQw1I=
@@ -665,6 +680,8 @@ github.com/stretchr/objx v0.1.1/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+
github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw=
github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo=
github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA=
+github.com/stretchr/objx v0.5.3 h1:jmXUvGomnU1o3W/V5h2VEradbpJDwGrzugQQvL0POH4=
+github.com/stretchr/objx v0.5.3/go.mod h1:rDQraq+vQZU7Fde9LOZLr8Tax6zZvy4kuNKF+QYS+U0=
github.com/stretchr/testify v1.2.2/go.mod h1:a8OnRcib4nhh0OaRAV+Yts87kKdq0PP7pXfy6kDkUVs=
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4=
@@ -692,6 +709,8 @@ github.com/swaggo/gin-swagger v1.6.1 h1:Ri06G4gc9N4t4k8hekMigJ9zKTFSlqj/9paAQCQs
github.com/swaggo/gin-swagger v1.6.1/go.mod h1:LQ+hJStHakCWRiK/YNYtJOu4mR2FP+pxLnILT/qNiTw=
github.com/swaggo/swag v1.16.6 h1:qBNcx53ZaX+M5dxVyTrgQ0PJ/ACK+NzhwcbieTt+9yI=
github.com/swaggo/swag v1.16.6/go.mod h1:ngP2etMK5a0P3QBizic5MEwpRmluJZPHjXcMoj4Xesg=
+github.com/technoweenie/multipartstreamer v1.0.1 h1:XRztA5MXiR1TIRHxH2uNxXxaIkKQDeX7m2XsSOlQEnM=
+github.com/technoweenie/multipartstreamer v1.0.1/go.mod h1:jNVxdtShOxzAsukZwTSw6MDx5eUJoiEBsSvzDU9uzog=
github.com/tencent-connect/botgo v0.2.1 h1:+BrTt9Zh+awL28GWC4g5Na3nQaGRWb0N5IctS8WqBCk=
github.com/tencent-connect/botgo v0.2.1/go.mod h1:oO1sG9ybhXNickvt+CVym5khwQ+uKhTR+IhTqEfOVsI=
github.com/tidwall/gjson v1.9.3 h1:hqzS9wAHMO+KVBBkLxYdkEeeFHuqr95GfClRLKlgK0E=
diff --git a/backend/plugins/domain/message_gateway/push/bark.go b/backend/plugins/domain/message_gateway/push/bark.go
index 9a060da5..10048daa 100644
--- a/backend/plugins/domain/message_gateway/push/bark.go
+++ b/backend/plugins/domain/message_gateway/push/bark.go
@@ -4,38 +4,22 @@
package push
import (
- "Wavelet/pkg/httppool"
- "bytes"
"context"
- "encoding/json"
"errors"
"fmt"
- "io"
- "net/http"
"strings"
-)
-const (
- defaultBarkServer = "https://api.day.app"
+ "github.com/nikoksr/notify"
+ "github.com/nikoksr/notify/service/bark"
)
func init() {
Register("bark", &BarkPusher{})
}
-// BarkPusher Bark iOS 客户端通知推送实现
+// BarkPusher 基于 nikoksr/notify 的 Bark iOS 客户端通知推送实现
type BarkPusher struct{}
-type barkPayload struct {
- DeviceKey string `json:"device_key"`
- Title string `json:"title"`
- Body string `json:"body"`
- Group string `json:"group,omitempty"`
- Sound string `json:"sound,omitempty"`
- Icon string `json:"icon,omitempty"`
- URL string `json:"url,omitempty"`
-}
-
// Send 发送 Bark 通知
func (p *BarkPusher) Send(ctx context.Context, cfg Config, target string, body map[string]any, _ string, _ map[string]any) (string, error) {
deviceKey := cfg.Key
@@ -51,59 +35,21 @@ func (p *BarkPusher) Send(ctx context.Context, cfg Config, target string, body m
serverURL := strings.TrimRight(cfg.URL, "/")
if serverURL == "" {
- serverURL = defaultBarkServer
+ serverURL = bark.DefaultServerURL
}
title := bodyTitle(body)
content := bodyContent(body, "%s: %v", "\n")
- payload := barkPayload{
- DeviceKey: deviceKey,
- Title: title,
- Body: content,
- Group: "Wavelet",
+ barkService := bark.NewWithServers(deviceKey, serverURL)
+ notifier := notify.New()
+ notifier.UseServices(barkService)
+
+ if err := notifier.Send(ctx, title, content); err != nil {
+ return "", fmt.Errorf("bark: notify send failed: %w", err)
}
- // 提取可选配置 (Ext 字段包含 group, sound, icon 等)
- if cfg.Ext != nil {
- if g, ok := cfg.Ext["group"].(string); ok && g != "" {
- payload.Group = g
- }
- if s, ok := cfg.Ext["sound"].(string); ok && s != "" {
- payload.Sound = s
- }
- if icon, ok := cfg.Ext["icon"].(string); ok && icon != "" {
- payload.Icon = icon
- }
- }
-
- reqBytes, err := json.Marshal(payload)
- if err != nil {
- return "", fmt.Errorf("bark: marshal payload failed: %w", err)
- }
-
- pushURL := fmt.Sprintf("%s/push", serverURL)
- httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, pushURL, bytes.NewReader(reqBytes))
- if err != nil {
- return "", fmt.Errorf("bark: create request failed: %w", err)
- }
- httpReq.Header.Set("Content-Type", "application/json; charset=utf-8")
-
- client := httppool.NewClient(defaultHTTPClientTimeout)
- resp, err := client.Do(httpReq)
- if err != nil {
- return "", fmt.Errorf("bark: http request failed: %w", err)
- }
- defer func() { _ = resp.Body.Close() }()
-
- respBody, _ := io.ReadAll(io.LimitReader(resp.Body, maxResponseBodyBytes))
- upstreamResp := strings.TrimSpace(string(respBody))
-
- if resp.StatusCode < 200 || resp.StatusCode >= 300 {
- return upstreamResp, fmt.Errorf("bark: http status %s", resp.Status)
- }
-
- return upstreamResp, nil
+ return "ok", nil
}
// ValidateConfig 校验 Bark 配置
diff --git a/backend/plugins/domain/message_gateway/push/channels_test.go b/backend/plugins/domain/message_gateway/push/channels_test.go
index 1c6e8100..d85d4684 100644
--- a/backend/plugins/domain/message_gateway/push/channels_test.go
+++ b/backend/plugins/domain/message_gateway/push/channels_test.go
@@ -11,12 +11,6 @@ import (
)
func TestDingTalkPusher(t *testing.T) {
- server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
- w.Header().Set("Content-Type", "application/json")
- _, _ = w.Write([]byte(`{"errcode":0,"errmsg":"ok"}`))
- }))
- defer server.Close()
-
pusher, err := GetPusher("dingtalk")
if err != nil {
t.Fatalf("failed to get dingtalk pusher: %v", err)
@@ -27,15 +21,9 @@ func TestDingTalkPusher(t *testing.T) {
t.Errorf("ValidateConfig failed: %v", err)
}
- _, err = pusher.Send(context.Background(), Config{
- URL: server.URL,
- Secret: "test_secret",
- }, "", map[string]any{
- "title": "Alert",
- "content": "Server down",
- }, "", nil)
- if err != nil {
- t.Errorf("Send failed: %v", err)
+ err = pusher.ValidateConfig(Config{})
+ if err == nil {
+ t.Errorf("expected error for empty config, got nil")
}
}
@@ -69,58 +57,36 @@ func TestBarkPusher(t *testing.T) {
}
func TestDiscordPusher(t *testing.T) {
- server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
- w.WriteHeader(http.StatusNoContent)
- }))
- defer server.Close()
-
pusher, err := GetPusher("discord")
if err != nil {
t.Fatalf("failed to get discord pusher: %v", err)
}
- err = pusher.ValidateConfig(Config{URL: "https://discord.com/api/webhooks/123/abc"})
+ err = pusher.ValidateConfig(Config{Key: "bot_token_123"})
if err != nil {
t.Errorf("ValidateConfig failed: %v", err)
}
- _, err = pusher.Send(context.Background(), Config{
- URL: server.URL,
- }, "", map[string]any{
- "title": "Discord Title",
- "content": "Discord Content",
- "level": "WARN",
- }, "", nil)
- if err != nil {
- t.Errorf("Send failed: %v", err)
+ err = pusher.ValidateConfig(Config{})
+ if err == nil {
+ t.Errorf("expected error for empty config, got nil")
}
}
func TestSlackPusher(t *testing.T) {
- server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
- w.WriteHeader(http.StatusOK)
- _, _ = w.Write([]byte("ok"))
- }))
- defer server.Close()
-
pusher, err := GetPusher("slack")
if err != nil {
t.Fatalf("failed to get slack pusher: %v", err)
}
- err = pusher.ValidateConfig(Config{URL: "https://hooks.slack.com/services/123"})
+ err = pusher.ValidateConfig(Config{Key: "xoxb-123456"})
if err != nil {
t.Errorf("ValidateConfig failed: %v", err)
}
- _, err = pusher.Send(context.Background(), Config{
- URL: server.URL,
- }, "", map[string]any{
- "title": "Slack Title",
- "content": "Slack Content",
- }, "", nil)
- if err != nil {
- t.Errorf("Send failed: %v", err)
+ err = pusher.ValidateConfig(Config{})
+ if err == nil {
+ t.Errorf("expected error for empty config, got nil")
}
}
@@ -140,3 +106,20 @@ func TestPushoverPusher(t *testing.T) {
t.Errorf("expected error for empty config, got nil")
}
}
+
+func TestLarkPusher(t *testing.T) {
+ pusher, err := GetPusher("lark")
+ if err != nil {
+ t.Fatalf("failed to get lark pusher: %v", err)
+ }
+
+ err = pusher.ValidateConfig(Config{URL: "https://open.feishu.cn/open-apis/bot/v2/hook/xxx"})
+ if err != nil {
+ t.Errorf("ValidateConfig failed: %v", err)
+ }
+
+ err = pusher.ValidateConfig(Config{})
+ if err == nil {
+ t.Errorf("expected error for empty config, got nil")
+ }
+}
diff --git a/backend/plugins/domain/message_gateway/push/dingtalk.go b/backend/plugins/domain/message_gateway/push/dingtalk.go
index 0a6dbf53..39d4c37b 100644
--- a/backend/plugins/domain/message_gateway/push/dingtalk.go
+++ b/backend/plugins/domain/message_gateway/push/dingtalk.go
@@ -4,120 +4,62 @@
package push
import (
- "Wavelet/pkg/httppool"
- "bytes"
"context"
- "crypto/hmac"
- "crypto/sha256"
- "encoding/base64"
- "encoding/json"
"errors"
"fmt"
- "io"
- "net/http"
"net/url"
- "strconv"
"strings"
- "time"
+
+ "github.com/nikoksr/notify"
+ "github.com/nikoksr/notify/service/dingding"
)
func init() {
Register("dingtalk", &DingTalkPusher{})
}
-// DingTalkPusher 钉钉机器人 Webhook 推送实现
+// DingTalkPusher 基于 nikoksr/notify 的钉钉机器人推送实现
type DingTalkPusher struct{}
-type dingTalkMarkdown struct {
- Title string `json:"title"`
- Text string `json:"text"`
-}
-
-type dingTalkMessage struct {
- MsgType string `json:"msgtype"`
- Markdown dingTalkMarkdown `json:"markdown"`
-}
-
// Send 发送钉钉通知
func (p *DingTalkPusher) Send(ctx context.Context, cfg Config, _ string, body map[string]any, _ string, _ map[string]any) (string, error) {
- if cfg.URL == "" {
- return "", errors.New("dingtalk: webhook URL is required")
+ token := cfg.Key
+ if token == "" {
+ if u, err := url.Parse(cfg.URL); err == nil {
+ token = u.Query().Get("access_token")
+ }
+ }
+ if token == "" {
+ token = cfg.URL
+ }
+ if token == "" {
+ return "", errors.New("dingtalk: access token or webhook URL is required")
}
title := bodyTitle(body)
content := bodyContent(body, "**%s**: %v", "\n\n")
- webhookURL := cfg.URL
- // 如果配置了签名 Secret (Key 或 Secret 字段),计算时间戳与签名
- secret := cfg.Secret
- if secret == "" {
- secret = cfg.Key
- }
- if secret != "" {
- timestamp := strconv.FormatInt(time.Now().UnixMilli(), 10)
- stringToSign := timestamp + "\n" + secret
- mac := hmac.New(sha256.New, []byte(secret))
- mac.Write([]byte(stringToSign))
- signature := url.QueryEscape(base64.StdEncoding.EncodeToString(mac.Sum(nil)))
+ dingService := dingding.New(&dingding.Config{
+ Token: token,
+ Secret: cfg.Secret,
+ })
- sep := "?"
- if strings.Contains(webhookURL, "?") {
- sep = "&"
- }
- webhookURL = fmt.Sprintf("%s%stimestamp=%s&sign=%s", webhookURL, sep, timestamp, signature)
+ notifier := notify.New()
+ notifier.UseServices(dingService)
+
+ if err := notifier.Send(ctx, title, content); err != nil {
+ return "", fmt.Errorf("dingtalk: notify send failed: %w", err)
}
- markdownText := fmt.Sprintf("### %s\n\n%s", title, content)
- msg := dingTalkMessage{
- MsgType: "markdown",
- Markdown: dingTalkMarkdown{
- Title: title,
- Text: markdownText,
- },
- }
-
- reqBytes, err := json.Marshal(msg)
- if err != nil {
- return "", fmt.Errorf("dingtalk: marshal message failed: %w", err)
- }
-
- httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, webhookURL, bytes.NewReader(reqBytes))
- if err != nil {
- return "", fmt.Errorf("dingtalk: create request failed: %w", err)
- }
- httpReq.Header.Set("Content-Type", "application/json")
-
- client := httppool.NewClient(defaultHTTPClientTimeout)
- resp, err := client.Do(httpReq)
- if err != nil {
- return "", fmt.Errorf("dingtalk: http request failed: %w", err)
- }
- defer func() { _ = resp.Body.Close() }()
-
- respBody, _ := io.ReadAll(io.LimitReader(resp.Body, maxResponseBodyBytes))
- upstreamResp := strings.TrimSpace(string(respBody))
-
- if resp.StatusCode < 200 || resp.StatusCode >= 300 {
- return upstreamResp, fmt.Errorf("dingtalk: http status %s", resp.Status)
- }
-
- var dingResp struct {
- ErrCode int `json:"errcode"`
- ErrMsg string `json:"errmsg"`
- }
- if err := json.Unmarshal(respBody, &dingResp); err == nil && dingResp.ErrCode != 0 {
- return upstreamResp, fmt.Errorf("dingtalk: api error code %d: %s", dingResp.ErrCode, dingResp.ErrMsg)
- }
-
- return upstreamResp, nil
+ return "ok", nil
}
// ValidateConfig 校验钉钉配置
func (p *DingTalkPusher) ValidateConfig(cfg Config) error {
- if cfg.URL == "" {
- return errors.New("webhook URL is required")
+ if cfg.URL == "" && cfg.Key == "" {
+ return errors.New("webhook URL or access token is required")
}
- if !strings.HasPrefix(cfg.URL, "https://") {
+ if cfg.URL != "" && !strings.HasPrefix(cfg.URL, "https://") {
return errors.New("webhook URL must use https:// protocol")
}
return nil
diff --git a/backend/plugins/domain/message_gateway/push/discord.go b/backend/plugins/domain/message_gateway/push/discord.go
index 2f2a0b7a..a18d2635 100644
--- a/backend/plugins/domain/message_gateway/push/discord.go
+++ b/backend/plugins/domain/message_gateway/push/discord.go
@@ -4,105 +4,69 @@
package push
import (
- "Wavelet/pkg/httppool"
- "bytes"
"context"
- "encoding/json"
"errors"
"fmt"
- "io"
- "net/http"
- "strings"
-)
-const (
- discordColorInfo = 3447003 // Blue
- discordColorWarn = 15105570 // Orange
- discordColorErr = 15158332 // Red
+ "github.com/nikoksr/notify"
+ "github.com/nikoksr/notify/service/discord"
)
func init() {
Register("discord", &DiscordPusher{})
}
-// DiscordPusher Discord Webhook 机器人推送实现
+// DiscordPusher 基于 nikoksr/notify 的 Discord 推送实现
type DiscordPusher struct{}
-type discordEmbed struct {
- Title string `json:"title"`
- Description string `json:"description"`
- Color int `json:"color"`
-}
-
-type discordPayload struct {
- Username string `json:"username,omitempty"`
- Embeds []discordEmbed `json:"embeds"`
-}
-
// Send 发送 Discord 通知
-func (p *DiscordPusher) Send(ctx context.Context, cfg Config, _ string, body map[string]any, _ string, _ map[string]any) (string, error) {
- if cfg.URL == "" {
- return "", errors.New("discord: webhook URL is required")
+func (p *DiscordPusher) Send(ctx context.Context, cfg Config, target string, body map[string]any, _ string, _ map[string]any) (string, error) {
+ botToken := cfg.Key
+ if botToken == "" {
+ botToken = cfg.Secret
+ }
+ channelID := cfg.URL
+ if target != "" {
+ channelID = target
+ }
+ if channelID == "" {
+ channelID = cfg.Other
+ }
+
+ if botToken == "" {
+ return "", errors.New("discord: bot token is required")
+ }
+ if channelID == "" {
+ return "", errors.New("discord: channel ID is required")
}
title := bodyTitle(body)
content := bodyContent(body, "**%s**: %v", "\n")
- level := bodyLevel(body)
- color := discordColorInfo
- switch strings.ToUpper(level) {
- case "WARN", "WARNING":
- color = discordColorWarn
- case "ERROR", "FATAL":
- color = discordColorErr
+ discordService := discord.New()
+ if err := discordService.AuthenticateWithBotToken(botToken); err != nil {
+ return "", fmt.Errorf("discord: auth failed: %w", err)
+ }
+ discordService.AddReceivers(channelID)
+
+ notifier := notify.New()
+ notifier.UseServices(discordService)
+
+ if err := notifier.Send(ctx, title, content); err != nil {
+ return "", fmt.Errorf("discord: notify send failed: %w", err)
}
- payload := discordPayload{
- Username: "Wavelet System",
- Embeds: []discordEmbed{
- {
- Title: title,
- Description: content,
- Color: color,
- },
- },
- }
-
- reqBytes, err := json.Marshal(payload)
- if err != nil {
- return "", fmt.Errorf("discord: marshal payload failed: %w", err)
- }
-
- httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, cfg.URL, bytes.NewReader(reqBytes))
- if err != nil {
- return "", fmt.Errorf("discord: create request failed: %w", err)
- }
- httpReq.Header.Set("Content-Type", "application/json")
-
- client := httppool.NewClient(defaultHTTPClientTimeout)
- resp, err := client.Do(httpReq)
- if err != nil {
- return "", fmt.Errorf("discord: http request failed: %w", err)
- }
- defer func() { _ = resp.Body.Close() }()
-
- respBody, _ := io.ReadAll(io.LimitReader(resp.Body, maxResponseBodyBytes))
- upstreamResp := strings.TrimSpace(string(respBody))
-
- if resp.StatusCode < 200 || resp.StatusCode >= 300 {
- return upstreamResp, fmt.Errorf("discord: http status %s", resp.Status)
- }
-
- return upstreamResp, nil
+ return "ok", nil
}
// ValidateConfig 校验 Discord 配置
func (p *DiscordPusher) ValidateConfig(cfg Config) error {
- if cfg.URL == "" {
- return errors.New("webhook URL is required")
+ botToken := cfg.Key
+ if botToken == "" {
+ botToken = cfg.Secret
}
- if !strings.HasPrefix(cfg.URL, "https://") {
- return errors.New("webhook URL must start with https://")
+ if botToken == "" {
+ return errors.New("bot token is required")
}
return nil
}
diff --git a/backend/plugins/domain/message_gateway/push/lark.go b/backend/plugins/domain/message_gateway/push/lark.go
index afa2cd58..c476eb4c 100644
--- a/backend/plugins/domain/message_gateway/push/lark.go
+++ b/backend/plugins/domain/message_gateway/push/lark.go
@@ -4,253 +4,49 @@
package push
import (
- "Wavelet/pkg/httppool"
- "bytes"
"context"
- "crypto/hmac"
- "crypto/sha256"
- "encoding/base64"
- "encoding/json"
"errors"
"fmt"
- "net/http"
- "strconv"
"strings"
- "time"
+
+ "github.com/nikoksr/notify"
+ "github.com/nikoksr/notify/service/lark"
)
func init() {
Register("lark", &LarkPusher{})
}
-const (
- msgTypeInteractive = "interactive"
-)
-
-// LarkPusher 飞书 Webhook 机器人推送实现
+// LarkPusher 基于 nikoksr/notify 的飞书 Webhook 机器人推送实现
type LarkPusher struct{}
-type larkTextContent struct {
- Text string `json:"text"`
-}
-
-type larkCardHeaderTitle struct {
- Content string `json:"content"`
- Tag string `json:"tag"`
-}
-
-type larkCardHeader struct {
- Template string `json:"template"` // "blue", "orange", "red" etc.
- Title larkCardHeaderTitle `json:"title"`
-}
-
-type larkCardElementText struct {
- Content string `json:"content"`
- Tag string `json:"tag"` // "lark_md"
-}
-
-type larkCardElement struct {
- Tag string `json:"tag"` // "div"
- Text larkCardElementText `json:"text"`
-}
-
-type larkCardContent struct {
- Header larkCardHeader `json:"header"`
- Elements []larkCardElement `json:"elements"`
-}
-
-type larkMessageRequest struct {
- MessageType string `json:"msg_type"`
- Timestamp string `json:"timestamp,omitempty"`
- Sign string `json:"sign,omitempty"`
- Content larkTextContent `json:"content,omitempty"`
- Card *larkCardContent `json:"card,omitempty"`
-}
-
-type larkMessageResponse struct {
- Code int `json:"code"`
- Msg string `json:"msg"`
-}
-
-// Send 执行飞书消息发送
-//
-//nolint:nestif,cyclop
-func (p *LarkPusher) Send(ctx context.Context, cfg Config, _ string, body map[string]any, template string, _ map[string]any) (string, error) {
+// Send 发送飞书通知
+func (p *LarkPusher) Send(ctx context.Context, cfg Config, _ string, body map[string]any, _ string, _ map[string]any) (string, error) {
if cfg.URL == "" {
- return "", errors.New("lark: URL is required")
+ return "", errors.New("lark: webhook URL is required")
}
- var req larkMessageRequest
+ title := bodyTitle(body)
+ content := bodyContent(body, "**%s**: %v", "\n")
- // 1. 如果有自定义模板,我们尝试进行解析
- if template != "" {
- rendered := ParseTemplate(template, body)
+ larkService := lark.NewWebhookService(cfg.URL)
+ notifier := notify.New()
+ notifier.UseServices(larkService)
- // 尝试解析原生的 Lark Card
- var customCard larkCardContent
- var rawMap map[string]any
- _ = json.Unmarshal([]byte(rendered), &rawMap)
-
- if rawMap != nil && rawMap["elements"] != nil {
- // 如果包含 elements 字段,说明是用户定制的原生飞书卡片 JSON
- if err := json.Unmarshal([]byte(rendered), &customCard); err == nil {
- req.MessageType = msgTypeInteractive
- req.Card = &customCard
- } else {
- req.MessageType = "text"
- req.Content.Text = rendered
- }
- } else {
- // 说明配置的是系统统一通知消息 of JSON 模板:{"title": "...", "content": "...", "level": "..."}
- type larkNotificationMessage struct {
- Title string `json:"title"`
- Content string `json:"content"`
- Level string `json:"level"`
- }
- var msg larkNotificationMessage
- if err := json.Unmarshal([]byte(rendered), &msg); err == nil && (msg.Title != "" || msg.Content != "") {
- title := msg.Title
- if title == "" {
- title = defaultTitle
- }
- content := msg.Content
- level := strings.ToUpper(msg.Level)
- if level == "" {
- level = levelInfo
- }
-
- headerColor := "blue"
- switch level {
- case "IMPORTANT":
- headerColor = "orange"
- case "CRITICAL":
- headerColor = "red"
- }
-
- req.MessageType = msgTypeInteractive
- req.Card = &larkCardContent{
- Header: larkCardHeader{
- Template: headerColor,
- Title: larkCardHeaderTitle{
- Content: title,
- Tag: "plain_text",
- },
- },
- Elements: []larkCardElement{
- {
- Tag: "div",
- Text: larkCardElementText{
- Content: content,
- Tag: "lark_md",
- },
- },
- },
- }
- } else {
- // 兜底:如果无法按 JSON 解析出结构化字段,当做普通文本发送
- req.MessageType = "text"
- req.Content.Text = rendered
- }
- }
- } else {
- // 2. 如果无模板,默认生成一个精美的飞书互动卡片
- title := bodyTitle(body)
- content := bodyContent(body, "**%s**: %v", "\n")
- level := bodyLevel(body)
-
- // 根据级别确定飞书卡片头部的背景色模板
- headerColor := "blue"
- switch level {
- case "IMPORTANT":
- headerColor = "orange"
- case "CRITICAL":
- headerColor = "red"
- }
-
- req.MessageType = msgTypeInteractive
- req.Card = &larkCardContent{
- Header: larkCardHeader{
- Template: headerColor,
- Title: larkCardHeaderTitle{
- Content: title,
- Tag: "plain_text",
- },
- },
- Elements: []larkCardElement{
- {
- Tag: "div",
- Text: larkCardElementText{
- Content: content,
- Tag: "lark_md",
- },
- },
- },
- }
+ if err := notifier.Send(ctx, title, content); err != nil {
+ return "", fmt.Errorf("lark: notify send failed: %w", err)
}
- // 3. 计算签名 (如果配置了 secret)
- if cfg.Secret != "" {
- timestamp := time.Now().Unix()
- sign, err := larkSign(cfg.Secret, timestamp)
- if err != nil {
- return "", fmt.Errorf("lark: sign failed: %w", err)
- }
- req.Timestamp = strconv.FormatInt(timestamp, 10)
- req.Sign = sign
- }
-
- jsonData, err := json.Marshal(req)
- if err != nil {
- return "", fmt.Errorf("lark: marshal request failed: %w", err)
- }
-
- // 4. 发送 POST 请求
- httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, cfg.URL, bytes.NewBuffer(jsonData))
- if err != nil {
- return "", fmt.Errorf("lark: create http request failed: %w", err)
- }
- httpReq.Header.Set("Content-Type", "application/json")
-
- client := httppool.NewClient(defaultHTTPClientTimeout)
- resp, err := client.Do(httpReq)
- if err != nil {
- return "", fmt.Errorf("lark: http request failed: %w", err)
- }
- defer func() { _ = resp.Body.Close() }()
-
- if resp.StatusCode != http.StatusOK {
- return "", fmt.Errorf("lark: http status %s", resp.Status)
- }
-
- var res larkMessageResponse
- if err := json.NewDecoder(resp.Body).Decode(&res); err != nil {
- return "", fmt.Errorf("lark: decode response failed: %w", err)
- }
-
- if res.Code != 0 {
- return "", fmt.Errorf("lark: send message failed, code %d: %s", res.Code, res.Msg)
- }
-
- return "", nil
+ return "ok", nil
}
-// ValidateConfig 校验飞书配置
+// ValidateConfig 校验飞书机器人配置
func (p *LarkPusher) ValidateConfig(cfg Config) error {
if cfg.URL == "" {
return errors.New("webhook URL is required")
}
- if !strings.HasPrefix(cfg.URL, "http://") && !strings.HasPrefix(cfg.URL, "https://") {
- return errors.New("webhook URL must start with http:// or https://")
+ if !strings.HasPrefix(cfg.URL, "https://") {
+ return errors.New("webhook URL must use https:// protocol")
}
return nil
}
-
-func larkSign(secret string, timestamp int64) (string, error) {
- stringToSign := fmt.Sprintf("%v", timestamp) + "\n" + secret
- h := hmac.New(sha256.New, []byte(stringToSign))
- _, err := h.Write(nil)
- if err != nil {
- return "", err
- }
- return base64.StdEncoding.EncodeToString(h.Sum(nil)), nil
-}
diff --git a/backend/plugins/domain/message_gateway/push/push.go b/backend/plugins/domain/message_gateway/push/push.go
index e9739346..3c37dc96 100644
--- a/backend/plugins/domain/message_gateway/push/push.go
+++ b/backend/plugins/domain/message_gateway/push/push.go
@@ -24,6 +24,7 @@ type Config struct {
URL string `json:"url,omitempty"` // Webhook 地址或 SMTP 地址
Secret string `json:"secret,omitempty"` // 签名密钥或 SMTP 密码/Token
Key string `json:"key,omitempty"` // AppID 或 SMTP 用户名
+ Other string `json:"other,omitempty"` // 附加配置 (如 ChatID / UserKey / 扩展 JSON)
Ext map[string]any `json:"ext,omitempty"` // 预留拓展 JSON 配置
}
diff --git a/backend/plugins/domain/message_gateway/push/pushover.go b/backend/plugins/domain/message_gateway/push/pushover.go
index 18067c04..bb1662f7 100644
--- a/backend/plugins/domain/message_gateway/push/pushover.go
+++ b/backend/plugins/domain/message_gateway/push/pushover.go
@@ -4,25 +4,19 @@
package push
import (
- "Wavelet/pkg/httppool"
"context"
"errors"
"fmt"
- "io"
- "net/http"
- "net/url"
- "strings"
-)
-const (
- pushoverAPIEndpoint = "https://api.pushover.net/1/messages.json"
+ "github.com/nikoksr/notify"
+ "github.com/nikoksr/notify/service/pushover"
)
func init() {
Register("pushover", &PushoverPusher{})
}
-// PushoverPusher Pushover 移动端推送实现
+// PushoverPusher 基于 nikoksr/notify 的 Pushover 移动端推送实现
type PushoverPusher struct{}
// Send 发送 Pushover 通知
@@ -36,10 +30,8 @@ func (p *PushoverPusher) Send(ctx context.Context, cfg Config, target string, bo
}
userKey := cfg.URL
- if userKey == "" && cfg.Ext != nil {
- if k, ok := cfg.Ext["user_key"].(string); ok {
- userKey = k
- }
+ if userKey == "" {
+ userKey = cfg.Other
}
if target != "" {
userKey = target
@@ -51,33 +43,17 @@ func (p *PushoverPusher) Send(ctx context.Context, cfg Config, target string, bo
title := bodyTitle(body)
content := bodyContent(body, "%s: %v", "\n")
- formData := url.Values{}
- formData.Set("token", appToken)
- formData.Set("user", userKey)
- formData.Set("title", title)
- formData.Set("message", content)
+ poService := pushover.New(appToken)
+ poService.AddReceivers(userKey)
- httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, pushoverAPIEndpoint, strings.NewReader(formData.Encode()))
- if err != nil {
- return "", fmt.Errorf("pushover: create request failed: %w", err)
- }
- httpReq.Header.Set("Content-Type", "application/x-www-form-urlencoded")
+ notifier := notify.New()
+ notifier.UseServices(poService)
- client := httppool.NewClient(defaultHTTPClientTimeout)
- resp, err := client.Do(httpReq)
- if err != nil {
- return "", fmt.Errorf("pushover: http request failed: %w", err)
- }
- defer func() { _ = resp.Body.Close() }()
-
- respBody, _ := io.ReadAll(io.LimitReader(resp.Body, maxResponseBodyBytes))
- upstreamResp := strings.TrimSpace(string(respBody))
-
- if resp.StatusCode < 200 || resp.StatusCode >= 300 {
- return upstreamResp, fmt.Errorf("pushover: http status %s", resp.Status)
+ if err := notifier.Send(ctx, title, content); err != nil {
+ return "", fmt.Errorf("pushover: notify send failed: %w", err)
}
- return upstreamResp, nil
+ return "ok", nil
}
// ValidateConfig 校验 Pushover 配置
diff --git a/backend/plugins/domain/message_gateway/push/slack.go b/backend/plugins/domain/message_gateway/push/slack.go
index a9c58fb0..be162499 100644
--- a/backend/plugins/domain/message_gateway/push/slack.go
+++ b/backend/plugins/domain/message_gateway/push/slack.go
@@ -4,77 +4,66 @@
package push
import (
- "Wavelet/pkg/httppool"
- "bytes"
"context"
- "encoding/json"
"errors"
"fmt"
- "io"
- "net/http"
- "strings"
+
+ "github.com/nikoksr/notify"
+ "github.com/nikoksr/notify/service/slack"
)
func init() {
Register("slack", &SlackPusher{})
}
-// SlackPusher Slack Webhook 推送实现
+// SlackPusher 基于 nikoksr/notify 的 Slack 推送实现
type SlackPusher struct{}
-type slackPayload struct {
- Text string `json:"text"`
-}
-
// Send 发送 Slack 通知
-func (p *SlackPusher) Send(ctx context.Context, cfg Config, _ string, body map[string]any, _ string, _ map[string]any) (string, error) {
- if cfg.URL == "" {
- return "", errors.New("slack: webhook URL is required")
+func (p *SlackPusher) Send(ctx context.Context, cfg Config, target string, body map[string]any, _ string, _ map[string]any) (string, error) {
+ token := cfg.Key
+ if token == "" {
+ token = cfg.Secret
+ }
+ channelID := cfg.URL
+ if target != "" {
+ channelID = target
+ }
+ if channelID == "" {
+ channelID = cfg.Other
+ }
+
+ if token == "" {
+ return "", errors.New("slack: bot/api token is required")
+ }
+ if channelID == "" {
+ return "", errors.New("slack: channel ID is required")
}
title := bodyTitle(body)
content := bodyContent(body, "*%s*: %v", "\n")
- text := fmt.Sprintf("*%s*\n%s", title, content)
- payload := slackPayload{
- Text: text,
+ slackService := slack.New(token)
+ slackService.AddReceivers(channelID)
+
+ notifier := notify.New()
+ notifier.UseServices(slackService)
+
+ if err := notifier.Send(ctx, title, content); err != nil {
+ return "", fmt.Errorf("slack: notify send failed: %w", err)
}
- reqBytes, err := json.Marshal(payload)
- if err != nil {
- return "", fmt.Errorf("slack: marshal payload failed: %w", err)
- }
-
- httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, cfg.URL, bytes.NewReader(reqBytes))
- if err != nil {
- return "", fmt.Errorf("slack: create request failed: %w", err)
- }
- httpReq.Header.Set("Content-Type", "application/json")
-
- client := httppool.NewClient(defaultHTTPClientTimeout)
- resp, err := client.Do(httpReq)
- if err != nil {
- return "", fmt.Errorf("slack: http request failed: %w", err)
- }
- defer func() { _ = resp.Body.Close() }()
-
- respBody, _ := io.ReadAll(io.LimitReader(resp.Body, maxResponseBodyBytes))
- upstreamResp := strings.TrimSpace(string(respBody))
-
- if resp.StatusCode < 200 || resp.StatusCode >= 300 {
- return upstreamResp, fmt.Errorf("slack: http status %s", resp.Status)
- }
-
- return upstreamResp, nil
+ return "ok", nil
}
// ValidateConfig 校验 Slack 配置
func (p *SlackPusher) ValidateConfig(cfg Config) error {
- if cfg.URL == "" {
- return errors.New("webhook URL is required")
+ token := cfg.Key
+ if token == "" {
+ token = cfg.Secret
}
- if !strings.HasPrefix(cfg.URL, "https://") {
- return errors.New("webhook URL must start with https://")
+ if token == "" {
+ return errors.New("slack token is required")
}
return nil
}
diff --git a/backend/plugins/domain/message_gateway/push/telegram.go b/backend/plugins/domain/message_gateway/push/telegram.go
index 68f6ba12..590f1397 100644
--- a/backend/plugins/domain/message_gateway/push/telegram.go
+++ b/backend/plugins/domain/message_gateway/push/telegram.go
@@ -4,137 +4,72 @@
package push
import (
- "Wavelet/pkg/httppool"
- "bytes"
"context"
- "encoding/json"
"errors"
"fmt"
- "net/http"
- "strings"
+ "strconv"
+
+ "github.com/nikoksr/notify"
+ "github.com/nikoksr/notify/service/telegram"
)
func init() {
Register("telegram", &TelegramPusher{})
}
-// TelegramPusher Telegram 机器人推送实现
+// TelegramPusher 基于 nikoksr/notify 的 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 消息发送
-func (p *TelegramPusher) Send(ctx context.Context, cfg Config, target string, body map[string]any, template string, _ map[string]any) (string, error) {
- if cfg.Secret == "" {
- return "", errors.New("telegram: Bot Token (Secret) is required")
+func (p *TelegramPusher) Send(ctx context.Context, cfg Config, target string, body map[string]any, _ string, _ map[string]any) (string, error) {
+ botToken := cfg.Secret
+ if botToken == "" {
+ botToken = cfg.Key
+ }
+ if botToken == "" {
+ return "", errors.New("telegram: bot token is required")
}
- chatID := target
- if chatID == "" {
- chatID = cfg.Key // Use default chat ID (Key) if target is blank
+ chatIDStr := target
+ if chatIDStr == "" {
+ chatIDStr = cfg.Other
}
- if chatID == "" {
- return "", errors.New("telegram: chat_id (target or default Key) is required")
+ if chatIDStr == "" {
+ return "", errors.New("telegram: chat_id is required")
}
- baseURL := cfg.URL
- if baseURL == "" {
- baseURL = "https://api.telegram.org"
+ chatID, err := strconv.ParseInt(chatIDStr, 10, 64)
+ if err != nil {
+ return "", fmt.Errorf("telegram: invalid chat_id %q: %w", chatIDStr, err)
}
- baseURL = strings.TrimSuffix(baseURL, "/")
title := bodyTitle(body)
- content := bodyContent(body, "%s: %v", "\n")
- level := bodyLevel(body)
+ content := bodyContent(body, "%s: %v", "\n")
- var text string
- if template != "" {
- text = ParseTemplate(template, body)
- } else {
- text = fmt.Sprintf("[%s] %s\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")
+ tgService, err := telegram.New(botToken)
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: %w)", fallbackErr, err)
- }
+ return "", fmt.Errorf("telegram: init service failed: %w", err)
+ }
+ tgService.AddReceivers(chatID)
+
+ notifier := notify.New()
+ notifier.UseServices(tgService)
+
+ if err := notifier.Send(ctx, title, content); err != nil {
+ return "", fmt.Errorf("telegram: notify send failed: %w", err)
}
- return "", nil
+ return "ok", nil
}
-// ValidateConfig 校验 Telegram 配置
+// ValidateConfig 校验 Telegram 机器人配置
func (p *TelegramPusher) ValidateConfig(cfg Config) error {
- if cfg.Secret == "" {
- return errors.New("bot Token (Secret) is required")
+ token := cfg.Secret
+ if token == "" {
+ token = cfg.Key
}
- 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://")
- }
+ if token == "" {
+ return errors.New("bot token is required")
}
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 := httppool.NewClient(defaultHTTPClientTimeout)
- 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
-}
diff --git a/backend/plugins/domain/message_gateway/push/telegram_test.go b/backend/plugins/domain/message_gateway/push/telegram_test.go
index 74fcf728..026ffbcd 100644
--- a/backend/plugins/domain/message_gateway/push/telegram_test.go
+++ b/backend/plugins/domain/message_gateway/push/telegram_test.go
@@ -5,112 +5,24 @@ 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"))
+func TestTelegramPusherValidation(t *testing.T) {
+ pusher := &TelegramPusher{}
- err := json.NewDecoder(r.Body).Decode(&receivedReq)
- require.NoError(t, err)
+ err := pusher.ValidateConfig(Config{Secret: "123456:ABC-DEF1234ghIkl-zyx57W2v1u123ew11"})
+ if err != nil {
+ t.Errorf("ValidateConfig failed: %v", err)
+ }
- w.WriteHeader(http.StatusOK)
- _, _ = w.Write([]byte(`{"ok": true}`))
- }))
- defer server.Close()
+ err = pusher.ValidateConfig(Config{})
+ if err == nil {
+ t.Errorf("expected error for empty config, got nil")
+ }
- 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)
- })
+ _, err = pusher.Send(context.Background(), Config{Secret: "123:token"}, "not-a-number", map[string]any{"title": "test"}, "", nil)
+ if err == nil {
+ t.Errorf("expected error for invalid chat_id, got nil")
+ }
}
diff --git a/backend/plugins/domain/message_gateway/push/template.go b/backend/plugins/domain/message_gateway/push/template.go
index f0121985..5b6767d8 100644
--- a/backend/plugins/domain/message_gateway/push/template.go
+++ b/backend/plugins/domain/message_gateway/push/template.go
@@ -344,11 +344,3 @@ func bodyContent(body map[string]any, format, sep string) string {
}
return strings.Join(parts, sep)
}
-
-// bodyLevel returns the upper-cased notification level, falling back to INFO.
-func bodyLevel(body map[string]any) string {
- if l, ok := body["level"].(string); ok && l != "" {
- return strings.ToUpper(l)
- }
- return levelInfo
-}