who_need_help/lib/who_need_help/external_boundary_drill.ex

747 lines
24 KiB
Elixir

defmodule WhoNeedHelp.ExternalBoundaryDrill do
@moduledoc """
Explicit local-only protocol drill for OAuth, SMTP, and the push boundary.
It is invoked by `scripts/external-boundaries-run.sh`; normal application
startup never calls it.
"""
import Ecto.Query, only: [where: 3]
import Swoosh.Email
alias WhoNeedHelp.Accounts
alias WhoNeedHelp.Accounts.{Scope, User}
alias WhoNeedHelp.{Catalog, Help, Messaging}
alias WhoNeedHelp.Push
alias WhoNeedHelp.Push.DeliveryWorker
alias WhoNeedHelp.Push.DisabledAdapter, as: PushDisabledAdapter
alias WhoNeedHelp.Push.HTTPAdapter, as: PushHTTPAdapter
alias WhoNeedHelp.Repo
alias WhoNeedHelp.GoogleAuth.AssentAdapter, as: GoogleAssentAdapter
alias WhoNeedHelp.SocialOAuth.AssentAdapter
@oauth_redirect_uri "http://boundary.local/auth/social/github/callback"
@google_redirect_uri "http://boundary.local/auth/google/callback"
@output_directory "/output"
@output_path "/output/summary.json"
def run! do
ensure_started!(:req)
ensure_started!(:swoosh)
ensure_started!(:gen_smtp)
ensure_started!(:who_need_help)
base_url = fetch_env!("EXTERNAL_MOCK_BASE_URL")
protocol_push = push_drill(base_url)
product_push = push_product_drill(base_url)
summary = %{
status: "passed",
oauth: oauth_drill(base_url),
google_oidc: google_oidc_drill(base_url),
smtp: smtp_drill(base_url),
push: Map.put(protocol_push, :domain_workflow_integration, product_push.status),
push_product_integration: product_push
}
File.mkdir_p!(@output_directory)
File.write!(@output_path, Jason.encode_to_iodata!(summary, pretty: true))
File.chmod!(@output_path, 0o600)
IO.puts("External boundary evidence: #{@output_path}")
:ok
end
defp google_oidc_drill(base_url) do
assert!(GoogleAssentAdapter.enabled?(), "runtime Google OIDC config is disabled")
control!(base_url, "google", "success")
{success_session, success_params} = google_authorization!()
{:ok, identity} =
GoogleAssentAdapter.callback(
@google_redirect_uri,
success_params,
success_session
)
assert!(
identity.provider_uid == "google-local-subject",
"Google subject was not normalized"
)
assert!(
identity.email == "google-helper@example.invalid",
"Google email was not normalized"
)
assert!(
not Map.has_key?(identity, :access_token),
"Google access token escaped adapter boundary"
)
success_state = state!(base_url)["google"]
assert!(success_state["token_requests"] == 1, "Google token request was not observed")
assert!(success_state["jwks_requests"] == 1, "Google JWKS request was not observed")
assert!(success_state["consumed_codes"] == 1, "Google code was not consumed")
assert_error!(
GoogleAssentAdapter.callback(
@google_redirect_uri,
success_params,
success_session
),
"Google authorization code replay unexpectedly succeeded"
)
replay_state = state!(base_url)["google"]
assert!(replay_state["token_requests"] == 2, "Google replay did not reach token endpoint")
assert!(replay_state["jwks_requests"] == 1, "Google replay reached the JWKS endpoint")
control!(base_url, "google", "success")
{mismatch_session, mismatch_params} = google_authorization!()
assert_error!(
GoogleAssentAdapter.callback(
@google_redirect_uri,
Map.put(mismatch_params, "state", "mismatched-state"),
mismatch_session
),
"Google state mismatch unexpectedly succeeded"
)
mismatch_state = state!(base_url)["google"]
assert!(mismatch_state["token_requests"] == 0, "Google state mismatch reached token endpoint")
control!(base_url, "google", "nonce_mismatch")
{nonce_session, nonce_params} = google_authorization!()
assert_error!(
GoogleAssentAdapter.callback(
@google_redirect_uri,
nonce_params,
nonce_session
),
"Google nonce mismatch unexpectedly succeeded"
)
nonce_state = state!(base_url)["google"]
assert!(nonce_state["token_requests"] == 1, "Google nonce test skipped token exchange")
assert!(nonce_state["jwks_requests"] == 1, "Google nonce test skipped signature verification")
control!(base_url, "google", "unverified_email")
{email_session, email_params} = google_authorization!()
assert!(
GoogleAssentAdapter.callback(
@google_redirect_uri,
email_params,
email_session
) == {:error, :email_not_verified},
"Google unverified email was accepted"
)
%{
success: "passed",
state_mismatch_blocked_before_token: true,
one_time_code_replay_rejected: true,
nonce_mismatch_rejected_after_signature_verification: true,
unverified_email_rejected: true,
access_token_returned_to_application: false
}
end
defp google_authorization! do
{:ok, %{url: url, session_params: session_params}} =
GoogleAssentAdapter.authorize_url(@google_redirect_uri)
response = Req.get!(url, redirect: false, retry: false)
assert!(response.status == 302, "Google OIDC mock did not redirect")
[location] = Req.Response.get_header(response, "location")
params = location |> URI.parse() |> Map.fetch!(:query) |> URI.decode_query()
{session_params, params}
end
defp oauth_drill(base_url) do
assert!(AssentAdapter.enabled?(:github), "runtime GitHub OAuth config is disabled")
control!(base_url, "oauth", "success")
{success_session, success_params} = oauth_authorization!()
{:ok, identity} =
AssentAdapter.callback(
:github,
@oauth_redirect_uri,
success_params,
success_session
)
assert!(identity.provider_uid == "4242", "OAuth identity UID was not normalized")
assert!(identity.handle == "@local-neighbor", "OAuth identity handle was not normalized")
assert!(not Map.has_key?(identity, :access_token), "OAuth token escaped adapter boundary")
success_state = state!(base_url)["oauth"]
assert!(success_state["token_requests"] == 1, "OAuth token request was not observed")
assert!(success_state["user_requests"] == 1, "OAuth user request was not observed")
assert!(success_state["consumed_codes"] == 1, "OAuth code was not consumed")
assert_error!(
AssentAdapter.callback(
:github,
@oauth_redirect_uri,
success_params,
success_session
),
"OAuth code replay unexpectedly succeeded"
)
replay_state = state!(base_url)["oauth"]
assert!(replay_state["token_requests"] == 2, "OAuth replay did not reach token endpoint")
assert!(replay_state["user_requests"] == 1, "OAuth replay reached the user endpoint")
control!(base_url, "oauth", "success")
{mismatch_session, mismatch_params} = oauth_authorization!()
assert_error!(
AssentAdapter.callback(
:github,
@oauth_redirect_uri,
Map.put(mismatch_params, "state", "mismatched-state"),
mismatch_session
),
"OAuth state mismatch unexpectedly succeeded"
)
mismatch_state = state!(base_url)["oauth"]
assert!(mismatch_state["token_requests"] == 0, "state mismatch reached token endpoint")
control!(base_url, "oauth", "deny_authorize")
{denied_session, denied_params} = oauth_authorization!()
assert_error!(
AssentAdapter.callback(
:github,
@oauth_redirect_uri,
denied_params,
denied_session
),
"provider authorization rejection unexpectedly succeeded"
)
denied_state = state!(base_url)["oauth"]
assert!(denied_state["token_requests"] == 0, "provider rejection reached token endpoint")
control!(base_url, "oauth", "token_temporary_once")
{temporary_session, temporary_params} = oauth_authorization!()
assert_error!(
AssentAdapter.callback(
:github,
@oauth_redirect_uri,
temporary_params,
temporary_session
),
"temporary OAuth token failure unexpectedly succeeded"
)
{retry_session, retry_params} = oauth_authorization!()
{:ok, _retried_identity} =
AssentAdapter.callback(
:github,
@oauth_redirect_uri,
retry_params,
retry_session
)
retry_state = state!(base_url)["oauth"]
assert!(retry_state["token_requests"] == 2, "fresh OAuth retry count was not two")
assert!(retry_state["user_requests"] == 1, "fresh OAuth retry did not fetch user")
control!(base_url, "oauth", "token_timeout_once")
{timeout_session, timeout_params} = oauth_authorization!()
assert_error!(
AssentAdapter.callback(
:github,
@oauth_redirect_uri,
timeout_params,
timeout_session
),
"OAuth timeout unexpectedly succeeded"
)
Process.sleep(400)
timeout_state = state!(base_url)["oauth"]
assert!(timeout_state["token_requests"] == 1, "OAuth timeout request was not observed")
assert!(timeout_state["user_requests"] == 0, "timed-out OAuth flow fetched user")
%{
success: "passed",
state_mismatch_blocked_before_token: true,
provider_rejection_blocked_before_token: true,
one_time_code_replay_rejected: true,
fresh_flow_retry_after_temporary_failure: "passed",
timeout_failed_closed: true,
access_token_returned_to_application: false
}
end
defp oauth_authorization! do
{:ok, %{url: url, session_params: session_params}} =
AssentAdapter.authorize_url(:github, @oauth_redirect_uri)
response = Req.get!(url, redirect: false, retry: false)
assert!(response.status == 302, "OAuth mock did not redirect")
[location] = Req.Response.get_header(response, "location")
params = location |> URI.parse() |> Map.fetch!(:query) |> URI.decode_query()
{session_params, params}
end
defp smtp_drill(base_url) do
relay = fetch_env!("EXTERNAL_SMTP_RELAY")
control!(base_url, "smtp", "success")
assert_smtp_success!(
smtp_email("success@example.invalid"),
smtp_options(relay, 2525, retries: 0)
)
success_state = state!(base_url)["smtp"]
assert!(success_state["messages"] == 1, "SMTP success message was not accepted")
control!(base_url, "smtp", "success")
assert_error!(
Swoosh.Adapters.SMTP.deliver(
smtp_email("reject@example.invalid"),
smtp_options(relay, 2525, retries: 1)
),
"SMTP permanent rejection unexpectedly succeeded"
)
rejection_state = state!(base_url)["smtp"]
assert!(rejection_state["rejections"] == 1, "SMTP rejection was not observed")
assert!(rejection_state["messages"] == 0, "rejected SMTP message was accepted")
assert!(
rejection_state["connections"]["2525"] == 1,
"permanent SMTP rejection was retried"
)
control!(base_url, "smtp", "success")
assert_smtp_success!(
smtp_email("retry@example.invalid"),
smtp_options(relay, 2526, retries: 1)
)
retry_state = state!(base_url)["smtp"]
assert!(retry_state["connections"]["2526"] == 2, "SMTP temporary failure was not retried")
assert!(retry_state["messages"] == 1, "retried SMTP message was not accepted")
control!(base_url, "smtp", "success")
assert_error!(
Swoosh.Adapters.SMTP.deliver(
smtp_email("timeout@example.invalid"),
smtp_options(relay, 2527, retries: 0, timeout: 100)
),
"SMTP timeout unexpectedly succeeded"
)
timeout_state = state!(base_url)["smtp"]
assert!(timeout_state["connections"]["2527"] == 1, "SMTP timeout connection was absent")
assert!(timeout_state["messages"] == 0, "timed-out SMTP message was accepted")
control!(base_url, "smtp", "success")
replay_email = smtp_email("replay@example.invalid")
replay_options = smtp_options(relay, 2525, retries: 0)
assert_smtp_success!(replay_email, replay_options)
assert_smtp_success!(replay_email, replay_options)
replay_state = state!(base_url)["smtp"]
assert!(replay_state["messages"] == 2, "SMTP replay behavior was not observed")
%{
success: "passed",
permanent_rejection_not_retried: true,
temporary_greeting_retried_once: true,
timeout_failed_closed: true,
repeated_submission_count: 2,
exactly_once_delivery_claimed: false
}
end
defp smtp_email(recipient) do
new()
|> to(recipient)
|> from({"Who Need Help boundary", "boundary@example.invalid"})
|> subject("External boundary drill")
|> text_body("Local protocol boundary payload")
end
defp smtp_options(relay, port, overrides) do
[
relay: relay,
port: port,
auth: :never,
tls: :never,
ssl: false,
no_mx_lookups: true,
timeout: 1_000,
retries: 0
]
|> Keyword.merge(overrides)
end
defp assert_smtp_success!(email, options) do
case Swoosh.Adapters.SMTP.deliver(email, options) do
{:ok, receipt} when is_binary(receipt) -> :ok
other -> raise "SMTP delivery failed: #{inspect(error_shape(other))}"
end
end
defp push_drill(base_url) do
endpoint = "#{base_url}/push"
bearer_token = fetch_env!("EXTERNAL_PUSH_BEARER_TOKEN")
assert_error!(
Push.deliver(push_notification("disabled"), adapter: PushDisabledAdapter),
"default push adapter unexpectedly delivered"
)
control!(base_url, "push", "success")
success = push_deliver!(endpoint, bearer_token, "push-success", max_attempts: 1)
assert!(success.duplicate == false, "first push was marked duplicate")
success_state = state!(base_url)["push"]
assert!(success_state["attempts"] == 1, "push success attempt count was not one")
assert!(success_state["deliveries"] == 1, "push success delivery count was not one")
control!(base_url, "push", "reject")
assert_error!(
push_deliver(endpoint, bearer_token, "push-reject", max_attempts: 2),
"push rejection unexpectedly succeeded"
)
rejection_state = state!(base_url)["push"]
assert!(rejection_state["attempts"] == 1, "permanent push rejection was retried")
assert!(rejection_state["deliveries"] == 0, "rejected push was delivered")
control!(base_url, "push", "temporary_once")
retry = push_deliver!(endpoint, bearer_token, "push-retry", max_attempts: 2)
assert!(retry.duplicate == false, "temporary push retry was marked duplicate")
retry_state = state!(base_url)["push"]
assert!(retry_state["attempts"] == 2, "temporary push failure was not retried once")
assert!(retry_state["deliveries"] == 1, "temporary push retry duplicated delivery")
control!(base_url, "push", "success")
first = push_deliver!(endpoint, bearer_token, "push-replay", max_attempts: 1)
replay = push_deliver!(endpoint, bearer_token, "push-replay", max_attempts: 1)
assert!(first.id == replay.id, "push replay returned a different receipt")
assert!(replay.duplicate, "push replay was not marked duplicate")
replay_state = state!(base_url)["push"]
assert!(replay_state["attempts"] == 2, "push replay attempt count was not two")
assert!(replay_state["deliveries"] == 1, "push replay created duplicate delivery")
control!(base_url, "push", "timeout_after_accept")
timeout =
push_deliver!(
endpoint,
bearer_token,
"push-timeout",
max_attempts: 2,
receive_timeout: 100
)
assert!(timeout.duplicate, "ambiguous push timeout retry was not deduplicated")
timeout_state = state!(base_url)["push"]
assert!(timeout_state["attempts"] == 2, "push timeout was not retried once")
assert!(timeout_state["deliveries"] == 1, "push timeout retry duplicated delivery")
%{
default_adapter_disabled: true,
success: "passed",
permanent_rejection_not_retried: true,
temporary_failure_retried_once: true,
replay_deduplicated: true,
timeout_after_accept_deduplicated: true
}
end
defp push_product_drill(base_url) do
expected_worker_replicas =
fetch_env!("EXTERNAL_EXPECTED_WORKER_REPLICAS")
|> parse_positive_integer!("EXTERNAL_EXPECTED_WORKER_REPLICAS")
requester = confirmed_user!("requester")
helper = confirmed_user!("helper")
requester_scope = Scope.for_user(requester)
helper_scope = Scope.for_user(helper)
category = Catalog.seed_defaults()
control!(base_url, "push", "success")
{:ok, request} =
Help.create_request(requester_scope, %{
"title" => "Boundary medicine pickup",
"description" => "The reserved legal medicine is ready for pickup.",
"pickup_instructions" => "Ask for the run-scoped reservation.",
"location_label" => "Boundary district",
"latitude" => "50.4501",
"longitude" => "30.5234",
"urgency" => "now",
"location_visibility" => "approximate_public",
"structured_data" => %{"pickup_status" => "reserved"},
"expires_at" => DateTime.utc_now(:second) |> DateTime.add(3, :hour),
"category_id" => category.id,
"safety_confirmed" => true
})
{:ok, assignment} = Help.accept_request(helper_scope, request.id)
acceptance_key = "request-accepted:#{assignment.id}:#{requester.id}"
acceptance_job = wait_for_completed_job!(acceptance_key, 1)
acceptance_record =
Repo.get_by!(WhoNeedHelp.Notifications.Notification, idempotency_key: acceptance_key)
acceptance_state = state!(base_url)["push"]
assert!(acceptance_job.attempt == 1, "acceptance push was not completed on its first job run")
assert!(acceptance_state["attempts"] == 1, "acceptance push attempt count was not one")
assert!(acceptance_state["deliveries"] == 1, "acceptance push delivery count was not one")
[acceptance_notification] = acceptance_state["notifications"]
assert!(
acceptance_notification["data"]["kind"] == "request_accepted",
"acceptance product event kind was not delivered"
)
assert!(
acceptance_notification["recipient"] == "user:#{requester.id}",
"acceptance product event targeted the wrong user"
)
{:ok, replay_notification} =
Push.enqueue_request_accepted(assignment.id, request.id, requester.id)
assert!(
replay_notification.id == acceptance_record.id,
"product event replay created another notification"
)
replay_job = wait_for_completed_job!(acceptance_key, 1)
assert!(replay_job.id == acceptance_job.id, "product event replay created another Oban job")
Process.sleep(250)
acceptance_replay_state = state!(base_url)["push"]
assert!(
acceptance_replay_state["attempts"] == 1,
"product event replay reached the transport boundary"
)
control!(base_url, "push", "temporary_once")
message_body = "Run-scoped chat text must not enter push"
{:ok, message} =
Messaging.send_message(requester_scope, assignment, %{"body" => message_body})
message_key = "message-created:#{message.id}:#{helper.id}"
message_job = wait_for_completed_job!(message_key, 2)
message_state = state!(base_url)["push"]
assert!(message_job.attempt == 2, "chat push did not complete on the second Oban attempt")
assert!(message_state["attempts"] == 2, "chat push temporary failure was not retried")
assert!(message_state["deliveries"] == 1, "chat push retry duplicated delivery")
[message_notification] = message_state["notifications"]
assert!(
message_notification["data"]["kind"] == "message_created",
"chat product event kind was not delivered"
)
assert!(
message_notification["recipient"] == "user:#{helper.id}",
"chat product event targeted the wrong user"
)
assert!(
message_notification["body"] != message_body and
not String.contains?(Jason.encode!(message_notification), message_body),
"chat message content escaped into the push payload"
)
%{
status: "passed",
acceptance_delivery: "passed",
chat_delivery: "passed",
oban_retry_observed: true,
replay_deduplicated_before_transport: true,
message_content_excluded: true,
expected_worker_replicas: expected_worker_replicas,
database_scope: "isolated_ephemeral_volume"
}
end
defp confirmed_user!(label) do
suffix = Ecto.UUID.generate()
{:ok, user} =
Accounts.register_user(%{
email: "#{label}-#{suffix}@boundary.invalid",
display_name: "Boundary #{label}",
terms_accepted: true
})
user
|> User.confirm_changeset()
|> Repo.update!()
end
defp wait_for_completed_job!(idempotency_key, minimum_attempts) do
deadline = System.monotonic_time(:millisecond) + 60_000
do_wait_for_completed_job!(idempotency_key, minimum_attempts, deadline)
end
defp do_wait_for_completed_job!(idempotency_key, minimum_attempts, deadline) do
worker_name = DeliveryWorker.__opts__() |> Keyword.fetch!(:worker)
job =
Oban.Job
|> where(
[job],
job.worker == ^worker_name and
fragment("?->>'idempotency_key' = ?", job.args, ^idempotency_key)
)
|> Repo.one()
cond do
job && job.state == "completed" && job.attempt >= minimum_attempts ->
job
job && job.state in ["cancelled", "discarded"] ->
raise "push product job ended in #{job.state}"
System.monotonic_time(:millisecond) >= deadline ->
observed =
case job do
nil ->
"missing"
job ->
inspect(%{
id: job.id,
state: job.state,
attempt: job.attempt,
max_attempts: job.max_attempts,
queue: job.queue,
scheduled_at: job.scheduled_at
})
end
raise "push product job did not complete within the local drill window: #{observed}"
true ->
Process.sleep(100)
do_wait_for_completed_job!(idempotency_key, minimum_attempts, deadline)
end
end
defp push_notification(idempotency_key) do
%{
idempotency_key: idempotency_key,
recipient: "future-device-token",
title: "Help update",
body: "A local boundary event",
data: %{"kind" => "boundary_drill"}
}
end
defp push_deliver!(endpoint, bearer_token, idempotency_key, overrides) do
case push_deliver(endpoint, bearer_token, idempotency_key, overrides) do
{:ok, receipt} -> receipt
other -> raise "push delivery failed: #{inspect(error_shape(other))}"
end
end
defp push_deliver(endpoint, bearer_token, idempotency_key, overrides) do
options =
[
adapter: PushHTTPAdapter,
endpoint: endpoint,
bearer_token: bearer_token,
max_attempts: 1,
receive_timeout: 1_000,
connect_timeout: 1_000,
retry_delay_ms: 0
]
|> Keyword.merge(overrides)
Push.deliver(push_notification(idempotency_key), options)
end
defp control!(base_url, component, mode) do
response =
Req.post!(
"#{base_url}/control",
json: %{component: component, mode: mode},
retry: false
)
assert!(response.status == 200, "boundary mock control failed")
end
defp state!(base_url) do
response = Req.get!("#{base_url}/state", retry: false)
assert!(response.status == 200, "boundary mock state request failed")
response.body
end
defp assert_error!({:error, _error}, _message), do: :ok
defp assert_error!(_result, message), do: raise(message)
defp assert!(true, _message), do: :ok
defp assert!(false, message), do: raise(message)
defp error_shape({:error, {kind, status, _body}}), do: {:error, kind, status}
defp error_shape({:error, error}) when is_atom(error), do: {:error, error}
defp error_shape({:error, %{__struct__: module}}), do: {:error, module}
defp error_shape(other), do: other
defp fetch_env!(name) do
case System.fetch_env(name) do
{:ok, value} when value != "" -> value
_other -> raise "#{name} is required for the external boundary drill."
end
end
defp parse_positive_integer!(value, name) do
case Integer.parse(value) do
{integer, ""} when integer > 0 -> integer
_other -> raise "#{name} must be a positive integer."
end
end
defp ensure_started!(application) do
case Application.ensure_all_started(application) do
{:ok, _started} -> :ok
{:error, reason} -> raise "could not start #{application}: #{inspect(reason)}"
end
end
end