Packages

phoenix_kit

1.7.13
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.process_sqs.ex
Raw

lib/mix/tasks/phoenix_kit.process_sqs.ex

defmodule Mix.Tasks.PhoenixKit.ProcessSqs do
@moduledoc """
Processes pending messages from AWS SQS queue for email events.
## Usage
# Process all messages in queue
mix phoenix_kit.process_sqs
# Process specific number of messages
mix phoenix_kit.process_sqs --count 10
# Show status without processing
mix phoenix_kit.process_sqs --status
## Options
* `--count` - Number of messages to process (default: all)
* `--status` - Show queue status without processing
* `--help` - Show this help
## Examples
# Process all pending messages
mix phoenix_kit.process_sqs
# Process up to 10 messages
mix phoenix_kit.process_sqs --count 10
# Check queue status
mix phoenix_kit.process_sqs --status
"""
use Mix.Task
alias PhoenixKit.Emails.SQSProcessor
alias PhoenixKit.Settings
@shortdoc "Process AWS SQS email event messages"
@switches [
count: :integer,
status: :boolean,
help: :boolean
]
@aliases [
c: :count,
s: :status,
h: :help
]
@impl Mix.Task
def run(args) do
{opts, _} = OptionParser.parse!(args, strict: @switches, aliases: @aliases)
cond do
opts[:help] ->
show_help()
opts[:status] ->
show_status()
true ->
count = opts[:count]
process_messages(count)
end
end
defp show_help do
Mix.shell().info(@moduledoc)
end
defp show_status do
Mix.Task.run("app.start")
queue_url = get_queue_url()
if queue_url do
Mix.shell().info("\n=== SQS Queue Status ===\n")
case get_queue_attributes(queue_url) do
{:ok, attrs} ->
available = attrs["ApproximateNumberOfMessages"] || "0"
in_flight = attrs["ApproximateNumberOfMessagesNotVisible"] || "0"
Mix.shell().info("Queue URL: #{queue_url}")
Mix.shell().info("Available messages: #{available}")
Mix.shell().info("In-flight messages: #{in_flight}")
Mix.shell().info("")
{:error, reason} ->
Mix.shell().error("Failed to get queue status: #{inspect(reason)}")
end
else
Mix.shell().error("AWS SQS configuration not found in Settings")
end
end
defp process_messages(count) do
Mix.Task.run("app.start")
queue_url = get_queue_url()
if queue_url do
Mix.shell().info("\n=== Processing SQS Messages ===\n")
max_count = count || 999_999
process_loop(queue_url, 0, max_count)
else
Mix.shell().error("AWS SQS configuration not found in Settings")
end
end
defp process_loop(_queue_url, processed, max_count) when processed >= max_count do
Mix.shell().info("\n✅ Processed #{processed} messages (limit reached)")
end
defp process_loop(queue_url, processed, max_count) do
case receive_message(queue_url) do
{:ok, nil} ->
if processed == 0 do
Mix.shell().info("No messages in queue")
else
Mix.shell().info("\n✅ Processed #{processed} messages total")
end
{:ok, message} ->
process_single_message(message, queue_url)
process_loop(queue_url, processed + 1, max_count)
{:error, reason} ->
Mix.shell().error("Failed to receive message: #{inspect(reason)}")
end
end
defp process_single_message(message, queue_url) do
body = message["Body"]
receipt_handle = message["ReceiptHandle"]
# Parse SNS message
case Jason.decode(body) do
{:ok, sns_body} ->
sns_message = Jason.decode!(sns_body["Message"])
event_type = sns_message["eventType"]
message_id = get_in(sns_message, ["mail", "messageId"])
Mix.shell().info("Processing: #{event_type} for #{message_id}")
# Process event
with {:ok, sns_data} <- SQSProcessor.parse_sns_message(%{"Body" => body}),
{:ok, result} <- SQSProcessor.process_email_event(sns_data) do
Mix.shell().info(" ✅ #{inspect(result)}")
# Delete message from queue
delete_message(queue_url, receipt_handle)
else
{:error, reason} ->
Mix.shell().error(" ❌ Failed: #{inspect(reason)}")
end
{:error, reason} ->
Mix.shell().error(" ❌ Failed to parse: #{inspect(reason)}")
end
end
defp get_queue_url do
region = Settings.get_setting("aws_region", "eu-north-1")
account_id = Settings.get_setting("aws_account_id")
queue_name = Settings.get_setting("aws_sqs_queue_name", "phoenixkit-email-queue")
if account_id do
"https://sqs.#{region}.amazonaws.com/#{account_id}/#{queue_name}"
else
nil
end
end
defp get_queue_attributes(queue_url) do
region = Settings.get_setting("aws_region", "eu-north-1")
case System.cmd(
"aws",
[
"sqs",
"get-queue-attributes",
"--queue-url",
queue_url,
"--attribute-names",
"ApproximateNumberOfMessages",
"ApproximateNumberOfMessagesNotVisible",
"--region",
region
],
stderr_to_stdout: true
) do
{output, 0} ->
case Jason.decode(output) do
{:ok, %{"Attributes" => attrs}} -> {:ok, attrs}
{:error, reason} -> {:error, reason}
end
{output, _code} ->
{:error, output}
end
end
defp receive_message(queue_url) do
region = Settings.get_setting("aws_region", "eu-north-1")
case System.cmd(
"aws",
[
"sqs",
"receive-message",
"--queue-url",
queue_url,
"--max-number-of-messages",
"1",
"--region",
region
],
stderr_to_stdout: true
) do
{output, 0} ->
case Jason.decode(output) do
{:ok, %{"Messages" => [message | _]}} -> {:ok, message}
{:ok, _} -> {:ok, nil}
{:error, reason} -> {:error, reason}
end
{output, _code} ->
{:error, output}
end
end
defp delete_message(queue_url, receipt_handle) do
region = Settings.get_setting("aws_region", "eu-north-1")
System.cmd(
"aws",
[
"sqs",
"delete-message",
"--queue-url",
queue_url,
"--receipt-handle",
receipt_handle,
"--region",
region
],
stderr_to_stdout: true
)
end
end