fix oban pro worker warnings
This commit is contained in:
parent
78c013e8b7
commit
58daf45b60
55 changed files with 129 additions and 142 deletions
|
|
@ -19,9 +19,9 @@ defmodule Towerops.Snmp.NeighborCleanupWorker do
|
|||
# Consider neighbors stale if not seen in 24 hours
|
||||
@stale_threshold_hours 24
|
||||
|
||||
@impl Oban.Worker
|
||||
@spec perform(Oban.Job.t()) :: :ok
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
@spec process(Oban.Job.t()) :: :ok
|
||||
def process(%Oban.Job{}) do
|
||||
cleanup_stale_records()
|
||||
cleanup_stale_topology_links()
|
||||
:ok
|
||||
|
|
|
|||
|
|
@ -29,8 +29,8 @@ defmodule Towerops.Workers.AgentLatencyEvaluator do
|
|||
# Look back 48 hours to cover multiple 8-hour probe cycles
|
||||
@hours_ago 48
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
start_time = System.monotonic_time(:millisecond)
|
||||
|
||||
candidates =
|
||||
|
|
|
|||
|
|
@ -10,8 +10,8 @@ defmodule Towerops.Workers.AgentReleaseWebhookWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
case Agents.broadcast_mass_update() do
|
||||
{:ok, result} ->
|
||||
Logger.info("Agent mass update: #{result.notified} notified, #{result.skipped} skipped")
|
||||
|
|
|
|||
|
|
@ -30,8 +30,8 @@ defmodule Towerops.Workers.AiNetworkInsightWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
org_ids =
|
||||
Organization
|
||||
|> select([o], o.id)
|
||||
|
|
|
|||
|
|
@ -18,8 +18,8 @@ defmodule Towerops.Workers.AlertDigestWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"user_id" => "__cron__"}}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: %{"user_id" => "__cron__"}}) do
|
||||
# Cron mode: find all pending digests and enqueue delivery
|
||||
|
||||
pending =
|
||||
|
|
@ -48,7 +48,7 @@ defmodule Towerops.Workers.AlertDigestWorker do
|
|||
:ok
|
||||
end
|
||||
|
||||
def perform(%Oban.Job{args: %{"user_id" => user_id}}) when user_id != "__cron__" do
|
||||
def process(%Oban.Job{args: %{"user_id" => user_id}}) when user_id != "__cron__" do
|
||||
case NotificationRateLimiter.get_pending_digest(user_id) do
|
||||
nil ->
|
||||
:ok
|
||||
|
|
|
|||
|
|
@ -19,8 +19,8 @@ defmodule Towerops.Workers.AlertNotificationWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"action" => "trigger", "alert_id" => alert_id}}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: %{"action" => "trigger", "alert_id" => alert_id}}) do
|
||||
with {:ok, alert} <- fetch_alert(alert_id),
|
||||
{:ok, device} <- fetch_device(alert.device_id) do
|
||||
handle_trigger_notification(alert, device, alert_id)
|
||||
|
|
@ -35,7 +35,7 @@ defmodule Towerops.Workers.AlertNotificationWorker do
|
|||
end
|
||||
end
|
||||
|
||||
def perform(%Oban.Job{args: %{"action" => "acknowledge", "alert_id" => alert_id}}) do
|
||||
def process(%Oban.Job{args: %{"action" => "acknowledge", "alert_id" => alert_id}}) do
|
||||
case fetch_alert(alert_id) do
|
||||
{:ok, alert} ->
|
||||
routing = alert_routing_for_alert(alert)
|
||||
|
|
@ -54,7 +54,7 @@ defmodule Towerops.Workers.AlertNotificationWorker do
|
|||
end
|
||||
end
|
||||
|
||||
def perform(%Oban.Job{args: %{"action" => "resolve", "alert_id" => alert_id}}) do
|
||||
def process(%Oban.Job{args: %{"action" => "resolve", "alert_id" => alert_id}}) do
|
||||
case fetch_alert(alert_id) do
|
||||
{:ok, alert} ->
|
||||
routing = alert_routing_for_alert(alert)
|
||||
|
|
|
|||
|
|
@ -12,9 +12,9 @@ defmodule Towerops.Workers.BackupSummaryWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
@spec perform(Oban.Job.t()) :: :ok
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
@spec process(Oban.Job.t()) :: :ok
|
||||
def process(%Oban.Job{}) do
|
||||
yesterday = DateTime.add(DateTime.utc_now(), -24 * 60 * 60, :second)
|
||||
|
||||
summary = BackupRequests.summary_since(yesterday)
|
||||
|
|
|
|||
|
|
@ -15,9 +15,9 @@ defmodule Towerops.Workers.BackupTimeoutWorker do
|
|||
|
||||
@timeout_minutes 5
|
||||
|
||||
@impl Oban.Worker
|
||||
@spec perform(Oban.Job.t()) :: :ok
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
@spec process(Oban.Job.t()) :: :ok
|
||||
def process(%Oban.Job{}) do
|
||||
cutoff = DateTime.add(DateTime.utc_now(), -@timeout_minutes * 60, :second)
|
||||
timeout_count = BackupRequests.mark_timed_out_requests(cutoff)
|
||||
|
||||
|
|
|
|||
|
|
@ -12,8 +12,8 @@ defmodule Towerops.Workers.BillingSyncWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(_job) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(_job) do
|
||||
Logger.info("Starting daily billing sync")
|
||||
|
||||
# Get all organizations with active paid subscriptions
|
||||
|
|
|
|||
|
|
@ -26,8 +26,8 @@ defmodule Towerops.Workers.CapacityInsightWorker do
|
|||
@warning_threshold 75
|
||||
@resolve_threshold 70
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
# Use streaming to avoid loading all org IDs into memory
|
||||
Repo.transaction(fn ->
|
||||
Organization
|
||||
|
|
|
|||
|
|
@ -50,8 +50,8 @@ defmodule Towerops.Workers.CheckExecutorWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"check_id" => check_id}}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: %{"check_id" => check_id}}) do
|
||||
case Monitoring.get_check(check_id) do
|
||||
nil ->
|
||||
Logger.debug("Check #{check_id} deleted, skipping execution")
|
||||
|
|
|
|||
|
|
@ -24,8 +24,8 @@ defmodule Towerops.Workers.CheckWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"check_id" => check_id}}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: %{"check_id" => check_id}}) do
|
||||
case Monitoring.get_check(check_id) do
|
||||
nil ->
|
||||
Logger.debug("Check #{check_id} no longer exists, skipping and not rescheduling")
|
||||
|
|
|
|||
|
|
@ -15,8 +15,8 @@ defmodule Towerops.Workers.CloudLatencyProbeWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
cloud_pollers = Agents.list_online_cloud_pollers()
|
||||
devices = Agents.list_cloud_polled_devices()
|
||||
|
||||
|
|
|
|||
|
|
@ -12,8 +12,8 @@ defmodule Towerops.Workers.CloudflareBanWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"ip_address" => ip_address}}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: %{"ip_address" => ip_address}}) do
|
||||
case CloudflareClient.block_ip(ip_address) do
|
||||
{:ok, _response} ->
|
||||
Logger.info("Successfully pushed permanent ban to Cloudflare: #{ip_address}")
|
||||
|
|
|
|||
|
|
@ -15,8 +15,8 @@ defmodule Towerops.Workers.CnMaestroSyncWorker do
|
|||
|
||||
@window_seconds 300
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: args}) when args == %{} do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: args}) when args == %{} do
|
||||
integrations = Integrations.list_enabled_integrations("cn_maestro")
|
||||
|
||||
now = DateTime.utc_now()
|
||||
|
|
@ -48,7 +48,7 @@ defmodule Towerops.Workers.CnMaestroSyncWorker do
|
|||
:ok
|
||||
end
|
||||
|
||||
def perform(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
def process(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
case Integrations.get_integration_by_id(id) do
|
||||
{:ok, integration} ->
|
||||
sync_integration(integration)
|
||||
|
|
|
|||
|
|
@ -56,8 +56,8 @@ defmodule Towerops.Workers.CoverageWorker do
|
|||
@sm_height_tiers_m [1.83, 3.05, 4.57, 6.10, 9.14, 12.19]
|
||||
@default_tier_m 3.05
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"coverage_id" => coverage_id, "organization_id" => organization_id}}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: %{"coverage_id" => coverage_id, "organization_id" => organization_id}}) do
|
||||
case Repo.get(Coverage, coverage_id) do
|
||||
nil ->
|
||||
Logger.warning("CoverageWorker: coverage #{coverage_id} not found, skipping")
|
||||
|
|
|
|||
|
|
@ -23,8 +23,8 @@ defmodule Towerops.Workers.DataRetentionWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(_job) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(_job) do
|
||||
retention_groups = fetch_retention_groups()
|
||||
|
||||
results =
|
||||
|
|
|
|||
|
|
@ -15,8 +15,8 @@ defmodule Towerops.Workers.DeviceHealthInsightWorker do
|
|||
|
||||
@poll_gap_hours 24
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
devices = list_stale_monitored_devices()
|
||||
|
||||
Enum.each(devices, &generate_poll_gap_insight/1)
|
||||
|
|
|
|||
|
|
@ -37,9 +37,9 @@ defmodule Towerops.Workers.DeviceMonitorWorker do
|
|||
|
||||
@monitor_interval 30
|
||||
|
||||
@impl Oban.Worker
|
||||
@spec perform(Oban.Job.t()) :: :ok
|
||||
def perform(%Oban.Job{args: %{"device_id" => device_id}}) do
|
||||
@impl Oban.Pro.Worker
|
||||
@spec process(Oban.Job.t()) :: :ok
|
||||
def process(%Oban.Job{args: %{"device_id" => device_id}}) do
|
||||
case Devices.get_device(device_id) do
|
||||
nil ->
|
||||
Logger.debug("Device #{device_id} no longer exists, skipping monitor and not rescheduling")
|
||||
|
|
|
|||
|
|
@ -46,9 +46,9 @@ defmodule Towerops.Workers.DevicePollerWorker do
|
|||
|
||||
@default_poll_interval 60
|
||||
|
||||
@impl Oban.Worker
|
||||
@spec perform(Oban.Job.t()) :: :ok
|
||||
def perform(%Oban.Job{args: %{"device_id" => device_id}} = job) do
|
||||
@impl Oban.Pro.Worker
|
||||
@spec process(Oban.Job.t()) :: :ok
|
||||
def process(%Oban.Job{args: %{"device_id" => device_id}} = job) do
|
||||
_ = Events.broadcast_job_event(job, :started)
|
||||
start_time = System.monotonic_time(:second)
|
||||
|
||||
|
|
|
|||
|
|
@ -38,9 +38,9 @@ defmodule Towerops.Workers.DiscoveryWorker do
|
|||
# Consider agents online if they checked in within last 10 minutes
|
||||
@agent_online_threshold_minutes 10
|
||||
|
||||
@impl Oban.Worker
|
||||
@spec perform(Oban.Job.t()) :: :ok | :discard | {:error, term()}
|
||||
def perform(%Oban.Job{args: %{"device_id" => device_id}} = job) do
|
||||
@impl Oban.Pro.Worker
|
||||
@spec process(Oban.Job.t()) :: :ok | :discard | {:error, term()}
|
||||
def process(%Oban.Job{args: %{"device_id" => device_id}} = job) do
|
||||
_ = Events.broadcast_job_event(job, :started)
|
||||
start_time = System.monotonic_time(:second)
|
||||
|
||||
|
|
|
|||
|
|
@ -10,8 +10,8 @@ defmodule Towerops.Workers.EscalationCheckWorker do
|
|||
|
||||
alias Towerops.OnCall.Escalation
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"incident_id" => incident_id}}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: %{"incident_id" => incident_id}}) do
|
||||
_ = Escalation.check_and_escalate(incident_id)
|
||||
:ok
|
||||
end
|
||||
|
|
|
|||
|
|
@ -10,8 +10,8 @@ defmodule Towerops.Workers.ExpiredBanCleanupWorker do
|
|||
|
||||
alias Towerops.Security.BruteForce
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(_job) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(_job) do
|
||||
BruteForce.delete_expired_bans()
|
||||
:ok
|
||||
end
|
||||
|
|
|
|||
|
|
@ -22,9 +22,9 @@ defmodule Towerops.Workers.FirmwareVersionFetcherWorker do
|
|||
@rss_url "https://cdn.mikrotik.com/routeros/latest-stable.rss"
|
||||
@changelog_url "https://mikrotik.com/download/changelogs"
|
||||
|
||||
@impl Oban.Worker
|
||||
@spec perform(Oban.Job.t()) :: :ok | {:error, term()}
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
@spec process(Oban.Job.t()) :: :ok | {:error, term()}
|
||||
def process(%Oban.Job{}) do
|
||||
Logger.info("Fetching MikroTik RouterOS firmware version from RSS feed")
|
||||
|
||||
with {:ok, rss_body} <- fetch_rss_feed(),
|
||||
|
|
|
|||
|
|
@ -18,8 +18,8 @@ defmodule Towerops.Workers.GaiiaInsightWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
integrations = Integrations.list_enabled_integrations("gaiia")
|
||||
|
||||
Enum.each(integrations, fn integration ->
|
||||
|
|
|
|||
|
|
@ -25,8 +25,8 @@ defmodule Towerops.Workers.GaiiaSyncWorker do
|
|||
@window_seconds 600
|
||||
|
||||
# Cron dispatcher: enqueues staggered individual sync jobs.
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: args}) when args == %{} do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: args}) when args == %{} do
|
||||
integrations = Integrations.list_enabled_integrations(@provider)
|
||||
|
||||
eligible = Enum.filter(integrations, &due_for_sync?/1)
|
||||
|
|
@ -54,7 +54,7 @@ defmodule Towerops.Workers.GaiiaSyncWorker do
|
|||
end
|
||||
|
||||
# Individual sync: loads integration by ID and syncs.
|
||||
def perform(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
def process(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
case Integrations.get_integration_by_id(id) do
|
||||
{:ok, integration} ->
|
||||
sync_integration(integration)
|
||||
|
|
|
|||
|
|
@ -8,8 +8,8 @@ defmodule Towerops.Workers.GaiiaWebhookWorker do
|
|||
|
||||
alias Towerops.Gaiia.Webhooks
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"organization_id" => org_id, "event" => event, "payload" => payload}}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: %{"organization_id" => org_id, "event" => event, "payload" => payload}}) do
|
||||
Webhooks.process_event(org_id, event, payload)
|
||||
:ok
|
||||
end
|
||||
|
|
|
|||
|
|
@ -27,8 +27,8 @@ defmodule Towerops.Workers.InsightExpiryWorker do
|
|||
|
||||
@default_max_age_days 7
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{} = _job) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{} = _job) do
|
||||
cutoff = DateTime.add(DateTime.utc_now(), -max_age_days() * 86_400, :second)
|
||||
|
||||
{count, _} =
|
||||
|
|
|
|||
|
|
@ -19,8 +19,8 @@ defmodule Towerops.Workers.InsightLlmEnrichmentWorker do
|
|||
|
||||
@batch_size 25
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
[limit: @batch_size]
|
||||
|> Insights.list_unenriched_insights()
|
||||
|> Enum.each(&enrich_one/1)
|
||||
|
|
|
|||
|
|
@ -24,9 +24,9 @@ defmodule Towerops.Workers.JobHealthCheckWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
@spec perform(Oban.Job.t()) :: :ok
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
@spec process(Oban.Job.t()) :: :ok
|
||||
def process(%Oban.Job{}) do
|
||||
check_recoveries = recover_missing_check_executor_jobs()
|
||||
|
||||
if check_recoveries > 0 do
|
||||
|
|
|
|||
|
|
@ -13,8 +13,8 @@ defmodule Towerops.Workers.LidarCatalogSyncWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
case run_sync() do
|
||||
{:ok, summary} ->
|
||||
Logger.info(
|
||||
|
|
|
|||
|
|
@ -15,9 +15,9 @@ defmodule Towerops.Workers.LoginHistoryCleanupWorker do
|
|||
alias Towerops.Accounts.LoginAttempt
|
||||
alias Towerops.Repo
|
||||
|
||||
@impl Oban.Worker
|
||||
@spec perform(Oban.Job.t()) :: {:ok, %{deleted: non_neg_integer(), anonymized_deleted: non_neg_integer()}}
|
||||
def perform(_job) do
|
||||
@impl Oban.Pro.Worker
|
||||
@spec process(Oban.Job.t()) :: {:ok, %{deleted: non_neg_integer(), anonymized_deleted: non_neg_integer()}}
|
||||
def process(_job) do
|
||||
# Delete login attempts older than retention period (365 days)
|
||||
retention_days = Application.get_env(:towerops, :login_history_retention_days, 365)
|
||||
cutoff = DateTime.add(DateTime.utc_now(), -retention_days, :day)
|
||||
|
|
|
|||
|
|
@ -19,9 +19,9 @@ defmodule Towerops.Workers.MikrotikBackupWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
@spec perform(Oban.Job.t()) :: :ok | {:error, String.t()}
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
@spec process(Oban.Job.t()) :: :ok | {:error, String.t()}
|
||||
def process(%Oban.Job{}) do
|
||||
Logger.info("Starting MikroTik configuration backup job")
|
||||
|
||||
devices = list_backup_eligible_devices()
|
||||
|
|
|
|||
|
|
@ -23,8 +23,8 @@ defmodule Towerops.Workers.MikrotikWebhookWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: args}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: args}) do
|
||||
organization_id = Map.fetch!(args, "organization_id")
|
||||
source = Map.fetch!(args, "source")
|
||||
event_type = Map.fetch!(args, "event_type")
|
||||
|
|
|
|||
|
|
@ -31,8 +31,8 @@ defmodule Towerops.Workers.MsBuildingsImportWorker do
|
|||
|
||||
@batch_size 500
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"path" => path} = args}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: %{"path" => path} = args}) do
|
||||
source = Map.get(args, "source", "ms_global_ml")
|
||||
|
||||
if File.exists?(path) do
|
||||
|
|
|
|||
|
|
@ -25,8 +25,8 @@ defmodule Towerops.Workers.NetBoxSyncWorker do
|
|||
@window_seconds 600
|
||||
|
||||
# Cron dispatcher: enqueues staggered individual sync jobs.
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: args}) when args == %{} do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: args}) when args == %{} do
|
||||
integrations = Integrations.list_enabled_integrations(@provider)
|
||||
|
||||
eligible = Enum.filter(integrations, &due_for_sync?/1)
|
||||
|
|
@ -54,7 +54,7 @@ defmodule Towerops.Workers.NetBoxSyncWorker do
|
|||
end
|
||||
|
||||
# Individual sync: loads integration by ID and syncs.
|
||||
def perform(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
def process(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
case Integrations.get_integration_by_id(id) do
|
||||
{:ok, integration} ->
|
||||
sync_integration(integration)
|
||||
|
|
|
|||
|
|
@ -10,8 +10,8 @@ defmodule Towerops.Workers.PagerdutyWebhookWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"event_type" => event_type, "data" => data, "organization_id" => org_id}}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: %{"event_type" => event_type, "data" => data, "organization_id" => org_id}}) do
|
||||
case event_type do
|
||||
"incident.resolved" ->
|
||||
resolve_from_pagerduty(data, org_id)
|
||||
|
|
|
|||
|
|
@ -12,8 +12,8 @@ defmodule Towerops.Workers.PreseemBaselineWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
integrations = Integrations.list_enabled_integrations("preseem")
|
||||
|
||||
Enum.each(integrations, fn integration ->
|
||||
|
|
|
|||
|
|
@ -22,8 +22,8 @@ defmodule Towerops.Workers.PreseemSyncWorker do
|
|||
@window_seconds 600
|
||||
|
||||
# Cron dispatcher: enqueues staggered individual sync jobs
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: args}) when args == %{} do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: args}) when args == %{} do
|
||||
integrations = Integrations.list_enabled_integrations("preseem")
|
||||
|
||||
eligible = Enum.filter(integrations, &should_sync?/1)
|
||||
|
|
@ -51,7 +51,7 @@ defmodule Towerops.Workers.PreseemSyncWorker do
|
|||
end
|
||||
|
||||
# Individual sync: loads integration by ID and syncs
|
||||
def perform(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
def process(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
case Integrations.get_integration_by_id(id) do
|
||||
{:ok, integration} ->
|
||||
sync_integration(integration)
|
||||
|
|
|
|||
|
|
@ -40,8 +40,8 @@ defmodule Towerops.Workers.RecommendationsRunWorker do
|
|||
OpticalRxLow
|
||||
]
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
Enum.each(list_organization_ids(), &run_rules_for_org/1)
|
||||
:ok
|
||||
end
|
||||
|
|
|
|||
|
|
@ -13,8 +13,8 @@ defmodule Towerops.Workers.ReportWorker do
|
|||
require Logger
|
||||
|
||||
# Cron dispatcher: find due reports and enqueue individual jobs
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: args}) when args == %{} do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: args}) when args == %{} do
|
||||
reports =
|
||||
Reports.Report
|
||||
|> Towerops.Repo.all()
|
||||
|
|
@ -42,7 +42,7 @@ defmodule Towerops.Workers.ReportWorker do
|
|||
end
|
||||
|
||||
# Individual report: generate and deliver
|
||||
def perform(%Oban.Job{args: %{"report_id" => id}}) do
|
||||
def process(%Oban.Job{args: %{"report_id" => id}}) do
|
||||
case Reports.get_report(id) do
|
||||
{:ok, report} ->
|
||||
run_report(report)
|
||||
|
|
|
|||
|
|
@ -9,9 +9,9 @@ defmodule Towerops.Workers.SessionCleanupWorker do
|
|||
|
||||
alias Towerops.Accounts
|
||||
|
||||
@impl Oban.Worker
|
||||
@spec perform(Oban.Job.t()) :: {:ok, %{sessions_deleted: non_neg_integer()}}
|
||||
def perform(_job) do
|
||||
@impl Oban.Pro.Worker
|
||||
@spec process(Oban.Job.t()) :: {:ok, %{sessions_deleted: non_neg_integer()}}
|
||||
def process(_job) do
|
||||
count = Accounts.delete_expired_browser_sessions()
|
||||
{:ok, %{sessions_deleted: count}}
|
||||
end
|
||||
|
|
|
|||
|
|
@ -25,8 +25,8 @@ defmodule Towerops.Workers.SonarSyncWorker do
|
|||
@window_seconds 300
|
||||
|
||||
# Cron dispatcher: enqueues staggered individual sync jobs.
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: args}) when args == %{} do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: args}) when args == %{} do
|
||||
integrations = Integrations.list_enabled_integrations(@provider)
|
||||
|
||||
eligible = Enum.filter(integrations, &due_for_sync?/1)
|
||||
|
|
@ -54,7 +54,7 @@ defmodule Towerops.Workers.SonarSyncWorker do
|
|||
end
|
||||
|
||||
# Individual sync: loads integration by ID and syncs.
|
||||
def perform(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
def process(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
case Integrations.get_integration_by_id(id) do
|
||||
{:ok, integration} ->
|
||||
sync_integration(integration)
|
||||
|
|
|
|||
|
|
@ -25,8 +25,8 @@ defmodule Towerops.Workers.SplynxSyncWorker do
|
|||
@window_seconds 300
|
||||
|
||||
# Cron dispatcher: enqueues staggered individual sync jobs.
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: args}) when args == %{} do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: args}) when args == %{} do
|
||||
integrations = Integrations.list_enabled_integrations(@provider)
|
||||
|
||||
eligible = Enum.filter(integrations, &due_for_sync?/1)
|
||||
|
|
@ -54,7 +54,7 @@ defmodule Towerops.Workers.SplynxSyncWorker do
|
|||
end
|
||||
|
||||
# Individual sync: loads integration by ID and syncs.
|
||||
def perform(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
def process(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
case Integrations.get_integration_by_id(id) do
|
||||
{:ok, integration} ->
|
||||
sync_integration(integration)
|
||||
|
|
|
|||
|
|
@ -27,8 +27,8 @@ defmodule Towerops.Workers.StaleAgentWorker do
|
|||
# Consider agents "newly stale" if not seen in 10-15 minutes (for warning logs)
|
||||
@newly_stale_threshold_minutes 15
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
stale_agents = find_stale_agents()
|
||||
|
||||
_ =
|
||||
|
|
|
|||
|
|
@ -10,8 +10,8 @@ defmodule Towerops.Workers.StaleViolationCleanupWorker do
|
|||
|
||||
alias Towerops.Security.BruteForce
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(_job) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(_job) do
|
||||
BruteForce.delete_stale_violations()
|
||||
:ok
|
||||
end
|
||||
|
|
|
|||
|
|
@ -10,8 +10,8 @@ defmodule Towerops.Workers.StripeWebhookWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"event" => event}}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: %{"event" => event}}) do
|
||||
case WebhookProcessor.process(event) do
|
||||
:ok ->
|
||||
:ok
|
||||
|
|
|
|||
|
|
@ -20,8 +20,8 @@ defmodule Towerops.Workers.SystemInsightWorker do
|
|||
|
||||
@stale_threshold_minutes 10
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
# Use streaming to avoid loading all org IDs into memory
|
||||
Repo.transaction(fn ->
|
||||
Organization
|
||||
|
|
|
|||
|
|
@ -22,8 +22,8 @@ defmodule Towerops.Workers.UispSyncWorker do
|
|||
@window_seconds 300
|
||||
|
||||
# Cron dispatcher: enqueues staggered individual sync jobs
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: args}) when args == %{} do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: args}) when args == %{} do
|
||||
integrations = Integrations.list_enabled_integrations("uisp")
|
||||
|
||||
eligible = Enum.filter(integrations, &should_sync?/1)
|
||||
|
|
@ -51,7 +51,7 @@ defmodule Towerops.Workers.UispSyncWorker do
|
|||
end
|
||||
|
||||
# Individual sync: loads integration by ID and syncs
|
||||
def perform(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
def process(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
case Integrations.get_integration_by_id(id) do
|
||||
{:ok, integration} ->
|
||||
sync_integration(integration)
|
||||
|
|
|
|||
|
|
@ -25,8 +25,8 @@ defmodule Towerops.Workers.VispSyncWorker do
|
|||
@window_seconds 300
|
||||
|
||||
# Cron dispatcher: enqueues staggered individual sync jobs.
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: args}) when args == %{} do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: args}) when args == %{} do
|
||||
integrations = Integrations.list_enabled_integrations(@provider)
|
||||
|
||||
eligible = Enum.filter(integrations, &due_for_sync?/1)
|
||||
|
|
@ -54,7 +54,7 @@ defmodule Towerops.Workers.VispSyncWorker do
|
|||
end
|
||||
|
||||
# Individual sync: loads integration by ID and syncs.
|
||||
def perform(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
def process(%Oban.Job{args: %{"integration_id" => id}}) do
|
||||
case Integrations.get_integration_by_id(id) do
|
||||
{:ok, integration} ->
|
||||
sync_integration(integration)
|
||||
|
|
|
|||
|
|
@ -16,8 +16,8 @@ defmodule Towerops.Workers.WeatherSyncWorker do
|
|||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
api_key = Application.get_env(:towerops, :openweathermap_api_key)
|
||||
|
||||
if is_nil(api_key) or api_key == "" do
|
||||
|
|
|
|||
|
|
@ -13,8 +13,8 @@ defmodule Towerops.Workers.WelcomeEmailWorker do
|
|||
|
||||
@delay_minutes 3
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"user_id" => user_id}}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{args: %{"user_id" => user_id}}) do
|
||||
case Accounts.get_user(user_id) do
|
||||
nil ->
|
||||
Logger.warning("Welcome email skipped: user #{user_id} not found")
|
||||
|
|
|
|||
|
|
@ -35,8 +35,8 @@ defmodule Towerops.Workers.WirelessInsightWorker do
|
|||
@ap_critical_threshold 75
|
||||
@ap_recovery_threshold 40
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
@impl Oban.Pro.Worker
|
||||
def process(%Oban.Job{}) do
|
||||
# Use streaming to avoid loading all org IDs into memory
|
||||
Repo.transaction(fn ->
|
||||
Organization
|
||||
|
|
|
|||
|
|
@ -105,8 +105,6 @@ defmodule ToweropsWeb.Org.IntegrationsLive do
|
|||
}
|
||||
]
|
||||
|
||||
@providers Enum.flat_map(@provider_categories, & &1.providers)
|
||||
|
||||
@impl true
|
||||
def mount(_params, _session, socket) do
|
||||
organization = socket.assigns.current_scope.organization
|
||||
|
|
|
|||
|
|
@ -18,8 +18,6 @@ defmodule ToweropsWeb.Org.SettingsLive do
|
|||
|
||||
require Logger
|
||||
|
||||
@billing_provider_ids ~w(gaiia sonar splynx visp)
|
||||
|
||||
@impl true
|
||||
def mount(_params, _session, socket) do
|
||||
organization = socket.assigns.current_scope.organization
|
||||
|
|
@ -560,15 +558,6 @@ defmodule ToweropsWeb.Org.SettingsLive do
|
|||
|> Map.new(fn integration -> {integration.provider, integration} end)
|
||||
end
|
||||
|
||||
defp active_billing_provider(integrations) do
|
||||
Enum.find_value(@billing_provider_ids, fn id ->
|
||||
case Map.get(integrations, id) do
|
||||
%{enabled: true} -> id
|
||||
_ -> nil
|
||||
end
|
||||
end)
|
||||
end
|
||||
|
||||
defp get_current_integration(socket) do
|
||||
provider = socket.assigns.configuring
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue