FILE / ScuroNeko/Laniakea
observer_async.go
Исходный файл и его история в репозитории.
(fix): harden concurrent lifecycle (tests): add regression coverage (doc): update v1.2 guidance
126 lines
2.7 KiB
Go
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() }
|