Packages

phoenix_kit

1.7.21
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
phoenix_kit lib mix tasks phoenix_kit.email.process_dlq.ex
Raw

lib/mix/tasks/phoenix_kit.email.process_dlq.ex

defmodule Mix.Tasks.PhoenixKit.Email.ProcessDlq do
@moduledoc """
Process accumulated messages from AWS SQS Dead Letter Queue (DLQ).
This task retrieves all messages from the DLQ, processes them through
the SQS processor to update email statuses, and optionally deletes
successfully processed messages.
## Usage
mix phoenix_kit.email.process_dlq [--batch-size 10] [--delete-after] [--dry-run]
## Options
* `--batch-size` - Number of messages to process in each batch (default: 10)
* `--delete-after` - Delete successfully processed messages from DLQ (default: false)
* `--dry-run` - Show what would be processed without making changes (default: false)
## Examples
# Process all DLQ messages without deleting them
mix phoenix_kit.email.process_dlq
# Process in small batches and delete successful ones
mix phoenix_kit.email.process_dlq --batch-size 5 --delete-after
# See what would be processed (no changes)
mix phoenix_kit.email.process_dlq --dry-run
## Requirements
- Email system must be enabled
- AWS credentials must be configured
- DLQ URL must be set in settings
"""
use Mix.Task
require Logger
alias PhoenixKit.Modules.Emails
alias PhoenixKit.Modules.Emails.SQSProcessor
alias PhoenixKit.Settings
@shortdoc "Process accumulated DLQ messages"
@impl Mix.Task
def run(args) do
# Start the application to ensure repo and settings are available
Mix.Task.run("app.start")
{options, [], []} =
OptionParser.parse(args,
strict: [
batch_size: :integer,
delete_after: :boolean,
dry_run: :boolean
],
aliases: [
b: :batch_size,
d: :delete_after,
n: :dry_run
]
)
batch_size = Keyword.get(options, :batch_size, 10)
delete_after = Keyword.get(options, :delete_after, false)
dry_run = Keyword.get(options, :dry_run, false)
if not Emails.enabled?() do
Mix.shell().error("❌ Email system is not enabled")
exit(:shutdown)
end
dlq_url = Settings.get_setting("aws_sqs_dlq_url")
unless is_binary(dlq_url) and dlq_url != "" do
Mix.shell().error("❌ DLQ URL not configured")
exit(:shutdown)
end
Mix.shell().info("🔄 Processing DLQ messages...")
Mix.shell().info("📋 Configuration:")
Mix.shell().info(" â€ĸ DLQ URL: #{dlq_url}")
Mix.shell().info(" â€ĸ Batch size: #{batch_size}")
Mix.shell().info(" â€ĸ Delete after: #{delete_after}")
Mix.shell().info(" â€ĸ Dry run: #{dry_run}")
Mix.shell().info("")
total_processed = 0
total_successful = 0
total_errors = 0
{final_processed, final_successful, final_errors} =
process_dlq_batches(
dlq_url,
batch_size,
delete_after,
dry_run,
total_processed,
total_successful,
total_errors
)
Mix.shell().info("")
Mix.shell().info("✅ DLQ processing completed!")
Mix.shell().info("📊 Summary:")
Mix.shell().info(" â€ĸ Total messages processed: #{final_processed}")
Mix.shell().info(" â€ĸ Successful: #{final_successful}")
Mix.shell().info(" â€ĸ Errors: #{final_errors}")
if dry_run do
Mix.shell().info("â„šī¸ This was a dry run - no changes were made")
end
end
# Recursively process message batches from DLQ
defp process_dlq_batches(
dlq_url,
batch_size,
delete_after,
dry_run,
total_processed,
total_successful,
total_errors
) do
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
Mix.shell().info("đŸ“Ļ No more messages in DLQ")
{total_processed, total_successful, total_errors}
else
Mix.shell().info("đŸ“Ļ Processing batch of #{length(messages)} messages...")
{batch_successful, batch_errors, processed_receipts} =
process_message_batch(messages, dry_run)
new_processed = total_processed + length(messages)
new_successful = total_successful + batch_successful
new_errors = total_errors + batch_errors
Mix.shell().info("✅ Batch completed: #{batch_successful}/#{length(messages)} successful")
# Delete successfully processed messages if required
if delete_after and not dry_run and not Enum.empty?(processed_receipts) do
delete_processed_messages(dlq_url, processed_receipts)
end
# Continue processing next batch
process_dlq_batches(
dlq_url,
batch_size,
delete_after,
dry_run,
new_processed,
new_successful,
new_errors
)
end
rescue
error ->
Mix.shell().error("❌ Error processing DLQ batch: #{inspect(error)}")
{total_processed, total_successful, total_errors + 1}
end
# Process one message batch
defp process_message_batch(messages, dry_run) do
results =
Enum.map(messages, fn message ->
if dry_run do
case analyze_message(message) do
{:ok, info} ->
Mix.shell().info(" Would process: #{info.event_type} for #{info.message_id}")
{:ok, message["ReceiptHandle"]}
{:error, reason} ->
Mix.shell().info(" Would skip: #{reason}")
{:error, reason}
end
else
process_single_message(message)
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
# Analyze message without processing it (for dry-run)
defp analyze_message(message) do
case SQSProcessor.parse_sns_message(message) do
{:ok, event_data} ->
message_id = get_in(event_data, ["mail", "messageId"])
event_type = event_data["eventType"]
{:ok, %{message_id: message_id, event_type: event_type}}
{:error, reason} ->
{:error, "Invalid message format: #{reason}"}
end
end
# Process single message
defp process_single_message(message) do
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} ->
Mix.shell().info(" ✅ #{event_type} for #{message_id}: #{inspect(result)}")
{:ok, message["ReceiptHandle"]}
{:error, reason} ->
Mix.shell().info(" ❌ Failed #{event_type} for #{message_id}: #{reason}")
{:error, reason}
end
{:error, reason} ->
Mix.shell().info(" ❌ Failed to parse message: #{reason}")
{:error, reason}
end
end
# Delete successfully processed messages from DLQ
defp delete_processed_messages(dlq_url, receipt_handles) do
Mix.shell().info("đŸ—‘ī¸ Deleting #{length(receipt_handles)} processed messages from DLQ...")
Enum.each(receipt_handles, fn receipt_handle ->
try do
ExAws.SQS.delete_message(dlq_url, receipt_handle)
|> ExAws.request()
rescue
error ->
Mix.shell().error(" ❌ Failed to delete message: #{inspect(error)}")
end
end)
Mix.shell().info(" ✅ Deleted #{length(receipt_handles)} messages")
end
end