| 1 |
|
defmodule WalletIntegrations.Workers.PaymentWorker do |
| 2 |
|
@moduledoc """ |
| 3 |
|
Async worker for outbound payment operations. |
| 4 |
|
|
| 5 |
|
Implements `perform/1` compatible with both `WalletIntegrations.JobQueue` and Oban. |
| 6 |
|
|
| 7 |
|
Supported operations (via `"operation"` arg key): |
| 8 |
|
- `"initiate"` — dispatch a new payment request to the provider. |
| 9 |
|
- `"get_status"` — poll provider status for an existing payment. |
| 10 |
|
- `"cancel"` — cancel a pending payment at the provider. |
| 11 |
|
- `"refund"` — initiate a refund for a completed payment. |
| 12 |
|
|
| 13 |
|
Resilience policy (ADR 0008): |
| 14 |
|
- Circuit breaker checked before each outbound call. |
| 15 |
|
- Success/failure recorded to circuit breaker after each call. |
| 16 |
|
- Unknown outcomes create IntegrationException for reconciliation. |
| 17 |
|
- No duplicate payment charges on retry (provider idempotency key enforced). |
| 18 |
|
|
| 19 |
|
Adapter resolved via Application.get_env(:wallet_integrations, :payment_adapter, StubAdapter). |
| 20 |
|
""" |
| 21 |
|
|
| 22 |
|
alias WalletIntegrations.{PaymentRequest, PaymentRequestStore} |
| 23 |
|
alias WalletIntegrations.{IntegrationException, IntegrationExceptionStore} |
| 24 |
|
alias WalletIntegrations.{AdapterRequest, CircuitBreaker} |
| 25 |
|
alias WalletIntegrations.Events.{PaymentCompleted, PaymentFailed} |
| 26 |
|
alias WalletObservability.AuditEvent |
| 27 |
|
|
| 28 |
|
@spec perform(args :: map()) :: :ok | {:error, term()} |
| 29 |
|
def perform(%{"request_id" => request_id, "operation" => operation} = args) do |
| 30 |
19 |
correlation_id = |
| 31 |
|
Map.get(args, "correlation_id", WalletSharedKernel.Correlation.new_correlation_id()) |
| 32 |
|
|
| 33 |
19 |
with {:ok, request} <- PaymentRequestStore.get(request_id) do |
| 34 |
18 |
execute_operation(operation, request, args, correlation_id) |
| 35 |
|
end |
| 36 |
|
end |
| 37 |
|
|
| 38 |
1 |
def perform(_args), do: {:error, :missing_required_args} |
| 39 |
|
|
| 40 |
|
# Private: dispatch by operation |
| 41 |
|
|
| 42 |
|
defp execute_operation("initiate", request, _args, correlation_id) do |
| 43 |
18 |
adapter = resolve_adapter(request.provider) |
| 44 |
18 |
operation_atom = :initiate_payment |
| 45 |
|
|
| 46 |
18 |
case CircuitBreaker.check(request.provider, operation_atom) do |
| 47 |
1 |
{:error, :circuit_open} -> |
| 48 |
|
{:error, :circuit_open} |
| 49 |
|
|
| 50 |
|
:ok -> |
| 51 |
17 |
adapter_request = |
| 52 |
|
AdapterRequest.new_initiate( |
| 53 |
17 |
request.idempotency_key || request.request_id, |
| 54 |
17 |
request.amount, |
| 55 |
17 |
request.currency, |
| 56 |
17 |
provider: request.provider, |
| 57 |
|
correlation_id: correlation_id |
| 58 |
|
) |
| 59 |
|
|
| 60 |
17 |
case PaymentRequest.dispatch(request) do |
| 61 |
1 |
{:error, :invalid_transition} -> |
| 62 |
|
# Request is in a terminal state (e.g. failed on prior attempt); nothing to do |
| 63 |
|
:ok |
| 64 |
|
|
| 65 |
|
{:ok, dispatched} -> |
| 66 |
16 |
PaymentRequestStore.update(dispatched) |
| 67 |
|
|
| 68 |
16 |
case apply(adapter, :initiate_payment, [adapter_request]) do |
| 69 |
|
{:ok, result} -> |
| 70 |
13 |
CircuitBreaker.record_success(request.provider, operation_atom) |
| 71 |
13 |
handle_success_result(dispatched, result, correlation_id) |
| 72 |
|
|
| 73 |
|
{:error, result} -> |
| 74 |
3 |
CircuitBreaker.record_failure(request.provider, operation_atom) |
| 75 |
3 |
handle_failure_result(dispatched, result, correlation_id) |
| 76 |
|
end |
| 77 |
|
end |
| 78 |
|
end |
| 79 |
|
end |
| 80 |
|
|
| 81 |
|
defp execute_operation("get_status", request, _args, correlation_id) do |
| 82 |
:-( |
adapter = resolve_adapter(request.provider) |
| 83 |
:-( |
operation_atom = :get_payment_status |
| 84 |
|
|
| 85 |
:-( |
case CircuitBreaker.check(request.provider, operation_atom) do |
| 86 |
:-( |
{:error, :circuit_open} -> |
| 87 |
|
{:error, :circuit_open} |
| 88 |
|
|
| 89 |
|
:ok -> |
| 90 |
:-( |
provider_ref = request.provider_reference |
| 91 |
|
|
| 92 |
:-( |
unless provider_ref do |
| 93 |
|
{:error, :no_provider_reference} |
| 94 |
|
else |
| 95 |
:-( |
case apply(adapter, :get_payment_status, [provider_ref]) do |
| 96 |
|
{:ok, result} -> |
| 97 |
:-( |
CircuitBreaker.record_success(request.provider, operation_atom) |
| 98 |
:-( |
handle_success_result(request, result, correlation_id) |
| 99 |
|
|
| 100 |
|
{:error, result} -> |
| 101 |
:-( |
CircuitBreaker.record_failure(request.provider, operation_atom) |
| 102 |
:-( |
handle_failure_result(request, result, correlation_id) |
| 103 |
|
end |
| 104 |
|
end |
| 105 |
|
end |
| 106 |
|
end |
| 107 |
|
|
| 108 |
|
defp execute_operation("cancel", request, _args, correlation_id) do |
| 109 |
:-( |
adapter = resolve_adapter(request.provider) |
| 110 |
:-( |
operation_atom = :cancel_payment |
| 111 |
|
|
| 112 |
:-( |
case CircuitBreaker.check(request.provider, operation_atom) do |
| 113 |
:-( |
{:error, :circuit_open} -> |
| 114 |
|
{:error, :circuit_open} |
| 115 |
|
|
| 116 |
|
:ok -> |
| 117 |
:-( |
provider_ref = request.provider_reference |
| 118 |
|
|
| 119 |
:-( |
unless provider_ref do |
| 120 |
|
{:error, :no_provider_reference} |
| 121 |
|
else |
| 122 |
:-( |
case apply(adapter, :cancel_payment, [provider_ref]) do |
| 123 |
|
{:ok, _result} -> |
| 124 |
:-( |
CircuitBreaker.record_success(request.provider, operation_atom) |
| 125 |
:-( |
{:ok, canceled} = PaymentRequest.cancel(request) |
| 126 |
:-( |
PaymentRequestStore.update(canceled) |
| 127 |
:-( |
emit_audit("cancel_completed", request.request_id, correlation_id, :success, %{}) |
| 128 |
|
:ok |
| 129 |
|
|
| 130 |
|
{:error, result} -> |
| 131 |
:-( |
CircuitBreaker.record_failure(request.provider, operation_atom) |
| 132 |
:-( |
handle_failure_result(request, result, correlation_id) |
| 133 |
|
end |
| 134 |
|
end |
| 135 |
|
end |
| 136 |
|
end |
| 137 |
|
|
| 138 |
|
defp execute_operation("refund", request, args, correlation_id) do |
| 139 |
:-( |
adapter = resolve_adapter(request.provider) |
| 140 |
:-( |
operation_atom = :refund_payment |
| 141 |
:-( |
amount = Map.get(args, "refund_amount", request.amount) |
| 142 |
|
|
| 143 |
:-( |
case CircuitBreaker.check(request.provider, operation_atom) do |
| 144 |
:-( |
{:error, :circuit_open} -> |
| 145 |
|
{:error, :circuit_open} |
| 146 |
|
|
| 147 |
|
:ok -> |
| 148 |
:-( |
provider_ref = request.provider_reference |
| 149 |
|
|
| 150 |
:-( |
unless provider_ref do |
| 151 |
|
{:error, :no_provider_reference} |
| 152 |
|
else |
| 153 |
:-( |
case apply(adapter, :refund_payment, [provider_ref, amount]) do |
| 154 |
|
{:ok, _result} -> |
| 155 |
:-( |
CircuitBreaker.record_success(request.provider, operation_atom) |
| 156 |
:-( |
emit_audit("refund_accepted", request.request_id, correlation_id, :success, %{amount: amount}) |
| 157 |
|
:ok |
| 158 |
|
|
| 159 |
|
{:error, result} -> |
| 160 |
:-( |
CircuitBreaker.record_failure(request.provider, operation_atom) |
| 161 |
:-( |
handle_failure_result(request, result, correlation_id) |
| 162 |
|
end |
| 163 |
|
end |
| 164 |
|
end |
| 165 |
|
end |
| 166 |
|
|
| 167 |
:-( |
defp execute_operation(op, _request, _args, _correlation_id) do |
| 168 |
|
{:error, {:unknown_operation, op}} |
| 169 |
|
end |
| 170 |
|
|
| 171 |
|
# Result handlers |
| 172 |
|
|
| 173 |
|
defp handle_success_result(request, result, correlation_id) do |
| 174 |
13 |
case result.status do |
| 175 |
|
s when s in [:accepted, :pending] -> |
| 176 |
9 |
{:ok, updated} = |
| 177 |
9 |
PaymentRequest.accept(request, result.provider_reference || "", |
| 178 |
9 |
provider_status_code: result.provider_status_code |
| 179 |
|
) |
| 180 |
|
|
| 181 |
9 |
PaymentRequestStore.update(updated) |
| 182 |
|
:ok |
| 183 |
|
|
| 184 |
|
:completed -> |
| 185 |
:-( |
{:ok, completed} = |
| 186 |
|
PaymentRequest.complete(request, |
| 187 |
:-( |
provider_reference: result.provider_reference, |
| 188 |
:-( |
provider_status_code: result.provider_status_code |
| 189 |
|
) |
| 190 |
|
|
| 191 |
:-( |
PaymentRequestStore.update(completed) |
| 192 |
|
|
| 193 |
:-( |
emit_event( |
| 194 |
:-( |
PaymentCompleted.build(request.request_id, |
| 195 |
:-( |
transfer_id: request.transfer_id, |
| 196 |
:-( |
provider_reference: result.provider_reference, |
| 197 |
:-( |
provider_status_code: result.provider_status_code, |
| 198 |
|
correlation_id: correlation_id |
| 199 |
|
) |
| 200 |
|
) |
| 201 |
|
|
| 202 |
:-( |
emit_audit("payment_completed", request.request_id, correlation_id, :success, %{ |
| 203 |
:-( |
provider_reference: result.provider_reference |
| 204 |
|
}) |
| 205 |
|
|
| 206 |
|
:ok |
| 207 |
|
|
| 208 |
|
:unknown -> |
| 209 |
4 |
create_reconciliation_exception(request, :unknown_outcome, "Provider returned unknown status", correlation_id) |
| 210 |
4 |
{:ok, unknown} = PaymentRequest.mark_unknown(request) |
| 211 |
4 |
PaymentRequestStore.update(unknown) |
| 212 |
|
:ok |
| 213 |
|
end |
| 214 |
|
end |
| 215 |
|
|
| 216 |
|
defp handle_failure_result(request, result, correlation_id) do |
| 217 |
3 |
if result.status == :unknown do |
| 218 |
:-( |
create_reconciliation_exception(request, :unknown_outcome, result.provider_status_code || "unknown", correlation_id) |
| 219 |
:-( |
{:ok, unknown} = PaymentRequest.mark_unknown(request) |
| 220 |
:-( |
PaymentRequestStore.update(unknown) |
| 221 |
|
:ok |
| 222 |
|
else |
| 223 |
3 |
reason = result.provider_status_code || "provider_failure" |
| 224 |
3 |
{:ok, failed} = PaymentRequest.fail(request, reason, retryable: result.retryable) |
| 225 |
3 |
PaymentRequestStore.update(failed) |
| 226 |
|
|
| 227 |
3 |
emit_event( |
| 228 |
3 |
PaymentFailed.build(request.request_id, |
| 229 |
3 |
transfer_id: request.transfer_id, |
| 230 |
3 |
provider_reference: result.provider_reference, |
| 231 |
|
failure_reason: reason, |
| 232 |
3 |
retryable: result.retryable, |
| 233 |
|
correlation_id: correlation_id |
| 234 |
|
) |
| 235 |
|
) |
| 236 |
|
|
| 237 |
3 |
emit_audit("payment_failed", request.request_id, correlation_id, :failure, %{ |
| 238 |
|
reason: reason, |
| 239 |
3 |
retryable: result.retryable |
| 240 |
|
}) |
| 241 |
|
|
| 242 |
3 |
if result.retryable do |
| 243 |
|
{:error, {:provider_failure, reason}} |
| 244 |
|
else |
| 245 |
|
:ok |
| 246 |
|
end |
| 247 |
|
end |
| 248 |
|
end |
| 249 |
|
|
| 250 |
4 |
defp create_reconciliation_exception(request, type, description, correlation_id) do |
| 251 |
4 |
exception = |
| 252 |
4 |
IntegrationException.new(type, request.provider, |
| 253 |
4 |
provider_reference: request.provider_reference, |
| 254 |
4 |
transfer_id: request.transfer_id, |
| 255 |
4 |
payment_request_id: request.request_id, |
| 256 |
|
description: description, |
| 257 |
|
correlation_id: correlation_id |
| 258 |
|
) |
| 259 |
|
|
| 260 |
4 |
IntegrationExceptionStore.store(exception) |
| 261 |
|
rescue |
| 262 |
:-( |
_ -> :ok |
| 263 |
|
end |
| 264 |
|
|
| 265 |
|
defp resolve_adapter(provider) do |
| 266 |
|
Application.get_env(:wallet_integrations, :adapters, %{}) |
| 267 |
|
|> Map.get(provider) |
| 268 |
18 |
|> Kernel.||(Application.get_env(:wallet_integrations, :payment_adapter, WalletIntegrations.Adapters.StubAdapter)) |
| 269 |
|
end |
| 270 |
|
|
| 271 |
3 |
defp emit_event(event) do |
| 272 |
3 |
pubsub = Application.get_env(:wallet_integrations, :pubsub, WalletWeb.PubSub) |
| 273 |
3 |
apply(Phoenix.PubSub, :broadcast, [pubsub, "wallet_integrations:events", {:domain_event, event}]) |
| 274 |
|
rescue |
| 275 |
:-( |
_ -> :ok |
| 276 |
|
end |
| 277 |
|
|
| 278 |
3 |
defp emit_audit(action, request_id, correlation_id, outcome, metadata) do |
| 279 |
3 |
audit = |
| 280 |
|
AuditEvent.build(:integrations, action, "payment_request", request_id, |
| 281 |
|
outcome, |
| 282 |
|
correlation_id: correlation_id, |
| 283 |
|
metadata: metadata |
| 284 |
|
) |
| 285 |
|
|
| 286 |
3 |
:telemetry.execute([:wallet_integrations, :audit], %{}, audit) |
| 287 |
|
rescue |
| 288 |
:-( |
_ -> :ok |
| 289 |
|
end |
| 290 |
|
end |