refactor(webhooks): enqueue Oban jobs immediately, ACK fast
All 4 webhook controllers now verify the signature, insert an Oban job, and return immediately. Previously, Stripe and PagerDuty did DB queries + writes inline; Gaiia ran upserts; Agent Release broadcast to all connected WebSockets — all in the request path. New workers: - StripeWebhookWorker — delegates to WebhookProcessor.process/1 - GaiiaWebhookWorker — delegates to Webhooks.process_event/3 - PagerdutyWebhookWorker — resolve/acknowledge alerts by dedup key - AgentReleaseWebhookWorker — broadcast mass update to agents Also fix: log error when resolve_alert fails in agent_channel (was silently discarded, causing "Resolving stuck device_down alert" to repeat in logs).
This commit is contained in:
parent
36edf54e7c
commit
3b0961dfc0
13 changed files with 336 additions and 314 deletions
25
lib/towerops/workers/agent_release_webhook_worker.ex
Normal file
25
lib/towerops/workers/agent_release_webhook_worker.ex
Normal file
|
|
@ -0,0 +1,25 @@
|
|||
defmodule Towerops.Workers.AgentReleaseWebhookWorker do
|
||||
@moduledoc """
|
||||
Broadcasts a mass update to all connected agents asynchronously so the
|
||||
webhook controller can ACK the CI sender immediately.
|
||||
"""
|
||||
|
||||
use Oban.Worker, queue: :default, max_attempts: 3
|
||||
|
||||
alias Towerops.Agents
|
||||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{}) do
|
||||
case Agents.broadcast_mass_update() do
|
||||
{:ok, result} ->
|
||||
Logger.info("Agent mass update: #{result.notified} notified, #{result.skipped} skipped")
|
||||
:ok
|
||||
|
||||
{:error, reason} ->
|
||||
Logger.error("Agent mass update broadcast failed: #{inspect(reason)}")
|
||||
{:error, reason}
|
||||
end
|
||||
end
|
||||
end
|
||||
18
lib/towerops/workers/gaiia_webhook_worker.ex
Normal file
18
lib/towerops/workers/gaiia_webhook_worker.ex
Normal file
|
|
@ -0,0 +1,18 @@
|
|||
defmodule Towerops.Workers.GaiiaWebhookWorker do
|
||||
@moduledoc """
|
||||
Processes a Gaiia webhook event asynchronously so the controller can ACK
|
||||
the webhook sender immediately.
|
||||
"""
|
||||
|
||||
use Oban.Worker, queue: :default, max_attempts: 3
|
||||
|
||||
alias Towerops.Gaiia.Webhooks
|
||||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"organization_id" => org_id, "event" => event, "payload" => payload}}) do
|
||||
Webhooks.process_event(org_id, event, payload)
|
||||
:ok
|
||||
end
|
||||
end
|
||||
79
lib/towerops/workers/pagerduty_webhook_worker.ex
Normal file
79
lib/towerops/workers/pagerduty_webhook_worker.ex
Normal file
|
|
@ -0,0 +1,79 @@
|
|||
defmodule Towerops.Workers.PagerdutyWebhookWorker do
|
||||
@moduledoc """
|
||||
Processes a PagerDuty webhook event asynchronously so the controller can ACK
|
||||
the webhook sender immediately.
|
||||
"""
|
||||
|
||||
use Oban.Worker, queue: :default, max_attempts: 3
|
||||
|
||||
alias Towerops.Alerts
|
||||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%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)
|
||||
|
||||
"incident.acknowledged" ->
|
||||
acknowledge_from_pagerduty(data, org_id)
|
||||
|
||||
other ->
|
||||
Logger.debug("Ignoring PagerDuty webhook event: #{other}")
|
||||
:ok
|
||||
end
|
||||
end
|
||||
|
||||
defp resolve_from_pagerduty(data, _org_id) do
|
||||
case extract_alert_id(data) do
|
||||
nil ->
|
||||
Logger.debug("PagerDuty resolve: no TowerOps alert ID found in incident")
|
||||
:ok
|
||||
|
||||
alert_id ->
|
||||
case Alerts.get_alert(alert_id) do
|
||||
nil ->
|
||||
Logger.warning("PagerDuty resolve: alert #{alert_id} not found")
|
||||
|
||||
%{resolved_at: resolved_at} when not is_nil(resolved_at) ->
|
||||
Logger.debug("PagerDuty resolve: alert #{alert_id} already resolved")
|
||||
|
||||
alert ->
|
||||
Logger.info("Resolving alert #{alert_id} from PagerDuty")
|
||||
Alerts.resolve_alert_silent(alert)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
defp acknowledge_from_pagerduty(data, _org_id) do
|
||||
case extract_alert_id(data) do
|
||||
nil ->
|
||||
:ok
|
||||
|
||||
alert_id ->
|
||||
case Alerts.get_alert(alert_id) do
|
||||
nil ->
|
||||
Logger.warning("PagerDuty acknowledge: alert #{alert_id} not found")
|
||||
|
||||
%{acknowledged_at: ack} when not is_nil(ack) ->
|
||||
Logger.debug("PagerDuty acknowledge: alert #{alert_id} already acknowledged")
|
||||
|
||||
alert ->
|
||||
Logger.info("Acknowledging alert #{alert_id} from PagerDuty")
|
||||
Alerts.acknowledge_alert_silent(alert)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
defp extract_alert_id(data) do
|
||||
alerts = data["alerts"] || []
|
||||
|
||||
Enum.find_value(alerts, fn alert ->
|
||||
case alert["alert_key"] do
|
||||
"towerops-alert-" <> id -> id
|
||||
_ -> nil
|
||||
end
|
||||
end)
|
||||
end
|
||||
end
|
||||
24
lib/towerops/workers/stripe_webhook_worker.ex
Normal file
24
lib/towerops/workers/stripe_webhook_worker.ex
Normal file
|
|
@ -0,0 +1,24 @@
|
|||
defmodule Towerops.Workers.StripeWebhookWorker do
|
||||
@moduledoc """
|
||||
Processes a Stripe webhook event asynchronously so the controller can ACK
|
||||
the webhook sender immediately.
|
||||
"""
|
||||
|
||||
use Oban.Worker, queue: :default, max_attempts: 3
|
||||
|
||||
alias Towerops.Billing.WebhookProcessor
|
||||
|
||||
require Logger
|
||||
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"event" => event}}) do
|
||||
case WebhookProcessor.process(event) do
|
||||
:ok ->
|
||||
:ok
|
||||
|
||||
{:error, reason} ->
|
||||
Logger.error("Stripe webhook processing failed: #{inspect(reason)}")
|
||||
{:error, reason}
|
||||
end
|
||||
end
|
||||
end
|
||||
|
|
@ -1754,13 +1754,20 @@ defmodule ToweropsWeb.AgentChannel do
|
|||
})
|
||||
|
||||
# Resolve the stuck alert
|
||||
Alerts.resolve_alert(alert)
|
||||
case Alerts.resolve_alert(alert) do
|
||||
{:ok, _} ->
|
||||
Phoenix.PubSub.broadcast(
|
||||
Towerops.PubSub,
|
||||
"alerts:org:#{device.organization_id}:resolved",
|
||||
{:alert_resolved, device.id, :device_down}
|
||||
)
|
||||
|
||||
Phoenix.PubSub.broadcast(
|
||||
Towerops.PubSub,
|
||||
"alerts:org:#{device.organization_id}:resolved",
|
||||
{:alert_resolved, device.id, :device_down}
|
||||
)
|
||||
{:error, reason} ->
|
||||
Logger.error(
|
||||
"Failed to resolve stuck device_down alert #{alert.id} for #{device.name}: #{inspect(reason)}",
|
||||
device_id: device.id
|
||||
)
|
||||
end
|
||||
end
|
||||
|
||||
# Also resolve any stuck device_up alerts (should have been created with resolved_at)
|
||||
|
|
|
|||
|
|
@ -2,21 +2,12 @@ defmodule ToweropsWeb.Api.V1.AgentReleaseWebhookController do
|
|||
@moduledoc """
|
||||
Webhook endpoint for CI to trigger mass agent updates.
|
||||
|
||||
Called by CI after publishing a new GitHub Release. Broadcasts an update
|
||||
command to all connected agents via their existing WebSocket channels.
|
||||
|
||||
## Authentication
|
||||
Uses Bearer token authentication via the WebhookAuth plug.
|
||||
The token must match the configured `agent_webhook_secret`.
|
||||
|
||||
## Optional Signature Verification
|
||||
For enhanced security, also supports HMAC-SHA256 signature verification
|
||||
via X-Agent-Webhook-Signature header (format: t=<timestamp>,v1=<signature>).
|
||||
Verifies the optional signature, enqueues an Oban job, and ACKs immediately.
|
||||
The actual broadcast happens in `Towerops.Workers.AgentReleaseWebhookWorker`.
|
||||
"""
|
||||
|
||||
use ToweropsWeb, :controller
|
||||
|
||||
alias Towerops.Agents
|
||||
alias Towerops.Workers.AgentReleaseWebhookWorker
|
||||
|
||||
require Logger
|
||||
|
||||
|
|
@ -25,7 +16,14 @@ defmodule ToweropsWeb.Api.V1.AgentReleaseWebhookController do
|
|||
def create(conn, _params) do
|
||||
case verify_optional_signature(conn) do
|
||||
:ok ->
|
||||
handle_update(conn)
|
||||
case Oban.insert(AgentReleaseWebhookWorker.new(%{})) do
|
||||
{:ok, _job} ->
|
||||
json(conn, %{status: "ok"})
|
||||
|
||||
{:error, changeset} ->
|
||||
Logger.error("Failed to enqueue agent release webhook: #{inspect(changeset.errors)}")
|
||||
conn |> put_status(500) |> json(%{status: "error", error: "Failed to enqueue"})
|
||||
end
|
||||
|
||||
{:error, reason} ->
|
||||
Logger.warning("Agent release webhook signature verification failed: #{inspect(reason)}")
|
||||
|
|
@ -36,25 +34,6 @@ defmodule ToweropsWeb.Api.V1.AgentReleaseWebhookController do
|
|||
end
|
||||
end
|
||||
|
||||
defp handle_update(conn) do
|
||||
case Agents.broadcast_mass_update() do
|
||||
{:ok, result} ->
|
||||
json(conn, %{
|
||||
status: "ok",
|
||||
notified: result.notified,
|
||||
skipped: result.skipped,
|
||||
version: result.version
|
||||
})
|
||||
|
||||
{:error, reason} ->
|
||||
Logger.error("Agent mass update broadcast failed: #{inspect(reason)}")
|
||||
|
||||
conn
|
||||
|> put_status(:bad_gateway)
|
||||
|> json(%{status: "error", error: "Broadcast failed"})
|
||||
end
|
||||
end
|
||||
|
||||
defp verify_optional_signature(conn) do
|
||||
sig_headers = Plug.Conn.get_req_header(conn, "x-agent-webhook-signature")
|
||||
|
||||
|
|
|
|||
|
|
@ -2,39 +2,38 @@ defmodule ToweropsWeb.Api.V1.GaiiaWebhookController do
|
|||
@moduledoc """
|
||||
Webhook endpoint for receiving real-time updates from Gaiia.
|
||||
|
||||
Implements signature verification per Gaiia's webhook security docs:
|
||||
https://app.gaiia.com/docs/webhooks/security
|
||||
|
||||
## Gaiia's verification algorithm:
|
||||
|
||||
1. **Extract** timestamp (`t`) and signature (`v1`) from the
|
||||
`X-Gaiia-Webhook-Signature` header. Format: `t=<unix>,v1=<hex>`
|
||||
|
||||
2. **Prepare** the signed payload: `<timestamp>.<raw_request_body>`
|
||||
|
||||
3. **Compute** HMAC-SHA256 of the signed payload using the endpoint's
|
||||
secret key. Hex-encode the result.
|
||||
|
||||
4. **Compare** the computed signature with `v1` from the header.
|
||||
Also check that the timestamp is within acceptable tolerance.
|
||||
Verifies the signature, enqueues an Oban job, and ACKs immediately.
|
||||
All processing happens in `Towerops.Workers.GaiiaWebhookWorker`.
|
||||
|
||||
Route: POST /api/v1/webhooks/gaiia/:organization_id
|
||||
"""
|
||||
use ToweropsWeb, :controller
|
||||
|
||||
alias Towerops.Gaiia.Webhooks
|
||||
alias Towerops.Integrations
|
||||
alias Towerops.Workers.GaiiaWebhookWorker
|
||||
|
||||
require Logger
|
||||
|
||||
@max_age_seconds 300
|
||||
|
||||
# ── Main action ──
|
||||
|
||||
def create(conn, %{"organization_id" => organization_id}) do
|
||||
with {:ok, integration} <- get_gaiia_integration(organization_id),
|
||||
:ok <- verify_webhook(conn, integration) do
|
||||
process_webhook(conn, organization_id)
|
||||
:ok <- verify_webhook(conn, integration),
|
||||
{:ok, event, payload} <- extract_event(conn) do
|
||||
case Oban.insert(
|
||||
GaiiaWebhookWorker.new(%{
|
||||
organization_id: organization_id,
|
||||
event: event,
|
||||
payload: payload
|
||||
})
|
||||
) do
|
||||
{:ok, _job} ->
|
||||
json(conn, %{status: "ok"})
|
||||
|
||||
{:error, changeset} ->
|
||||
Logger.error("Failed to enqueue Gaiia webhook: #{inspect(changeset.errors)}")
|
||||
conn |> put_status(500) |> json(%{error: "Failed to enqueue"})
|
||||
end
|
||||
else
|
||||
{:error, :not_found} ->
|
||||
conn |> put_status(:not_found) |> json(%{error: "Gaiia integration not found"})
|
||||
|
|
@ -45,41 +44,33 @@ defmodule ToweropsWeb.Api.V1.GaiiaWebhookController do
|
|||
{:error, :invalid_signature} ->
|
||||
conn |> put_status(:unauthorized) |> json(%{error: "Invalid webhook signature"})
|
||||
|
||||
{:error, :missing_event} ->
|
||||
conn |> put_status(:bad_request) |> json(%{error: "Missing eventName field"})
|
||||
|
||||
{:error, reason} ->
|
||||
conn |> put_status(:unauthorized) |> json(%{error: "Webhook verification failed: #{reason}"})
|
||||
end
|
||||
end
|
||||
|
||||
defp process_webhook(conn, organization_id) do
|
||||
defp extract_event(conn) do
|
||||
event = Map.get(conn.body_params, "eventName") || Map.get(conn.body_params, "event")
|
||||
|
||||
if is_nil(event) do
|
||||
Logger.warning("Gaiia webhook missing eventName: #{inspect(Map.keys(conn.body_params))}")
|
||||
conn |> put_status(:bad_request) |> json(%{error: "Missing eventName field"})
|
||||
{:error, :missing_event}
|
||||
else
|
||||
data = Map.get(conn.body_params, "payload", Map.get(conn.body_params, "data", %{}))
|
||||
entity_id = Map.get(conn.body_params, "objectId", Map.get(data, "id"))
|
||||
|
||||
Webhooks.process_event(organization_id, event, %{
|
||||
"entity_id" => entity_id,
|
||||
"data" => data
|
||||
})
|
||||
|
||||
json(conn, %{status: "ok"})
|
||||
{:ok, event, %{"entity_id" => entity_id, "data" => data}}
|
||||
end
|
||||
end
|
||||
|
||||
# ── Integration lookup ──
|
||||
|
||||
defp get_gaiia_integration(organization_id) do
|
||||
case Integrations.get_integration(organization_id, "gaiia") do
|
||||
{:ok, integration} -> {:ok, integration}
|
||||
{:ok, _integration} = ok -> ok
|
||||
{:error, :not_found} -> {:error, :not_found}
|
||||
end
|
||||
end
|
||||
|
||||
# ── Signature verification per Gaiia docs ──
|
||||
|
||||
defp verify_webhook(conn, integration) do
|
||||
secret = get_in(integration.credentials, ["webhook_secret"])
|
||||
sig_headers = Plug.Conn.get_req_header(conn, "x-gaiia-webhook-signature")
|
||||
|
|
@ -99,28 +90,20 @@ defmodule ToweropsWeb.Api.V1.GaiiaWebhookController do
|
|||
end
|
||||
end
|
||||
|
||||
# Step 1: Parse the header to extract t and v1
|
||||
# Step 2-4: Build signed payload, compute HMAC, compare
|
||||
defp verify_gaiia_signature(header, raw_body, secret) do
|
||||
with {:ok, timestamp, v1_signature} <- parse_signature_header(header),
|
||||
:ok <- check_timestamp(timestamp) do
|
||||
# Step 2: signedPayload = "<timestamp>.<raw_body>"
|
||||
signed_payload = timestamp <> "." <> raw_body
|
||||
|
||||
# Step 3: HMAC-SHA256, hex-encoded lowercase
|
||||
expected =
|
||||
:hmac
|
||||
|> :crypto.mac(:sha256, secret, signed_payload)
|
||||
|> Base.encode16(case: :lower)
|
||||
|
||||
# Step 4: Constant-time comparison
|
||||
if Plug.Crypto.secure_compare(expected, v1_signature) do
|
||||
:ok
|
||||
else
|
||||
Logger.warning(
|
||||
"Gaiia webhook signature mismatch — " <>
|
||||
"ts=#{timestamp} body_len=#{byte_size(raw_body)}"
|
||||
)
|
||||
Logger.warning("Gaiia webhook signature mismatch — ts=#{timestamp} body_len=#{byte_size(raw_body)}")
|
||||
|
||||
{:error, :invalid_signature}
|
||||
end
|
||||
|
|
@ -135,7 +118,6 @@ defmodule ToweropsWeb.Api.V1.GaiiaWebhookController do
|
|||
end
|
||||
end
|
||||
|
||||
# Parse "t=1492774577,v1=5257a869..." into {:ok, "1492774577", "5257a869..."}
|
||||
defp parse_signature_header(header) do
|
||||
parts =
|
||||
header
|
||||
|
|
@ -153,13 +135,9 @@ defmodule ToweropsWeb.Api.V1.GaiiaWebhookController do
|
|||
end
|
||||
end
|
||||
|
||||
# Step 4 (from docs): "calculate the difference between the current
|
||||
# timestamp and the received timestamp, then decide if the difference
|
||||
# is within your tolerance"
|
||||
defp check_timestamp(timestamp_str) do
|
||||
case Integer.parse(timestamp_str) do
|
||||
{ts, ""} ->
|
||||
# Gaiia sends timestamps in milliseconds — normalize to seconds
|
||||
ts_seconds = if ts > 9_999_999_999, do: div(ts, 1000), else: ts
|
||||
age = abs(System.system_time(:second) - ts_seconds)
|
||||
if age <= @max_age_seconds, do: :ok, else: {:error, :expired}
|
||||
|
|
@ -168,22 +146,4 @@ defmodule ToweropsWeb.Api.V1.GaiiaWebhookController do
|
|||
{:error, :malformed}
|
||||
end
|
||||
end
|
||||
|
||||
# ── Error handling ──
|
||||
|
||||
def action(conn, _opts) do
|
||||
case apply(__MODULE__, action_name(conn), [conn, conn.params]) do
|
||||
{:error, :not_found} ->
|
||||
conn |> put_status(:not_found) |> json(%{error: "Gaiia integration not found"})
|
||||
|
||||
{:error, :invalid_signature} ->
|
||||
conn |> put_status(:unauthorized) |> json(%{error: "Invalid webhook signature"})
|
||||
|
||||
%Plug.Conn{} = conn ->
|
||||
conn
|
||||
|
||||
other ->
|
||||
other
|
||||
end
|
||||
end
|
||||
end
|
||||
|
|
|
|||
|
|
@ -2,16 +2,13 @@ defmodule ToweropsWeb.Api.V1.PagerdutyWebhookController do
|
|||
@moduledoc """
|
||||
Handles inbound PagerDuty V3 webhook events.
|
||||
|
||||
When an incident is resolved or acknowledged in PagerDuty, this controller
|
||||
updates the corresponding TowerOps alert so we don't re-alert on it.
|
||||
|
||||
PagerDuty signs webhooks with HMAC-SHA256 using a per-integration secret.
|
||||
The signature is in the `X-PagerDuty-Signature` header as `v1=<hex_digest>`.
|
||||
Verifies the signature, enqueues an Oban job, and ACKs immediately.
|
||||
All processing happens in `Towerops.Workers.PagerdutyWebhookWorker`.
|
||||
"""
|
||||
use ToweropsWeb, :controller
|
||||
|
||||
alias Towerops.Alerts
|
||||
alias Towerops.Integrations
|
||||
alias Towerops.Workers.PagerdutyWebhookWorker
|
||||
|
||||
require Logger
|
||||
|
||||
|
|
@ -19,14 +16,27 @@ defmodule ToweropsWeb.Api.V1.PagerdutyWebhookController do
|
|||
with {:ok, raw_body} <- get_raw_body(conn),
|
||||
{:ok, integration} <- get_pagerduty_integration(organization_id),
|
||||
:ok <- validate_organization_match(integration, organization_id),
|
||||
:ok <- verify_signature(conn, raw_body, integration) do
|
||||
# params already parsed by Phoenix JSON parser
|
||||
process_webhook(params, organization_id)
|
||||
json(conn, %{status: "accepted"})
|
||||
else
|
||||
{:error, :organization_mismatch} ->
|
||||
Logger.warning("PagerDuty webhook: rejected - organization ID mismatch (URL vs integration)")
|
||||
:ok <- verify_signature(conn, raw_body, integration),
|
||||
{:ok, event_type, data} <- extract_event(params) do
|
||||
case Oban.insert(
|
||||
PagerdutyWebhookWorker.new(%{
|
||||
organization_id: organization_id,
|
||||
event_type: event_type,
|
||||
data: data
|
||||
})
|
||||
) do
|
||||
{:ok, _job} ->
|
||||
json(conn, %{status: "accepted"})
|
||||
|
||||
{:error, changeset} ->
|
||||
Logger.error("Failed to enqueue PagerDuty webhook: #{inspect(changeset.errors)}")
|
||||
conn |> put_status(500) |> json(%{error: "Failed to enqueue"})
|
||||
end
|
||||
else
|
||||
{:error, :no_body} ->
|
||||
conn |> put_status(400) |> json(%{error: "Missing body"})
|
||||
|
||||
{:error, :organization_mismatch} ->
|
||||
conn |> put_status(403) |> json(%{error: "Organization mismatch"})
|
||||
|
||||
{:error, :not_configured} ->
|
||||
|
|
@ -38,11 +48,26 @@ defmodule ToweropsWeb.Api.V1.PagerdutyWebhookController do
|
|||
{:error, :no_secret_configured} ->
|
||||
conn |> put_status(403) |> json(%{error: "Webhook secret not configured"})
|
||||
|
||||
{:error, :unhandled_event} ->
|
||||
json(conn, %{status: "accepted"})
|
||||
|
||||
{:error, _reason} ->
|
||||
conn |> put_status(500) |> json(%{error: "Internal error"})
|
||||
end
|
||||
end
|
||||
|
||||
defp extract_event(%{"event" => %{"event_type" => event_type, "data" => data}}) do
|
||||
case event_type do
|
||||
type when type in ~w(incident.resolved incident.acknowledged) ->
|
||||
{:ok, type, data}
|
||||
|
||||
_ ->
|
||||
{:error, :unhandled_event}
|
||||
end
|
||||
end
|
||||
|
||||
defp extract_event(_), do: {:error, :unhandled_event}
|
||||
|
||||
defp get_raw_body(conn) do
|
||||
case conn.private[:raw_body] do
|
||||
body when is_binary(body) and body != "" -> {:ok, body}
|
||||
|
|
@ -102,75 +127,4 @@ defmodule ToweropsWeb.Api.V1.PagerdutyWebhookController do
|
|||
_ -> false
|
||||
end)
|
||||
end
|
||||
|
||||
defp process_webhook(%{"event" => %{"event_type" => event_type, "data" => data}}, organization_id) do
|
||||
case event_type do
|
||||
"incident.resolved" ->
|
||||
resolve_from_pagerduty(data, organization_id)
|
||||
|
||||
"incident.acknowledged" ->
|
||||
acknowledge_from_pagerduty(data, organization_id)
|
||||
|
||||
other ->
|
||||
Logger.debug("Ignoring PagerDuty webhook event: #{other}")
|
||||
:ok
|
||||
end
|
||||
end
|
||||
|
||||
defp process_webhook(_payload, _organization_id), do: :ok
|
||||
|
||||
defp resolve_from_pagerduty(data, _organization_id) do
|
||||
case extract_alert_id(data) do
|
||||
nil ->
|
||||
Logger.debug("PagerDuty resolve: no TowerOps alert ID found in incident")
|
||||
:ok
|
||||
|
||||
alert_id ->
|
||||
case Alerts.get_alert(alert_id) do
|
||||
nil ->
|
||||
Logger.warning("PagerDuty resolve: alert #{alert_id} not found")
|
||||
|
||||
%{resolved_at: resolved_at} when not is_nil(resolved_at) ->
|
||||
Logger.debug("PagerDuty resolve: alert #{alert_id} already resolved")
|
||||
|
||||
alert ->
|
||||
Logger.info("Resolving alert #{alert_id} from PagerDuty")
|
||||
Alerts.resolve_alert_silent(alert)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
defp acknowledge_from_pagerduty(data, _organization_id) do
|
||||
case extract_alert_id(data) do
|
||||
nil ->
|
||||
:ok
|
||||
|
||||
alert_id ->
|
||||
case Alerts.get_alert(alert_id) do
|
||||
nil ->
|
||||
Logger.warning("PagerDuty acknowledge: alert #{alert_id} not found")
|
||||
|
||||
%{acknowledged_at: ack} when not is_nil(ack) ->
|
||||
Logger.debug("PagerDuty acknowledge: alert #{alert_id} already acknowledged")
|
||||
|
||||
alert ->
|
||||
Logger.info("Acknowledging alert #{alert_id} from PagerDuty")
|
||||
Alerts.acknowledge_alert_silent(alert)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
# PagerDuty includes our dedup_key as the alert key.
|
||||
# Our dedup_key format is "towerops-alert-<uuid>"
|
||||
defp extract_alert_id(data) do
|
||||
# Try the first alert's dedup key
|
||||
alerts = data["alerts"] || []
|
||||
|
||||
Enum.find_value(alerts, fn alert ->
|
||||
case alert["alert_key"] do
|
||||
"towerops-alert-" <> id -> id
|
||||
_ -> nil
|
||||
end
|
||||
end)
|
||||
end
|
||||
end
|
||||
|
|
|
|||
|
|
@ -1,64 +1,41 @@
|
|||
defmodule ToweropsWeb.Api.V1.StripeWebhookController do
|
||||
use ToweropsWeb, :controller
|
||||
|
||||
alias Towerops.Billing.WebhookProcessor
|
||||
alias Towerops.Workers.StripeWebhookWorker
|
||||
|
||||
require Logger
|
||||
|
||||
@doc """
|
||||
Stripe webhook endpoint.
|
||||
|
||||
Receives events like payment_intent.succeeded, customer.subscription.updated.
|
||||
Verifies signature before processing.
|
||||
"""
|
||||
def create(conn, _params) do
|
||||
signature = conn |> get_req_header("stripe-signature") |> List.first()
|
||||
raw_body = conn.private[:raw_body] || ""
|
||||
|
||||
with :ok <- verify_signature(raw_body, signature),
|
||||
{:ok, event} <- decode_payload(raw_body) do
|
||||
process_webhook(conn, event)
|
||||
case Oban.insert(StripeWebhookWorker.new(%{event: event})) do
|
||||
{:ok, _job} ->
|
||||
json(conn, %{received: true})
|
||||
|
||||
{:error, changeset} ->
|
||||
Logger.error("Failed to enqueue Stripe webhook: #{inspect(changeset.errors)}")
|
||||
conn |> put_status(500) |> json(%{error: "Failed to enqueue"})
|
||||
end
|
||||
else
|
||||
{:error, :invalid_signature} ->
|
||||
Logger.warning("Invalid webhook signature")
|
||||
|
||||
conn
|
||||
|> put_status(401)
|
||||
|> json(%{error: "Invalid signature"})
|
||||
conn |> put_status(401) |> json(%{error: "Invalid signature"})
|
||||
|
||||
{:error, :invalid_json} ->
|
||||
conn
|
||||
|> put_status(400)
|
||||
|> json(%{error: "Invalid JSON payload"})
|
||||
conn |> put_status(400) |> json(%{error: "Invalid JSON payload"})
|
||||
|
||||
{:error, reason} ->
|
||||
Logger.error("Webhook verification failed: #{inspect(reason)}")
|
||||
|
||||
conn
|
||||
|> put_status(400)
|
||||
|> json(%{error: "Verification failed"})
|
||||
Logger.error("Stripe webhook verification failed: #{inspect(reason)}")
|
||||
conn |> put_status(400) |> json(%{error: "Verification failed"})
|
||||
end
|
||||
end
|
||||
|
||||
defp decode_payload(raw_body) do
|
||||
case Jason.decode(raw_body) do
|
||||
{:ok, event} ->
|
||||
{:ok, event}
|
||||
|
||||
{:error, decode_error} ->
|
||||
Logger.error("Failed to decode webhook payload: #{inspect(decode_error)}")
|
||||
{:error, :invalid_json}
|
||||
end
|
||||
end
|
||||
|
||||
defp process_webhook(conn, event) do
|
||||
case WebhookProcessor.process(event) do
|
||||
:ok ->
|
||||
json(conn, %{received: true})
|
||||
|
||||
{:error, reason} ->
|
||||
Logger.error("Webhook processing failed: #{inspect(reason)}")
|
||||
json(conn, %{received: true})
|
||||
{:ok, event} -> {:ok, event}
|
||||
{:error, _} -> {:error, :invalid_json}
|
||||
end
|
||||
end
|
||||
|
||||
|
|
|
|||
|
|
@ -1,7 +1,8 @@
|
|||
defmodule ToweropsWeb.Api.V1.AgentReleaseWebhookControllerTest do
|
||||
use ToweropsWeb.ConnCase, async: true
|
||||
use Oban.Testing, repo: Towerops.Repo
|
||||
|
||||
alias Towerops.Agents.ReleaseChecker
|
||||
alias Towerops.Workers.AgentReleaseWebhookWorker
|
||||
|
||||
setup do
|
||||
conn =
|
||||
|
|
@ -34,47 +35,11 @@ defmodule ToweropsWeb.Api.V1.AgentReleaseWebhookControllerTest do
|
|||
end
|
||||
|
||||
describe "create/2 response" do
|
||||
@tag :integration
|
||||
test "returns 502 when GitHub API is unavailable", %{conn: conn} do
|
||||
# In test environment, GitHub API returns 404 (private repo, no auth token).
|
||||
# The controller should surface this as a 502.
|
||||
test "enqueues job and returns 200", %{conn: conn} do
|
||||
conn = post(conn, ~p"/api/v1/webhooks/agent-release")
|
||||
|
||||
response = json_response(conn, 502)
|
||||
assert response["status"] == "error"
|
||||
end
|
||||
|
||||
test "returns 200 ok with notified/skipped/version on success", %{conn: conn} do
|
||||
# Stub ReleaseChecker so broadcast_mass_update returns {:ok, ...}
|
||||
Req.Test.stub(ReleaseChecker, fn r_conn ->
|
||||
Req.Test.json(r_conn, %{
|
||||
"tag_name" => "v9.9.9",
|
||||
"assets" => []
|
||||
})
|
||||
end)
|
||||
|
||||
ReleaseChecker.invalidate_cache()
|
||||
|
||||
conn = post(conn, ~p"/api/v1/webhooks/agent-release")
|
||||
|
||||
response = json_response(conn, 200)
|
||||
assert response["status"] == "ok"
|
||||
assert response["version"] == "9.9.9"
|
||||
assert is_integer(response["notified"])
|
||||
assert is_integer(response["skipped"])
|
||||
end
|
||||
|
||||
test "returns 502 when broadcast_mass_update returns error", %{conn: conn} do
|
||||
# Stub the underlying API call to fail with a non-200 status
|
||||
Req.Test.stub(ReleaseChecker, fn r_conn ->
|
||||
r_conn |> Plug.Conn.put_status(503) |> Req.Test.json(%{})
|
||||
end)
|
||||
|
||||
ReleaseChecker.invalidate_cache()
|
||||
|
||||
conn = post(conn, ~p"/api/v1/webhooks/agent-release")
|
||||
response = json_response(conn, 502)
|
||||
assert response["status"] == "error"
|
||||
assert json_response(conn, 200) == %{"status" => "ok"}
|
||||
assert_enqueued(worker: AgentReleaseWebhookWorker)
|
||||
end
|
||||
end
|
||||
|
||||
|
|
@ -183,8 +148,7 @@ defmodule ToweropsWeb.Api.V1.AgentReleaseWebhookControllerTest do
|
|||
assert json_response(conn, 401)["error"] == "Signature verification failed"
|
||||
end
|
||||
|
||||
@tag :integration
|
||||
test "passes signature verification with correct hmac", %{conn: conn} do
|
||||
test "enqueues job with valid signature and correct hmac", %{conn: conn} do
|
||||
ts = :second |> System.system_time() |> to_string()
|
||||
secret = Application.get_env(:towerops, :agent_webhook_secret)
|
||||
raw_body = ""
|
||||
|
|
@ -199,7 +163,8 @@ defmodule ToweropsWeb.Api.V1.AgentReleaseWebhookControllerTest do
|
|||
|> put_req_header("x-agent-webhook-signature", "t=#{ts},v1=#{sig}")
|
||||
|> post(~p"/api/v1/webhooks/agent-release")
|
||||
|
||||
assert conn.status in [200, 502]
|
||||
assert json_response(conn, 200) == %{"status" => "ok"}
|
||||
assert_enqueued(worker: AgentReleaseWebhookWorker)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
|
|
|||
|
|
@ -1,11 +1,13 @@
|
|||
defmodule ToweropsWeb.Api.V1.GaiiaWebhookControllerTest do
|
||||
use ToweropsWeb.ConnCase, async: true
|
||||
use Oban.Testing, repo: Towerops.Repo
|
||||
|
||||
import Towerops.AccountsFixtures
|
||||
import Towerops.OrganizationsFixtures
|
||||
|
||||
alias Towerops.Gaiia
|
||||
alias Towerops.Integrations
|
||||
alias Towerops.Workers.GaiiaWebhookWorker
|
||||
|
||||
@webhook_secret "test-gaiia-webhook-secret-abc123"
|
||||
|
||||
|
|
@ -42,6 +44,17 @@ defmodule ToweropsWeb.Api.V1.GaiiaWebhookControllerTest do
|
|||
|> post(~p"/api/v1/webhooks/gaiia/#{org_id}", body)
|
||||
end
|
||||
|
||||
defp worker_args(org, event, data) do
|
||||
%{
|
||||
"organization_id" => org.id,
|
||||
"event" => event,
|
||||
"payload" => %{
|
||||
"entity_id" => data["id"],
|
||||
"data" => data
|
||||
}
|
||||
}
|
||||
end
|
||||
|
||||
describe "authentication" do
|
||||
test "returns 401 without signature header", %{conn: conn, org: org} do
|
||||
conn =
|
||||
|
|
@ -87,7 +100,7 @@ defmodule ToweropsWeb.Api.V1.GaiiaWebhookControllerTest do
|
|||
end
|
||||
|
||||
describe "event processing" do
|
||||
test "processes account.created event", %{conn: conn, org: org} do
|
||||
test "enqueues job for account.created event", %{conn: conn, org: org} do
|
||||
payload = %{
|
||||
"id" => "gaiia-acct-1",
|
||||
"readable_id" => "ACCT-001",
|
||||
|
|
@ -99,11 +112,17 @@ defmodule ToweropsWeb.Api.V1.GaiiaWebhookControllerTest do
|
|||
conn = post_webhook(conn, org.id, "account.created", payload)
|
||||
assert json_response(conn, 200)["status"] == "ok"
|
||||
|
||||
args = worker_args(org, "account.created", payload)
|
||||
assert_enqueued(worker: GaiiaWebhookWorker)
|
||||
[job] = all_enqueued(worker: GaiiaWebhookWorker)
|
||||
assert ^args = job.args
|
||||
|
||||
perform_job(GaiiaWebhookWorker, args)
|
||||
assert account = Gaiia.get_account(org.id, "gaiia-acct-1")
|
||||
assert account.name == "Jane Doe"
|
||||
end
|
||||
|
||||
test "processes billing_subscription.activated event", %{conn: conn, org: org} do
|
||||
test "enqueues job for billing_subscription.activated event", %{conn: conn, org: org} do
|
||||
payload = %{
|
||||
"id" => "gaiia-sub-1",
|
||||
"account_id" => "gaiia-acct-1",
|
||||
|
|
@ -116,11 +135,13 @@ defmodule ToweropsWeb.Api.V1.GaiiaWebhookControllerTest do
|
|||
conn = post_webhook(conn, org.id, "billing_subscription.activated", payload)
|
||||
assert json_response(conn, 200)["status"] == "ok"
|
||||
|
||||
args = worker_args(org, "billing_subscription.activated", payload)
|
||||
perform_job(GaiiaWebhookWorker, args)
|
||||
assert sub = Gaiia.get_billing_subscription(org.id, "gaiia-sub-1")
|
||||
assert sub.product_name == "50 Mbps Plan"
|
||||
end
|
||||
|
||||
test "processes inventory_item.ip_address_assigned event", %{conn: conn, org: org} do
|
||||
test "enqueues job for inventory_item.ip_address_assigned event", %{conn: conn, org: org} do
|
||||
{:ok, _} =
|
||||
Gaiia.upsert_inventory_item(org.id, %{
|
||||
gaiia_id: "gaiia-inv-1",
|
||||
|
|
@ -136,11 +157,13 @@ defmodule ToweropsWeb.Api.V1.GaiiaWebhookControllerTest do
|
|||
conn = post_webhook(conn, org.id, "inventory_item.ip_address_assigned", payload)
|
||||
assert json_response(conn, 200)["status"] == "ok"
|
||||
|
||||
args = worker_args(org, "inventory_item.ip_address_assigned", payload)
|
||||
perform_job(GaiiaWebhookWorker, args)
|
||||
assert item = Gaiia.get_inventory_item(org.id, "gaiia-inv-1")
|
||||
assert item.ip_address == "10.0.0.50"
|
||||
end
|
||||
|
||||
test "returns ok for unknown events", %{conn: conn, org: org} do
|
||||
test "returns ok for unknown events (still enqueued)", %{conn: conn, org: org} do
|
||||
conn = post_webhook(conn, org.id, "some.future.event", %{})
|
||||
assert json_response(conn, 200)["status"] == "ok"
|
||||
end
|
||||
|
|
@ -165,7 +188,6 @@ defmodule ToweropsWeb.Api.V1.GaiiaWebhookControllerTest do
|
|||
test "returns 404 when no Gaiia integration is configured for org", %{conn: conn} do
|
||||
orphan_user = user_fixture()
|
||||
orphan_org = organization_fixture(orphan_user.id)
|
||||
# No integration created for orphan_org
|
||||
|
||||
conn = post_webhook(conn, orphan_org.id, "test.event", %{})
|
||||
assert json_response(conn, 404)["error"] =~ "not found"
|
||||
|
|
@ -195,7 +217,6 @@ defmodule ToweropsWeb.Api.V1.GaiiaWebhookControllerTest do
|
|||
|
||||
test "returns 401 when signature timestamp has expired", %{conn: conn, org: org} do
|
||||
body = Jason.encode!(%{"event" => "test", "data" => %{}})
|
||||
# 1 hour ago — well past the 5-minute window
|
||||
timestamp = DateTime.to_unix(DateTime.utc_now()) - 3600
|
||||
signed_payload = "#{timestamp}.#{body}"
|
||||
|
||||
|
|
|
|||
|
|
@ -1,10 +1,12 @@
|
|||
defmodule ToweropsWeb.Api.V1.PagerdutyWebhookControllerTest do
|
||||
use ToweropsWeb.ConnCase, async: true
|
||||
use Oban.Testing, repo: Towerops.Repo
|
||||
|
||||
import Towerops.AccountsFixtures
|
||||
import Towerops.OrganizationsFixtures
|
||||
|
||||
alias Towerops.Integrations
|
||||
alias Towerops.Workers.PagerdutyWebhookWorker
|
||||
|
||||
@webhook_secret "test-pagerduty-webhook-secret"
|
||||
|
||||
|
|
@ -144,8 +146,7 @@ defmodule ToweropsWeb.Api.V1.PagerdutyWebhookControllerTest do
|
|||
end
|
||||
|
||||
describe "incident.resolved event" do
|
||||
test "resolves an active alert matching the dedup key", %{conn: conn, org: org} do
|
||||
# Create a device and alert
|
||||
test "enqueues job and resolves an active alert", %{conn: conn, org: org} do
|
||||
site = site_fixture(org.id)
|
||||
|
||||
{:ok, device} =
|
||||
|
|
@ -180,10 +181,15 @@ defmodule ToweropsWeb.Api.V1.PagerdutyWebhookControllerTest do
|
|||
}
|
||||
|
||||
conn = post_webhook(conn, org.id, payload)
|
||||
|
||||
assert json_response(conn, 200)["status"] == "accepted"
|
||||
|
||||
# Verify the alert was resolved
|
||||
# Run the worker
|
||||
perform_job(PagerdutyWebhookWorker, %{
|
||||
"organization_id" => org.id,
|
||||
"event_type" => "incident.resolved",
|
||||
"data" => %{"alerts" => [%{"alert_key" => "towerops-alert-#{alert.id}"}]}
|
||||
})
|
||||
|
||||
updated_alert = Towerops.Alerts.get_alert(alert.id)
|
||||
assert updated_alert.resolved_at
|
||||
end
|
||||
|
|
@ -205,7 +211,7 @@ defmodule ToweropsWeb.Api.V1.PagerdutyWebhookControllerTest do
|
|||
end
|
||||
|
||||
describe "incident.acknowledged event" do
|
||||
test "acknowledges an active alert matching the dedup key", %{conn: conn, org: org, user: _user} do
|
||||
test "enqueues job and acknowledges an active alert", %{conn: conn, org: org, user: _user} do
|
||||
site = site_fixture(org.id)
|
||||
|
||||
{:ok, device} =
|
||||
|
|
@ -240,9 +246,15 @@ defmodule ToweropsWeb.Api.V1.PagerdutyWebhookControllerTest do
|
|||
}
|
||||
|
||||
conn = post_webhook(conn, org.id, payload)
|
||||
|
||||
assert json_response(conn, 200)["status"] == "accepted"
|
||||
|
||||
# Run the worker
|
||||
perform_job(PagerdutyWebhookWorker, %{
|
||||
"organization_id" => org.id,
|
||||
"event_type" => "incident.acknowledged",
|
||||
"data" => %{"alerts" => [%{"alert_key" => "towerops-alert-#{alert.id}"}]}
|
||||
})
|
||||
|
||||
updated_alert = Towerops.Alerts.get_alert(alert.id)
|
||||
assert updated_alert.acknowledged_at
|
||||
end
|
||||
|
|
|
|||
|
|
@ -1,5 +1,6 @@
|
|||
defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
||||
use ToweropsWeb.ConnCase
|
||||
use Oban.Testing, repo: Towerops.Repo
|
||||
|
||||
import Towerops.AccountsFixtures
|
||||
import Towerops.BillingFixtures
|
||||
|
|
@ -7,6 +8,7 @@ defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
|||
|
||||
alias Towerops.Billing.StripeWebhookEvent
|
||||
alias Towerops.Repo
|
||||
alias Towerops.Workers.StripeWebhookWorker
|
||||
|
||||
setup do
|
||||
user = user_fixture()
|
||||
|
|
@ -16,7 +18,7 @@ defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
|||
end
|
||||
|
||||
describe "POST /api/v1/webhooks/stripe" do
|
||||
test "processes valid webhook with correct signature", %{conn: conn, organization: organization} do
|
||||
test "enqueues job with correct signature", %{conn: conn, organization: organization} do
|
||||
subscription =
|
||||
stripe_subscription_object(%{
|
||||
id: "sub_webhook_123",
|
||||
|
|
@ -36,7 +38,15 @@ defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
|||
|
||||
assert json_response(conn, 200) == %{"received" => true}
|
||||
|
||||
# Verify event was processed
|
||||
# Verify job was enqueued with correct args
|
||||
assert_enqueued(worker: StripeWebhookWorker)
|
||||
|
||||
# Verify args match
|
||||
[job] = all_enqueued(worker: StripeWebhookWorker)
|
||||
assert %{"event" => ^event} = job.args
|
||||
|
||||
# Run the job and verify processing
|
||||
perform_job(StripeWebhookWorker, %{event: event})
|
||||
assert Repo.get(StripeWebhookEvent, event["id"])
|
||||
end
|
||||
|
||||
|
|
@ -44,7 +54,6 @@ defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
|||
event = stripe_webhook_event_object("customer.subscription.updated", %{})
|
||||
payload = Jason.encode!(event)
|
||||
|
||||
# Use current timestamp with invalid signature
|
||||
timestamp = System.system_time(:second)
|
||||
invalid_signature = "t=#{timestamp},v1=invalid_signature_hash"
|
||||
|
||||
|
|
@ -55,9 +64,7 @@ defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
|||
|> post(~p"/api/v1/webhooks/stripe", payload)
|
||||
|
||||
assert json_response(conn, 401) == %{"error" => "Invalid signature"}
|
||||
|
||||
# Verify event was not processed
|
||||
refute Repo.get(StripeWebhookEvent, event["id"])
|
||||
refute_enqueued(worker: StripeWebhookWorker)
|
||||
end
|
||||
|
||||
test "returns 400 for malformed signature header", %{conn: conn} do
|
||||
|
|
@ -77,7 +84,6 @@ defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
|||
event = stripe_webhook_event_object("customer.subscription.updated", %{})
|
||||
payload = Jason.encode!(event)
|
||||
|
||||
# Generate signature with old timestamp (> 5 minutes ago)
|
||||
old_timestamp = System.system_time(:second) - 400
|
||||
signature = generate_signature(payload, old_timestamp)
|
||||
|
||||
|
|
@ -90,8 +96,7 @@ defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
|||
assert json_response(conn, 400) == %{"error" => "Verification failed"}
|
||||
end
|
||||
|
||||
test "returns 200 even when event processing fails", %{conn: conn} do
|
||||
# Event with invalid data that will fail processing
|
||||
test "returns 200 even with bad event data (enqueue succeeds)", %{conn: conn} do
|
||||
event = %{
|
||||
"id" => "evt_fail_test",
|
||||
"type" => "customer.subscription.updated",
|
||||
|
|
@ -111,11 +116,10 @@ defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
|||
|> put_req_header("content-type", "application/json")
|
||||
|> post(~p"/api/v1/webhooks/stripe", payload)
|
||||
|
||||
# Should still return 200 to Stripe (idempotency)
|
||||
assert json_response(conn, 200) == %{"received" => true}
|
||||
end
|
||||
|
||||
test "handles duplicate webhook deliveries correctly", %{conn: conn, organization: organization} do
|
||||
test "duplicate webhook: enqueues twice, dedup handled by processor", %{conn: conn, organization: organization} do
|
||||
subscription =
|
||||
stripe_subscription_object(%{
|
||||
customer: organization.stripe_customer_id,
|
||||
|
|
@ -126,7 +130,6 @@ defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
|||
payload = Jason.encode!(event)
|
||||
signature = generate_signature(payload)
|
||||
|
||||
# Send same event twice
|
||||
conn1 =
|
||||
conn
|
||||
|> put_req_header("stripe-signature", signature)
|
||||
|
|
@ -142,29 +145,27 @@ defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
|||
assert json_response(conn1, 200) == %{"received" => true}
|
||||
assert json_response(conn2, 200) == %{"received" => true}
|
||||
|
||||
# Verify only one event record
|
||||
# Run both jobs — processor deduplicates
|
||||
perform_job(StripeWebhookWorker, %{event: event})
|
||||
perform_job(StripeWebhookWorker, %{event: event})
|
||||
assert Repo.aggregate(StripeWebhookEvent, :count) == 1
|
||||
end
|
||||
|
||||
test "processes different event types", %{conn: _conn, organization: organization} do
|
||||
event_types = [
|
||||
"customer.subscription.created",
|
||||
"customer.subscription.updated",
|
||||
"customer.subscription.deleted",
|
||||
"invoice.payment_succeeded",
|
||||
"invoice.payment_failed",
|
||||
"customer.updated"
|
||||
]
|
||||
test "enqueues jobs for different event types", %{conn: _conn, organization: organization} do
|
||||
event_types = ~w(
|
||||
customer.subscription.created
|
||||
customer.subscription.updated
|
||||
customer.subscription.deleted
|
||||
invoice.payment_succeeded
|
||||
invoice.payment_failed
|
||||
customer.updated
|
||||
)
|
||||
|
||||
for event_type <- event_types do
|
||||
data =
|
||||
case event_type do
|
||||
"invoice." <> _ ->
|
||||
stripe_invoice_object(%{customer: organization.stripe_customer_id})
|
||||
|
||||
_ ->
|
||||
stripe_subscription_object(%{customer: organization.stripe_customer_id})
|
||||
end
|
||||
if String.starts_with?(event_type, "invoice."),
|
||||
do: stripe_invoice_object(%{customer: organization.stripe_customer_id}),
|
||||
else: stripe_subscription_object(%{customer: organization.stripe_customer_id})
|
||||
|
||||
event = stripe_webhook_event_object(event_type, data)
|
||||
payload = Jason.encode!(event)
|
||||
|
|
@ -177,6 +178,9 @@ defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
|||
|> post(~p"/api/v1/webhooks/stripe", payload)
|
||||
|
||||
assert json_response(conn_result, 200) == %{"received" => true}
|
||||
|
||||
# Run the job
|
||||
perform_job(StripeWebhookWorker, %{event: event})
|
||||
assert Repo.get(StripeWebhookEvent, event["id"])
|
||||
end
|
||||
end
|
||||
|
|
@ -204,7 +208,6 @@ defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
|||
event = stripe_webhook_event_object("customer.subscription.updated", %{})
|
||||
payload = Jason.encode!(event)
|
||||
|
||||
# Timestamp 10 minutes in future
|
||||
future_timestamp = System.system_time(:second) + 600
|
||||
signature = generate_signature(payload, future_timestamp)
|
||||
|
||||
|
|
@ -224,7 +227,6 @@ defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
|||
timestamp = System.system_time(:second)
|
||||
signed_payload = "#{timestamp}.#{payload}"
|
||||
|
||||
# Sign with wrong secret
|
||||
wrong_signature =
|
||||
:hmac
|
||||
|> :crypto.mac(:sha256, "wrong_secret", signed_payload)
|
||||
|
|
@ -242,7 +244,6 @@ defmodule ToweropsWeb.Api.V1.StripeWebhookControllerTest do
|
|||
end
|
||||
end
|
||||
|
||||
# Helper to generate valid Stripe signature
|
||||
defp generate_signature(payload, timestamp \\ nil) do
|
||||
timestamp = timestamp || System.system_time(:second)
|
||||
webhook_secret = Application.fetch_env!(:towerops, :stripe_webhook_secret)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue