diff --git a/client/publisher.go b/client/publisher.go index 05717d7..188ac39 100644 --- a/client/publisher.go +++ b/client/publisher.go @@ -2,9 +2,12 @@ package client import ( "encoding/binary" + "errors" "fmt" + "log/slog" "net" "sync" + "syscall" "github.com/willdot/messagebroker/internal/server" ) @@ -13,10 +16,23 @@ import ( type Publisher struct { conn net.Conn connMu sync.Mutex + addr string } // NewPublisher connects to the server at the given address and registers as a publisher func NewPublisher(addr string) (*Publisher, error) { + conn, err := connect(addr) + if err != nil { + return nil, fmt.Errorf("failed to connect to server: %w", err) + } + + return &Publisher{ + conn: conn, + addr: addr, + }, nil +} + +func connect(addr string) (net.Conn, error) { conn, err := net.Dial("tcp", addr) if err != nil { return nil, fmt.Errorf("failed to dial: %w", err) @@ -27,10 +43,7 @@ func NewPublisher(addr string) (*Publisher, error) { conn.Close() return nil, fmt.Errorf("failed to register publish to server: %w", err) } - - return &Publisher{ - conn: conn, - }, nil + return conn, nil } // Close cleanly shuts down the publisher @@ -40,6 +53,10 @@ func (p *Publisher) Close() error { // Publish will publish the given message to the server func (p *Publisher) PublishMessage(message *Message) error { + return p.publishMessageWithRetry(message, 0) +} + +func (p *Publisher) publishMessageWithRetry(message *Message, attempt int) error { op := func(conn net.Conn) error { // send topic first topic := fmt.Sprintf("topic:%s", message.Topic) @@ -60,7 +77,32 @@ func (p *Publisher) PublishMessage(message *Message) error { return nil } - return p.connOperation(op) + err := p.connOperation(op) + if err == nil { + return nil + } + + // we can handle a broken pipe by trying to reconnect, but if it's a different error return it + if !errors.Is(err, syscall.EPIPE) { + return err + } + + slog.Info("error is broken pipe") + + if attempt >= 5 { + return fmt.Errorf("failed to publish message after max attempts to reconnect (%d): %w", attempt, err) + } + + slog.Error("failed to publish message", "error", err) + + conn, connectErr := connect(p.addr) + if connectErr != nil { + return fmt.Errorf("failed to reconnect after failing to publish message: %w", connectErr) + } + + p.conn = conn + + return p.publishMessageWithRetry(message, attempt+1) } func (p *Publisher) connOperation(op connOpp) error { diff --git a/client/subscriber.go b/client/subscriber.go index 1177dd8..d2d78df 100644 --- a/client/subscriber.go +++ b/client/subscriber.go @@ -190,6 +190,7 @@ func (s *Subscriber) consume(ctx context.Context, consumer *Consumer) { err := s.readMessage(ctx, consumer.msgs) if err != nil { + // TODO: if broken pipe, we need to somehow reconnect and subscribe again....YIKES consumer.Err = err return } diff --git a/example/main.go b/example/main.go index 74e7afb..d7fdb85 100644 --- a/example/main.go +++ b/example/main.go @@ -57,6 +57,8 @@ func main() { msg.Ack(true) } + time.Sleep(time.Second * 30) + } func sendMessages() { @@ -81,6 +83,8 @@ func sendMessages() { continue } + slog.Info("message sent") + time.Sleep(time.Millisecond * 500) } }