Packages
phoenix_kit
1.7.74
1.7.208
1.7.207
1.7.206
1.7.205
1.7.204
1.7.203
1.7.202
1.7.201
1.7.200
1.7.199
1.7.198
1.7.197
1.7.196
1.7.194
1.7.193
1.7.192
1.7.191
1.7.190
1.7.189
1.7.187
1.7.186
1.7.185
1.7.184
1.7.183
1.7.182
1.7.181
1.7.180
1.7.179
1.7.178
1.7.177
1.7.176
1.7.175
1.7.174
1.7.173
1.7.172
1.7.171
1.7.170
1.7.169
1.7.168
1.7.167
1.7.166
1.7.165
1.7.164
1.7.162
1.7.161
1.7.160
1.7.159
1.7.157
1.7.156
1.7.155
1.7.154
1.7.153
1.7.152
1.7.151
1.7.150
1.7.149
1.7.146
1.7.145
1.7.144
1.7.143
1.7.138
1.7.133
1.7.132
1.7.131
1.7.130
1.7.128
1.7.126
1.7.125
1.7.121
1.7.120
1.7.119
1.7.118
1.7.117
1.7.116
1.7.115
1.7.114
1.7.113
1.7.112
1.7.111
1.7.110
1.7.109
1.7.108
1.7.107
1.7.106
1.7.105
1.7.104
1.7.103
1.7.102
1.7.101
1.7.100
1.7.99
1.7.98
1.7.97
1.7.96
1.7.95
1.7.94
1.7.93
1.7.92
1.7.91
1.7.90
1.7.89
1.7.88
1.7.87
1.7.86
1.7.85
1.7.84
1.7.83
1.7.82
1.7.81
1.7.80
1.7.79
1.7.78
1.7.77
1.7.76
1.7.75
1.7.74
1.7.71
1.7.70
1.7.69
1.7.66
1.7.65
1.7.64
1.7.63
1.7.62
1.7.61
1.7.59
1.7.58
1.7.57
1.7.56
1.7.55
1.7.54
1.7.53
1.7.52
1.7.51
1.7.49
1.7.44
1.7.43
1.7.42
1.7.41
1.7.39
1.7.38
1.7.37
1.7.36
1.7.34
1.7.33
1.7.31
1.7.30
1.7.29
1.7.28
1.7.27
1.7.26
1.7.25
1.7.24
1.7.23
1.7.22
1.7.21
1.7.20
1.7.19
1.7.18
1.7.17
1.7.16
1.7.15
1.7.14
1.7.13
1.7.12
1.7.11
1.7.10
1.7.9
1.7.8
1.7.7
1.7.6
1.7.5
1.7.4
1.7.3
1.7.2
1.7.1
1.7.0
1.6.20
1.6.19
1.6.18
1.6.17
1.6.16
1.6.15
1.6.14
1.6.13
1.6.12
1.6.11
1.6.10
1.6.9
1.6.8
1.6.7
1.6.6
1.6.5
1.6.4
1.6.3
1.5.2
1.5.1
1.5.0
1.4.9
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.2
1.3.1
1.3.0
1.2.10
1.2.9
1.2.8
1.2.7
1.2.5
1.2.4
1.2.2
1.2.1
1.2.0
1.1.0
1.0.0
A foundation for building Elixir Phoenix apps — SaaS, social networks, ERP systems, marketplaces, and more
Current section
Files
Jump to
Current section
Files
lib/modules/emails/sqs_worker.ex
defmodule PhoenixKit.Modules.Emails.SQSWorker do
@moduledoc """
SQS Worker for processing email events from AWS SQS Queue.
## ⚠️ DEPRECATION NOTICE
This GenServer-based worker is deprecated in favor of the Oban-based approach.
For new installations, use `PhoenixKit.Modules.Emails.SQSPollingManager` instead.
The GenServer approach has been replaced with Oban jobs to allow:
- Dynamic enabling/disabling without application restart
- Better job monitoring and failure tracking
- Automatic retries via Oban's built-in retry mechanism
- Integration with existing Oban infrastructure
## Migration Path
If you're currently using SQSWorker, you can migrate to the new approach:
# Old approach (GenServer)
PhoenixKit.Modules.Emails.SQSWorker.status()
PhoenixKit.Modules.Emails.SQSWorker.pause()
PhoenixKit.Modules.Emails.SQSWorker.resume()
# New approach (Oban-based)
PhoenixKit.Modules.Emails.SQSPollingManager.status()
PhoenixKit.Modules.Emails.SQSPollingManager.disable_polling()
PhoenixKit.Modules.Emails.SQSPollingManager.enable_polling()
## Backward Compatibility
This module maintains backward compatibility by delegating to SQSPollingManager
where appropriate. The GenServer will still work but is not recommended for
new installations.
## Architecture
```
AWS SES → SNS Topic → SQS Queue → SQS Worker → Database
```
## Configuration
All settings are retrieved from PhoenixKit Settings and checked dynamically:
- `email_ses_events` - master switch for AWS SES events processing
- `sqs_polling_enabled` - enable/disable polling
- `sqs_polling_interval_ms` - interval between polling cycles
- `sqs_max_messages_per_poll` - maximum messages per batch
- `aws_sqs_queue_url` - SQS queue URL
- `aws_region` - AWS region
## Usage (Legacy)
# In supervision tree (deprecated)
{PhoenixKit.Modules.Emails.SQSWorker, []}
# Worker management (delegates to new API)
PhoenixKit.Modules.Emails.SQSWorker.status()
PhoenixKit.Modules.Emails.SQSWorker.process_now()
PhoenixKit.Modules.Emails.SQSWorker.pause()
PhoenixKit.Modules.Emails.SQSWorker.resume()
"""
use GenServer
require Logger
alias PhoenixKit.Modules.Emails
alias PhoenixKit.Modules.Emails.SQSPollingManager
alias PhoenixKit.Modules.Emails.SQSProcessor
alias PhoenixKit.Utils.Date, as: UtilsDate
# 20 seconds
@default_long_poll_timeout 20
## --- Client API ---
@doc """
Starts the SQS Worker process.
## Options
- `:name` - process name (defaults to `__MODULE__`)
## Examples
{:ok, pid} = PhoenixKit.Modules.Emails.SQSWorker.start_link()
"""
def start_link(opts \\ []) do
name = Keyword.get(opts, :name, __MODULE__)
GenServer.start_link(__MODULE__, opts, name: name)
end
@doc """
Returns the current status of the worker process.
**Note**: This function now delegates to SQSPollingManager for consistency.
## Examples
iex> PhoenixKit.Modules.Emails.SQSWorker.status()
%{
polling_enabled: true,
messages_processed: 150,
errors_count: 2,
last_poll: ~U[2025-09-20 15:30:45.123456Z],
queue_url: "https://sqs.eu-north-1.amazonaws.com/123456789012/phoenixkit-email-queue",
average_processing_time_ms: 45.2
}
"""
def status(worker \\ __MODULE__) do
# Try new API first, fallback to GenServer if needed
case get_status_from_manager() do
{:ok, status} -> status
:error -> get_status_from_genserver(worker)
end
end
defp get_status_from_manager do
{:ok, SQSPollingManager.status()}
catch
_, _ -> :error
end
defp get_status_from_genserver(worker) do
GenServer.call(worker, :status)
catch
_, _ ->
%{
error: "Worker not responding",
message: "Consider using PhoenixKit.Modules.Emails.SQSPollingManager.status() instead"
}
end
@doc """
Forces a polling cycle to start immediately.
**Note**: This function now delegates to SQSPollingManager.
## Examples
iex> PhoenixKit.Modules.Emails.SQSWorker.process_now()
:ok
"""
def process_now(_worker \\ __MODULE__) do
Logger.info(
"SQSWorker.process_now/1 is deprecated - delegating to SQSPollingManager.poll_now/0"
)
case SQSPollingManager.poll_now() do
{:ok, _job} -> :ok
{:error, _reason} -> :ok
end
end
@doc """
Pauses polling (temporarily).
**Note**: This function now delegates to SQSPollingManager.
## Examples
iex> PhoenixKit.Modules.Emails.SQSWorker.pause()
:ok
"""
def pause(_worker \\ __MODULE__) do
Logger.info(
"SQSWorker.pause/1 is deprecated - delegating to SQSPollingManager.disable_polling/0"
)
case SQSPollingManager.disable_polling() do
:ok -> :ok
{:error, _reason} -> :ok
end
end
@doc """
Resumes polling after pause.
**Note**: This function now delegates to SQSPollingManager.
## Examples
iex> PhoenixKit.Modules.Emails.SQSWorker.resume()
:ok
"""
def resume(_worker \\ __MODULE__) do
Logger.info(
"SQSWorker.resume/1 is deprecated - delegating to SQSPollingManager.enable_polling/0"
)
case SQSPollingManager.enable_polling() do
{:ok, _job} -> :ok
{:error, _reason} -> :ok
end
end
@doc """
Processes all messages from DLQ (Dead Letter Queue).
This function retrieves all messages from DLQ, processes them through
SQSProcessor, and optionally deletes successfully processed messages.
## Parameters
- `opts` - Processing options:
- `:batch_size` - Batch size (default 10)
- `:delete_after` - Delete successfully processed messages (default false)
- `:max_batches` - Maximum number of batches (default 100)
## Returns
- `{:ok, result}` - Successful processing with results
- `{:error, reason}` - Processing error
## Examples
iex> PhoenixKit.Modules.Emails.SQSWorker.process_dlq_messages()
{:ok, %{total_processed: 15, successful: 12, errors: 3}}
iex> PhoenixKit.Modules.Emails.SQSWorker.process_dlq_messages(delete_after: true)
{:ok, %{total_processed: 8, successful: 8, errors: 0, deleted: 8}}
"""
def process_dlq_messages(opts \\ []) do
GenServer.call(__MODULE__, {:process_dlq, opts}, 30_000)
end
@doc """
Deletes processed messages from DLQ.
## Parameters
- `receipt_handles` - List of receipt handles to delete
## Returns
- `{:ok, deleted_count}` - Number of deleted messages
- `{:error, reason}` - Deletion error
## Examples
iex> PhoenixKit.Modules.Emails.SQSWorker.delete_dlq_messages(["receipt1", "receipt2"])
{:ok, 2}
"""
def delete_dlq_messages(receipt_handles) when is_list(receipt_handles) do
GenServer.call(__MODULE__, {:delete_dlq_messages, receipt_handles}, 10_000)
end
## --- Server Callbacks ---
@doc false
def init(_opts) do
# Get configuration at startup
config = Emails.get_sqs_config()
state = %{
queue_url: config.queue_url,
polling_enabled: config.polling_enabled,
polling_interval_ms: config.polling_interval_ms,
max_messages_per_poll: config.max_messages_per_poll,
visibility_timeout: config.visibility_timeout,
paused: false,
# AWS configuration
aws_config: build_aws_config(config),
# Metrics
messages_processed: 0,
errors_count: 0,
last_poll: nil,
total_processing_time_ms: 0,
# Internal state
poll_timer_ref: nil
}
# Check configuration validity
case validate_configuration(state) do
:ok ->
if state.polling_enabled and not state.paused do
# Start immediately
{:ok, schedule_next_poll(state, 0)}
else
Logger.info("SQS Worker started but polling is disabled")
{:ok, state}
end
{:error, reason} ->
Logger.error("SQS Worker failed to start: #{reason}")
{:ok, %{state | polling_enabled: false}}
end
end
@doc false
def handle_info(:poll_sqs, state) do
# Check current polling status from database (dynamic check)
ses_events_enabled = Emails.ses_events_enabled?()
current_polling_enabled = Emails.sqs_polling_enabled?()
cond do
# AWS SES events processing is disabled - stop polling
not ses_events_enabled ->
if state.polling_enabled do
Logger.info("SQS Worker: AWS SES events disabled via settings, stopping polling cycle")
end
# Check status again after 30 seconds to detect if it was re-enabled
new_state =
%{state | polling_enabled: false}
|> schedule_next_poll(30_000, true)
{:noreply, new_state}
# Polling is enabled and not paused - perform normal polling cycle
ses_events_enabled and current_polling_enabled and not state.paused ->
# Log if polling was just re-enabled
if not state.polling_enabled do
Logger.info("SQS Worker: Polling enabled via settings, resuming polling cycle")
end
new_state =
%{state | polling_enabled: true}
|> perform_polling_cycle()
|> schedule_next_poll(state.polling_interval_ms)
{:noreply, new_state}
# Polling is disabled - log once and schedule status check
not current_polling_enabled ->
if state.polling_enabled do
Logger.info("SQS Worker: Polling disabled via settings, stopping polling cycle")
end
# Check status again after 30 seconds to detect if polling was re-enabled
new_state =
%{state | polling_enabled: false}
|> schedule_next_poll(30_000, true)
{:noreply, new_state}
# Worker is paused - keep checking
state.paused ->
new_state = schedule_next_poll(state, state.polling_interval_ms)
{:noreply, new_state}
end
end
@doc false
def handle_cast(:process_now, state) do
Logger.info("SQS Worker: Processing triggered manually")
new_state = perform_polling_cycle(state)
{:noreply, new_state}
end
def handle_cast(:pause, state) do
Logger.info("SQS Worker: Paused")
new_state = cancel_timer(%{state | paused: true})
{:noreply, new_state}
end
def handle_cast(:resume, state) do
Logger.info("SQS Worker: Resumed")
new_state =
%{state | paused: false}
# Resume immediately
|> schedule_next_poll(0)
{:noreply, new_state}
end
@doc false
def handle_call(:status, _from, state) do
average_processing_time =
if state.messages_processed > 0 do
state.total_processing_time_ms / state.messages_processed
else
0.0
end
status = %{
polling_enabled: state.polling_enabled,
paused: state.paused,
messages_processed: state.messages_processed,
errors_count: state.errors_count,
last_poll: state.last_poll,
queue_url: state.queue_url,
average_processing_time_ms: Float.round(average_processing_time, 2)
}
{:reply, status, state}
end
@doc false
def handle_call({:process_dlq, opts}, _from, state) do
batch_size = Keyword.get(opts, :batch_size, 10)
delete_after = Keyword.get(opts, :delete_after, false)
max_batches = Keyword.get(opts, :max_batches, 100)
dlq_url = PhoenixKit.Settings.get_setting("aws_sqs_dlq_url")
if dlq_url do
result = process_dlq_batches_sync(dlq_url, batch_size, delete_after, max_batches)
{:reply, {:ok, result}, state}
else
{:reply, {:error, :dlq_url_not_configured}, state}
end
rescue
error ->
Logger.error("DLQ processing failed: #{inspect(error)}")
{:reply, {:error, error}, state}
end
@doc false
def handle_call({:delete_dlq_messages, receipt_handles}, _from, state) do
dlq_url = PhoenixKit.Settings.get_setting("aws_sqs_dlq_url")
if dlq_url do
deleted_count = delete_messages_from_dlq(dlq_url, receipt_handles)
{:reply, {:ok, deleted_count}, state}
else
{:reply, {:error, :dlq_url_not_configured}, state}
end
rescue
error ->
Logger.error("DLQ message deletion failed: #{inspect(error)}")
{:reply, {:error, error}, state}
end
@doc false
def terminate(reason, state) do
Logger.info("SQS Worker shutting down", %{reason: inspect(reason)})
cancel_timer(state)
:ok
end
## --- Private Helper Functions ---
# Performs one SQS polling cycle
defp perform_polling_cycle(state) do
_start_time = System.monotonic_time(:millisecond)
case receive_messages(state) do
{:ok, [_ | _] = messages} ->
Logger.info("SQS Worker: Received #{length(messages)} messages")
processing_start = System.monotonic_time(:millisecond)
processed_count = process_messages(messages, state)
processing_time = System.monotonic_time(:millisecond) - processing_start
%{
state
| messages_processed: state.messages_processed + processed_count,
total_processing_time_ms: state.total_processing_time_ms + processing_time,
last_poll: UtilsDate.utc_now()
}
{:ok, []} ->
%{state | last_poll: UtilsDate.utc_now()}
{:error, reason} ->
Logger.error("SQS Worker: Failed to receive messages", %{
reason: inspect(reason),
queue_url: state.queue_url
})
%{state | errors_count: state.errors_count + 1, last_poll: UtilsDate.utc_now()}
end
end
# Retrieves messages from SQS queue
defp receive_messages(state) do
case state.queue_url do
nil ->
{:error, :queue_url_not_configured}
queue_url when is_binary(queue_url) ->
request =
ExAws.SQS.receive_message(
queue_url,
max_number_of_messages: state.max_messages_per_poll,
wait_time_seconds: @default_long_poll_timeout,
visibility_timeout: state.visibility_timeout,
message_attribute_names: [:all],
attribute_names: [:all]
)
case ExAws.request(request, state.aws_config) do
# Handle all possible response formats from ExAws
{:ok, %{"Messages" => messages}} when is_list(messages) ->
{:ok, messages}
{:ok, %{"messages" => messages}} when is_list(messages) ->
{:ok, messages}
{:ok, %{body: %{messages: messages}}} when is_list(messages) ->
{:ok, messages}
{:ok, %{body: %{"Messages" => messages}}} when is_list(messages) ->
{:ok, messages}
{:ok, %{body: %{"messages" => messages}}} when is_list(messages) ->
{:ok, messages}
{:ok, _response} ->
# No messages (any other successful response)
{:ok, []}
{:error, error} ->
error_info = safe_extract_error_info(error)
Logger.error("SQS Worker: ExAws request failed", %{
error: inspect(error),
error_type: error_info.type,
error_details: error_info.details,
queue_url: queue_url
})
{:error, error}
end
end
end
# Processes message list in parallel
defp process_messages(messages, state) do
# Create tasks for parallel processing
tasks =
Enum.map(messages, fn message ->
Task.async(fn ->
process_single_message(message, state.queue_url, state.aws_config)
end)
end)
# Wait for all tasks to complete
results = Task.await_many(tasks, 30_000)
# Count successfully processed
Enum.count(results, & &1)
end
# Processes a single message
defp process_single_message(message, queue_url, aws_config) do
message_id = message["MessageId"]
receipt_handle = message["ReceiptHandle"]
with {:ok, event_data} <- SQSProcessor.parse_sns_message(message),
{:ok, _result} <- SQSProcessor.process_email_event(event_data),
:ok <- delete_message(queue_url, receipt_handle, aws_config) do
true
else
{:error, reason} ->
Logger.error("Failed to process SQS message", %{
message_id: message_id,
reason: inspect(reason)
})
false
end
end
# Deletes processed message from queue
defp delete_message(queue_url, receipt_handle, aws_config) do
ExAws.SQS.delete_message(queue_url, receipt_handle)
|> ExAws.request(aws_config)
|> case do
{:ok, _} ->
:ok
{:error, error} ->
Logger.error("Failed to delete SQS message", %{
error: inspect(error),
error_type: if(is_map(error), do: error.__struct__, else: :unknown),
queue_url: queue_url,
receipt_handle: String.slice(receipt_handle, 0, 50) <> "...",
has_aws_config: not Enum.empty?(aws_config)
})
# Non-critical error, message will return to queue
:ok
end
end
# Schedules the next polling cycle
defp schedule_next_poll(state, delay_ms, force \\ false) do
new_state = cancel_timer(state)
# Schedule next cycle if:
# - polling is enabled and not paused, OR
# - force is true (for status checks when polling is disabled)
if (new_state.polling_enabled and not new_state.paused) or force do
timer_ref = Process.send_after(self(), :poll_sqs, delay_ms)
%{new_state | poll_timer_ref: timer_ref}
else
new_state
end
end
# Cancels the current timer
defp cancel_timer(state) do
if state.poll_timer_ref do
Process.cancel_timer(state.poll_timer_ref)
end
%{state | poll_timer_ref: nil}
end
# Validates configuration correctness
defp validate_configuration(state) do
cond do
is_nil(state.queue_url) or state.queue_url == "" ->
{:error, "SQS queue URL not configured"}
Enum.empty?(state.aws_config) ->
Logger.warning(
"AWS credentials not configured - will use default credential provider chain"
)
validate_other_configuration(state)
not is_list(state.aws_config) ->
{:error, "Invalid AWS configuration format"}
true ->
validate_other_configuration(state)
end
end
# Validates other configuration (non-AWS)
defp validate_other_configuration(state) do
cond do
not is_integer(state.polling_interval_ms) or state.polling_interval_ms <= 0 ->
{:error, "Invalid polling interval"}
not is_integer(state.max_messages_per_poll) or
state.max_messages_per_poll <= 0 or
state.max_messages_per_poll > 10 ->
{:error, "Invalid max messages per poll (must be 1-10)"}
not is_integer(state.visibility_timeout) or state.visibility_timeout <= 0 ->
{:error, "Invalid visibility timeout"}
true ->
:ok
end
end
# Builds AWS configuration from settings
defp build_aws_config(config) do
if is_binary(config.aws_access_key_id) and config.aws_access_key_id != "" and
is_binary(config.aws_secret_access_key) and config.aws_secret_access_key != "" and
is_binary(config.aws_region) and config.aws_region != "" do
[
access_key_id: String.trim(config.aws_access_key_id),
secret_access_key: String.trim(config.aws_secret_access_key),
region: String.trim(config.aws_region)
]
else
# If AWS credentials are not configured, use empty list
# ExAws will use default credential provider chain
[]
end
end
# Safely extracts error information
defp safe_extract_error_info(error) do
if is_map(error) and Map.has_key?(error, :__struct__) do
%{
type: error.__struct__,
details: Map.from_struct(error)
}
else
%{
type: :string_error,
details: %{message: to_string(error)}
}
end
rescue
_ ->
%{
type: :unknown_error,
details: %{message: "Failed to extract error details"}
}
end
# Synchronous processing of message batches from DLQ
defp process_dlq_batches_sync(dlq_url, batch_size, delete_after, max_batches) do
Logger.info("Processing DLQ messages", %{
dlq_url: dlq_url,
batch_size: batch_size,
delete_after: delete_after,
max_batches: max_batches
})
process_dlq_batches_recursive(dlq_url, batch_size, delete_after, max_batches, 0, %{
total_processed: 0,
successful: 0,
errors: 0,
deleted: 0
})
end
# Recursive DLQ batch processing
defp process_dlq_batches_recursive(
dlq_url,
batch_size,
delete_after,
max_batches,
current_batch,
stats
) do
if current_batch >= max_batches do
Logger.warning("Reached max batches limit", %{
current_batch: current_batch,
max_batches: max_batches
})
stats
else
messages =
ExAws.SQS.receive_message(dlq_url,
max_number_of_messages: batch_size,
wait_time_seconds: 1
)
|> ExAws.request()
|> case do
{:ok, %{body: %{messages: messages}}} -> messages
_ -> []
end
if Enum.empty?(messages) do
Logger.info("No more messages in DLQ", %{batches_processed: current_batch})
stats
else
Logger.info("Processing DLQ batch #{current_batch + 1}: #{length(messages)} messages")
{successful_count, error_count, processed_receipts} = process_dlq_message_batch(messages)
new_stats = %{
total_processed: stats.total_processed + length(messages),
successful: stats.successful + successful_count,
errors: stats.errors + error_count,
deleted: stats.deleted
}
# Delete successfully processed messages if required
final_stats =
if delete_after and not Enum.empty?(processed_receipts) do
deleted_count = delete_messages_from_dlq(dlq_url, processed_receipts)
%{new_stats | deleted: new_stats.deleted + deleted_count}
else
new_stats
end
# Continue processing next batch
process_dlq_batches_recursive(
dlq_url,
batch_size,
delete_after,
max_batches,
current_batch + 1,
final_stats
)
end
end
end
# Processing one batch of messages from DLQ
defp process_dlq_message_batch(messages) do
results =
Enum.map(messages, fn message ->
case SQSProcessor.parse_sns_message(message) do
{:ok, event_data} ->
message_id = get_in(event_data, ["mail", "messageId"])
event_type = event_data["eventType"]
case SQSProcessor.process_email_event(event_data) do
{:ok, _result} ->
{:ok, message["ReceiptHandle"]}
{:error, reason} ->
Logger.warning(
"Failed to process #{event_type} for #{message_id}: #{inspect(reason)}"
)
{:error, reason}
end
{:error, reason} ->
Logger.warning("Failed to parse DLQ message: #{inspect(reason)}")
{:error, reason}
end
end)
successful_results =
Enum.filter(results, fn
{:ok, _} -> true
_ -> false
end)
error_results =
Enum.filter(results, fn
{:error, _} -> true
_ -> false
end)
processed_receipts = Enum.map(successful_results, fn {:ok, receipt} -> receipt end)
{length(successful_results), length(error_results), processed_receipts}
end
# Deleting messages from DLQ
defp delete_messages_from_dlq(dlq_url, receipt_handles) do
Logger.info("Deleting #{length(receipt_handles)} messages from DLQ")
successful_deletes =
Enum.count(receipt_handles, fn receipt_handle ->
try do
ExAws.SQS.delete_message(dlq_url, receipt_handle)
|> ExAws.request()
|> case do
{:ok, _} ->
true
{:error, reason} ->
Logger.error("Failed to delete DLQ message: #{inspect(reason)}")
false
end
rescue
error ->
Logger.error("Exception while deleting DLQ message: #{inspect(error)}")
false
end
end)
Logger.info(
"Successfully deleted #{successful_deletes}/#{length(receipt_handles)} messages from DLQ"
)
successful_deletes
end
end