2022-07-24 20:51:42 +00:00
|
|
|
package rest
|
|
|
|
|
|
|
|
import (
|
2022-12-10 16:21:22 +00:00
|
|
|
"context"
|
|
|
|
|
2022-11-09 19:53:01 +00:00
|
|
|
v2 "github.com/waku-org/go-waku/waku/v2"
|
|
|
|
"github.com/waku-org/go-waku/waku/v2/protocol"
|
2022-07-24 20:51:42 +00:00
|
|
|
)
|
|
|
|
|
|
|
|
type Adder func(msg *protocol.Envelope)
|
|
|
|
|
|
|
|
type runnerService struct {
|
|
|
|
broadcaster v2.Broadcaster
|
|
|
|
ch chan *protocol.Envelope
|
2022-12-10 16:21:22 +00:00
|
|
|
cancel context.CancelFunc
|
2022-07-24 20:51:42 +00:00
|
|
|
adder Adder
|
|
|
|
}
|
|
|
|
|
|
|
|
func newRunnerService(broadcaster v2.Broadcaster, adder Adder) *runnerService {
|
|
|
|
return &runnerService{
|
|
|
|
broadcaster: broadcaster,
|
|
|
|
adder: adder,
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2022-12-10 16:21:22 +00:00
|
|
|
func (r *runnerService) Start(ctx context.Context) {
|
|
|
|
ctx, cancel := context.WithCancel(ctx)
|
2022-07-24 20:51:42 +00:00
|
|
|
r.ch = make(chan *protocol.Envelope, 1024)
|
2022-12-10 16:21:22 +00:00
|
|
|
r.cancel = cancel
|
2022-07-24 20:51:42 +00:00
|
|
|
r.broadcaster.Register(nil, r.ch)
|
|
|
|
for {
|
|
|
|
select {
|
2022-12-10 16:21:22 +00:00
|
|
|
case <-ctx.Done():
|
2022-07-24 20:51:42 +00:00
|
|
|
return
|
|
|
|
case envelope := <-r.ch:
|
|
|
|
r.adder(envelope)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (r *runnerService) Stop() {
|
2023-01-06 22:37:57 +00:00
|
|
|
if r.cancel == nil {
|
|
|
|
return
|
|
|
|
}
|
2022-12-10 16:21:22 +00:00
|
|
|
r.cancel()
|
2022-07-24 20:51:42 +00:00
|
|
|
r.broadcaster.Unregister(nil, r.ch)
|
|
|
|
close(r.ch)
|
|
|
|
}
|