cover/Elixir.WalletIntegrations.Workers.PaymentWorker.html

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
Line Hits Source