173 lines
5.0 KiB
Go
173 lines
5.0 KiB
Go
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
|
|
}
|
|
}
|