Skip to content

Commit 582797e

Browse files
committed
Refactor nack handling and update docs
1 parent 2358a2d commit 582797e

4 files changed

Lines changed: 95 additions & 12 deletions

File tree

lib/broadway_sqs/ex_aws_client.ex

Lines changed: 42 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -37,33 +37,67 @@ defmodule BroadwaySQS.ExAwsClient do
3737
def ack(ack_ref, successful, failed) do
3838
ack_options = :persistent_term.get(ack_ref)
3939

40-
messages =
41-
Enum.filter(successful, &ack?(&1, ack_options, :on_success)) ++
42-
Enum.filter(failed, &ack?(&1, ack_options, :on_failure))
40+
success_messages_by_ack_action =
41+
Enum.group_by(successful, &ack_action(&1, ack_options, :on_success))
4342

44-
messages
43+
failed_messages_by_ack_action =
44+
Enum.group_by(failed, &ack_action(&1, ack_options, :on_failure))
45+
46+
delete_messages =
47+
Map.get(success_messages_by_ack_action, :ack, []) ++
48+
Map.get(failed_messages_by_ack_action, :ack, [])
49+
50+
delete_messages
4551
|> Enum.chunk_every(@max_num_messages_allowed_by_aws)
46-
|> Enum.each(fn messages -> delete_messages(messages, ack_options) end)
52+
|> Enum.each(&delete_messages_batch(&1, ack_options))
53+
54+
change_visibility_entries =
55+
collect_nack_entries(success_messages_by_ack_action) ++
56+
collect_nack_entries(failed_messages_by_ack_action)
57+
58+
change_visibility_entries
59+
|> Enum.chunk_every(@max_num_messages_allowed_by_aws)
60+
|> Enum.each(&change_message_visibility_batch(&1, ack_options))
4761
end
4862

49-
defp ack?(message, ack_options, option) do
63+
defp collect_nack_entries(messages_by_ack_action) do
64+
Enum.flat_map(messages_by_ack_action, fn
65+
{{:nack, timeout}, messages} -> Enum.map(messages, &{&1, timeout})
66+
_ -> []
67+
end)
68+
end
69+
70+
defp ack_action(message, ack_options, option) do
5071
{_, _, message_ack_options} = message.acknowledger
51-
(message_ack_options[option] || Map.fetch!(ack_options, option)) == :ack
72+
message_ack_options[option] || Map.fetch!(ack_options, option)
5273
end
5374

5475
@impl Acknowledger
5576
def configure(_ack_ref, ack_data, options) do
5677
{:ok, Map.merge(ack_data, Map.new(options))}
5778
end
5879

59-
defp delete_messages(messages, ack_options) do
80+
defp delete_messages_batch(messages, ack_options) do
6081
receipts = Enum.map(messages, &extract_message_receipt/1)
6182

6283
ack_options.queue_url
6384
|> ExAws.SQS.delete_message_batch(receipts)
6485
|> ExAws.request!(ack_options.config)
6586
end
6687

88+
defp change_message_visibility_batch(messages, ack_options) do
89+
entries =
90+
Enum.map(messages, fn {message, timeout} ->
91+
message
92+
|> extract_message_receipt()
93+
|> Map.put(:visibility_timeout, timeout)
94+
end)
95+
96+
ack_options.queue_url
97+
|> ExAws.SQS.change_message_visibility_batch(entries)
98+
|> ExAws.request!(ack_options.config)
99+
end
100+
67101
defp wrap_received_messages({:ok, %{body: body}}, %{ack_ref: ack_ref}) do
68102
Enum.map(body.messages, fn message ->
69103
metadata = Map.delete(message, :body)

lib/broadway_sqs/options.ex

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -34,15 +34,15 @@ defmodule BroadwaySQS.Options do
3434
default: 5000
3535
],
3636
on_success: [
37-
type: :atom,
37+
type: {:custom, __MODULE__, :type_ack_action, [[{:name, :on_success}]]},
3838
default: :ack,
3939
doc: """
4040
configures the acking behaviour for successful messages. See the
4141
"Acknowledgments" section below for all the possible values.
4242
"""
4343
],
4444
on_failure: [
45-
type: :atom,
45+
type: {:custom, __MODULE__, :type_ack_action, [[{:name, :on_failure}]]},
4646
default: :noop,
4747
doc: """
4848
configures the acking behaviour for failed messages. See the
@@ -221,4 +221,16 @@ defmodule BroadwaySQS.Options do
221221
"expected :#{name} to be a list with possible members #{inspect(allowed_members)}, got: #{inspect(value)}"}
222222
end
223223
end
224+
225+
def type_ack_action(:ack, _opts), do: {:ok, :ack}
226+
def type_ack_action(:noop, _opts), do: {:ok, :noop}
227+
def type_ack_action({:nack, timeout}, [{:name, name}])
228+
when is_integer(timeout) and timeout >= 0 and timeout <= 43_200 do
229+
{:ok, {:nack, timeout}}
230+
end
231+
232+
def type_ack_action(value, [{:name, name}]) do
233+
{:error,
234+
"expected :#{name} to be :ack, :noop, or {:nack, timeout}, got: #{inspect(value)}"}
235+
end
224236
end

lib/broadway_sqs/producer.ex

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -31,8 +31,13 @@ defmodule BroadwaySQS.Producer do
3131
and will not redeliver it to any other consumer.
3232
3333
* `:noop` - do not acknowledge the message. SQS will eventually redeliver the message
34-
or remove it based on the "Visibility Timeout" and "Max Receive Count"
35-
configurations. For more information, see:
34+
or remove it based on the "Visibility Timeout" and "Max Receive Count" configurations.
35+
36+
* `{:nack, timeout}` - change the message visibility timeout to `timeout` seconds
37+
(0 to 43_200). The message will become available again for processing after
38+
the given amount of time.
39+
40+
For more information, see:
3641
3742
* ["Visibility Timeout" page on Amazon SQS](https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/sqs-visibility-timeout.html)
3843
* ["Dead Letter Queue" page on Amazon SQS](https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/sqs-dead-letter-queues.html)

test/broadway_sqs/ex_aws_client_test.exs

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,12 @@ defmodule BroadwaySQS.ExAwsClientTest do
5252

5353
{:ok, %{status_code: 200, body: "<DeleteMessageBatchResponse />"}}
5454
end
55+
56+
def request(:post, url, "Action=ChangeMessageVisibilityBatch" <> _ = body, _, _) do
57+
send(self(), {:http_request_called, %{url: url, body: body}})
58+
59+
{:ok, %{status_code: 200, body: "<ChangeMessageVisibilityBatchResponse />"}}
60+
end
5561
end
5662

5763
defmodule FakeHttpClientWithError do
@@ -181,6 +187,7 @@ defmodule BroadwaySQS.ExAwsClientTest do
181187
assert_received {:http_request_called, %{url: url}}
182188
assert url == "http://localhost:9324/"
183189
end
190+
184191
end
185192

186193
describe "ack/3" do
@@ -292,6 +299,31 @@ defmodule BroadwaySQS.ExAwsClientTest do
292299
assert_received {:http_request_called, %{url: url}}
293300
assert url == "http://localhost:9324/"
294301
end
302+
303+
test "request with :nack strategy", %{opts: base_opts} do
304+
{:ok, opts} = ExAwsClient.init(base_opts ++ [on_failure: {:nack, 10}])
305+
306+
:persistent_term.put(opts.ack_ref, %{
307+
queue_url: opts[:queue_url],
308+
config: opts[:config],
309+
on_success: opts[:on_success],
310+
on_failure: opts[:on_failure]
311+
})
312+
313+
ack_data = %{receipt: %{id: "1", receipt_handle: "abc"}}
314+
message = %Message{acknowledger: {ExAwsClient, opts.ack_ref, ack_data}, data: nil}
315+
316+
ExAwsClient.ack(opts.ack_ref, [], [message])
317+
318+
assert_received {:http_request_called, %{body: body}}
319+
320+
assert body ==
321+
"Action=ChangeMessageVisibilityBatch" <>
322+
"&ChangeMessageVisibilityBatchRequestEntry.1.Id=1" <>
323+
"&ChangeMessageVisibilityBatchRequestEntry.1.ReceiptHandle=abc" <>
324+
"&ChangeMessageVisibilityBatchRequestEntry.1.VisibilityTimeout=10" <>
325+
"&QueueUrl=my_queue"
326+
end
295327
end
296328

297329
defp fill_persistent_term(ack_ref, base_opts) do

0 commit comments

Comments
 (0)