consumer constructor
This commit is contained in:
parent
9366cbbf6d
commit
4c5b1c7030
2 changed files with 22 additions and 8 deletions
17
consumer.go
17
consumer.go
|
|
@ -7,6 +7,23 @@ import (
|
|||
"log"
|
||||
)
|
||||
|
||||
type Consumer interface {
|
||||
Start() chan []byte
|
||||
}
|
||||
|
||||
type consumeHandler struct {
|
||||
ctx context.Context
|
||||
client *Client
|
||||
chanLen int
|
||||
}
|
||||
|
||||
func (c *consumeHandler) Start() chan []byte {
|
||||
msgCh := make(chan []byte, c.chanLen)
|
||||
go runConsumer(c.ctx, c.client, msgCh)
|
||||
|
||||
return msgCh
|
||||
}
|
||||
|
||||
func runConsumer(ctx context.Context, client *Client, msgCh chan []byte) {
|
||||
runCtx, cancel := context.WithCancel(ctx)
|
||||
defer cancel()
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue