FILE / ScuroNeko/Laniakea

observer_async.go

Исходный файл и его история в репозитории.
FILE 24040fe164649c3e4a7ef2f9d37b9747957b0cc1
Files
Laniakea/observer_async.go
T
ScuroNeko 24040fe164
Golang lint / lint (push) Successful in 58s
Golang lint / lint (pull_request) Successful in 13m10s
(new): expand runtime APIs
(fix): harden concurrent lifecycle
(tests): add regression coverage
(doc): update v1.2 guidance
2026-08-20 11:08:45 +03:00

126 lines
2.7 KiB
Go

package laniakea
import (
"context"
"fmt"
"sync"
"sync/atomic"
"time"
"git.scuroneko.dev/scuroneko/sneklog/v2"
)
const observerQueueSize = 1024
type queuedObserverEvent struct {
ctx context.Context
event Event
}
type observerDispatcher struct {
observer Observer
logger *sneklog.Logger
ctx context.Context
cancel context.CancelFunc
queue chan queuedObserverEvent
stop chan struct{}
mu sync.RWMutex
closed bool
wg sync.WaitGroup
stopOnce sync.Once
dropped atomic.Uint64
}
func newObserverDispatcher(observer Observer, logger *sneklog.Logger) *observerDispatcher {
ctx, cancel := context.WithCancel(context.Background())
dispatcher := &observerDispatcher{
observer: observer,
logger: logger,
ctx: ctx,
cancel: cancel,
queue: make(chan queuedObserverEvent, observerQueueSize),
stop: make(chan struct{}),
}
dispatcher.wg.Add(1)
go dispatcher.run()
return dispatcher
}
func (d *observerDispatcher) enqueue(ctx context.Context, event Event) {
if ctx == nil {
ctx = context.Background()
}
ctx = observerEventContext{Context: context.WithoutCancel(ctx), lifecycle: d.ctx}
d.mu.RLock()
defer d.mu.RUnlock()
if d.closed {
return
}
select {
case d.queue <- queuedObserverEvent{ctx: ctx, event: event}:
default:
dropped := d.dropped.Add(1)
if d.logger != nil && (dropped == 1 || dropped&(dropped-1) == 0) {
d.logger.Warnf("observer queue full; dropped %d events", dropped)
}
}
}
func (d *observerDispatcher) run() {
defer d.wg.Done()
for {
select {
case queued := <-d.queue:
d.dispatch(queued)
case <-d.stop:
for {
select {
case queued := <-d.queue:
d.dispatch(queued)
default:
return
}
}
}
}
}
func (d *observerDispatcher) dispatch(queued queuedObserverEvent) {
defer func() {
if recovered := recover(); recovered != nil && d.logger != nil {
d.logger.Errorln(fmt.Sprintf("panic in observer: %v", recovered))
}
}()
emitObserverEvent(d.observer, queued.ctx, queued.event)
}
func (d *observerDispatcher) close(ctx context.Context) error {
d.stopOnce.Do(func() {
d.mu.Lock()
d.closed = true
d.cancel()
close(d.stop)
d.mu.Unlock()
})
done := make(chan struct{})
go func() {
d.wg.Wait()
close(done)
}()
select {
case <-done:
return nil
case <-ctx.Done():
return fmt.Errorf("%w: %v", ErrObserverShutdownTimeout, ctx.Err())
}
}
type observerEventContext struct {
context.Context
lifecycle context.Context
}
func (ctx observerEventContext) Deadline() (time.Time, bool) { return ctx.lifecycle.Deadline() }
func (ctx observerEventContext) Done() <-chan struct{} { return ctx.lifecycle.Done() }
func (ctx observerEventContext) Err() error { return ctx.lifecycle.Err() }