mirror of
https://github.com/yusing/godoxy.git
synced 2026-04-11 03:06:51 +02:00
refactor(events): move event_queue.go to goutils/eventqueue package
This commit is contained in:
@@ -17,6 +17,7 @@ import (
|
||||
W "github.com/yusing/godoxy/internal/watcher"
|
||||
"github.com/yusing/godoxy/internal/watcher/events"
|
||||
gperr "github.com/yusing/goutils/errs"
|
||||
"github.com/yusing/goutils/eventqueue"
|
||||
"github.com/yusing/goutils/task"
|
||||
)
|
||||
|
||||
@@ -115,19 +116,19 @@ func (p *Provider) Start(parent task.Parent) error {
|
||||
|
||||
err := errs.Wait().Error()
|
||||
|
||||
eventQueue := events.NewEventQueue(
|
||||
t.Subtask("event_queue", false),
|
||||
providerEventFlushInterval,
|
||||
func(events []events.Event) {
|
||||
opts := eventqueue.Options[events.Event]{
|
||||
FlushInterval: providerEventFlushInterval,
|
||||
OnFlush: func(events []events.Event) {
|
||||
handler := p.newEventHandler()
|
||||
// routes' lifetime should follow the provider's lifetime
|
||||
handler.Handle(t, events)
|
||||
handler.Log()
|
||||
},
|
||||
func(err error) {
|
||||
OnError: func(err error) {
|
||||
p.Logger().Err(err).Msg("event error")
|
||||
},
|
||||
)
|
||||
}
|
||||
eventQueue := eventqueue.New(t.Subtask("event_queue", false), opts)
|
||||
eventQueue.Start(p.watcher.Events(t.Context()))
|
||||
|
||||
if err != nil {
|
||||
|
||||
Reference in New Issue
Block a user