mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-30 22:26:38 +08:00
refactor(push): clean up legacy push_config and w_push_events hardcoded insertions
- Remove ConfigKeyPushConfig and delete legacy push_config query logic from EventTrigger.Trigger. - Simplify the push event notification dispatching engine to rely purely on database custom channels. - Remove w_push_events hardcoded INSERT statement from migration files, letting SyncEvents handle default event registration. - Rewrite unit tests to use model.PushChannel instead of push_config.
This commit is contained in:
@@ -124,40 +124,16 @@ func (t *EventTrigger) Trigger(ctx context.Context, meta EventMetadata, body map
|
||||
return
|
||||
}
|
||||
|
||||
// 2. Read push configs
|
||||
configs, err := t.getPushConfigs(asyncCtx)
|
||||
if err != nil {
|
||||
logger.ErrorF(asyncCtx, "push_event_trigger: getPushConfigs failed: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
// 3. Build and render notification message
|
||||
// 2. Build and render notification message
|
||||
flatBody := getFlatBody(body)
|
||||
msg, renderedTemplate := t.buildMessage(&event, meta, flatBody, body)
|
||||
msg, _ := t.buildMessage(&event, meta, flatBody, body)
|
||||
|
||||
// 4. Enqueue tasks for each matching channel
|
||||
t.enqueuePushTasks(asyncCtx, meta, &event, configs, msg, renderedTemplate, flatBody)
|
||||
// 3. Enqueue tasks for each matching channel
|
||||
t.enqueuePushTasks(asyncCtx, meta, &event, msg, flatBody)
|
||||
}()
|
||||
}
|
||||
|
||||
func (t *EventTrigger) getPushConfigs(ctx context.Context) ([]pkgpush.Config, error) {
|
||||
var configVal string
|
||||
var sc model.SystemConfig
|
||||
if err := sc.GetByKey(ctx, model.ConfigKeyPushConfig); err == nil {
|
||||
configVal = sc.Value
|
||||
}
|
||||
|
||||
if configVal == "" || configVal == "[]" {
|
||||
return nil, errors.New("push_config is empty or not configured")
|
||||
}
|
||||
|
||||
var configs []pkgpush.Config
|
||||
if err := json.Unmarshal([]byte(configVal), &configs); err != nil {
|
||||
return nil, fmt.Errorf("unmarshal push_config failed: %w", err)
|
||||
}
|
||||
|
||||
return configs, nil
|
||||
}
|
||||
|
||||
func (t *EventTrigger) buildMessage(event *model.PushEvent, meta EventMetadata, flatBody map[string]any, body map[string]any) (NotificationMessage, string) {
|
||||
var msg NotificationMessage
|
||||
@@ -244,13 +220,8 @@ func (t *EventTrigger) parseDefaultTemplate(meta EventMetadata, flatBody map[str
|
||||
return msg
|
||||
}
|
||||
|
||||
func (t *EventTrigger) enqueuePushTasks(ctx context.Context, meta EventMetadata, event *model.PushEvent, configs []pkgpush.Config, msg NotificationMessage, renderedTemplate string, flatBody map[string]any) {
|
||||
func (t *EventTrigger) enqueuePushTasks(ctx context.Context, meta EventMetadata, event *model.PushEvent, msg NotificationMessage, flatBody map[string]any) {
|
||||
for _, channelName := range event.Channels {
|
||||
if channelName == channelEmail {
|
||||
t.enqueueEmailPushTasks(ctx, meta, event, configs, msg, renderedTemplate, flatBody)
|
||||
continue
|
||||
}
|
||||
|
||||
// 检查是不是自定义数据库渠道
|
||||
var customChannel model.PushChannel
|
||||
err := db.DB(ctx).Where("name = ? AND enabled = ?", channelName, true).First(&customChannel).Error
|
||||
@@ -259,74 +230,7 @@ func (t *EventTrigger) enqueuePushTasks(ctx context.Context, meta EventMetadata,
|
||||
continue
|
||||
}
|
||||
|
||||
// 数据库中不存在。我们核对它是不是通过代码内置注册的 Pusher 渠道
|
||||
if _, errPusher := pkgpush.GetPusher(channelName); errPusher == nil {
|
||||
t.enqueueBuiltinPushTasks(ctx, meta, event, channelName, configs, msg, renderedTemplate, flatBody)
|
||||
continue
|
||||
}
|
||||
|
||||
logger.WarnF(ctx, "push_event_trigger: channel %q not found in DB and not registered as built-in: %v", channelName, err)
|
||||
}
|
||||
}
|
||||
|
||||
func (t *EventTrigger) enqueueEmailPushTasks(ctx context.Context, meta EventMetadata, event *model.PushEvent, configs []pkgpush.Config, msg NotificationMessage, renderedTemplate string, flatBody map[string]any) {
|
||||
var matchedConfigs []pkgpush.Config
|
||||
for _, cfg := range configs {
|
||||
if cfg.Channel == channelEmail {
|
||||
matchedConfigs = append(matchedConfigs, cfg)
|
||||
}
|
||||
}
|
||||
|
||||
if len(matchedConfigs) == 0 {
|
||||
logger.WarnF(ctx, "push_event_trigger: no active settings for channel %q", channelEmail)
|
||||
return
|
||||
}
|
||||
|
||||
for _, cfg := range matchedConfigs {
|
||||
if cfg.URL == "" || cfg.Key == "" {
|
||||
var smtpHost, smtpPort, smtpUser, smtpPass model.SystemConfig
|
||||
_ = smtpHost.GetByKey(ctx, model.ConfigKeySMTPHost)
|
||||
_ = smtpPort.GetByKey(ctx, model.ConfigKeySMTPPort)
|
||||
_ = smtpUser.GetByKey(ctx, model.ConfigKeySMTPUsername)
|
||||
_ = smtpPass.GetByKey(ctx, model.ConfigKeySMTPPassword)
|
||||
|
||||
if smtpHost.Value != "" && smtpUser.Value != "" {
|
||||
port := smtpPort.Value
|
||||
if port == "" {
|
||||
port = "587"
|
||||
}
|
||||
cfg.URL = smtpHost.Value + ":" + port
|
||||
cfg.Key = smtpUser.Value
|
||||
cfg.Secret = smtpPass.Value
|
||||
}
|
||||
}
|
||||
|
||||
if len(event.Targets) > 0 {
|
||||
for _, target := range event.Targets {
|
||||
resolvedTarget := resolveTarget(ctx, target, flatBody, channelEmail)
|
||||
payload := SendPayload{
|
||||
EventKey: meta.Key,
|
||||
Config: cfg,
|
||||
Target: resolvedTarget,
|
||||
Body: msg,
|
||||
Template: renderedTemplate,
|
||||
}
|
||||
if err := enqueuePushTask(ctx, payload); err != nil {
|
||||
logger.ErrorF(ctx, "push_event_trigger: enqueuePushTask failed for %s -> %s: %v", channelEmail, resolvedTarget, err)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
payload := SendPayload{
|
||||
EventKey: meta.Key,
|
||||
Config: cfg,
|
||||
Target: "",
|
||||
Body: msg,
|
||||
Template: renderedTemplate,
|
||||
}
|
||||
if err := enqueuePushTask(ctx, payload); err != nil {
|
||||
logger.ErrorF(ctx, "push_event_trigger: enqueuePushTask failed for %s: %v", channelEmail, err)
|
||||
}
|
||||
}
|
||||
logger.WarnF(ctx, "push_event_trigger: channel %q not found in DB or disabled: %v", channelName, err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -398,48 +302,7 @@ func (t *EventTrigger) enqueueSingleCustomPushChannelTask(ctx context.Context, m
|
||||
}
|
||||
}
|
||||
|
||||
func (t *EventTrigger) enqueueBuiltinPushTasks(ctx context.Context, meta EventMetadata, event *model.PushEvent, channelName string, configs []pkgpush.Config, msg NotificationMessage, renderedTemplate string, flatBody map[string]any) {
|
||||
var matchedConfigs []pkgpush.Config
|
||||
for _, cfg := range configs {
|
||||
if cfg.Channel == channelName {
|
||||
matchedConfigs = append(matchedConfigs, cfg)
|
||||
}
|
||||
}
|
||||
|
||||
if len(matchedConfigs) == 0 {
|
||||
logger.WarnF(ctx, "push_event_trigger: no active settings for built-in channel %q", channelName)
|
||||
return
|
||||
}
|
||||
|
||||
for _, cfg := range matchedConfigs {
|
||||
if len(event.Targets) > 0 {
|
||||
for _, target := range event.Targets {
|
||||
resolvedTarget := resolveTarget(ctx, target, flatBody, channelName)
|
||||
payload := SendPayload{
|
||||
EventKey: meta.Key,
|
||||
Config: cfg,
|
||||
Target: resolvedTarget,
|
||||
Body: msg,
|
||||
Template: renderedTemplate,
|
||||
}
|
||||
if err := enqueuePushTask(ctx, payload); err != nil {
|
||||
logger.ErrorF(ctx, "push_event_trigger: enqueuePushTask failed for builtin %s -> %s: %v", channelName, resolvedTarget, err)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
payload := SendPayload{
|
||||
EventKey: meta.Key,
|
||||
Config: cfg,
|
||||
Target: "",
|
||||
Body: msg,
|
||||
Template: renderedTemplate,
|
||||
}
|
||||
if err := enqueuePushTask(ctx, payload); err != nil {
|
||||
logger.ErrorF(ctx, "push_event_trigger: enqueuePushTask failed for builtin %s: %v", channelName, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func enqueuePushTask(ctx context.Context, payload SendPayload) error {
|
||||
payloadBytes, err := json.Marshal(payload)
|
||||
|
||||
@@ -148,10 +148,6 @@ func TestEventTrigger(t *testing.T) {
|
||||
dbConn, _, cleanup := setupPushTest(t)
|
||||
defer cleanup()
|
||||
|
||||
// Register mock pusher
|
||||
mPusher := &mockPusher{}
|
||||
pkgpush.Register("mock_channel", mPusher)
|
||||
|
||||
// SyncEvents
|
||||
err := SyncEvents(context.Background())
|
||||
require.NoError(t, err)
|
||||
@@ -173,14 +169,17 @@ func TestEventTrigger(t *testing.T) {
|
||||
})
|
||||
|
||||
t.Run("trigger enabled event enqueues task", func(t *testing.T) {
|
||||
// Set push_config system configuration
|
||||
cfgJson := `[{"channel": "mock_channel", "url": "http://mock"}]`
|
||||
sysConfig := &model.SystemConfig{
|
||||
Key: model.ConfigKeyPushConfig,
|
||||
Value: cfgJson,
|
||||
// Create an enabled custom channel in GORM
|
||||
customChan := &model.PushChannel{
|
||||
Name: "mock_channel",
|
||||
Type: "custom",
|
||||
URL: "https://webhook.site/trigger",
|
||||
Other: `{"text": "$content"}`,
|
||||
Enabled: true,
|
||||
}
|
||||
err = dbConn.Create(sysConfig).Error
|
||||
err = dbConn.Create(customChan).Error
|
||||
require.NoError(t, err)
|
||||
defer dbConn.Delete(customChan)
|
||||
|
||||
// Enable the push event in DB using struct to trigger JSON serializer
|
||||
var event model.PushEvent
|
||||
@@ -216,7 +215,8 @@ func TestEventTrigger(t *testing.T) {
|
||||
err = json.Unmarshal([]byte(execution.Payload), &payload)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, "admin_login", payload.EventKey)
|
||||
assert.Equal(t, "mock_channel", payload.Config.Channel)
|
||||
assert.Equal(t, "custom", payload.Config.Channel)
|
||||
assert.Equal(t, "https://webhook.site/trigger", payload.Config.URL)
|
||||
assert.Equal(t, "admin_user", payload.Target)
|
||||
assert.Equal(t, "管理员登录提醒", payload.Body.Title)
|
||||
assert.Contains(t, payload.Body.Content, "super_admin")
|
||||
@@ -224,15 +224,17 @@ func TestEventTrigger(t *testing.T) {
|
||||
})
|
||||
|
||||
t.Run("trigger without user injects virtual system user", func(t *testing.T) {
|
||||
// Set push_config system configuration
|
||||
dbConn.Where("key = ?", model.ConfigKeyPushConfig).Delete(&model.SystemConfig{})
|
||||
cfgJson := `[{"channel": "mock_channel", "url": "http://mock"}]`
|
||||
sysConfig := &model.SystemConfig{
|
||||
Key: model.ConfigKeyPushConfig,
|
||||
Value: cfgJson,
|
||||
// Create an enabled custom channel in GORM
|
||||
customChan := &model.PushChannel{
|
||||
Name: "mock_channel",
|
||||
Type: "custom",
|
||||
URL: "https://webhook.site/trigger",
|
||||
Other: `{"text": "$content"}`,
|
||||
Enabled: true,
|
||||
}
|
||||
err = dbConn.Create(sysConfig).Error
|
||||
err = dbConn.Create(customChan).Error
|
||||
require.NoError(t, err)
|
||||
defer dbConn.Delete(customChan)
|
||||
|
||||
// Enable the push event in DB
|
||||
var event model.PushEvent
|
||||
|
||||
@@ -1,27 +0,0 @@
|
||||
-- +goose Up
|
||||
INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at)
|
||||
VALUES (
|
||||
'push_config',
|
||||
'[]',
|
||||
'system',
|
||||
0,
|
||||
'通知推送渠道配置',
|
||||
CURRENT_TIMESTAMP,
|
||||
CURRENT_TIMESTAMP
|
||||
) ON CONFLICT (key) DO NOTHING;
|
||||
|
||||
INSERT INTO w_push_events (event_key, name, channels, targets, template, enabled, created_at, updated_at)
|
||||
VALUES (
|
||||
'admin_login',
|
||||
'管理员登录',
|
||||
'[]',
|
||||
'[]',
|
||||
'{"title": "管理员登录提醒", "content": "管理员 {{user.username}} 于 {{time}} 从 IP: {{ip}} 登录成功。", "level": "INFO"}',
|
||||
false,
|
||||
CURRENT_TIMESTAMP,
|
||||
CURRENT_TIMESTAMP
|
||||
) ON CONFLICT (event_key) DO NOTHING;
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_system_configs WHERE key = 'push_config';
|
||||
DELETE FROM w_push_events WHERE event_key = 'admin_login';
|
||||
@@ -0,0 +1,5 @@
|
||||
-- +goose Up
|
||||
DELETE FROM w_system_configs WHERE key = 'push_config';
|
||||
DELETE FROM w_push_events WHERE event_key = 'admin_login';
|
||||
|
||||
-- +goose Down
|
||||
@@ -1,27 +0,0 @@
|
||||
-- +goose Up
|
||||
INSERT OR IGNORE INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at)
|
||||
VALUES (
|
||||
'push_config',
|
||||
'[]',
|
||||
'system',
|
||||
0,
|
||||
'通知推送渠道配置',
|
||||
CURRENT_TIMESTAMP,
|
||||
CURRENT_TIMESTAMP
|
||||
);
|
||||
|
||||
INSERT OR IGNORE INTO w_push_events (event_key, name, channels, targets, template, enabled, created_at, updated_at)
|
||||
VALUES (
|
||||
'admin_login',
|
||||
'管理员登录',
|
||||
'[]',
|
||||
'[]',
|
||||
'{"title": "管理员登录提醒", "content": "管理员 {{user.username}} 于 {{time}} 从 IP: {{ip}} 登录成功。", "level": "INFO"}',
|
||||
false,
|
||||
CURRENT_TIMESTAMP,
|
||||
CURRENT_TIMESTAMP
|
||||
);
|
||||
|
||||
-- +goose Down
|
||||
DELETE FROM w_system_configs WHERE key = 'push_config';
|
||||
DELETE FROM w_push_events WHERE event_key = 'admin_login';
|
||||
@@ -0,0 +1,5 @@
|
||||
-- +goose Up
|
||||
DELETE FROM w_system_configs WHERE key = 'push_config';
|
||||
DELETE FROM w_push_events WHERE event_key = 'admin_login';
|
||||
|
||||
-- +goose Down
|
||||
@@ -17,7 +17,7 @@ import (
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
const expectedMigratedSystemConfigCount = 32
|
||||
const expectedMigratedSystemConfigCount = 31
|
||||
|
||||
func TestMigrateInitializesSQLiteDatabase(t *testing.T) {
|
||||
sqliteDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{
|
||||
|
||||
@@ -49,7 +49,6 @@ const (
|
||||
ConfigKeyLoginSessionTTLHours = "login_session_ttl_hours" // 登录会话过期时间 (小时,0表示浏览器关闭后自动退出登录,-1表示永不过期)
|
||||
ConfigKeyUpdateUpstreamRepository = "update_upstream_repository" // GitHub Actions Release 上游仓库
|
||||
ConfigKeyStorageConfig = "storage_config" // 文件存储配置 (JSON)
|
||||
ConfigKeyPushConfig = "push_config" // 通知推送渠道配置 (JSON)
|
||||
)
|
||||
|
||||
const (
|
||||
|
||||
Reference in New Issue
Block a user