-
Notifications
You must be signed in to change notification settings - Fork 0
/
blocking_producer.go
36 lines (29 loc) · 1.05 KB
/
blocking_producer.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
package producer
import (
"context"
"fmt"
"github.com/rs/zerolog/log"
"google.golang.org/protobuf/proto"
"github.com/areknoster/public-distributed-commit-log/sentinel/sentinelpb"
"github.com/areknoster/public-distributed-commit-log/storage"
)
type BlockingProducer struct {
storage storage.MessageWriter
sentinelClient sentinelpb.SentinelClient
}
func NewBlockingProducer(writer storage.MessageWriter, sentinelClient sentinelpb.SentinelClient) *BlockingProducer {
return &BlockingProducer{storage: writer, sentinelClient: sentinelClient}
}
func (m *BlockingProducer) Produce(ctx context.Context, message proto.Message) error {
cid, err := m.storage.Write(ctx, message)
if err != nil {
return fmt.Errorf("save message to storage: %w", err)
}
log.Debug().Stringer("cid", cid).Msg("message stored")
_, err = m.sentinelClient.Publish(ctx, &sentinelpb.PublishRequest{Cid: cid.String()})
if err != nil {
return fmt.Errorf("publish MessageBuf to sentinel: %w", err)
}
log.Info().Stringer("cid", cid).Msg("message accepted by sentinel")
return nil
}