@@ -0,0 +1,100 @@
|
||||
package laniakea
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
|
||||
"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
|
||||
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 {
|
||||
dispatcher := &observerDispatcher{
|
||||
observer: observer,
|
||||
logger: logger,
|
||||
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()
|
||||
} else {
|
||||
ctx = context.WithoutCancel(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() {
|
||||
d.stopOnce.Do(func() {
|
||||
d.mu.Lock()
|
||||
d.closed = true
|
||||
close(d.stop)
|
||||
d.mu.Unlock()
|
||||
})
|
||||
d.wg.Wait()
|
||||
}
|
||||
Reference in New Issue
Block a user