feat: 完成一般消息通知能力的开发
This commit is contained in:
@@ -0,0 +1,172 @@
|
||||
package notice
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"sort"
|
||||
"sync"
|
||||
|
||||
"go.uber.org/zap"
|
||||
"go.uber.org/zap/zapcore"
|
||||
)
|
||||
|
||||
type zapNoticeCore struct {
|
||||
senders []Sender `json:"-" dc:"日志通知使用的消息发送实例列表"`
|
||||
levels map[zapcore.Level]struct{} `json:"-" dc:"触发消息通知的精确日志等级集合"`
|
||||
fields []zapcore.Field `json:"-" dc:"通过 Logger.With 附加的上下文字段"`
|
||||
pending *sync.WaitGroup `json:"-" dc:"等待尚未完成的异步消息发送"`
|
||||
onError func(error) `json:"-" dc:"异步消息发送失败处理函数"`
|
||||
}
|
||||
|
||||
func (f *Factory) WrapZap(base *zap.Logger, config ZapConfig) (*zap.Logger, error) {
|
||||
if base == nil {
|
||||
return nil, fmt.Errorf("%w: base logger is nil", ErrInvalidZapConfig)
|
||||
}
|
||||
channels, levels, err := normalizeZapConfig(config)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
senders := make([]Sender, 0, len(channels))
|
||||
for _, channel := range channels {
|
||||
sender, getErr := f.Get(channel)
|
||||
if getErr != nil {
|
||||
return nil, getErr
|
||||
}
|
||||
senders = append(senders, sender)
|
||||
}
|
||||
core := &zapNoticeCore{
|
||||
senders: senders,
|
||||
levels: levels,
|
||||
fields: []zapcore.Field{},
|
||||
pending: &sync.WaitGroup{},
|
||||
onError: config.OnError,
|
||||
}
|
||||
return base.WithOptions(zap.WrapCore(func(original zapcore.Core) zapcore.Core {
|
||||
return zapcore.NewTee(original, core)
|
||||
})), nil
|
||||
}
|
||||
|
||||
func normalizeZapConfig(config ZapConfig) ([]Channel, map[zapcore.Level]struct{}, error) {
|
||||
if len(config.Channels) == 0 {
|
||||
return nil, nil, fmt.Errorf("%w: channels are empty", ErrInvalidZapConfig)
|
||||
}
|
||||
seen := make(map[Channel]struct{}, len(config.Channels))
|
||||
channels := make([]Channel, 0, len(config.Channels))
|
||||
for _, channel := range config.Channels {
|
||||
if !validChannel(channel) {
|
||||
return nil, nil, fmt.Errorf("%w: unsupported channel %q", ErrInvalidZapConfig, channel)
|
||||
}
|
||||
if _, exists := seen[channel]; exists {
|
||||
continue
|
||||
}
|
||||
seen[channel] = struct{}{}
|
||||
channels = append(channels, channel)
|
||||
}
|
||||
configuredLevels := config.Levels
|
||||
if len(configuredLevels) == 0 {
|
||||
configuredLevels = []zapcore.Level{zapcore.WarnLevel, zapcore.ErrorLevel}
|
||||
}
|
||||
levels := make(map[zapcore.Level]struct{}, len(configuredLevels))
|
||||
for _, level := range configuredLevels {
|
||||
if level < zapcore.DebugLevel || level > zapcore.FatalLevel {
|
||||
return nil, nil, fmt.Errorf("%w: invalid level %d", ErrInvalidZapConfig, level)
|
||||
}
|
||||
levels[level] = struct{}{}
|
||||
}
|
||||
return channels, levels, nil
|
||||
}
|
||||
|
||||
func (c *zapNoticeCore) Enabled(level zapcore.Level) bool {
|
||||
_, exists := c.levels[level]
|
||||
return exists
|
||||
}
|
||||
|
||||
func (c *zapNoticeCore) With(fields []zapcore.Field) zapcore.Core {
|
||||
cloned := *c
|
||||
cloned.fields = make([]zapcore.Field, 0, len(c.fields)+len(fields))
|
||||
cloned.fields = append(cloned.fields, c.fields...)
|
||||
cloned.fields = append(cloned.fields, fields...)
|
||||
return &cloned
|
||||
}
|
||||
|
||||
func (c *zapNoticeCore) Check(entry zapcore.Entry, checked *zapcore.CheckedEntry) *zapcore.CheckedEntry {
|
||||
if c.Enabled(entry.Level) {
|
||||
return checked.AddCore(entry, c)
|
||||
}
|
||||
return checked
|
||||
}
|
||||
|
||||
func (c *zapNoticeCore) Write(entry zapcore.Entry, fields []zapcore.Field) error {
|
||||
message := zapEntryMessage(entry, append(append([]zapcore.Field{}, c.fields...), fields...))
|
||||
c.pending.Add(1)
|
||||
go func() {
|
||||
defer c.pending.Done()
|
||||
var sendErrors []error
|
||||
for index, sender := range c.senders {
|
||||
if err := sender.Send(context.Background(), message); err != nil {
|
||||
sendErrors = append(sendErrors, fmt.Errorf("notice sender %d: %w", index, err))
|
||||
}
|
||||
}
|
||||
if len(sendErrors) > 0 && c.onError != nil {
|
||||
c.onError(errors.Join(sendErrors...))
|
||||
}
|
||||
}()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *zapNoticeCore) Sync() error {
|
||||
c.pending.Wait()
|
||||
return nil
|
||||
}
|
||||
|
||||
func zapEntryMessage(entry zapcore.Entry, fields []zapcore.Field) Message {
|
||||
titleSuffix := entry.LoggerName
|
||||
if titleSuffix == "" {
|
||||
titleSuffix = "日志告警"
|
||||
}
|
||||
cardFields := encodeZapFields(fields)
|
||||
cardFields = append(cardFields, CardField{Name: "timestamp", Value: entry.Time.Format("2006-01-02T15:04:05.000Z07:00")})
|
||||
if entry.Caller.Defined {
|
||||
cardFields = append(cardFields, CardField{Name: "caller", Value: entry.Caller.TrimmedPath()})
|
||||
}
|
||||
return Card(CardContent{
|
||||
Title: fmt.Sprintf("[%s] %s", entry.Level.CapitalString(), titleSuffix),
|
||||
Theme: zapLevelTheme(entry.Level),
|
||||
Markdown: entry.Message,
|
||||
Fields: cardFields,
|
||||
})
|
||||
}
|
||||
|
||||
func encodeZapFields(fields []zapcore.Field) []CardField {
|
||||
encoder := zapcore.NewMapObjectEncoder()
|
||||
for _, field := range fields {
|
||||
field.AddTo(encoder)
|
||||
}
|
||||
keys := make([]string, 0, len(encoder.Fields))
|
||||
for key := range encoder.Fields {
|
||||
keys = append(keys, key)
|
||||
}
|
||||
sort.Strings(keys)
|
||||
result := make([]CardField, 0, len(keys))
|
||||
for _, key := range keys {
|
||||
data, err := json.Marshal(encoder.Fields[key])
|
||||
if err != nil {
|
||||
data = []byte(fmt.Sprint(encoder.Fields[key]))
|
||||
}
|
||||
result = append(result, CardField{Name: key, Value: string(data)})
|
||||
}
|
||||
return result
|
||||
}
|
||||
|
||||
func zapLevelTheme(level zapcore.Level) CardTheme {
|
||||
switch {
|
||||
case level >= zapcore.ErrorLevel:
|
||||
return CardThemeRed
|
||||
case level == zapcore.WarnLevel:
|
||||
return CardThemeOrange
|
||||
default:
|
||||
return CardThemeBlue
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user