diff --git a/lib/towerops/workers/agent_release_webhook_worker.ex b/lib/towerops/workers/agent_release_webhook_worker.ex new file mode 100644 index 00000000..c1182f3b --- /dev/null +++ b/lib/towerops/workers/agent_release_webhook_worker.ex @@ -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 diff --git a/lib/towerops/workers/gaiia_webhook_worker.ex b/lib/towerops/workers/gaiia_webhook_worker.ex new file mode 100644 index 00000000..627f314d --- /dev/null +++ b/lib/towerops/workers/gaiia_webhook_worker.ex @@ -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 diff --git a/lib/towerops/workers/pagerduty_webhook_worker.ex b/lib/towerops/workers/pagerduty_webhook_worker.ex new file mode 100644 index 00000000..952ff86d --- /dev/null +++ b/lib/towerops/workers/pagerduty_webhook_worker.ex @@ -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 diff --git a/lib/towerops/workers/stripe_webhook_worker.ex b/lib/towerops/workers/stripe_webhook_worker.ex new file mode 100644 index 00000000..ea9f030a --- /dev/null +++ b/lib/towerops/workers/stripe_webhook_worker.ex @@ -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 diff --git a/lib/towerops_web/channels/agent_channel.ex b/lib/towerops_web/channels/agent_channel.ex index 744f2173..c72c9dad 100644 --- a/lib/towerops_web/channels/agent_channel.ex +++ b/lib/towerops_web/channels/agent_channel.ex @@ -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) diff --git a/lib/towerops_web/controllers/api/v1/agent_release_webhook_controller.ex b/lib/towerops_web/controllers/api/v1/agent_release_webhook_controller.ex index 4da884f0..8120b595 100644 --- a/lib/towerops_web/controllers/api/v1/agent_release_webhook_controller.ex +++ b/lib/towerops_web/controllers/api/v1/agent_release_webhook_controller.ex @@ -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=,v1=). + 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") diff --git a/lib/towerops_web/controllers/api/v1/gaiia_webhook_controller.ex b/lib/towerops_web/controllers/api/v1/gaiia_webhook_controller.ex index e82f3d8a..1e3534d0 100644 --- a/lib/towerops_web/controllers/api/v1/gaiia_webhook_controller.ex +++ b/lib/towerops_web/controllers/api/v1/gaiia_webhook_controller.ex @@ -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=,v1=` - - 2. **Prepare** the signed payload: `.` - - 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 = "." 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 diff --git a/lib/towerops_web/controllers/api/v1/pagerduty_webhook_controller.ex b/lib/towerops_web/controllers/api/v1/pagerduty_webhook_controller.ex index f16d9704..b66f7a4e 100644 --- a/lib/towerops_web/controllers/api/v1/pagerduty_webhook_controller.ex +++ b/lib/towerops_web/controllers/api/v1/pagerduty_webhook_controller.ex @@ -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=`. + 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-" - 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 diff --git a/lib/towerops_web/controllers/api/v1/stripe_webhook_controller.ex b/lib/towerops_web/controllers/api/v1/stripe_webhook_controller.ex index b6dd039f..610409c5 100644 --- a/lib/towerops_web/controllers/api/v1/stripe_webhook_controller.ex +++ b/lib/towerops_web/controllers/api/v1/stripe_webhook_controller.ex @@ -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 diff --git a/test/towerops_web/controllers/api/v1/agent_release_webhook_controller_test.exs b/test/towerops_web/controllers/api/v1/agent_release_webhook_controller_test.exs index 6a3fdb7d..5951799f 100644 --- a/test/towerops_web/controllers/api/v1/agent_release_webhook_controller_test.exs +++ b/test/towerops_web/controllers/api/v1/agent_release_webhook_controller_test.exs @@ -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 diff --git a/test/towerops_web/controllers/api/v1/gaiia_webhook_controller_test.exs b/test/towerops_web/controllers/api/v1/gaiia_webhook_controller_test.exs index fafba339..71f64504 100644 --- a/test/towerops_web/controllers/api/v1/gaiia_webhook_controller_test.exs +++ b/test/towerops_web/controllers/api/v1/gaiia_webhook_controller_test.exs @@ -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}" diff --git a/test/towerops_web/controllers/api/v1/pagerduty_webhook_controller_test.exs b/test/towerops_web/controllers/api/v1/pagerduty_webhook_controller_test.exs index 63ffbdd6..585c8912 100644 --- a/test/towerops_web/controllers/api/v1/pagerduty_webhook_controller_test.exs +++ b/test/towerops_web/controllers/api/v1/pagerduty_webhook_controller_test.exs @@ -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 diff --git a/test/towerops_web/controllers/api/v1/stripe_webhook_controller_test.exs b/test/towerops_web/controllers/api/v1/stripe_webhook_controller_test.exs index 2b101174..d9d4ff62 100644 --- a/test/towerops_web/controllers/api/v1/stripe_webhook_controller_test.exs +++ b/test/towerops_web/controllers/api/v1/stripe_webhook_controller_test.exs @@ -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)