Skip to content

Commit ec1d758

Browse files
Reconnect only while waiting for a worker job
Keep retry around consumer setup and the delivery receive. Once a real AMQP message arrives, handle it as a one-shot job and do not reconnect, so a broker flap cannot duplicate execution. Co-authored-by: Max Schmitt <max@schmitt.mx>
1 parent eb7ee7a commit ec1d758

2 files changed

Lines changed: 70 additions & 54 deletions

File tree

‎internal/worker/worker.go‎

Lines changed: 55 additions & 45 deletions
Original file line numberDiff line numberDiff line change
@@ -62,61 +62,75 @@ func (w *Worker) Run() {
6262
}
6363

6464
for {
65-
conn, err := amqp.Dial(os.Getenv("AMQP_URL"))
65+
conn, msgs, err := w.openConsumer()
6666
if err != nil {
67-
log.Printf("could not dial to amqp: %v", err)
67+
log.Printf("%v", err)
6868
time.Sleep(time.Second)
6969
continue
7070
}
7171

72-
w.channel, err = conn.Channel()
72+
incomingMessage, err := readDelivery(msgs)
7373
if err != nil {
7474
conn.Close()
75-
log.Printf("could not open a channel: %v", err)
76-
time.Sleep(time.Second)
77-
continue
78-
}
79-
if _, err := w.channel.QueueDeclare(
80-
queue_name,
81-
false, // durable
82-
true, // delete when unused
83-
false, // exclusive
84-
false, // noWait
85-
nil, // args
86-
); err != nil {
87-
conn.Close()
88-
log.Printf("could not declare queue: %v", err)
75+
log.Printf("amqp delivery channel closed, reconnecting")
8976
time.Sleep(time.Second)
9077
continue
9178
}
92-
msgs, err := w.channel.Consume(
93-
queue_name,
94-
"", // consumer
95-
false, // auto-ack
96-
false, // exclusive
97-
false, // no-local
98-
false, // no-wait
99-
nil, // args
100-
)
79+
80+
err = w.handleDelivery(incomingMessage)
81+
conn.Close()
10182
if err != nil {
102-
conn.Close()
103-
log.Printf("could not consume channel messages: %v", err)
104-
time.Sleep(time.Second)
105-
continue
83+
log.Fatalf("could not consume messages: %v", err)
10684
}
85+
return
86+
}
87+
}
10788

108-
err = w.consumeMessage(msgs)
89+
func (w *Worker) openConsumer() (*amqp.Connection, <-chan amqp.Delivery, error) {
90+
conn, err := amqp.Dial(os.Getenv("AMQP_URL"))
91+
if err != nil {
92+
return nil, nil, fmt.Errorf("could not dial to amqp: %w", err)
93+
}
94+
95+
channel, err := conn.Channel()
96+
if err != nil {
10997
conn.Close()
110-
if err == nil {
111-
return
112-
}
113-
if errors.Is(err, errAMQPChannelClosed) {
114-
log.Printf("amqp delivery channel closed, reconnecting")
115-
time.Sleep(time.Second)
116-
continue
117-
}
118-
log.Fatalf("could not consume messages: %v", err)
98+
return nil, nil, fmt.Errorf("could not open a channel: %w", err)
99+
}
100+
if _, err := channel.QueueDeclare(
101+
queue_name,
102+
false, // durable
103+
true, // delete when unused
104+
false, // exclusive
105+
false, // noWait
106+
nil, // args
107+
); err != nil {
108+
conn.Close()
109+
return nil, nil, fmt.Errorf("could not declare queue: %w", err)
110+
}
111+
msgs, err := channel.Consume(
112+
queue_name,
113+
"", // consumer
114+
false, // auto-ack
115+
false, // exclusive
116+
false, // no-local
117+
false, // no-wait
118+
nil, // args
119+
)
120+
if err != nil {
121+
conn.Close()
122+
return nil, nil, fmt.Errorf("could not consume channel messages: %w", err)
119123
}
124+
w.channel = channel
125+
return conn, msgs, nil
126+
}
127+
128+
func readDelivery(incomingMessages <-chan amqp.Delivery) (amqp.Delivery, error) {
129+
incomingMessage, ok := <-incomingMessages
130+
if !ok {
131+
return amqp.Delivery{}, errAMQPChannelClosed
132+
}
133+
return incomingMessage, nil
120134
}
121135

122136
func (w *Worker) AddEnv(key, value string) {
@@ -168,11 +182,7 @@ func (w *Worker) ExecCommand(name string, args ...string) error {
168182
return nil
169183
}
170184

171-
func (w *Worker) consumeMessage(incomingMessages <-chan amqp.Delivery) error {
172-
incomingMessage, ok := <-incomingMessages
173-
if !ok {
174-
return errAMQPChannelClosed
175-
}
185+
func (w *Worker) handleDelivery(incomingMessage amqp.Delivery) error {
176186
var incomingMessageParsed *workertypes.WorkerRequestPayload
177187
if err := json.Unmarshal(incomingMessage.Body, &incomingMessageParsed); err != nil {
178188
return fmt.Errorf("could not parse incoming amqp message: %w", err)

‎internal/worker/worker_test.go‎

Lines changed: 15 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -8,22 +8,28 @@ import (
88
amqp "github.com/rabbitmq/amqp091-go"
99
)
1010

11-
func TestConsumeMessageClosedDeliveryChannel(t *testing.T) {
12-
w := NewWorker(&WorkerExecutionOptions{
13-
Handler: func(worker *Worker, code string) error {
14-
t.Fatal("handler should not run when the delivery channel is closed")
15-
return nil
16-
},
17-
})
18-
11+
func TestReadDeliveryClosedChannel(t *testing.T) {
1912
incoming := make(chan amqp.Delivery)
2013
close(incoming)
2114

22-
err := w.consumeMessage(incoming)
15+
_, err := readDelivery(incoming)
2316
if !errors.Is(err, errAMQPChannelClosed) {
2417
t.Fatalf("expected errAMQPChannelClosed, got %v", err)
2518
}
2619
if strings.Contains(err.Error(), "unexpected end of JSON input") {
2720
t.Fatalf("closed delivery channel was treated as a JSON payload: %v", err)
2821
}
2922
}
23+
24+
func TestReadDeliveryReturnsMessage(t *testing.T) {
25+
incoming := make(chan amqp.Delivery, 1)
26+
incoming <- amqp.Delivery{Body: []byte(`{"code":"x"}`)}
27+
28+
msg, err := readDelivery(incoming)
29+
if err != nil {
30+
t.Fatal(err)
31+
}
32+
if string(msg.Body) != `{"code":"x"}` {
33+
t.Fatalf("unexpected body %q", msg.Body)
34+
}
35+
}

0 commit comments

Comments
 (0)