go-waku/cmd/waku/server/rest/runner.go

49 lines
938 B
Go
Raw Normal View History

package rest
import (
"context"
"github.com/waku-org/go-waku/waku/v2/protocol"
2023-05-05 15:19:15 +05:30
"github.com/waku-org/go-waku/waku/v2/protocol/relay"
)
type Adder func(msg *protocol.Envelope)
type runnerService struct {
2023-05-05 15:19:15 +05:30
broadcaster relay.Broadcaster
sub *relay.Subscription
cancel context.CancelFunc
adder Adder
}
2023-05-05 15:19:15 +05:30
func newRunnerService(broadcaster relay.Broadcaster, adder Adder) *runnerService {
return &runnerService{
broadcaster: broadcaster,
adder: adder,
}
}
func (r *runnerService) Start(ctx context.Context) {
ctx, cancel := context.WithCancel(ctx)
r.cancel = cancel
r.sub = r.broadcaster.RegisterForAll(relay.WithBufferSize(relay.DefaultRelaySubscriptionBufferSize))
for {
select {
case <-ctx.Done():
return
2023-05-05 15:19:15 +05:30
case envelope, ok := <-r.sub.Ch:
2023-05-05 17:33:44 +05:30
if ok {
2023-05-05 15:19:15 +05:30
r.adder(envelope)
}
}
}
}
func (r *runnerService) Stop() {
2023-01-06 18:37:57 -04:00
if r.cancel == nil {
return
}
2023-05-05 15:19:15 +05:30
r.sub.Unsubscribe()
r.cancel()
}