| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278 |
- package loghub
- import (
- "context"
- "log/slog"
- "sort"
- "strings"
- "sync"
- "time"
- )
- // Entry is the stable, secret-neutral representation exposed by the log API.
- // Callers remain responsible for never adding credentials or keying material
- // to slog attributes.
- type Entry struct {
- Time time.Time `json:"time"`
- Level string `json:"level"`
- Message string `json:"message"`
- Caller string `json:"caller,omitempty"`
- Fields map[string]any `json:"fields,omitempty"`
- }
- type core struct {
- mu sync.RWMutex
- capacity int
- entries []Entry
- subscribers map[uint64]chan Entry
- nextID uint64
- }
- // Hub is both a slog.Handler and a bounded live log source.
- type Hub struct {
- next slog.Handler
- core *core
- attrs []slog.Attr
- groups []string
- }
- func New(next slog.Handler, capacity int) *Hub {
- if next == nil {
- next = slog.NewTextHandler(discardWriter{}, nil)
- }
- if capacity < 100 {
- capacity = 100
- }
- return &Hub{
- next: next,
- core: &core{
- capacity: capacity,
- entries: make([]Entry, 0, capacity),
- subscribers: make(map[uint64]chan Entry),
- },
- }
- }
- func (h *Hub) Enabled(ctx context.Context, level slog.Level) bool {
- return h.next.Enabled(ctx, level)
- }
- func (h *Hub) Handle(ctx context.Context, record slog.Record) error {
- err := h.next.Handle(ctx, record)
- fields := make(map[string]any)
- for _, attr := range h.attrs {
- appendAttribute(fields, h.groups, attr)
- }
- record.Attrs(func(attr slog.Attr) bool {
- appendAttribute(fields, h.groups, attr)
- return true
- })
- entry := Entry{
- Time: record.Time.UTC(),
- Level: levelName(record.Level),
- Message: record.Message,
- Fields: fields,
- }
- if len(fields) == 0 {
- entry.Fields = nil
- }
- h.publish(entry)
- return err
- }
- func (h *Hub) WithAttrs(attrs []slog.Attr) slog.Handler {
- nextAttrs := append(append([]slog.Attr(nil), h.attrs...), attrs...)
- return &Hub{
- next: h.next.WithAttrs(attrs),
- core: h.core,
- attrs: nextAttrs,
- groups: append([]string(nil), h.groups...),
- }
- }
- func (h *Hub) WithGroup(name string) slog.Handler {
- name = strings.TrimSpace(name)
- groups := append([]string(nil), h.groups...)
- if name != "" {
- groups = append(groups, name)
- }
- return &Hub{
- next: h.next.WithGroup(name),
- core: h.core,
- attrs: append([]slog.Attr(nil), h.attrs...),
- groups: groups,
- }
- }
- func (h *Hub) publish(entry Entry) {
- h.core.mu.Lock()
- if len(h.core.entries) == h.core.capacity {
- copy(h.core.entries, h.core.entries[1:])
- h.core.entries[len(h.core.entries)-1] = cloneEntry(entry)
- } else {
- h.core.entries = append(h.core.entries, cloneEntry(entry))
- }
- for _, subscriber := range h.core.subscribers {
- select {
- case subscriber <- cloneEntry(entry):
- default:
- select {
- case <-subscriber:
- default:
- }
- select {
- case subscriber <- cloneEntry(entry):
- default:
- }
- }
- }
- h.core.mu.Unlock()
- }
- // History returns the newest matching entries in chronological order.
- func (h *Hub) History(limit int, minimum slog.Level, search string) []Entry {
- if limit < 1 {
- limit = 1
- }
- if limit > h.core.capacity {
- limit = h.core.capacity
- }
- search = strings.ToLower(strings.TrimSpace(search))
- h.core.mu.RLock()
- result := make([]Entry, 0, limit)
- for index := len(h.core.entries) - 1; index >= 0 && len(result) < limit; index-- {
- entry := h.core.entries[index]
- if parseLevel(entry.Level) < minimum {
- continue
- }
- if search != "" && !entryContains(entry, search) {
- continue
- }
- result = append(result, cloneEntry(entry))
- }
- h.core.mu.RUnlock()
- sort.SliceStable(result, func(i, j int) bool { return result[i].Time.Before(result[j].Time) })
- return result
- }
- func (h *Hub) Subscribe(buffer int) (<-chan Entry, func()) {
- if buffer < 1 {
- buffer = 1
- }
- if buffer > 1000 {
- buffer = 1000
- }
- channel := make(chan Entry, buffer)
- h.core.mu.Lock()
- id := h.core.nextID
- h.core.nextID++
- h.core.subscribers[id] = channel
- h.core.mu.Unlock()
- var once sync.Once
- cancel := func() {
- once.Do(func() {
- h.core.mu.Lock()
- delete(h.core.subscribers, id)
- close(channel)
- h.core.mu.Unlock()
- })
- }
- return channel, cancel
- }
- func appendAttribute(fields map[string]any, groups []string, attr slog.Attr) {
- attr.Value = attr.Value.Resolve()
- if attr.Equal(slog.Attr{}) {
- return
- }
- target := fields
- for _, group := range groups {
- next, ok := target[group].(map[string]any)
- if !ok {
- next = make(map[string]any)
- target[group] = next
- }
- target = next
- }
- if attr.Value.Kind() == slog.KindGroup {
- group := make(map[string]any)
- for _, child := range attr.Value.Group() {
- appendAttribute(group, nil, child)
- }
- target[attr.Key] = group
- return
- }
- target[attr.Key] = attr.Value.Any()
- }
- func levelName(level slog.Level) string {
- switch {
- case level >= slog.LevelError:
- return "error"
- case level >= slog.LevelWarn:
- return "warn"
- case level >= slog.LevelInfo:
- return "info"
- default:
- return "debug"
- }
- }
- func parseLevel(value string) slog.Level {
- switch strings.ToLower(strings.TrimSpace(value)) {
- case "error":
- return slog.LevelError
- case "warn", "warning":
- return slog.LevelWarn
- case "info", "":
- return slog.LevelInfo
- default:
- return slog.LevelDebug
- }
- }
- func entryContains(entry Entry, search string) bool {
- if strings.Contains(strings.ToLower(entry.Message), search) ||
- strings.Contains(strings.ToLower(entry.Caller), search) {
- return true
- }
- for key, value := range entry.Fields {
- if strings.Contains(strings.ToLower(key), search) ||
- strings.Contains(strings.ToLower(toString(value)), search) {
- return true
- }
- }
- return false
- }
- func cloneEntry(entry Entry) Entry {
- if entry.Fields != nil {
- entry.Fields = cloneMap(entry.Fields)
- }
- return entry
- }
- func cloneMap(source map[string]any) map[string]any {
- result := make(map[string]any, len(source))
- for key, value := range source {
- if nested, ok := value.(map[string]any); ok {
- result[key] = cloneMap(nested)
- } else {
- result[key] = value
- }
- }
- return result
- }
- func toString(value any) string {
- if stringValue, ok := value.(string); ok {
- return stringValue
- }
- return slog.AnyValue(value).String()
- }
- type discardWriter struct{}
- func (discardWriter) Write(data []byte) (int, error) {
- return len(data), nil
- }
|