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 } }