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
Current section
Files
lib/modules/emails/sqs_polling_manager.ex
defmodule PhoenixKit.Modules.Emails.SQSPollingManager do
@moduledoc """
Manager module for SQS polling via Oban jobs.
This module provides a unified API for managing SQS polling that can be
enabled/disabled dynamically without application restart.
## Features
- **Enable/Disable Polling**: Start or stop polling without restart
- **Manual Triggering**: Force immediate polling when needed
- **Status Monitoring**: Get current polling status and job information
- **Settings Integration**: Automatically uses PhoenixKit Settings
- **Interval Control**: Dynamically adjust polling frequency
## Architecture
Instead of using a GenServer, this manager uses Oban jobs for polling:
- Each job polls SQS once and schedules the next job
- Jobs check settings before executing (dynamic control)
- No need to restart GenServer when settings change
## Usage
# Enable polling
iex> PhoenixKit.Modules.Emails.SQSPollingManager.enable_polling()
{:ok, %Oban.Job{}}
# Disable polling
iex> PhoenixKit.Modules.Emails.SQSPollingManager.disable_polling()
:ok
# Check status
iex> PhoenixKit.Modules.Emails.SQSPollingManager.status()
%{
enabled: true,
interval_ms: 5000,
pending_jobs: 1,
last_run: ~U[2025-09-20 15:30:45Z],
queue_url: "https://sqs.eu-north-1.amazonaws.com/..."
}
# Trigger immediate poll
iex> PhoenixKit.Modules.Emails.SQSPollingManager.poll_now()
{:ok, %Oban.Job{}}
# Change polling interval
iex> PhoenixKit.Modules.Emails.SQSPollingManager.set_polling_interval(3000)
{:ok, %Setting{}}
## Integration
This manager works alongside the existing SQSWorker for backward compatibility.
The SQSWorker can delegate to this manager when needed.
"""
require Logger
alias PhoenixKit.Modules.Emails
alias PhoenixKit.Modules.Emails.SQSPollingJob
@doc """
Enables SQS polling by setting the configuration and starting the first job.
## Returns
- `{:ok, job}` - Successfully enabled and started first job
- `{:error, reason}` - Failed to enable polling
## Examples
iex> PhoenixKit.Modules.Emails.SQSPollingManager.enable_polling()
{:ok, %Oban.Job{id: 1, queue: "sqs_polling"}}
"""
def enable_polling do
Logger.info("SQS Polling Manager: Enabling polling")
with {:ok, _setting} <- Emails.set_sqs_polling(true),
{:ok, job} <- start_initial_job() do
Logger.info("SQS Polling Manager: Polling enabled and first job started")
{:ok, job}
else
{:error, reason} = error ->
Logger.error("SQS Polling Manager: Failed to enable polling", %{
reason: inspect(reason)
})
error
end
end
@doc """
Disables SQS polling by updating the configuration.
Note: Existing scheduled jobs will check this setting and skip execution.
## Returns
- `:ok` - Successfully disabled
## Examples
iex> PhoenixKit.Modules.Emails.SQSPollingManager.disable_polling()
:ok
"""
def disable_polling do
Logger.info("SQS Polling Manager: Disabling polling")
case Emails.set_sqs_polling(false) do
{:ok, _setting} ->
# Cancel any scheduled polling jobs
SQSPollingJob.cancel_scheduled()
Logger.info("SQS Polling Manager: Polling disabled")
:ok
{:error, reason} ->
Logger.error("SQS Polling Manager: Failed to disable polling", %{
reason: inspect(reason)
})
{:error, reason}
end
end
@doc """
Sets the polling interval in milliseconds.
The new interval will be used for subsequent job scheduling.
## Parameters
- `interval_ms` - Interval in milliseconds (minimum 1000ms)
## Returns
- `{:ok, setting}` - Successfully updated
- `{:error, reason}` - Failed to update
## Examples
iex> PhoenixKit.Modules.Emails.SQSPollingManager.set_polling_interval(3000)
{:ok, %Setting{}}
"""
def set_polling_interval(interval_ms) when is_integer(interval_ms) and interval_ms >= 1000 do
Logger.info("SQS Polling Manager: Setting polling interval to #{interval_ms}ms")
Emails.set_sqs_polling_interval(interval_ms)
end
def set_polling_interval(interval_ms) do
{:error, "Invalid interval: #{interval_ms}. Must be >= 1000ms"}
end
@doc """
Triggers an immediate polling job.
This creates a new job that will execute as soon as possible,
regardless of the normal polling schedule.
## Returns
- `{:ok, job}` - Successfully created immediate job
- `{:error, reason}` - Failed to create job
## Examples
iex> PhoenixKit.Modules.Emails.SQSPollingManager.poll_now()
{:ok, %Oban.Job{}}
"""
def poll_now do
Logger.info("SQS Polling Manager: Triggering immediate poll")
unless polling_enabled?() do
Logger.warning("SQS Polling Manager: Polling is disabled, but executing manual poll")
end
case start_immediate_job() do
{:ok, job} ->
Logger.info("SQS Polling Manager: Immediate poll job created", %{job_id: job.id})
{:ok, job}
{:error, reason} = error ->
Logger.error("SQS Polling Manager: Failed to create immediate poll job", %{
reason: inspect(reason)
})
error
end
end
@doc """
Returns the current status of SQS polling.
## Returns
A map with:
- `enabled` - Whether polling is enabled
- `interval_ms` - Current polling interval
- `pending_jobs` - Number of scheduled jobs
- `last_run` - Timestamp of last completed job (if any)
- `queue_url` - Configured SQS queue URL
## Examples
iex> PhoenixKit.Modules.Emails.SQSPollingManager.status()
%{
enabled: true,
interval_ms: 5000,
pending_jobs: 1,
last_run: ~U[2025-09-20 15:30:45Z],
queue_url: "https://sqs.eu-north-1.amazonaws.com/..."
}
"""
def status do
config = Emails.get_sqs_config()
pending_jobs = count_pending_jobs()
last_completed = get_last_completed_job()
%{
enabled: config.polling_enabled,
interval_ms: config.polling_interval_ms,
pending_jobs: pending_jobs,
last_run: last_completed && last_completed.completed_at,
queue_url: config.queue_url,
aws_region: config.aws_region,
max_messages_per_poll: config.max_messages_per_poll,
system_enabled: Emails.enabled?(),
ses_events_enabled: Emails.ses_events_enabled?()
}
end
## --- Private Functions ---
# Start the initial polling job
defp start_initial_job do
%{}
|> SQSPollingJob.new()
|> Oban.insert()
end
# Start an immediate polling job
defp start_immediate_job do
%{}
|> SQSPollingJob.new()
|> Oban.insert()
end
# Check if polling is currently enabled
defp polling_enabled? do
Emails.sqs_polling_enabled?()
end
# Count pending/scheduled SQS polling jobs
defp count_pending_jobs do
repo = PhoenixKit.RepoHelper.repo()
import Ecto.Query
from(j in Oban.Job,
where: j.worker == "PhoenixKit.Modules.Emails.SQSPollingJob",
where: j.state in ["available", "scheduled", "executing"],
select: count(j.id)
)
|> repo.one()
rescue
_ -> 0
end
# Get the last completed job
defp get_last_completed_job do
repo = PhoenixKit.RepoHelper.repo()
import Ecto.Query
from(j in Oban.Job,
where: j.worker == "PhoenixKit.Modules.Emails.SQSPollingJob",
where: j.state == "completed",
order_by: [desc: j.completed_at],
limit: 1
)
|> repo.one()
rescue
_ -> nil
end
end