Packages

phoenix_kit

1.7.21
1.7.210 1.7.209 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 modules emails sqs_polling_job.ex
Raw

lib/modules/emails/sqs_polling_job.ex

defmodule PhoenixKit.Modules.Emails.SQSPollingJob do
@moduledoc """
Oban worker for polling AWS SQS queue for email events.
This worker replaces the GenServer-based SQSWorker with an Oban-based
approach that allows dynamic enabling/disabling without application restart.
## Architecture
```
AWS SES → SNS Topic → SQS Queue → SQSPollingJob (Oban) → SQSProcessor → Database
```
## Features
- **Dynamic Configuration**: Automatically responds to settings changes without restart
- **Oban Integration**: Uses Oban's job system for reliable background processing
- **Self-Scheduling**: Each job schedules the next polling cycle
- **Batch Processing**: Process up to 10 messages at a time
- **Error Handling**: Retry logic with Dead Letter Queue
- **Settings-Based Control**: Polling can be enabled/disabled via Settings
## Configuration
All settings are retrieved from PhoenixKit Settings:
- `sqs_polling_enabled` - enable/disable polling (checked before each cycle)
- `sqs_polling_interval_ms` - interval between polling cycles
- `sqs_max_messages_per_poll` - maximum messages per batch
- `sqs_visibility_timeout` - time for message processing
- `aws_sqs_queue_url` - SQS queue URL
- `aws_region` - AWS region
## Usage
# Enable polling (starts first job)
PhoenixKit.Modules.Emails.SQSPollingManager.enable_polling()
# Disable polling (stops scheduling new jobs)
PhoenixKit.Modules.Emails.SQSPollingManager.disable_polling()
# Trigger immediate polling
PhoenixKit.Modules.Emails.SQSPollingManager.poll_now()
# Check status
PhoenixKit.Modules.Emails.SQSPollingManager.status()
## Oban Queue Configuration
Add to your `config/config.exs`:
config :your_app, Oban,
repo: YourApp.Repo,
queues: [
sqs_polling: 1 # Only one concurrent polling job
]
## Implementation Notes
- Uses `unique: [period: 60]` to prevent duplicate jobs
- Schedules next job only if polling is enabled
- Uses existing SQSProcessor for event processing
- Compatible with existing SQSWorker API
"""
use Oban.Worker,
queue: :sqs_polling,
max_attempts: 3,
# Short unique period to prevent duplicate submissions while allowing self-scheduling
# 10 seconds is enough to prevent accidental double-clicks but allows 5s polling interval
unique: [period: 10, states: [:scheduled, :available, :executing]]
require Logger
import Ecto.Query
alias PhoenixKit.Modules.Emails
alias PhoenixKit.Modules.Emails.SQSProcessor
@default_long_poll_timeout 20
@impl Oban.Worker
def perform(%Oban.Job{}) do
# Check if polling is enabled before processing
if should_poll?() do
Logger.debug("SQS Polling Job: Starting polling cycle")
config = Emails.get_sqs_config()
case validate_configuration(config) do
:ok ->
result = perform_polling_cycle(config)
schedule_next_poll(config.polling_interval_ms)
result
{:error, reason} ->
Logger.error("SQS Polling Job: Invalid configuration - #{reason}")
{:error, reason}
end
else
Logger.debug("SQS Polling Job: Polling disabled, skipping cycle")
:ok
end
end
@doc """
Cancels all scheduled SQS polling jobs.
Called when polling is disabled to immediately clean up pending jobs.
## Returns
- `{:ok, count}` - Number of cancelled jobs
## Examples
iex> PhoenixKit.Modules.Emails.SQSPollingJob.cancel_scheduled()
{:ok, 2}
"""
@spec cancel_scheduled() :: {:ok, non_neg_integer()}
def cancel_scheduled do
worker_name = inspect(__MODULE__)
{count, _} =
Oban.Job
|> where([j], j.worker == ^worker_name)
|> where([j], j.state in ["available", "scheduled"])
|> get_repo().delete_all()
Logger.info("SQSPollingJob: Cancelled #{count} scheduled jobs")
{:ok, count}
end
defp get_repo do
PhoenixKit.RepoHelper.repo()
end
## --- Private Functions ---
# Check if polling should be performed
defp should_poll? do
Emails.enabled?() and
Emails.ses_events_enabled?() and
Emails.sqs_polling_enabled?()
end
# Validate SQS configuration
defp validate_configuration(config) do
cond do
is_nil(config.queue_url) or config.queue_url == "" ->
{:error, "SQS queue URL not configured"}
not is_integer(config.polling_interval_ms) or config.polling_interval_ms <= 0 ->
{:error, "Invalid polling interval"}
not is_integer(config.max_messages_per_poll) or
config.max_messages_per_poll <= 0 or
config.max_messages_per_poll > 10 ->
{:error, "Invalid max messages per poll (must be 1-10)"}
not is_integer(config.visibility_timeout) or config.visibility_timeout <= 0 ->
{:error, "Invalid visibility timeout"}
true ->
:ok
end
end
# Perform one polling cycle
defp perform_polling_cycle(config) do
_start_time = System.monotonic_time(:millisecond)
case receive_messages(config) do
{:ok, [_ | _] = messages} ->
Logger.info("SQS Polling Job: Received #{length(messages)} messages")
processing_start = System.monotonic_time(:millisecond)
processed_count = process_messages(messages, config)
processing_time = System.monotonic_time(:millisecond) - processing_start
Logger.info(
"SQS Polling Job: Processed #{processed_count}/#{length(messages)} messages in #{processing_time}ms"
)
{:ok,
%{processed: processed_count, total: length(messages), duration_ms: processing_time}}
{:ok, []} ->
Logger.debug("SQS Polling Job: No messages in queue")
{:ok, %{processed: 0, total: 0}}
{:error, reason} ->
Logger.error("SQS Polling Job: Failed to receive messages", %{
reason: inspect(reason),
queue_url: config.queue_url
})
{:error, reason}
end
end
# Receive messages from SQS queue
defp receive_messages(config) do
aws_config = build_aws_config(config)
request =
ExAws.SQS.receive_message(
config.queue_url,
max_number_of_messages: config.max_messages_per_poll,
wait_time_seconds: @default_long_poll_timeout,
visibility_timeout: config.visibility_timeout,
message_attribute_names: [:all],
attribute_names: [:all]
)
case ExAws.request(request, aws_config) do
{: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} ->
{:ok, []}
{:error, error} ->
Logger.error("SQS Polling Job: ExAws request failed", %{
error: inspect(error),
queue_url: config.queue_url
})
{:error, error}
end
end
# Process message list in parallel
defp process_messages(messages, config) do
aws_config = build_aws_config(config)
tasks =
Enum.map(messages, fn message ->
Task.async(fn ->
process_single_message(message, config.queue_url, aws_config)
end)
end)
results = Task.await_many(tasks, 30_000)
Enum.count(results, & &1)
end
# Process 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("SQS Polling Job: Failed to process message", %{
message_id: message_id,
reason: inspect(reason)
})
false
end
end
# Delete 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("SQS Polling Job: Failed to delete message", %{
error: inspect(error),
queue_url: queue_url
})
:ok
end
end
# Schedule next polling job
defp schedule_next_poll(interval_ms) do
if should_poll?() do
%{}
|> __MODULE__.new(schedule_in: div(interval_ms, 1000))
|> Oban.insert()
|> case do
{:ok, _job} ->
Logger.debug("SQS Polling Job: Next poll scheduled in #{interval_ms}ms")
:ok
{:error, reason} ->
Logger.error("SQS Polling Job: Failed to schedule next poll", %{
reason: inspect(reason)
})
:ok
end
else
Logger.debug("SQS Polling Job: Polling disabled, not scheduling next poll")
:ok
end
end
# Build AWS configuration
defp build_aws_config(config) do
if config.aws_access_key_id && config.aws_secret_access_key && config.aws_region do
[
access_key_id: config.aws_access_key_id,
secret_access_key: config.aws_secret_access_key,
region: config.aws_region
]
else
[]
end
end
end