diff --git a/CHANGELOG.md b/CHANGELOG.md index b7adf5cf..927b79f0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -30,3 +30,29 @@ manually. architecture; gRPC/OpenAPI/NuGet contract packages marked out of scope. - `docs/architecture.md`: corrected the Service Bus queue inventory to match `AppHost.cs`. - `docs/deploy-ordersphere.md`: fixed step numbering. +- `Order` status transitions (`Confirm`, `MarkShipped`, `MarkDelivered`, `Cancel`) and + `PaymentRecord` transitions return `Result` and reject invalid transitions instead of throwing + or applying them. +- `OrderStatusChangedIntegrationEvent` carries an optional `CustomerId` (additive). +- Stripe client: 30 s timeout with 2 network retries (SDK default 80 s), keeping one payment run + inside the Service Bus lock renewal window. + +### Fixed +- Stripe calls carry deterministic idempotency keys; a redelivered payment request no longer + creates a second PaymentIntent. Only declines and invalid requests are reported as failures; + transient faults are retried through Service Bus redelivery. +- A failed capture releases the authorization and keeps the PaymentIntent id on the payment record. +- The Stripe webhook is reachable (`/webhooks/stripe` on the BFF), finds payments by order or intent, + asks Stripe to retry while the record does not exist yet, rejects a missing signature with 400 + instead of 500, and applies only valid status transitions. +- A payment result for an already cancelled order no longer re-confirms it; a captured payment on a + cancelled order is refunded. +- Admin coupon management is routed through the API Gateway (was 404). + +### Security +- Webhook subscriptions only receive events of their own customer (previously every subscriber of + an event type received all customers' events). +- Webhook target URLs are restricted to public HTTPS hosts, checked on save and again for every + resolved address at connect time; redirects are not followed; failed deliveries no longer store + the target's response body. +- `/bff/login` only accepts a local `returnUrl` (open redirect). diff --git a/docs/operations.md b/docs/operations.md index 17e9f251..57ea1f35 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -379,6 +379,32 @@ the platform metric alert (`Service Bus DLQ`, row above) remains the source of t inbox (`EfInboxStore` only marks on success), so the replayed message is reprocessed normally. 5. If the root cause is not fixed, the message dead-letters again — replay batches are capped (`DlqAdminOptions.ReplayBatchLimit`, default 50) to avoid a replay storm. +6. **Payment queues with Stripe** (`payment-requests`, `order-confirmation-failed`, + `refund-requested`): every Stripe call carries an idempotency key derived from the order or + intent, and Stripe replays the stored result for at least 24 hours. Replay within 24 hours of the + original failure. After that, check the order's PaymentIntent in the Stripe dashboard first + (search by metadata `orderId`): Stripe no longer deduplicates, and a create/refund that already + succeeded would run a second time. + +### Stripe webhook + +Stripe delivers events to the BFF at `POST https:///webhooks/stripe` — the endpoint to +register in the Stripe dashboard (or `stripe listen --forward-to https://localhost:/webhooks/stripe` +locally). The route is anonymous on the BFF and the API Gateway and outside `/api`, so neither the +session policy nor the CSRF check applies; the Payment API authenticates each request by its +`Stripe-Signature` against `Stripe:WebhookSecret`. + +The endpoint answers: + +| Status | Meaning | +|---|---| +| 200 | Reconciled, already reconciled, out of date, or not an OrderSphere intent. | +| 400 | Missing or invalid signature. | +| 503 | The intent carries an OrderSphere `orderId`, but the payment worker has not stored the record yet. Stripe retries later. | + +A contradiction that no automatic transition resolves (for example `payment_intent.succeeded` for a +payment recorded as failed) is logged at `Error` as EventId 5001 *Stripe reconciliation required*. +Resolve it in the Stripe dashboard (refund or cancel the intent) and in the order. --- diff --git a/src/BuildingBlocks/OrderSphere.BuildingBlocks.Contracts/Events/OrderStatusChangedIntegrationEvent.cs b/src/BuildingBlocks/OrderSphere.BuildingBlocks.Contracts/Events/OrderStatusChangedIntegrationEvent.cs index 01ae3736..374dcdc2 100644 --- a/src/BuildingBlocks/OrderSphere.BuildingBlocks.Contracts/Events/OrderStatusChangedIntegrationEvent.cs +++ b/src/BuildingBlocks/OrderSphere.BuildingBlocks.Contracts/Events/OrderStatusChangedIntegrationEvent.cs @@ -8,4 +8,10 @@ public sealed record OrderStatusChangedIntegrationEvent : IntegrationEvent public required string PreviousStatus { get; init; } public required string NewStatus { get; init; } public required string CustomerEmail { get; init; } + + /// + /// Owner of the order. Webhook delivery is scoped to this customer's subscriptions; + /// an event without it (published before the property existed) is delivered to nobody. + /// + public Guid? CustomerId { get; init; } } diff --git a/src/Frontend/OrderSphere.Web/Services/LoginRedirect.cs b/src/Frontend/OrderSphere.Web/Services/LoginRedirect.cs index efd31b01..a969fded 100644 --- a/src/Frontend/OrderSphere.Web/Services/LoginRedirect.cs +++ b/src/Frontend/OrderSphere.Web/Services/LoginRedirect.cs @@ -13,8 +13,12 @@ public static class LoginRedirect { public static string Url(string returnUrl) => $"/bff/login?returnUrl={Uri.EscapeDataString(returnUrl)}"; - /// Full-page navigation to login, returning to the current page afterwards. - public static void Go(NavigationManager navigation) => navigation.NavigateTo(Url(navigation.Uri), forceLoad: true); + /// + /// Full-page navigation to login, returning to the current page afterwards. The return + /// target is sent as a local path: the BFF rejects absolute URLs to prevent open redirects. + /// + public static void Go(NavigationManager navigation) => + navigation.NavigateTo(Url("/" + navigation.ToBaseRelativePath(navigation.Uri)), forceLoad: true); /// True when signed in; otherwise redirects to login and returns false. public static async Task EnsureSignedInAsync(Task? authState, NavigationManager navigation) diff --git a/src/Gateways/OrderSphere.ApiGateway/appsettings.json b/src/Gateways/OrderSphere.ApiGateway/appsettings.json index 27b1adee..b5ec9bf8 100644 --- a/src/Gateways/OrderSphere.ApiGateway/appsettings.json +++ b/src/Gateways/OrderSphere.ApiGateway/appsettings.json @@ -88,6 +88,13 @@ }, "AuthorizationPolicy": "default" }, + "ordering-admin-coupons": { + "ClusterId": "ordering", + "Match": { + "Path": "/api/v1/admin/coupons/{**catch-all}" + }, + "AuthorizationPolicy": "default" + }, "ordering-worker-dlq": { "ClusterId": "ordering-worker", "Match": { @@ -179,6 +186,15 @@ }, "AuthorizationPolicy": "default" }, + "payment-stripe-webhook": { + "ClusterId": "payment", + "Order": -1, + "Match": { + "Path": "/api/v1/payments/webhooks/stripe", + "Methods": [ "POST" ] + }, + "AuthorizationPolicy": "anonymous" + }, "payment": { "ClusterId": "payment", "Match": { diff --git a/src/Gateways/OrderSphere.Bff/Auth/LocalReturnUrl.cs b/src/Gateways/OrderSphere.Bff/Auth/LocalReturnUrl.cs new file mode 100644 index 00000000..437054ee --- /dev/null +++ b/src/Gateways/OrderSphere.Bff/Auth/LocalReturnUrl.cs @@ -0,0 +1,21 @@ +namespace OrderSphere.Bff.Auth; + +/// +/// Restricts the post-login redirect of /bff/login to paths on this origin. +/// returnUrl is a query parameter any link can set; accepting an absolute or +/// protocol-relative URL would send the user to a foreign site after a genuine sign-in. +/// +public static class LocalReturnUrl +{ + /// Returns when it is a local path, otherwise /. + public static string Sanitize(string? returnUrl) => IsLocal(returnUrl) ? returnUrl : "/"; + + // Same rule as ASP.NET Core's IsLocalUrl: one leading '/', not followed by '/' or '\' + // (browsers treat "//host" and "/\host" as protocol-relative), and no control + // characters (browsers strip e.g. a tab, so "/\t/host" collapses to "//host"). + private static bool IsLocal([System.Diagnostics.CodeAnalysis.NotNullWhen(true)] string? url) => + !string.IsNullOrEmpty(url) + && url[0] == '/' + && (url.Length == 1 || (url[1] != '/' && url[1] != '\\')) + && !url.Any(char.IsControl); +} diff --git a/src/Gateways/OrderSphere.Bff/Program.cs b/src/Gateways/OrderSphere.Bff/Program.cs index 3d01acd7..dc367ce7 100644 --- a/src/Gateways/OrderSphere.Bff/Program.cs +++ b/src/Gateways/OrderSphere.Bff/Program.cs @@ -100,12 +100,9 @@ app.UseOrderSphereRequestLogging(); app.MapGet("/bff/login", (HttpContext ctx, string? returnUrl) => -{ - var redirect = string.IsNullOrEmpty(returnUrl) ? "/" : returnUrl; - return Results.Challenge( - new AuthenticationProperties { RedirectUri = redirect }, - [OpenIdConnectDefaults.AuthenticationScheme]); -}); + Results.Challenge( + new AuthenticationProperties { RedirectUri = LocalReturnUrl.Sanitize(returnUrl) }, + [OpenIdConnectDefaults.AuthenticationScheme])); app.MapPost("/bff/logout", (HttpContext _) => Results.SignOut( diff --git a/src/Gateways/OrderSphere.Bff/appsettings.Development.json b/src/Gateways/OrderSphere.Bff/appsettings.Development.json index e7c929c5..be241ea4 100644 --- a/src/Gateways/OrderSphere.Bff/appsettings.Development.json +++ b/src/Gateways/OrderSphere.Bff/appsettings.Development.json @@ -52,6 +52,18 @@ }, "AuthorizationPolicy": "anonymous" }, + "stripe-webhook": { + "ClusterId": "api-gateway", + "Order": -1, + "Match": { + "Path": "/webhooks/stripe", + "Methods": [ "POST" ] + }, + "AuthorizationPolicy": "anonymous", + "Transforms": [ + { "PathSet": "/api/v1/payments/webhooks/stripe" } + ] + }, "to-gateway": { "ClusterId": "api-gateway", "Match": { diff --git a/src/Gateways/OrderSphere.Bff/appsettings.json b/src/Gateways/OrderSphere.Bff/appsettings.json index a5595fbf..4d0d82ac 100644 --- a/src/Gateways/OrderSphere.Bff/appsettings.json +++ b/src/Gateways/OrderSphere.Bff/appsettings.json @@ -61,6 +61,18 @@ }, "AuthorizationPolicy": "anonymous" }, + "stripe-webhook": { + "ClusterId": "api-gateway", + "Order": -1, + "Match": { + "Path": "/webhooks/stripe", + "Methods": [ "POST" ] + }, + "AuthorizationPolicy": "anonymous", + "Transforms": [ + { "PathSet": "/api/v1/payments/webhooks/stripe" } + ] + }, "to-gateway": { "ClusterId": "api-gateway", "Match": { diff --git a/src/Services/Ordering/OrderSphere.Ordering.Application/Features/Order/Admin/CancelOrderCommandHandler.cs b/src/Services/Ordering/OrderSphere.Ordering.Application/Features/Order/Admin/CancelOrderCommandHandler.cs index be0814d0..1cde144e 100644 --- a/src/Services/Ordering/OrderSphere.Ordering.Application/Features/Order/Admin/CancelOrderCommandHandler.cs +++ b/src/Services/Ordering/OrderSphere.Ordering.Application/Features/Order/Admin/CancelOrderCommandHandler.cs @@ -31,11 +31,11 @@ public async Task Handle(CancelOrderCommand request, CancellationToken c // Paid/Shipped → the reservation was confirmed (on-hand stock decremented); restore it. var wasConfirmed = order.Status is not OrderStatus.Created; - try { order.Cancel(); } - catch (InvalidOperationException ex) + var cancel = order.Cancel(); + if (cancel.IsFailure) { - logger.LogWarning(ex, "Cannot cancel order {OrderId} in current status", request.OrderId); - return Result.Failure(OrderErrors.InvalidStatusTransition); + logger.LogWarning("Cannot cancel order {OrderId} in status {Status}", request.OrderId, order.Status); + return cancel; } if (wasConfirmed) diff --git a/src/Services/Ordering/OrderSphere.Ordering.Application/Features/Order/Admin/UpdateOrderStatusCommandHandler.cs b/src/Services/Ordering/OrderSphere.Ordering.Application/Features/Order/Admin/UpdateOrderStatusCommandHandler.cs index 2e0fc768..39e2d2cf 100644 --- a/src/Services/Ordering/OrderSphere.Ordering.Application/Features/Order/Admin/UpdateOrderStatusCommandHandler.cs +++ b/src/Services/Ordering/OrderSphere.Ordering.Application/Features/Order/Admin/UpdateOrderStatusCommandHandler.cs @@ -25,25 +25,18 @@ public async Task Handle(UpdateOrderStatusCommand request, CancellationT if (order is null) return Result.Failure(OrderErrors.OrderNotFoundError); - try + var transition = request.NewStatus switch { - switch (request.NewStatus) - { - case OrderStatus.Shipped: - order.MarkShipped(); - break; - case OrderStatus.Delivered: - order.MarkDelivered(); - break; - default: - return Result.Failure(OrderErrors.InvalidStatusTransition); - } - } - catch (InvalidOperationException ex) + OrderStatus.Shipped => order.MarkShipped(), + OrderStatus.Delivered => order.MarkDelivered(), + _ => Result.Failure(OrderErrors.InvalidStatusTransition) + }; + + if (transition.IsFailure) { - logger.LogWarning(ex, "Invalid status transition for order {OrderId} to {NewStatus}", - request.OrderId, request.NewStatus); - return Result.Failure(OrderErrors.InvalidStatusTransition); + logger.LogWarning("Invalid status transition for order {OrderId} from {Status} to {NewStatus}", + request.OrderId, order.Status, request.NewStatus); + return transition; } await eventStore.AppendAsync(order, cancellationToken); diff --git a/src/Services/Ordering/OrderSphere.Ordering.Domain/Entities/Order.cs b/src/Services/Ordering/OrderSphere.Ordering.Domain/Entities/Order.cs index 6330800c..2990561d 100644 --- a/src/Services/Ordering/OrderSphere.Ordering.Domain/Entities/Order.cs +++ b/src/Services/Ordering/OrderSphere.Ordering.Domain/Entities/Order.cs @@ -1,7 +1,9 @@ using OrderSphere.BuildingBlocks.Abstraction; +using OrderSphere.BuildingBlocks.Primitives; using OrderSphere.BuildingBlocks.StronglyTypedIds; using OrderSphere.BuildingBlocks.ValueObjects; using OrderSphere.Ordering.Domain.Enums; +using OrderSphere.Ordering.Domain.Errors; using OrderSphere.Ordering.Domain.OrderEvents; using OrderSphere.Ordering.Domain.ValueObjects; @@ -93,34 +95,44 @@ public void ApplyDiscount(string couponCode, decimal amount) public void SetShippingCost(decimal amount) => Raise(new ShippingCostSet(amount, DateTime.UtcNow)); - public void Confirm(string trackingNumber) - => Raise(new OrderConfirmed(trackingNumber, DateTime.UtcNow)); + // Transitions are guarded here and only here. Apply stays unguarded because it also folds + // persisted streams, which may already contain sequences these guards now reject. - public void MarkShipped() + /// Marks the order as paid. Only a freshly created order can be confirmed. + public Result Confirm(string trackingNumber) + { + if (Status is not OrderStatus.Created) + return Result.Failure(OrderErrors.InvalidStatusTransition); + + Raise(new OrderConfirmed(trackingNumber, DateTime.UtcNow)); + return Result.Success(); + } + + public Result MarkShipped() { if (Status is not OrderStatus.Paid) - throw new InvalidOperationException( - $"Order can only be marked as shipped when status is Paid (current: {Status})."); + return Result.Failure(OrderErrors.InvalidStatusTransition); Raise(new OrderShipped(DateTime.UtcNow)); + return Result.Success(); } - public void MarkDelivered() + public Result MarkDelivered() { if (Status is not OrderStatus.Shipped) - throw new InvalidOperationException( - $"Order can only be marked as delivered when status is Shipped (current: {Status})."); + return Result.Failure(OrderErrors.InvalidStatusTransition); Raise(new OrderDelivered(DateTime.UtcNow)); + return Result.Success(); } - public void Cancel() + public Result Cancel() { if (Status is OrderStatus.Delivered or OrderStatus.Cancelled) - throw new InvalidOperationException( - $"Order in status {Status} cannot be cancelled."); + return Result.Failure(OrderErrors.InvalidStatusTransition); Raise(new OrderCancelled(DateTime.UtcNow)); + return Result.Success(); } /// Clears the uncommitted buffer once the store has persisted the events. diff --git a/src/Services/Ordering/OrderSphere.Ordering.Worker/Workers/PaymentResultProcessor.cs b/src/Services/Ordering/OrderSphere.Ordering.Worker/Workers/PaymentResultProcessor.cs index a70ac4d0..3a75d709 100644 --- a/src/Services/Ordering/OrderSphere.Ordering.Worker/Workers/PaymentResultProcessor.cs +++ b/src/Services/Ordering/OrderSphere.Ordering.Worker/Workers/PaymentResultProcessor.cs @@ -121,6 +121,12 @@ internal async Task ProcessPaymentResultAsync( // Tracked, so the SaveChanges inside MarkAsProcessedAsync persists it atomically. var saga = await context.OrderSagas.FirstOrDefaultAsync(s => s.CorrelationId == evt.CorrelationId, ct); + // The order may have left Created before the payment result arrived (admin cancellation, + // or a second result for an already settled order). Decide before touching the reservation: + // Catalog's confirm answers 204 even when no active hold is left to confirm. + if (order.Status is not OrderStatus.Created) + return await HandleResultForSettledOrderAsync(evt, order, saga, catalogClient, context, eventStore, inboxStore, ct); + if (evt.Succeeded) { // Confirm the stock reservation (decrements on-hand stock) before persisting the @@ -139,10 +145,17 @@ internal async Task ProcessPaymentResultAsync( throw new InvalidOperationException( $"Reservation confirm failed transiently (delivery {deliveryCount}) for order {order.Id} (correlation {order.CorrelationId}): {confirm.Error.Code}"); - return await CompensateConfirmationFailureAsync(evt, order, saga, catalogClient, context, eventStore, inboxStore, ct); + return await CompensateConfirmationFailureAsync( + evt, order, saga, catalogClient, context, eventStore, inboxStore, + "Reservation confirm conflict (stock can no longer cover the reservation); refunding payment.", + cancelOrder: true, ct); } - order.Confirm(TrackingNumberGenerator.Generate()); + // Status was checked above; a failure here is a programming error, not a business outcome. + var confirmed = order.Confirm(TrackingNumberGenerator.Generate()); + if (confirmed.IsFailure) + throw new InvalidOperationException($"Order {order.Id} could not be confirmed from status {order.Status}."); + saga?.MarkConfirmed(); OrderingMetrics.OrdersConfirmed.Add(1); OrderingMetrics.RecordSagaTransition(nameof(SagaState.Confirmed)); @@ -151,7 +164,10 @@ internal async Task ProcessPaymentResultAsync( } else { - order.Cancel(); + var cancelled = order.Cancel(); + if (cancelled.IsFailure) + throw new InvalidOperationException($"Order {order.Id} could not be cancelled from status {order.Status}."); + saga?.MarkCancelled(evt.FailureReason); OrderingMetrics.OrdersCancelled.Add(1); OrderingMetrics.RecordSagaTransition(nameof(SagaState.Cancelled)); @@ -214,7 +230,8 @@ internal async Task ProcessPaymentResultAsync( OrderId = order.Id.Value, PreviousStatus = "Pending", NewStatus = evt.Succeeded ? "Confirmed" : "Cancelled", - CustomerEmail = evt.CustomerEmail + CustomerEmail = evt.CustomerEmail, + CustomerId = order.CustomerId.Value })); // Stage the new order events and their read projection, then commit everything in one @@ -226,8 +243,54 @@ internal async Task ProcessPaymentResultAsync( } /// - /// Compensates a captured-but-unconfirmable order: cancels the order, releases the reservation - /// (best-effort), advances the saga to , and queues an + /// Handles a payment result for an order that is no longer . + /// A captured payment on a cancelled order is refunded; a failed payment on a cancelled order + /// only closes the saga. For a paid or later order the result is recorded but never changes + /// the order: a success is a duplicate, a failure needs an operator. + /// + private async Task HandleResultForSettledOrderAsync( + PaymentProcessedIntegrationEvent evt, + Order order, + OrderSaga? saga, + ICatalogClient catalogClient, + OrderingDbContext context, + OrderEventStore eventStore, + IInboxStore inboxStore, + CancellationToken ct) + { + switch (order.Status, evt.Succeeded) + { + case (OrderStatus.Cancelled, true): + return await CompensateConfirmationFailureAsync( + evt, order, saga, catalogClient, context, eventStore, inboxStore, + "Order was cancelled before the payment was captured; refunding payment.", + cancelOrder: false, ct); + + case (OrderStatus.Cancelled, false): + saga?.MarkCancelled(evt.FailureReason); + logger.LogInformation("Payment failure for already cancelled order {OrderId} recorded.", order.Id); + break; + + case (_, true): + logger.LogInformation("Payment success for order {OrderId} in status {Status} ignored; already confirmed.", + order.Id, order.Status); + break; + + default: + logger.LogError( + "Payment failure reported for order {OrderId} in status {Status}; order is not cancelled automatically and needs review.", + order.Id, order.Status); + break; + } + + await inboxStore.MarkAsProcessedAsync(evt.Id, nameof(PaymentProcessedIntegrationEvent), ct); + return PaymentResultOutcome.Processed; + } + + /// + /// Compensates a captured payment that cannot be kept: cancels the order (unless it already is), + /// releases the reservation (best-effort), advances the saga to + /// , and queues an /// so Payment refunds the capture. All writes /// commit atomically with the inbox mark in the single SaveChanges inside MarkAsProcessedAsync. /// @@ -239,17 +302,25 @@ private async Task CompensateConfirmationFailureAsync( OrderingDbContext context, OrderEventStore eventStore, IInboxStore inboxStore, + string reason, + bool cancelOrder, CancellationToken ct) { - var reason = "Reservation confirm conflict (stock can no longer cover the reservation); refunding payment."; + if (cancelOrder) + { + var cancelled = order.Cancel(); + if (cancelled.IsFailure) + throw new InvalidOperationException($"Order {order.Id} could not be cancelled from status {order.Status}."); + + OrderingMetrics.OrdersCancelled.Add(1); + } - order.Cancel(); saga?.MarkCompensationPending(reason); - OrderingMetrics.OrdersCancelled.Add(1); OrderingMetrics.RecordSagaTransition(nameof(SagaState.CompensationPending)); - logger.LogError( - "Order {OrderId} confirmation failed after payment was captured; requesting refund. CorrelationId: {CorrelationId}", - order.Id, evt.CorrelationId); + // Warning, not Error: the refund below is the designed outcome and needs no operator. + logger.LogWarning( + "Order {OrderId} cannot keep its captured payment; requesting refund. Reason: {Reason}", + order.Id, reason); // Release the reservation that was never committed (best-effort; the TTL sweeper backstops). var release = await catalogClient.ReleaseReservationAsync(order.CorrelationId, ct); @@ -284,16 +355,20 @@ private async Task CompensateConfirmationFailureAsync( OrderId = order.Id.Value })); - context.AddOutboxMessage( - nameof(OrderStatusChangedIntegrationEvent), - JsonSerializer.Serialize(new OrderStatusChangedIntegrationEvent - { - CorrelationId = evt.CorrelationId, - OrderId = order.Id.Value, - PreviousStatus = "Pending", - NewStatus = "Cancelled", - CustomerEmail = evt.CustomerEmail - })); + if (cancelOrder) + { + context.AddOutboxMessage( + nameof(OrderStatusChangedIntegrationEvent), + JsonSerializer.Serialize(new OrderStatusChangedIntegrationEvent + { + CorrelationId = evt.CorrelationId, + OrderId = order.Id.Value, + PreviousStatus = "Pending", + NewStatus = "Cancelled", + CustomerEmail = evt.CustomerEmail, + CustomerId = order.CustomerId.Value + })); + } await eventStore.AppendAsync(order, ct); await inboxStore.MarkAsProcessedAsync(evt.Id, nameof(PaymentProcessedIntegrationEvent), ct); diff --git a/src/Services/Payment/OrderSphere.Payment.Api/Endpoints/StripeWebhookEndpoints.cs b/src/Services/Payment/OrderSphere.Payment.Api/Endpoints/StripeWebhookEndpoints.cs index a647999c..a94f4ac5 100644 --- a/src/Services/Payment/OrderSphere.Payment.Api/Endpoints/StripeWebhookEndpoints.cs +++ b/src/Services/Payment/OrderSphere.Payment.Api/Endpoints/StripeWebhookEndpoints.cs @@ -1,5 +1,8 @@ using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Options; +using OrderSphere.BuildingBlocks.Security; +using OrderSphere.BuildingBlocks.StronglyTypedIds; +using OrderSphere.Payment.Api.Logging; using OrderSphere.Payment.Domain.Enums; using OrderSphere.Payment.Infrastructure.Persistence; using OrderSphere.Payment.Infrastructure.Providers; @@ -12,6 +15,12 @@ namespace OrderSphere.Payment.Api.Endpoints; /// source of truth for settlement, so these events reconcile the local PaymentRecord /// with the provider. The endpoint is anonymous (authenticated by Stripe signature) and /// idempotent — a re-delivered event is a no-op once the record reached the target state. +/// +/// Only valid status transitions are applied. An event that is merely out of date (a failure +/// notice for a payment that later succeeded) is ignored; one that contradicts the record +/// (a capture for a payment recorded as failed, a different intent) is logged as +/// for an operator. +/// /// public static class StripeWebhookEndpoints { @@ -35,13 +44,22 @@ public static void MapStripeWebhookEndpoints(this WebApplication app) return Results.Ok(); } + // The route is anonymous; without a signature header the SDK does not raise a + // StripeException, so reject it here instead of surfacing a 500. + var signature = request.Headers["Stripe-Signature"].ToString(); + if (string.IsNullOrEmpty(signature)) + { + logger.LogWarning("Stripe webhook rejected: Stripe-Signature header missing."); + return Results.BadRequest(); + } + using var reader = new StreamReader(request.Body); var json = await reader.ReadToEndAsync(ct); Event stripeEvent; try { - stripeEvent = EventUtility.ConstructEvent(json, request.Headers["Stripe-Signature"], secret); + stripeEvent = EventUtility.ConstructEvent(json, signature, secret); } catch (StripeException ex) { @@ -49,13 +67,19 @@ public static void MapStripeWebhookEndpoints(this WebApplication app) return Results.BadRequest(); } - // payment_intent.* carry a PaymentIntent; charge.refunded carries a Charge that - // references the originating PaymentIntent. Both resolve to a PaymentIntent id. - var paymentIntentId = stripeEvent.Data.Object switch + if (stripeEvent.Type is not (PaymentIntentSucceeded or PaymentIntentPaymentFailed or ChargeRefunded)) { - PaymentIntent intent => intent.Id, - Charge charge => charge.PaymentIntentId, - _ => null + logger.LogDebug("Stripe webhook {Type} ignored (not reconciled).", stripeEvent.Type); + return Results.Ok(); + } + + // payment_intent.* carry the PaymentIntent with the metadata the provider set; + // charge.refunded carries a Charge that references the originating PaymentIntent. + var (paymentIntentId, orderId) = stripeEvent.Data.Object switch + { + PaymentIntent intent => (intent.Id, ReadOrderId(intent)), + Charge charge => (charge.PaymentIntentId, (Guid?)null), + _ => (null, null) }; if (paymentIntentId is null) @@ -64,46 +88,100 @@ public static void MapStripeWebhookEndpoints(this WebApplication app) return Results.Ok(); } - return await ReconcileAsync(context, logger, paymentIntentId, stripeEvent.Type, ct); - }); + return await ReconcileAsync(context, logger, stripeEvent.Type, paymentIntentId, orderId, ct); + }) + .AllowAnonymous(); } private static async Task ReconcileAsync( PaymentDbContext context, ILogger logger, - string paymentIntentId, string eventType, + string paymentIntentId, + Guid? orderId, CancellationToken ct) { - var record = await context.Payments.FirstOrDefaultAsync(p => p.TransactionId == paymentIntentId, ct); + var record = await FindRecordAsync(context, paymentIntentId, orderId, ct); if (record is null) { - logger.LogWarning("Stripe webhook {Type}: no PaymentRecord for intent {IntentId}.", eventType, paymentIntentId); + if (orderId is null) + { + logger.LogInformation("Stripe webhook {Type} for intent {IntentId} matches no payment; ignored.", + eventType, paymentIntentId); + return Results.Ok(); + } + + // Created by the payment worker, which stores the record only once the payment has run. + // A non-2xx makes Stripe redeliver later, when the record exists. + logger.LogWarning("Stripe webhook {Type} for order {OrderId} arrived before its payment record; Stripe will retry.", + eventType, orderId); + return Results.StatusCode(StatusCodes.Status503ServiceUnavailable); + } + + if (record.TransactionId is not null && record.TransactionId != paymentIntentId) + { + logger.StripeReconciliationRequired(eventType, paymentIntentId, record.OrderId.Value, record.Status, record.TransactionId); return Results.Ok(); } - switch (eventType) + // Stripe calls anonymously; act in the record's tenant so audit rows are stamped with it. + using var tenantScope = AmbientTenantContext.BeginScope(record.TenantId); + + var transition = eventType switch { - case PaymentIntentSucceeded when record.Status != PaymentStatus.Captured: - record.MarkCaptured(paymentIntentId); - break; - case PaymentIntentPaymentFailed when record.Status is not (PaymentStatus.Failed or PaymentStatus.Captured): - record.MarkFailed("Stripe reported payment failure."); - break; - case ChargeRefunded when record.Status != PaymentStatus.Refunded: - record.MarkRefunded(); - break; - default: - // Already reconciled or an event we don't act on. - return Results.Ok(); + PaymentIntentSucceeded => record.MarkCaptured(paymentIntentId), + PaymentIntentPaymentFailed => record.MarkFailed("Stripe reported payment failure."), + _ => record.MarkRefunded() + }; + + if (transition.IsFailure) + { + if (eventType == PaymentIntentSucceeded && record.Status == PaymentStatus.Failed) + logger.StripeReconciliationRequired(eventType, paymentIntentId, record.OrderId.Value, record.Status, record.TransactionId); + else + logger.LogDebug("Stripe webhook {Type} for intent {IntentId} is out of date; payment is {Status}.", + eventType, paymentIntentId, record.Status); + return Results.Ok(); + } + + if (!context.ChangeTracker.HasChanges()) + { + logger.LogDebug("Stripe webhook {Type} for intent {IntentId} already reconciled.", eventType, paymentIntentId); + return Results.Ok(); } await context.SaveChangesAsync(ct); - logger.LogInformation("Stripe webhook {Type} reconciled PaymentRecord for intent {IntentId}.", - eventType, paymentIntentId); + logger.LogInformation("Stripe webhook {Type} reconciled the payment for order {OrderId}.", + eventType, record.OrderId); return Results.Ok(); } + /// + /// Finds the record by the orderId metadata (unique), falling back to the intent id. + /// The request carries no tenant, so the tenant query filter is bypassed; soft-deleted rows + /// stay excluded explicitly because IgnoreQueryFilters drops that filter too. + /// + private static async Task FindRecordAsync( + PaymentDbContext context, string paymentIntentId, Guid? orderId, CancellationToken ct) + { + var payments = context.Payments.IgnoreQueryFilters().Where(p => !p.IsDeleted); + + if (orderId is { } id + && await payments.FirstOrDefaultAsync(p => p.OrderId == OrderId.From(id), ct) is { } byOrder) + { + return byOrder; + } + + return await payments.FirstOrDefaultAsync(p => p.TransactionId == paymentIntentId, ct); + } + + private static Guid? ReadOrderId(PaymentIntent intent) => + intent.Metadata is not null + && intent.Metadata.TryGetValue("orderId", out var raw) + && Guid.TryParse(raw, out var orderId) + ? orderId + : null; + /// Logger category marker for the webhook endpoint. public sealed class StripeWebhookMarker; } diff --git a/src/Services/Payment/OrderSphere.Payment.Api/Logging/PaymentApiLog.cs b/src/Services/Payment/OrderSphere.Payment.Api/Logging/PaymentApiLog.cs new file mode 100644 index 00000000..951cc7fa --- /dev/null +++ b/src/Services/Payment/OrderSphere.Payment.Api/Logging/PaymentApiLog.cs @@ -0,0 +1,27 @@ +using OrderSphere.Payment.Domain.Enums; + +namespace OrderSphere.Payment.Api.Logging; + +/// +/// Source-generated log methods for the Payment API. EventIds 5001-5099 (Payment range +/// 5000-5999, see docs/logging.md). +/// +internal static partial class PaymentApiLog +{ + /// + /// Stripe and the local record disagree in a way no automatic transition resolves — for + /// example Stripe reports a capture for a payment recorded as failed, so the customer may have + /// been charged for a cancelled order. An operator must reconcile it in the Stripe dashboard. + /// + [LoggerMessage( + EventId = 5001, + Level = LogLevel.Error, + Message = "Stripe reconciliation required: {EventType} for intent {IntentId} contradicts the payment for order {OrderId} (status {Status}, transaction {TransactionId}).")] + public static partial void StripeReconciliationRequired( + this ILogger logger, + string eventType, + string intentId, + Guid orderId, + PaymentStatus status, + string? transactionId); +} diff --git a/src/Services/Payment/OrderSphere.Payment.Domain/Entities/PaymentRecord.cs b/src/Services/Payment/OrderSphere.Payment.Domain/Entities/PaymentRecord.cs index f62f322e..64446ecb 100644 --- a/src/Services/Payment/OrderSphere.Payment.Domain/Entities/PaymentRecord.cs +++ b/src/Services/Payment/OrderSphere.Payment.Domain/Entities/PaymentRecord.cs @@ -1,8 +1,10 @@ using OrderSphere.BuildingBlocks.Abstraction; +using OrderSphere.BuildingBlocks.Primitives; using OrderSphere.BuildingBlocks.StronglyTypedIds; using OrderSphere.BuildingBlocks.ValueObjects; using OrderSphere.Payment.Domain.DomainEvents; using OrderSphere.Payment.Domain.Enums; +using OrderSphere.Payment.Domain.Errors; namespace OrderSphere.Payment.Domain.Entities; @@ -41,30 +43,61 @@ public PaymentRecord( Status = PaymentStatus.Pending; } - public void MarkAuthorized(string transactionId) + // Valid transitions: Pending → Authorized → Captured → Refunded, and Pending|Authorized → Failed. + // Repeating the current state (same transaction id) succeeds without a new domain event, so + // redelivered messages and webhooks are idempotent. Capture may replace the authorization id: + // some providers issue a separate capture reference. + + public Result MarkAuthorized(string transactionId) { + if (Status == PaymentStatus.Authorized && TransactionId == transactionId) + return Result.Success(); + if (Status != PaymentStatus.Pending) + return Result.Failure(PaymentErrors.InvalidStatusTransition); + TransactionId = transactionId; Status = PaymentStatus.Authorized; RaiseDomainEvent(new PaymentAuthorizedDomainEvent(Id, OrderId, transactionId)); + return Result.Success(); } - public void MarkCaptured(string transactionId) + public Result MarkCaptured(string transactionId) { + if (Status == PaymentStatus.Captured) + return TransactionId == transactionId + ? Result.Success() + : Result.Failure(PaymentErrors.InvalidStatusTransition); + if (Status is not (PaymentStatus.Pending or PaymentStatus.Authorized)) + return Result.Failure(PaymentErrors.InvalidStatusTransition); + TransactionId = transactionId; Status = PaymentStatus.Captured; RaiseDomainEvent(new PaymentCapturedDomainEvent(Id, OrderId, transactionId)); + return Result.Success(); } - public void MarkFailed(string reason) + public Result MarkFailed(string reason) { + if (Status == PaymentStatus.Failed) + return Result.Success(); + if (Status is not (PaymentStatus.Pending or PaymentStatus.Authorized)) + return Result.Failure(PaymentErrors.InvalidStatusTransition); + FailureReason = reason; Status = PaymentStatus.Failed; RaiseDomainEvent(new PaymentFailedDomainEvent(Id, OrderId, reason)); + return Result.Success(); } - public void MarkRefunded() + public Result MarkRefunded() { + if (Status == PaymentStatus.Refunded) + return Result.Success(); + if (Status != PaymentStatus.Captured) + return Result.Failure(PaymentErrors.InvalidStatusTransition); + Status = PaymentStatus.Refunded; + return Result.Success(); } /// GDPR right-to-erasure: overwrites the customer email, keeping the payment diff --git a/src/Services/Payment/OrderSphere.Payment.Domain/Errors/PaymentErrors.cs b/src/Services/Payment/OrderSphere.Payment.Domain/Errors/PaymentErrors.cs index 1eed3af6..9da44988 100644 --- a/src/Services/Payment/OrderSphere.Payment.Domain/Errors/PaymentErrors.cs +++ b/src/Services/Payment/OrderSphere.Payment.Domain/Errors/PaymentErrors.cs @@ -9,5 +9,7 @@ public static class PaymentErrors public static readonly Error AuthorizationFailed = new("Payment.AuthorizationFailed", "Payment authorization failed.", ErrorType.Failure); public static readonly Error CaptureFailed = new("Payment.CaptureFailed", "Payment capture failed.", ErrorType.Failure); public static readonly Error RefundFailed = new("Payment.RefundFailed", "Payment refund failed.", ErrorType.Failure); + public static readonly Error VoidFailed = new("Payment.VoidFailed", "Releasing the payment authorization failed.", ErrorType.Failure); + public static readonly Error InvalidStatusTransition = new("Payment.InvalidStatusTransition", "The payment's current status does not allow this transition.", ErrorType.Conflict); public static readonly Error DuplicatePayment = new("Payment.Duplicate", "A payment for this order already exists.", ErrorType.Conflict); } diff --git a/src/Services/Payment/OrderSphere.Payment.Infrastructure/DependencyInjection.cs b/src/Services/Payment/OrderSphere.Payment.Infrastructure/DependencyInjection.cs index 1696458a..20c7eeef 100644 --- a/src/Services/Payment/OrderSphere.Payment.Infrastructure/DependencyInjection.cs +++ b/src/Services/Payment/OrderSphere.Payment.Infrastructure/DependencyInjection.cs @@ -37,7 +37,14 @@ public static IServiceCollection AddPaymentInfrastructure( var stripeApiKey = configuration.GetSection(StripeOptions.SectionName)["ApiKey"]; if (!string.IsNullOrWhiteSpace(stripeApiKey)) { - services.AddSingleton(new Stripe.StripeClient(stripeApiKey)); + // The SDK default (80 s timeout, 2 retries) lets one payment run far past the 5-minute + // Service Bus lock renewal. 30 s × 3 attempts keeps create + capture + cancel inside it; + // the idempotency keys make the SDK's own retries safe. + services.AddSingleton(new Stripe.StripeClient( + stripeApiKey, + httpClient: new Stripe.SystemNetHttpClient( + new HttpClient { Timeout = TimeSpan.FromSeconds(30) }, + maxNetworkRetries: 2))); services.AddSingleton(); } else diff --git a/src/Services/Payment/OrderSphere.Payment.Infrastructure/EntityConfigurations/PaymentRecordConfiguration.cs b/src/Services/Payment/OrderSphere.Payment.Infrastructure/EntityConfigurations/PaymentRecordConfiguration.cs index 5d647551..fdfe4534 100644 --- a/src/Services/Payment/OrderSphere.Payment.Infrastructure/EntityConfigurations/PaymentRecordConfiguration.cs +++ b/src/Services/Payment/OrderSphere.Payment.Infrastructure/EntityConfigurations/PaymentRecordConfiguration.cs @@ -15,6 +15,8 @@ public void Configure(EntityTypeBuilder builder) builder.HasIndex(p => p.OrderId).IsUnique(); builder.HasIndex(p => p.CorrelationId); + // Stripe webhooks for charges reference only the PaymentIntent id. + builder.HasIndex(p => p.TransactionId); // Money mapped onto the existing "Amount"/"Currency" columns — no schema change. builder.ComplexProperty(p => p.Amount, b => diff --git a/src/Services/Payment/OrderSphere.Payment.Infrastructure/Migrations/20260923194654_AddPaymentTransactionIdIndex.Designer.cs b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Migrations/20260923194654_AddPaymentTransactionIdIndex.Designer.cs new file mode 100644 index 00000000..c8d81509 --- /dev/null +++ b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Migrations/20260923194654_AddPaymentTransactionIdIndex.Designer.cs @@ -0,0 +1,217 @@ +// +using System; +using System.Collections.Generic; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Migrations; +using Microsoft.EntityFrameworkCore.Storage.ValueConversion; +using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata; +using OrderSphere.Payment.Infrastructure.Persistence; + +#nullable disable + +namespace OrderSphere.Payment.Infrastructure.Migrations +{ + [DbContext(typeof(PaymentDbContext))] + [Migration("20260923194654_AddPaymentTransactionIdIndex")] + partial class AddPaymentTransactionIdIndex + { + /// + protected override void BuildTargetModel(ModelBuilder modelBuilder) + { +#pragma warning disable 612, 618 + modelBuilder + .HasAnnotation("ProductVersion", "10.0.11") + .HasAnnotation("Relational:MaxIdentifierLength", 63); + + NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder); + + modelBuilder.Entity("OrderSphere.BuildingBlocks.Auditing.AuditLogEntry", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("uuid"); + + b.Property("Action") + .IsRequired() + .HasMaxLength(20) + .HasColumnType("character varying(20)"); + + b.Property("ChangedBy") + .HasMaxLength(256) + .HasColumnType("character varying(256)"); + + b.Property("Changes") + .IsRequired() + .HasColumnType("text"); + + b.Property("EntityId") + .IsRequired() + .HasMaxLength(256) + .HasColumnType("character varying(256)"); + + b.Property("EntityType") + .IsRequired() + .HasMaxLength(256) + .HasColumnType("character varying(256)"); + + b.Property("OccurredAt") + .HasColumnType("timestamp with time zone"); + + b.HasKey("Id"); + + b.HasIndex("EntityType", "EntityId", "OccurredAt"); + + b.ToTable("audit_log_entries", (string)null); + }); + + modelBuilder.Entity("OrderSphere.BuildingBlocks.EventBus.Inbox.InboxMessage", b => + { + b.Property("EventId") + .ValueGeneratedOnAdd() + .HasColumnType("uuid"); + + b.Property("EventType") + .IsRequired() + .HasMaxLength(256) + .HasColumnType("character varying(256)"); + + b.Property("ProcessedAt") + .HasColumnType("timestamp with time zone"); + + b.HasKey("EventId"); + + b.HasIndex("ProcessedAt"); + + b.ToTable("inbox_messages", (string)null); + }); + + modelBuilder.Entity("OrderSphere.BuildingBlocks.EventBus.Outbox.OutboxMessage", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("uuid"); + + b.Property("Content") + .IsRequired() + .HasColumnType("text"); + + b.Property("CorrelationId") + .HasMaxLength(128) + .HasColumnType("character varying(128)"); + + b.Property("Error") + .HasColumnType("text"); + + b.Property("OccurredAt") + .HasColumnType("timestamp with time zone"); + + b.Property("ProcessedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("RetryCount") + .ValueGeneratedOnAdd() + .HasColumnType("integer") + .HasDefaultValue(0); + + b.Property("TraceParent") + .HasMaxLength(55) + .HasColumnType("character varying(55)"); + + b.Property("Type") + .IsRequired() + .HasMaxLength(256) + .HasColumnType("character varying(256)"); + + b.Property("xmin") + .IsConcurrencyToken() + .ValueGeneratedOnAddOrUpdate() + .HasColumnType("xid") + .HasColumnName("xmin"); + + b.HasKey("Id"); + + b.HasIndex("RetryCount"); + + b.HasIndex("ProcessedAt", "OccurredAt"); + + b.ToTable("outbox_messages", (string)null); + }); + + modelBuilder.Entity("OrderSphere.Payment.Domain.Entities.PaymentRecord", b => + { + b.Property("Id") + .HasColumnType("uuid"); + + b.Property("CorrelationId") + .HasColumnType("uuid"); + + b.Property("CreatedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("CustomerEmail") + .IsRequired() + .HasMaxLength(256) + .HasColumnType("character varying(256)"); + + b.Property("FailureReason") + .HasMaxLength(1024) + .HasColumnType("character varying(1024)"); + + b.Property("IsDeleted") + .HasColumnType("boolean"); + + b.Property("OrderId") + .HasColumnType("uuid"); + + b.Property("PaymentMethod") + .IsRequired() + .HasMaxLength(50) + .HasColumnType("character varying(50)"); + + b.Property("Status") + .IsRequired() + .HasMaxLength(20) + .HasColumnType("character varying(20)"); + + b.Property("TenantId") + .HasColumnType("uuid"); + + b.Property("TransactionId") + .HasMaxLength(256) + .HasColumnType("character varying(256)"); + + b.Property("UpdatedAt") + .HasColumnType("timestamp with time zone"); + + b.ComplexProperty(typeof(Dictionary), "Amount", "OrderSphere.Payment.Domain.Entities.PaymentRecord.Amount#Money", b1 => + { + b1.IsRequired(); + + b1.Property("Amount") + .HasPrecision(18, 2) + .HasColumnType("numeric(18,2)") + .HasColumnName("Amount"); + + b1.Property("Currency") + .IsRequired() + .HasMaxLength(3) + .HasColumnType("character varying(3)") + .HasColumnName("Currency"); + }); + + b.HasKey("Id"); + + b.HasIndex("CorrelationId"); + + b.HasIndex("OrderId") + .IsUnique(); + + b.HasIndex("TransactionId"); + + b.ToTable("payments", (string)null); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/src/Services/Payment/OrderSphere.Payment.Infrastructure/Migrations/20260923194654_AddPaymentTransactionIdIndex.cs b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Migrations/20260923194654_AddPaymentTransactionIdIndex.cs new file mode 100644 index 00000000..62ff2cb0 --- /dev/null +++ b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Migrations/20260923194654_AddPaymentTransactionIdIndex.cs @@ -0,0 +1,27 @@ +using Microsoft.EntityFrameworkCore.Migrations; + +#nullable disable + +namespace OrderSphere.Payment.Infrastructure.Migrations +{ + /// + public partial class AddPaymentTransactionIdIndex : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.CreateIndex( + name: "IX_payments_TransactionId", + table: "payments", + column: "TransactionId"); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropIndex( + name: "IX_payments_TransactionId", + table: "payments"); + } + } +} diff --git a/src/Services/Payment/OrderSphere.Payment.Infrastructure/Migrations/PaymentDbContextModelSnapshot.cs b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Migrations/PaymentDbContextModelSnapshot.cs index 931be44e..17eee1fe 100644 --- a/src/Services/Payment/OrderSphere.Payment.Infrastructure/Migrations/PaymentDbContextModelSnapshot.cs +++ b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Migrations/PaymentDbContextModelSnapshot.cs @@ -204,6 +204,8 @@ protected override void BuildModel(ModelBuilder modelBuilder) b.HasIndex("OrderId") .IsUnique(); + b.HasIndex("TransactionId"); + b.ToTable("payments", (string)null); }); #pragma warning restore 612, 618 diff --git a/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/CreditCardPaymentProvider.cs b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/CreditCardPaymentProvider.cs index 2a42db8b..d7e3d17f 100644 --- a/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/CreditCardPaymentProvider.cs +++ b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/CreditCardPaymentProvider.cs @@ -25,7 +25,14 @@ public Task> CaptureAsync(string transactionId, de return Task.FromResult(Result.Success(new PaymentProviderResult(transactionId))); } - public Task RefundAsync(string transactionId, decimal amount, CancellationToken ct = default) + public Task VoidAsync(string transactionId, CancellationToken ct = default) + { + logger.LogInformation("CreditCard authorization voided. TransactionId: {TransactionId}", transactionId); + + return Task.FromResult(Result.Success()); + } + + public Task RefundAsync(string transactionId, decimal amount, string refundReference, CancellationToken ct = default) { logger.LogInformation("CreditCard payment refunded. TransactionId: {TransactionId}, Amount: {Amount}", transactionId, amount); diff --git a/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/IPaymentProvider.cs b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/IPaymentProvider.cs index d29eeba6..4c9b970f 100644 --- a/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/IPaymentProvider.cs +++ b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/IPaymentProvider.cs @@ -2,18 +2,36 @@ namespace OrderSphere.Payment.Infrastructure.Providers; +/// +/// A payment service provider. Business outcomes (declined card, rejected request) are returned +/// as failures. Transient faults — network errors, timeouts, rate limits, +/// provider outages — are thrown, so the calling worker abandons the message and Service Bus +/// redelivers it; providers must make every call idempotent so that the redelivery completes or +/// replays the original request instead of repeating it. +/// public interface IPaymentProvider { string MethodName { get; } Task> AuthorizeAsync(PaymentRequest request, CancellationToken ct = default); Task> CaptureAsync(string transactionId, decimal amount, CancellationToken ct = default); - Task RefundAsync(string transactionId, decimal amount, CancellationToken ct = default); + + /// Releases an authorization that will not be captured. + Task VoidAsync(string transactionId, CancellationToken ct = default); + + /// + /// Refunds of a captured payment. + /// identifies the business reason (a return request, a failed order confirmation) and makes the + /// refund idempotent per reason, so partial refunds for different returns stay distinct. + /// + Task RefundAsync(string transactionId, decimal amount, string refundReference, CancellationToken ct = default); } public sealed record PaymentRequest( Guid OrderId, decimal Amount, string Currency, - string CustomerEmail); + string CustomerEmail, + Guid TenantId, + Guid CorrelationId); public sealed record PaymentProviderResult(string TransactionId); diff --git a/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/InvoicePaymentProvider.cs b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/InvoicePaymentProvider.cs index a02d7a63..327bc7a1 100644 --- a/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/InvoicePaymentProvider.cs +++ b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/InvoicePaymentProvider.cs @@ -24,7 +24,14 @@ public Task> CaptureAsync(string transactionId, de return Task.FromResult(Result.Success(new PaymentProviderResult(transactionId))); } - public Task RefundAsync(string transactionId, decimal amount, CancellationToken ct = default) + public Task VoidAsync(string transactionId, CancellationToken ct = default) + { + logger.LogInformation("Invoice authorization voided. TransactionId: {TransactionId}", transactionId); + + return Task.FromResult(Result.Success()); + } + + public Task RefundAsync(string transactionId, decimal amount, string refundReference, CancellationToken ct = default) { logger.LogInformation("Invoice payment refunded. TransactionId: {TransactionId}, Amount: {Amount}", transactionId, amount); diff --git a/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/PayPalPaymentProvider.cs b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/PayPalPaymentProvider.cs index ed0404f7..5fa9e154 100644 --- a/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/PayPalPaymentProvider.cs +++ b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/PayPalPaymentProvider.cs @@ -25,7 +25,14 @@ public Task> CaptureAsync(string transactionId, de return Task.FromResult(Result.Success(new PaymentProviderResult(transactionId))); } - public Task RefundAsync(string transactionId, decimal amount, CancellationToken ct = default) + public Task VoidAsync(string transactionId, CancellationToken ct = default) + { + logger.LogInformation("PayPal authorization voided. TransactionId: {TransactionId}", transactionId); + + return Task.FromResult(Result.Success()); + } + + public Task RefundAsync(string transactionId, decimal amount, string refundReference, CancellationToken ct = default) { logger.LogInformation("PayPal payment refunded. TransactionId: {TransactionId}, Amount: {Amount}", transactionId, amount); diff --git a/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/StripePaymentProvider.cs b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/StripePaymentProvider.cs index 9baf3d4c..775419b9 100644 --- a/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/StripePaymentProvider.cs +++ b/src/Services/Payment/OrderSphere.Payment.Infrastructure/Providers/StripePaymentProvider.cs @@ -1,3 +1,4 @@ +using System.Net; using Microsoft.Extensions.Logging; using OrderSphere.BuildingBlocks.Primitives; using OrderSphere.Payment.Domain.Errors; @@ -6,24 +7,40 @@ namespace OrderSphere.Payment.Infrastructure.Providers; /// -/// Real payment provider backed by Stripe (test mode). Maps the three-stage -/// Authorize → Capture → Refund flow onto Stripe PaymentIntents with manual capture. -/// Registered under the "CreditCard" method name so existing checkout routes here without a -/// UI contract change. Stripe SDK exceptions are mapped to failures — -/// business outcomes never surface as exceptions. +/// Real payment provider backed by Stripe (test mode). Maps the Authorize → Capture → Refund flow +/// onto PaymentIntents with manual capture; cancels an uncaptured intent. +/// Registered under the "CreditCard" method name so existing checkout routes here without a UI +/// contract change. +/// +/// Every call carries a deterministic derived from the +/// order or intent. Stripe stores the first response per key for at least 24 hours and replays it, +/// so a Service Bus redelivery — including one after a timeout on a request that did succeed — +/// converges on the original intent, capture, cancellation or refund instead of repeating it. +/// +/// +/// Error mapping: declines and invalid requests (card_error, invalid_request_error +/// with 400/402/404) are business outcomes and return a failure. Everything +/// else — rate limits, authentication, idempotency conflicts, 5xx, network errors and timeouts — +/// propagates, so the worker abandons the message and the redelivery retries under the same key. +/// /// internal sealed class StripePaymentProvider( IStripeClient stripeClient, ILogger logger) : IPaymentProvider { + private const string RequiresCapture = "requires_capture"; + private const string Succeeded = "succeeded"; + private const string Canceled = "canceled"; + public string MethodName => "CreditCard"; public async Task> AuthorizeAsync(PaymentRequest request, CancellationToken ct = default) { var service = new PaymentIntentService(stripeClient); + PaymentIntent intent; try { - var intent = await service.CreateAsync(new PaymentIntentCreateOptions + intent = await service.CreateAsync(new PaymentIntentCreateOptions { Amount = ToMinorUnits(request.Amount), Currency = request.Currency.ToLowerInvariant(), @@ -38,20 +55,37 @@ public async Task> AuthorizeAsync(PaymentRequest r AllowRedirects = "never" }, ReceiptEmail = request.CustomerEmail, - Metadata = new Dictionary { ["orderId"] = request.OrderId.ToString() } - }, cancellationToken: ct); - - logger.LogInformation( - "Stripe PaymentIntent {IntentId} authorized for order {OrderId} (status {Status}).", - intent.Id, request.OrderId, intent.Status); - - return Result.Success(new PaymentProviderResult(intent.Id)); + // orderId lets the webhook find the record even before the worker has stored the + // intent id; tenantId and correlationId tie the intent back to the saga. + Metadata = new Dictionary + { + ["orderId"] = request.OrderId.ToString(), + ["tenantId"] = request.TenantId.ToString(), + ["correlationId"] = request.CorrelationId.ToString() + } + }, Idempotent($"os-pi-create-{request.OrderId:N}"), ct); + } + catch (StripeException ex) when (IsPermanent(ex)) + { + logger.LogWarning(ex, "Stripe authorization declined for order {OrderId}. Code: {StripeErrorCode}", + request.OrderId, ex.StripeError?.Code); + return Result.Failure(PaymentErrors.AuthorizationFailed); } - catch (StripeException ex) + + if (intent.Status != RequiresCapture) { - logger.LogWarning(ex, "Stripe authorization failed for order {OrderId}.", request.OrderId); + // Confirmed with manual capture, anything but requires_capture (e.g. requires_action + // for 3-D Secure) means no funds are held. Release the intent and report a decline. + await VoidAsync(intent.Id, ct); + logger.LogWarning( + "Stripe PaymentIntent {IntentId} for order {OrderId} ended in status {Status} instead of requires_capture.", + intent.Id, request.OrderId, intent.Status); return Result.Failure(PaymentErrors.AuthorizationFailed); } + + logger.LogInformation("Stripe PaymentIntent {IntentId} authorized for order {OrderId}.", + intent.Id, request.OrderId); + return Result.Success(new PaymentProviderResult(intent.Id)); } public async Task> CaptureAsync(string transactionId, decimal amount, CancellationToken ct = default) @@ -59,19 +93,58 @@ public async Task> CaptureAsync(string transaction var service = new PaymentIntentService(stripeClient); try { - var intent = await service.CaptureAsync(transactionId, new PaymentIntentCaptureOptions(), cancellationToken: ct); - logger.LogInformation("Stripe PaymentIntent {IntentId} capture requested (status {Status}).", - intent.Id, intent.Status); + var intent = await service.CaptureAsync( + transactionId, new PaymentIntentCaptureOptions(), Idempotent($"os-pi-capture-{transactionId}"), ct); + logger.LogInformation("Stripe PaymentIntent {IntentId} captured (status {Status}).", intent.Id, intent.Status); return Result.Success(new PaymentProviderResult(intent.Id)); } - catch (StripeException ex) + catch (StripeException ex) when (IsUnexpectedState(ex)) { - logger.LogWarning(ex, "Stripe capture failed for intent {IntentId}.", transactionId); + // The intent moved on without this call — typically captured by an earlier attempt + // whose idempotency key has expired. Stripe's current state decides. + var current = await service.GetAsync(transactionId, cancellationToken: ct); + if (current.Status == Succeeded) + return Result.Success(new PaymentProviderResult(current.Id)); + + logger.LogWarning("Stripe capture of intent {IntentId} rejected; intent is {Status}.", transactionId, current.Status); + return Result.Failure(PaymentErrors.CaptureFailed); + } + catch (StripeException ex) when (IsPermanent(ex)) + { + logger.LogWarning(ex, "Stripe capture declined for intent {IntentId}. Code: {StripeErrorCode}", + transactionId, ex.StripeError?.Code); return Result.Failure(PaymentErrors.CaptureFailed); } } - public async Task RefundAsync(string transactionId, decimal amount, CancellationToken ct = default) + public async Task VoidAsync(string transactionId, CancellationToken ct = default) + { + var service = new PaymentIntentService(stripeClient); + try + { + await service.CancelAsync( + transactionId, new PaymentIntentCancelOptions(), Idempotent($"os-pi-cancel-{transactionId}"), ct); + logger.LogInformation("Stripe PaymentIntent {IntentId} canceled; authorization released.", transactionId); + return Result.Success(); + } + catch (StripeException ex) when (IsUnexpectedState(ex)) + { + var current = await service.GetAsync(transactionId, cancellationToken: ct); + if (current.Status == Canceled) + return Result.Success(); + + logger.LogWarning("Stripe cancel of intent {IntentId} rejected; intent is {Status}.", transactionId, current.Status); + return Result.Failure(PaymentErrors.VoidFailed); + } + catch (StripeException ex) when (IsPermanent(ex)) + { + logger.LogWarning(ex, "Stripe cancel rejected for intent {IntentId}. Code: {StripeErrorCode}", + transactionId, ex.StripeError?.Code); + return Result.Failure(PaymentErrors.VoidFailed); + } + } + + public async Task RefundAsync(string transactionId, decimal amount, string refundReference, CancellationToken ct = default) { var service = new RefundService(stripeClient); try @@ -80,17 +153,35 @@ await service.CreateAsync(new RefundCreateOptions { PaymentIntent = transactionId, Amount = ToMinorUnits(amount) - }, cancellationToken: ct); + }, Idempotent($"os-refund-{transactionId}-{refundReference}"), ct); logger.LogInformation("Stripe refund issued for intent {IntentId}, amount {Amount}.", transactionId, amount); return Result.Success(); } - catch (StripeException ex) + catch (StripeException ex) when (ex.StripeError?.Code == "charge_already_refunded") + { + // Refunded by an earlier attempt whose idempotency key has expired: the goal is met. + logger.LogWarning("Stripe intent {IntentId} was already fully refunded.", transactionId); + return Result.Success(); + } + catch (StripeException ex) when (IsPermanent(ex)) { - logger.LogWarning(ex, "Stripe refund failed for intent {IntentId}.", transactionId); + logger.LogWarning(ex, "Stripe refund rejected for intent {IntentId}. Code: {StripeErrorCode}", + transactionId, ex.StripeError?.Code); return Result.Failure(PaymentErrors.RefundFailed); } } + private static RequestOptions Idempotent(string key) => new() { IdempotencyKey = key }; + + // 429 (rate limit, lock timeout) also uses invalid_request_error; the status check keeps it + // on the retry path together with 401/403 and 409. + private static bool IsPermanent(StripeException ex) => + ex.StripeError?.Type is "card_error" or "invalid_request_error" + && ex.HttpStatusCode is HttpStatusCode.BadRequest or HttpStatusCode.PaymentRequired or HttpStatusCode.NotFound; + + private static bool IsUnexpectedState(StripeException ex) => + ex.StripeError?.Code == "payment_intent_unexpected_state"; + // Stripe expects amounts in the currency's minor unit (e.g. cents). Two-decimal // currencies (EUR/USD) are the showcase scope; rounding away from zero matches Money. private static long ToMinorUnits(decimal amount) => diff --git a/src/Services/Payment/OrderSphere.Payment.Worker/Workers/OrderConfirmationFailedProcessor.cs b/src/Services/Payment/OrderSphere.Payment.Worker/Workers/OrderConfirmationFailedProcessor.cs index 7b62fab5..7542b628 100644 --- a/src/Services/Payment/OrderSphere.Payment.Worker/Workers/OrderConfirmationFailedProcessor.cs +++ b/src/Services/Payment/OrderSphere.Payment.Worker/Workers/OrderConfirmationFailedProcessor.cs @@ -132,7 +132,8 @@ internal async Task ProcessConfirmationFailureAsync( throw new InvalidOperationException( $"No provider for method '{payment.PaymentMethod}' to refund order {evt.OrderId}."); - var refund = await provider.RefundAsync(payment.TransactionId!, payment.Amount, ct); + // One confirmation failure per order, so a fixed reference makes the refund idempotent. + var refund = await provider.RefundAsync(payment.TransactionId!, payment.Amount, "ocf", ct); if (refund.IsFailure) // Throw so Service Bus redelivers — a transient refund failure must not silently drop. throw new InvalidOperationException( diff --git a/src/Services/Payment/OrderSphere.Payment.Worker/Workers/PaymentProcessor.cs b/src/Services/Payment/OrderSphere.Payment.Worker/Workers/PaymentProcessor.cs index 7f608a98..df63b002 100644 --- a/src/Services/Payment/OrderSphere.Payment.Worker/Workers/PaymentProcessor.cs +++ b/src/Services/Payment/OrderSphere.Payment.Worker/Workers/PaymentProcessor.cs @@ -6,8 +6,10 @@ using OrderSphere.BuildingBlocks.Contracts.Events; using OrderSphere.BuildingBlocks.EventBus.AzureServiceBus; using OrderSphere.BuildingBlocks.EventBus.Inbox; +using OrderSphere.BuildingBlocks.Primitives; using OrderSphere.BuildingBlocks.StronglyTypedIds; using OrderSphere.Payment.Domain.Entities; +using OrderSphere.Payment.Domain.Enums; using OrderSphere.Payment.Infrastructure.Persistence; using OrderSphere.Payment.Infrastructure.Providers; @@ -80,7 +82,8 @@ await args.DeadLetterMessageAsync(args.Message, } var sw = Stopwatch.StartNew(); - var succeeded = await ProcessPaymentAsync(evt, context, providerFactory, args.CancellationToken); + var record = await ProcessPaymentAsync(evt, context, providerFactory, args.CancellationToken); + var succeeded = record.Status == PaymentStatus.Captured; sw.Stop(); PaymentMetrics.Duration.Record(sw.Elapsed.TotalMilliseconds, @@ -93,7 +96,7 @@ await args.DeadLetterMessageAsync(args.Message, // SaveChangesAsync below — a single PostgreSQL transaction guarantees atomicity. // The OutboxDispatcher publishes to Service Bus asynchronously, so a crash // between Save and publish does not lose the event. - EnqueuePaymentProcessedOutboxMessage(context, evt, succeeded); + EnqueuePaymentProcessedOutboxMessage(context, evt, record); await inboxStore.MarkAsProcessedAsync(evt.Id, nameof(PaymentRequestedIntegrationEvent)); await context.SaveChangesAsync(args.CancellationToken); @@ -108,7 +111,12 @@ await args.DeadLetterMessageAsync(args.Message, } } - internal async Task ProcessPaymentAsync( + /// + /// Runs authorize → capture for the order and returns the resulting record, staged but not + /// saved. Provider exceptions (transient faults) propagate: nothing is persisted, the message + /// is abandoned, and the redelivery replays the provider calls under the same idempotency keys. + /// + internal async Task ProcessPaymentAsync( PaymentRequestedIntegrationEvent evt, PaymentDbContext context, IPaymentProviderFactory providerFactory, @@ -121,7 +129,7 @@ internal async Task ProcessPaymentAsync( { logger.LogInformation("Payment for order {OrderId} already exists with status {Status}.", evt.OrderId, existing.Status); - return existing.Status is Domain.Enums.PaymentStatus.Captured or Domain.Enums.PaymentStatus.Authorized; + return existing; } var record = new PaymentRecord( @@ -131,65 +139,84 @@ internal async Task ProcessPaymentAsync( evt.PaymentMethod, evt.CustomerEmail, evt.CorrelationId); + await context.Payments.AddAsync(record, ct); if (options.Value.BypassProviders) { var devTransactionId = $"DEV-{Guid.CreateVersion7():N}"; - record.MarkCaptured(devTransactionId); - await context.Payments.AddAsync(record, ct); + Transition(record.MarkCaptured(devTransactionId)); logger.LogInformation( "Provider bypass active — marking order {OrderId} as captured without contacting a provider. TransactionId: {TransactionId}", evt.OrderId, devTransactionId); - return true; + return record; } var provider = providerFactory.GetProvider(evt.PaymentMethod); if (provider is null) { - record.MarkFailed($"Unsupported payment method: {evt.PaymentMethod}"); - await context.Payments.AddAsync(record, ct); - return false; + Transition(record.MarkFailed($"Unsupported payment method: {evt.PaymentMethod}")); + return record; } - var request = new PaymentRequest(evt.OrderId, evt.Amount, evt.Currency, evt.CustomerEmail); + var request = new PaymentRequest( + evt.OrderId, evt.Amount, evt.Currency, evt.CustomerEmail, evt.TenantId, evt.CorrelationId); var authResult = await provider.AuthorizeAsync(request, ct); if (authResult.IsFailure) { - record.MarkFailed(authResult.Error.Description ?? "Authorization failed."); - await context.Payments.AddAsync(record, ct); - return false; + Transition(record.MarkFailed(authResult.Error.Description ?? "Authorization failed.")); + return record; } - var captureResult = await provider.CaptureAsync(authResult.Value.TransactionId, evt.Amount, ct); + // Stored before capture so a failed capture still leaves the provider reference on the + // record — the Stripe webhook reconciles by it. + var authorizationId = authResult.Value.TransactionId; + Transition(record.MarkAuthorized(authorizationId)); + + var captureResult = await provider.CaptureAsync(authorizationId, evt.Amount, ct); if (captureResult.IsFailure) { - record.MarkFailed(captureResult.Error.Description ?? "Capture failed."); - await context.Payments.AddAsync(record, ct); - return false; + // Release the hold so the customer's funds are not blocked until the authorization lapses. + var voided = await provider.VoidAsync(authorizationId, ct); + if (voided.IsFailure) + logger.LogWarning( + "Releasing authorization {TransactionId} for order {OrderId} failed; it lapses at the provider.", + authorizationId, evt.OrderId); + + Transition(record.MarkFailed(captureResult.Error.Description ?? "Capture failed.")); + return record; } - record.MarkCaptured(captureResult.Value.TransactionId); - await context.Payments.AddAsync(record, ct); + Transition(record.MarkCaptured(captureResult.Value.TransactionId)); logger.LogInformation("Payment captured for order {OrderId}. TransactionId: {TransactionId}", evt.OrderId, captureResult.Value.TransactionId); - return true; + return record; + } + + // Every transition above starts from a record this method just created, so a rejected + // transition is a programming error, not a business outcome. + private static void Transition(Result result) + { + if (result.IsFailure) + throw new InvalidOperationException($"Invalid payment status transition: {result.Error.Code}"); } - private static void EnqueuePaymentProcessedOutboxMessage( + internal static void EnqueuePaymentProcessedOutboxMessage( PaymentDbContext context, PaymentRequestedIntegrationEvent source, - bool succeeded) + PaymentRecord record) { + var succeeded = record.Status == PaymentStatus.Captured; var processed = new PaymentProcessedIntegrationEvent { CorrelationId = source.CorrelationId, OrderId = source.OrderId, Succeeded = succeeded, FailureReason = succeeded ? null : "Payment processing failed.", + TransactionId = record.TransactionId, CustomerEmail = source.CustomerEmail, PaymentMethod = source.PaymentMethod }; diff --git a/src/Services/Payment/OrderSphere.Payment.Worker/Workers/RefundRequestedProcessor.cs b/src/Services/Payment/OrderSphere.Payment.Worker/Workers/RefundRequestedProcessor.cs index 9df297d6..504579fa 100644 --- a/src/Services/Payment/OrderSphere.Payment.Worker/Workers/RefundRequestedProcessor.cs +++ b/src/Services/Payment/OrderSphere.Payment.Worker/Workers/RefundRequestedProcessor.cs @@ -135,7 +135,8 @@ internal async Task ProcessRefundRequestAsync( throw new InvalidOperationException( $"No provider for method '{payment.PaymentMethod}' to refund order {evt.OrderId}."); - var refund = await provider.RefundAsync(payment.TransactionId!, evt.Amount, ct); + // Keyed per return request, so a refund for one return never replays another's. + var refund = await provider.RefundAsync(payment.TransactionId!, evt.Amount, evt.ReturnRequestId.ToString("N"), ct); if (refund.IsFailure) // Throw so Service Bus redelivers — a transient refund failure must not silently drop. throw new InvalidOperationException( diff --git a/src/Services/Webhooks/OrderSphere.Webhooks.Application/Features/Subscriptions/CreateSubscription/CreateSubscriptionCommandValidator.cs b/src/Services/Webhooks/OrderSphere.Webhooks.Application/Features/Subscriptions/CreateSubscription/CreateSubscriptionCommandValidator.cs index e793aaa1..fe1dd07e 100644 --- a/src/Services/Webhooks/OrderSphere.Webhooks.Application/Features/Subscriptions/CreateSubscription/CreateSubscriptionCommandValidator.cs +++ b/src/Services/Webhooks/OrderSphere.Webhooks.Application/Features/Subscriptions/CreateSubscription/CreateSubscriptionCommandValidator.cs @@ -1,3 +1,5 @@ +using OrderSphere.Webhooks.Application.Security; + namespace OrderSphere.Webhooks.Application.Features.Subscriptions.CreateSubscription; public sealed class CreateSubscriptionCommandValidator : AbstractValidator @@ -6,10 +8,17 @@ public CreateSubscriptionCommandValidator() { RuleFor(x => x.Url) .NotEmpty().WithMessage("A URL is required.") - .Must(url => Uri.TryCreate(url, UriKind.Absolute, out var uri) && uri.Scheme == "https") - .WithMessage("Only absolute HTTPS URLs are accepted."); + .MaximumLength(2048) + .Must(WebhookTargetPolicy.IsAllowedUrl) + .WithMessage("Only absolute HTTPS URLs to publicly reachable hosts are accepted."); + + RuleFor(x => x.Secret) + .MaximumLength(256); RuleFor(x => x.Events) .NotEmpty().WithMessage("At least one event type is required."); + + RuleForEach(x => x.Events) + .IsInEnum(); } } diff --git a/src/Services/Webhooks/OrderSphere.Webhooks.Application/Features/Subscriptions/UpdateSubscription/UpdateSubscriptionCommandValidator.cs b/src/Services/Webhooks/OrderSphere.Webhooks.Application/Features/Subscriptions/UpdateSubscription/UpdateSubscriptionCommandValidator.cs index 45aa61d1..73b25df6 100644 --- a/src/Services/Webhooks/OrderSphere.Webhooks.Application/Features/Subscriptions/UpdateSubscription/UpdateSubscriptionCommandValidator.cs +++ b/src/Services/Webhooks/OrderSphere.Webhooks.Application/Features/Subscriptions/UpdateSubscription/UpdateSubscriptionCommandValidator.cs @@ -1,3 +1,5 @@ +using OrderSphere.Webhooks.Application.Security; + namespace OrderSphere.Webhooks.Application.Features.Subscriptions.UpdateSubscription; public sealed class UpdateSubscriptionCommandValidator : AbstractValidator @@ -6,10 +8,17 @@ public UpdateSubscriptionCommandValidator() { RuleFor(x => x.Url) .NotEmpty().WithMessage("A URL is required.") - .Must(url => Uri.TryCreate(url, UriKind.Absolute, out var uri) && uri.Scheme == "https") - .WithMessage("Only absolute HTTPS URLs are accepted."); + .MaximumLength(2048) + .Must(WebhookTargetPolicy.IsAllowedUrl) + .WithMessage("Only absolute HTTPS URLs to publicly reachable hosts are accepted."); + + RuleFor(x => x.Secret) + .MaximumLength(256); RuleFor(x => x.Events) .NotEmpty().WithMessage("At least one event type is required."); + + RuleForEach(x => x.Events) + .IsInEnum(); } } diff --git a/src/Services/Webhooks/OrderSphere.Webhooks.Application/Security/WebhookTargetPolicy.cs b/src/Services/Webhooks/OrderSphere.Webhooks.Application/Security/WebhookTargetPolicy.cs new file mode 100644 index 00000000..ea32b486 --- /dev/null +++ b/src/Services/Webhooks/OrderSphere.Webhooks.Application/Security/WebhookTargetPolicy.cs @@ -0,0 +1,124 @@ +using System.Net; +using System.Net.Sockets; + +namespace OrderSphere.Webhooks.Application.Security; + +/// +/// Decides which targets a webhook may be delivered to. The subscriber chooses the URL and the +/// worker that posts to it runs inside the service network, so an unrestricted URL would let a +/// customer make the platform call internal services or cloud metadata endpoints (SSRF). +/// +/// The policy is applied twice. rejects obviously internal targets +/// when a subscription is saved. is applied to every resolved +/// address when the worker connects; that check is authoritative, because a public host name +/// can resolve — or later re-resolve — to an internal address. +/// +/// +public static class WebhookTargetPolicy +{ + private static readonly string[] InternalSuffixes = + [".localhost", ".local", ".internal", ".home.arpa"]; + + // Azure platform endpoint (DHCP, DNS, health probes); reachable from inside a VNet. + private static readonly IPAddress AzureWireServer = IPAddress.Parse("168.63.129.16"); + + private static ReadOnlySpan Nat64Prefix => [0x00, 0x64, 0xff, 0x9b, 0, 0, 0, 0, 0, 0, 0, 0]; + + /// + /// Syntax check for a subscription URL: absolute https, no user info, and a host that is + /// neither a blocked IP literal nor a name that only resolves inside the platform network. + /// + public static bool IsAllowedUrl(string? url) + { + if (!Uri.TryCreate(url, UriKind.Absolute, out var uri) + || uri.Scheme != Uri.UriSchemeHttps + || !string.IsNullOrEmpty(uri.UserInfo)) + { + return false; + } + + switch (uri.HostNameType) + { + case UriHostNameType.IPv4: + case UriHostNameType.IPv6: + return IPAddress.TryParse(uri.DnsSafeHost, out var address) && !IsBlockedAddress(address); + + case UriHostNameType.Dns: + var host = uri.IdnHost.TrimEnd('.'); + // Single-label names ("localhost", Aspire service names such as "ordersphere-catalog") + // resolve through service discovery or the cluster's DNS search domains. + return host.Contains('.') + && !InternalSuffixes.Any(s => host.EndsWith(s, StringComparison.OrdinalIgnoreCase)); + + default: + return false; + } + } + + /// + /// True for addresses a webhook must never connect to: loopback, private (RFC 1918), CGNAT, + /// link-local (including cloud metadata at 169.254.169.254), unspecified, multicast and + /// reserved ranges, IPv6 unique-local/site-local/link-local, and IPv4 addresses embedded in + /// IPv6 (mapped, NAT64, 6to4). + /// + public static bool IsBlockedAddress(IPAddress address) + { + if (address.IsIPv4MappedToIPv6) + address = address.MapToIPv4(); + + return address.AddressFamily switch + { + AddressFamily.InterNetwork => IsBlockedIPv4(address), + AddressFamily.InterNetworkV6 => IsBlockedIPv6(address), + _ => true + }; + } + + private static bool IsBlockedIPv4(IPAddress address) + { + if (address.Equals(AzureWireServer)) + return true; + + Span b = stackalloc byte[4]; + address.TryWriteBytes(b, out _); + + return b[0] switch + { + 0 => true, // 0.0.0.0/8 "this network" + 10 => true, // 10.0.0.0/8 private + 100 => b[1] is >= 64 and <= 127, // 100.64.0.0/10 carrier-grade NAT + 127 => true, // 127.0.0.0/8 loopback + 169 => b[1] == 254, // 169.254.0.0/16 link-local, cloud metadata + 172 => b[1] is >= 16 and <= 31, // 172.16.0.0/12 private + 192 => b[1] == 168 || (b[1] == 0 && b[2] is 0 or 2), // 192.168/16 private, 192.0.0/24, 192.0.2/24 + 198 => b[1] is 18 or 19 || (b[1] == 51 && b[2] == 100), // 198.18/15 benchmarking, 198.51.100/24 + 203 => b[1] == 0 && b[2] == 113, // 203.0.113.0/24 documentation + >= 224 => true, // multicast, reserved, broadcast + _ => false + }; + } + + private static bool IsBlockedIPv6(IPAddress address) + { + if (address.Equals(IPAddress.IPv6Any) || address.Equals(IPAddress.IPv6Loopback) + || address.IsIPv6LinkLocal || address.IsIPv6SiteLocal + || address.IsIPv6UniqueLocal || address.IsIPv6Multicast) + { + return true; + } + + Span b = stackalloc byte[16]; + address.TryWriteBytes(b, out _); + + // 64:ff9b::/96 (NAT64) carries the IPv4 target in the last four bytes. + if (b[..12].SequenceEqual(Nat64Prefix)) + return IsBlockedIPv4(new IPAddress(b[12..])); + + // 2002::/16 (6to4) carries the IPv4 relay in bytes 2..5. + if (b[0] == 0x20 && b[1] == 0x02) + return IsBlockedIPv4(new IPAddress(b[2..6])); + + // ::/96 (deprecated IPv4-compatible) — nothing public lives there. + return b[..12].IndexOfAnyExcept((byte)0) < 0; + } +} diff --git a/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Delivery/WebhookTargetConnector.cs b/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Delivery/WebhookTargetConnector.cs new file mode 100644 index 00000000..5313d20a --- /dev/null +++ b/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Delivery/WebhookTargetConnector.cs @@ -0,0 +1,40 @@ +using System.Net; +using System.Net.Sockets; +using OrderSphere.Webhooks.Application.Security; + +namespace OrderSphere.Webhooks.Worker.Delivery; + +/// +/// for webhook deliveries. Resolves the target +/// itself and connects only when every resolved address passes +/// . Checking at connect time — not only when the +/// subscription is saved — covers host names that resolve (or re-resolve) to internal addresses, +/// and hosts rewritten by service discovery before the request reaches this handler. +/// +internal static class WebhookTargetConnector +{ + public static async ValueTask ConnectAsync(SocketsHttpConnectionContext context, CancellationToken ct) + { + var host = context.DnsEndPoint.Host; + IPAddress[] addresses = IPAddress.TryParse(host, out var literal) + ? [literal] + : await Dns.GetHostAddressesAsync(host, ct); + + // Reject the host if any address is blocked, rather than connecting to the allowed + // subset: mixed answers are a known way to slip an internal address past a filter. + if (addresses.Length == 0 || addresses.Any(WebhookTargetPolicy.IsBlockedAddress)) + throw new HttpRequestException($"Webhook target host '{host}' resolves to a blocked address."); + + var socket = new Socket(SocketType.Stream, ProtocolType.Tcp) { NoDelay = true }; + try + { + await socket.ConnectAsync(addresses, context.DnsEndPoint.Port, ct); + return new NetworkStream(socket, ownsSocket: true); + } + catch + { + socket.Dispose(); + throw; + } + } +} diff --git a/src/Services/Webhooks/OrderSphere.Webhooks.Worker/OrderSphere.Webhooks.Worker.csproj b/src/Services/Webhooks/OrderSphere.Webhooks.Worker/OrderSphere.Webhooks.Worker.csproj index bafd3aed..b4e91aca 100644 --- a/src/Services/Webhooks/OrderSphere.Webhooks.Worker/OrderSphere.Webhooks.Worker.csproj +++ b/src/Services/Webhooks/OrderSphere.Webhooks.Worker/OrderSphere.Webhooks.Worker.csproj @@ -24,4 +24,8 @@ + + + + diff --git a/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Program.cs b/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Program.cs index 2ed9936c..6333fb35 100644 --- a/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Program.cs +++ b/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Program.cs @@ -5,6 +5,7 @@ using OrderSphere.Webhooks.Application; using OrderSphere.Webhooks.Infrastructure; using OrderSphere.Webhooks.Infrastructure.Persistence; +using OrderSphere.Webhooks.Worker.Delivery; using OrderSphere.Webhooks.Worker.Workers; var builder = WebApplication.CreateBuilder(args); @@ -24,10 +25,19 @@ builder.AddAzureServiceBusClient("azure-service-bus"); +// Targets are chosen by customers: every connection is checked against the SSRF policy, redirects +// are not followed (a public host could bounce to an internal one), and no proxy is used so the +// check always applies to the real destination. builder.Services.AddHttpClient("WebhookDelivery", client => { client.Timeout = TimeSpan.FromSeconds(10); client.DefaultRequestHeaders.Add("User-Agent", "OrderSphere-Webhooks/1.0"); +}) +.ConfigurePrimaryHttpMessageHandler(() => new SocketsHttpHandler +{ + AllowAutoRedirect = false, + UseProxy = false, + ConnectCallback = WebhookTargetConnector.ConnectAsync, }); builder.Services.AddHostedService(); diff --git a/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Workers/WebhookDeliveryProcessor.cs b/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Workers/WebhookDeliveryProcessor.cs index da54dab7..b03ce869 100644 --- a/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Workers/WebhookDeliveryProcessor.cs +++ b/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Workers/WebhookDeliveryProcessor.cs @@ -149,9 +149,10 @@ private async Task DeliverAsync( } else { - var body = await response.Content.ReadAsStringAsync(ct); - var error = $"HTTP {statusCode}: {body}"; - delivery.RecordFailure(statusCode, error); + // The response body is deliberately not stored: deliveries are readable by the + // subscriber through the API, so persisting it would echo whatever the target + // returned — including content of a host the subscriber should not reach. + delivery.RecordFailure(statusCode, $"HTTP {statusCode}"); WebhookMetrics.Failed.Add(1, new KeyValuePair("status", statusCode)); logger.LogWarning( "Webhook delivery {DeliveryId} to {Url} failed with {StatusCode}. Attempt {Attempt}/{Max}.", diff --git a/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Workers/WebhookEventProcessor.cs b/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Workers/WebhookEventProcessor.cs index 57a6fe4b..2ff1c6e8 100644 --- a/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Workers/WebhookEventProcessor.cs +++ b/src/Services/Webhooks/OrderSphere.Webhooks.Worker/Workers/WebhookEventProcessor.cs @@ -112,33 +112,12 @@ await args.DeadLetterMessageAsync(args.Message, return; } - // Find all active subscriptions that listen to this event type. - var eventTypeName = webhookEventType.Value.ToString(); - var subscriptions = await db.Subscriptions - .Where(s => s.IsActive && s.Events.Contains(eventTypeName)) - .ToListAsync(args.CancellationToken); - - // Filter precisely (Contains is a substring match; verify exact enum membership). - var matchingSubscriptions = subscriptions - .Where(s => s.ListensTo(webhookEventType.Value)) - .ToList(); - // "No subscriber" is a normal outcome, not a special case: it takes the same path and // produces the same Information record with Count = 0. Previously it returned early // with only a Debug line, which made a completed message indistinguishable from a // processor that never ran — see docs/logging.md, one Information record per message. - // - // Create a delivery record for each matching subscription. - foreach (var sub in matchingSubscriptions) - { - var delivery = new WebhookDelivery( - sub.Id, - eventTypeName, - eventId, - body); - - db.Deliveries.Add(delivery); - } + var created = await StageDeliveriesAsync( + db, webhookEventType.Value, eventId, body, args.CancellationToken); await db.SaveChangesAsync(args.CancellationToken); await inboxStore.MarkAsProcessedAsync(eventId, eventType, args.CancellationToken); @@ -146,7 +125,7 @@ await args.DeadLetterMessageAsync(args.Message, logger.LogInformation( "Created {Count} webhook deliveries for event {EventId} ({EventType}).", - matchingSubscriptions.Count, eventId, eventTypeName); + created, eventId, webhookEventType.Value); } catch (Exception ex) { @@ -155,6 +134,36 @@ await args.DeadLetterMessageAsync(args.Message, } } + /// + /// Stages one delivery per active subscription that listens to the event type and belongs to + /// the customer the event is about. Subscriptions are owned by a customer and the payload is + /// that customer's data, so an event without a customer id (published before the property + /// existed) is delivered to nobody rather than to every subscriber. + /// + internal static async Task StageDeliveriesAsync( + WebhooksDbContext db, + WebhookEventType webhookEventType, + Guid eventId, + string body, + CancellationToken ct) + { + if (ExtractCustomerId(body) is not { } customerId) + return 0; + + var owner = CustomerId.From(customerId); + var eventTypeName = webhookEventType.ToString(); + var subscriptions = await db.Subscriptions + .Where(s => s.IsActive && s.CustomerId == owner && s.Events.Contains(eventTypeName)) + .ToListAsync(ct); + + // Contains is a substring match; verify exact enum membership. + var matching = subscriptions.Where(s => s.ListensTo(webhookEventType)).ToList(); + foreach (var sub in matching) + db.Deliveries.Add(new WebhookDelivery(sub.Id, eventTypeName, eventId, body)); + + return matching.Count; + } + private Task ProcessErrorAsync(ProcessErrorEventArgs args) { logger.ProcessorError(args.Exception, args.EntityPath, args.ErrorSource.ToString()); @@ -214,6 +223,25 @@ private static Guid ExtractTenantId(string body) return TenantId.Default; } + /// Reads the optional CustomerId the publishing service puts on the event. + private static Guid? ExtractCustomerId(string body) + { + try + { + using var doc = JsonDocument.Parse(body); + if (doc.RootElement.TryGetProperty("CustomerId", out var prop) + && prop.ValueKind == JsonValueKind.String + && prop.TryGetGuid(out var customerId) + && customerId != Guid.Empty) + { + return customerId; + } + } + catch { /* Body already validated as JSON by ExtractEventId; be defensive anyway. */ } + + return null; + } + private static WebhookEventType? MapToWebhookEventType(string eventType) => eventType switch { nameof(OrderPlacedIntegrationEvent) or "OrderPlaced" => WebhookEventType.OrderPlaced, diff --git a/tests/OrderSphere.Bff.Tests/LocalReturnUrlTests.cs b/tests/OrderSphere.Bff.Tests/LocalReturnUrlTests.cs new file mode 100644 index 00000000..a8ed547f --- /dev/null +++ b/tests/OrderSphere.Bff.Tests/LocalReturnUrlTests.cs @@ -0,0 +1,36 @@ +using FluentAssertions; +using OrderSphere.Bff.Auth; +using Xunit; + +namespace OrderSphere.Bff.Tests; + +/// +/// /bff/login?returnUrl= must only redirect within this origin after sign-in; +/// anything else falls back to the start page. +/// +public sealed class LocalReturnUrlTests +{ + [Theory] + [InlineData("/")] + [InlineData("/products/some-slug")] + [InlineData("/checkout?step=2#payment")] + public void Sanitize_KeepsLocalPaths(string returnUrl) + { + LocalReturnUrl.Sanitize(returnUrl).Should().Be(returnUrl); + } + + [Theory] + [InlineData(null)] + [InlineData("")] + [InlineData("https://evil.example/")] + [InlineData("http://evil.example")] + [InlineData("//evil.example")] + [InlineData("/\\evil.example")] + [InlineData("/\t/evil.example")] + [InlineData("evil.example")] + [InlineData("javascript:alert(1)")] + public void Sanitize_ReplacesNonLocalTargetsWithRoot(string? returnUrl) + { + LocalReturnUrl.Sanitize(returnUrl).Should().Be("/"); + } +} diff --git a/tests/OrderSphere.Bff.Tests/StripeWebhookRouteTests.cs b/tests/OrderSphere.Bff.Tests/StripeWebhookRouteTests.cs new file mode 100644 index 00000000..44a88f17 --- /dev/null +++ b/tests/OrderSphere.Bff.Tests/StripeWebhookRouteTests.cs @@ -0,0 +1,49 @@ +using System.Net; +using System.Text; +using FluentAssertions; +using Microsoft.AspNetCore.Mvc.Testing; +using Xunit; + +namespace OrderSphere.Bff.Tests; + +/// +/// Stripe delivers webhooks without a session or antiforgery token; the request is +/// authenticated downstream by its Stripe-Signature. The BFF therefore exposes +/// exactly one anonymous POST route outside /api, where neither +/// BffUserPolicy nor the CSRF middleware applies. +/// +/// No gateway runs behind the test host, so a request that clears the BFF ends as a +/// proxy error rather than a payload (see ). +/// +/// +public sealed class StripeWebhookRouteTests(BffWebApplicationFactory factory) + : IClassFixture +{ + private HttpClient Client() => factory.CreateClient(new WebApplicationFactoryClientOptions + { + AllowAutoRedirect = false, + HandleCookies = true, + }); + + private static StringContent StripePayload() => + new("{\"id\":\"evt_test\"}", Encoding.UTF8, "application/json"); + + [Fact] + public async Task AnonymousPost_ToStripeWebhook_IsProxiedWithoutSessionOrCsrfToken() + { + var response = await Client().PostAsync("/webhooks/stripe", StripePayload()); + + // Without the route the request falls through to the SPA fallback (404 for POST); + // with it, YARP forwards and fails on the missing gateway with a 5xx proxy error. + ((int)response.StatusCode).Should().BeGreaterThanOrEqualTo(500, + "the request must reach the proxy instead of being rejected or unrouted"); + } + + [Fact] + public async Task AnonymousPost_ToPaymentApi_StillRequiresASession() + { + var response = await Client().PostAsync("/api/v1/payments/webhooks/stripe", StripePayload()); + + response.StatusCode.Should().BeOneOf(HttpStatusCode.Unauthorized, HttpStatusCode.Forbidden); + } +} diff --git a/tests/OrderSphere.Domain.Tests/Aggregates/OrderTests.cs b/tests/OrderSphere.Domain.Tests/Aggregates/OrderTests.cs index a380e135..bb400fce 100644 --- a/tests/OrderSphere.Domain.Tests/Aggregates/OrderTests.cs +++ b/tests/OrderSphere.Domain.Tests/Aggregates/OrderTests.cs @@ -3,6 +3,7 @@ using OrderSphere.BuildingBlocks.ValueObjects; using OrderSphere.Ordering.Domain.Entities; using OrderSphere.Ordering.Domain.Enums; +using OrderSphere.Ordering.Domain.Errors; using OrderSphere.Ordering.Domain.OrderEvents; using OrderSphere.Ordering.Domain.ValueObjects; using Xunit; @@ -102,13 +103,32 @@ public void MarkShipped_FromPaid_SetsStatusShipped() } [Fact] - public void MarkShipped_FromCreated_Throws() + public void MarkShipped_FromCreated_FailsWithoutRaisingEvent() { var order = CreateOrder(); - var act = () => order.MarkShipped(); + var result = order.MarkShipped(); - act.Should().Throw(); + result.IsFailure.Should().BeTrue(); + result.Error.Should().Be(OrderErrors.InvalidStatusTransition); + order.Status.Should().Be(OrderStatus.Created); + order.UncommittedEvents.OfType().Should().BeEmpty(); + } + + [Theory] + [InlineData(OrderStatus.Paid)] + [InlineData(OrderStatus.Cancelled)] + public void Confirm_FromNonCreatedStatus_FailsAndKeepsStatus(OrderStatus status) + { + var order = CreateOrder(); + if (status is OrderStatus.Paid) order.Confirm("T"); + if (status is OrderStatus.Cancelled) order.Cancel(); + + var result = order.Confirm("T-2"); + + result.IsFailure.Should().BeTrue(); + result.Error.Should().Be(OrderErrors.InvalidStatusTransition); + order.Status.Should().Be(status); } [Fact] @@ -124,14 +144,15 @@ public void MarkDelivered_FromShipped_SetsStatusDelivered() } [Fact] - public void MarkDelivered_FromPaid_Throws() + public void MarkDelivered_FromPaid_Fails() { var order = CreateOrder(); order.Confirm("T"); - var act = () => order.MarkDelivered(); + var result = order.MarkDelivered(); - act.Should().Throw(); + result.Error.Should().Be(OrderErrors.InvalidStatusTransition); + order.Status.Should().Be(OrderStatus.Paid); } [Fact] @@ -146,27 +167,46 @@ public void Cancel_FromCreated_SetsStatusCancelled() } [Fact] - public void Cancel_FromDelivered_Throws() + public void Cancel_FromDelivered_Fails() { var order = CreateOrder(); order.Confirm("T"); order.MarkShipped(); order.MarkDelivered(); - var act = () => order.Cancel(); + var result = order.Cancel(); - act.Should().Throw(); + result.Error.Should().Be(OrderErrors.InvalidStatusTransition); + order.Status.Should().Be(OrderStatus.Delivered); } [Fact] - public void Cancel_AlreadyCancelled_Throws() + public void Cancel_AlreadyCancelled_Fails() { var order = CreateOrder(); order.Cancel(); - var act = () => order.Cancel(); + var result = order.Cancel(); + + result.Error.Should().Be(OrderErrors.InvalidStatusTransition); + order.UncommittedEvents.OfType().Should().ContainSingle(); + } + + [Fact] + public void Rehydrate_FoldsHistoricSequencesTheGuardsNowReject() + { + // Streams written before the guards existed can contain Cancelled → Confirmed. + // Loading them must still work; only new transitions are guarded. + var source = CreateOrder(); + var stream = source.UncommittedEvents + .Append(new OrderCancelled(DateTime.UtcNow)) + .Append(new OrderConfirmed("T", DateTime.UtcNow)) + .ToList(); + + var rebuilt = Order.Rehydrate(source.Id, stream); - act.Should().Throw(); + rebuilt.Status.Should().Be(OrderStatus.Paid); + rebuilt.Version.Should().Be(3); } [Fact] diff --git a/tests/OrderSphere.Domain.Tests/Aggregates/PaymentRecordTests.cs b/tests/OrderSphere.Domain.Tests/Aggregates/PaymentRecordTests.cs index 9340d576..5abac747 100644 --- a/tests/OrderSphere.Domain.Tests/Aggregates/PaymentRecordTests.cs +++ b/tests/OrderSphere.Domain.Tests/Aggregates/PaymentRecordTests.cs @@ -3,6 +3,7 @@ using OrderSphere.Payment.Domain.DomainEvents; using OrderSphere.Payment.Domain.Entities; using OrderSphere.Payment.Domain.Enums; +using OrderSphere.Payment.Domain.Errors; using Xunit; namespace OrderSphere.Domain.Tests.Aggregates; @@ -135,4 +136,88 @@ public void MarkAuthorized_EventContainsOrderId() var @event = record.PopDomainEvents().OfType().Single(); @event.OrderId.Should().Be(Order); } + + [Fact] + public void MarkCaptured_FromFailed_IsRejected() + { + var record = CreateRecord(); + record.MarkAuthorized("TXN-001"); + record.MarkFailed("Capture rejected."); + + var result = record.MarkCaptured("TXN-001"); + + result.Error.Should().Be(PaymentErrors.InvalidStatusTransition); + record.Status.Should().Be(PaymentStatus.Failed); + } + + [Fact] + public void MarkCaptured_Repeated_WithSameTransaction_SucceedsWithoutNewEvent() + { + var record = CreateRecord(); + record.MarkCaptured("TXN-001"); + record.PopDomainEvents(); + + var result = record.MarkCaptured("TXN-001"); + + result.IsSuccess.Should().BeTrue(); + record.PopDomainEvents().Should().BeEmpty(); + } + + [Fact] + public void MarkCaptured_Repeated_WithOtherTransaction_IsRejected() + { + var record = CreateRecord(); + record.MarkCaptured("TXN-001"); + + var result = record.MarkCaptured("TXN-002"); + + result.IsFailure.Should().BeTrue(); + record.TransactionId.Should().Be("TXN-001"); + } + + [Fact] + public void MarkCaptured_MayReplaceTheAuthorizationReference() + { + var record = CreateRecord(); + record.MarkAuthorized("AUTH-1"); + + var result = record.MarkCaptured("CAP-1"); + + result.IsSuccess.Should().BeTrue(); + record.TransactionId.Should().Be("CAP-1"); + } + + [Fact] + public void MarkFailed_FromCaptured_IsRejected() + { + var record = CreateRecord(); + record.MarkCaptured("TXN-001"); + + var result = record.MarkFailed("late failure notice"); + + result.IsFailure.Should().BeTrue(); + record.Status.Should().Be(PaymentStatus.Captured); + } + + [Fact] + public void MarkRefunded_BeforeCapture_IsRejected() + { + var record = CreateRecord(); + record.MarkAuthorized("TXN-001"); + + var result = record.MarkRefunded(); + + result.IsFailure.Should().BeTrue(); + record.Status.Should().Be(PaymentStatus.Authorized); + } + + [Fact] + public void MarkRefunded_Repeated_Succeeds() + { + var record = CreateRecord(); + record.MarkCaptured("TXN-001"); + record.MarkRefunded(); + + record.MarkRefunded().IsSuccess.Should().BeTrue(); + } } diff --git a/tests/OrderSphere.IntegrationTests/Api/StripeWebhookTests.cs b/tests/OrderSphere.IntegrationTests/Api/StripeWebhookTests.cs new file mode 100644 index 00000000..e3bf2dc9 --- /dev/null +++ b/tests/OrderSphere.IntegrationTests/Api/StripeWebhookTests.cs @@ -0,0 +1,149 @@ +using System.Globalization; +using System.Net; +using System.Security.Cryptography; +using System.Text; +using System.Text.Json; +using FluentAssertions; +using Microsoft.AspNetCore.Hosting; +using Microsoft.AspNetCore.Mvc.Testing; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; +using OrderSphere.BuildingBlocks.StronglyTypedIds; +using OrderSphere.Payment.Api; +using OrderSphere.Payment.Application.Abstractions; +using OrderSphere.Payment.Domain.Entities; +using OrderSphere.Payment.Domain.Enums; +using StripeConfiguration = Stripe.StripeConfiguration; +using Xunit; + +namespace OrderSphere.IntegrationTests.Api; + +/// +/// The Stripe webhook reconciles PaymentRecords from signed Stripe events. It must find the +/// record by the orderId metadata or the intent id, ask Stripe to retry (503) while the worker +/// has not yet stored a record for one of our intents, ignore foreign intents, and never move a +/// record through an invalid transition. +/// +public sealed class StripeWebhookTests : IClassFixture +{ + private const string Secret = "whsec_test_secret"; + private const string Path = "api/v1/payments/webhooks/stripe"; + + private readonly WebApplicationFactory _factory; + + public StripeWebhookTests(PaymentApiFactory factory) => + _factory = factory.WithWebHostBuilder(b => b.UseSetting("Stripe:WebhookSecret", Secret)); + + [Fact] + public async Task Unsigned_request_is_rejected() + { + var response = await _factory.CreateClient().PostAsync(Path, + new StringContent(PaymentIntentEvent("payment_intent.succeeded", "pi_x", null), Encoding.UTF8, "application/json")); + + response.StatusCode.Should().Be(HttpStatusCode.BadRequest); + } + + [Fact] + public async Task Intent_of_ours_without_record_yet_asks_Stripe_to_retry() + { + var response = await PostSignedAsync(PaymentIntentEvent("payment_intent.succeeded", "pi_new", Guid.NewGuid())); + + response.StatusCode.Should().Be(HttpStatusCode.ServiceUnavailable); + } + + [Fact] + public async Task Foreign_intent_is_acknowledged_and_ignored() + { + var response = await PostSignedAsync(PaymentIntentEvent("payment_intent.succeeded", "pi_foreign", null)); + + response.StatusCode.Should().Be(HttpStatusCode.OK); + } + + [Fact] + public async Task Capture_notice_for_failed_payment_does_not_resurrect_it() + { + var orderId = await SeedAsync(r => { r.MarkAuthorized("pi_failed"); r.MarkFailed("Capture rejected."); }); + + var response = await PostSignedAsync(PaymentIntentEvent("payment_intent.succeeded", "pi_failed", orderId)); + + response.StatusCode.Should().Be(HttpStatusCode.OK); + (await StatusOfAsync(orderId)).Should().Be(PaymentStatus.Failed); + } + + [Fact] + public async Task Charge_refund_is_matched_by_intent_id_and_marks_the_payment_refunded() + { + var orderId = await SeedAsync(r => r.MarkCaptured("pi_refund")); + + var response = await PostSignedAsync(ChargeRefundedEvent("pi_refund")); + + response.StatusCode.Should().Be(HttpStatusCode.OK); + (await StatusOfAsync(orderId)).Should().Be(PaymentStatus.Refunded); + } + + private async Task SeedAsync(Action arrange) + { + var orderId = Guid.NewGuid(); + using var scope = _factory.Services.CreateScope(); + var context = scope.ServiceProvider.GetRequiredService(); + var record = new PaymentRecord(OrderId.From(orderId), 49.99m, "EUR", "CreditCard", "payer@example.com", Guid.NewGuid()); + arrange(record); + context.Payments.Add(record); + await context.SaveChangesAsync(CancellationToken.None); + return orderId; + } + + private async Task StatusOfAsync(Guid orderId) + { + using var scope = _factory.Services.CreateScope(); + var context = scope.ServiceProvider.GetRequiredService(); + return (await context.Payments.AsNoTracking().SingleAsync(p => p.OrderId == OrderId.From(orderId))).Status; + } + + private Task PostSignedAsync(string json) + { + var timestamp = DateTimeOffset.UtcNow.ToUnixTimeSeconds().ToString(CultureInfo.InvariantCulture); + var signature = Convert.ToHexStringLower( + HMACSHA256.HashData(Encoding.UTF8.GetBytes(Secret), Encoding.UTF8.GetBytes($"{timestamp}.{json}"))); + + var request = new HttpRequestMessage(HttpMethod.Post, Path) + { + Content = new StringContent(json, Encoding.UTF8, "application/json") + }; + request.Headers.Add("Stripe-Signature", $"t={timestamp},v1={signature}"); + return _factory.CreateClient().SendAsync(request); + } + + private static string PaymentIntentEvent(string type, string intentId, Guid? orderId) => + Envelope(type, new Dictionary + { + ["id"] = intentId, + ["object"] = "payment_intent", + ["status"] = type == "payment_intent.succeeded" ? "succeeded" : "requires_payment_method", + ["metadata"] = orderId is null + ? new Dictionary() + : new Dictionary { ["orderId"] = orderId.Value.ToString() } + }); + + private static string ChargeRefundedEvent(string intentId) => + Envelope("charge.refunded", new Dictionary + { + ["id"] = "ch_1", + ["object"] = "charge", + ["payment_intent"] = intentId, + ["refunded"] = true + }); + + private static string Envelope(string type, Dictionary dataObject) => + JsonSerializer.Serialize(new Dictionary + { + ["id"] = $"evt_{Guid.NewGuid():N}", + ["object"] = "event", + ["api_version"] = StripeConfiguration.ApiVersion, + ["created"] = DateTimeOffset.UtcNow.ToUnixTimeSeconds(), + ["livemode"] = false, + ["pending_webhooks"] = 1, + ["type"] = type, + ["data"] = new Dictionary { ["object"] = dataObject } + }); +} diff --git a/tests/OrderSphere.IntegrationTests/PaymentResultToOrderFlowTests.cs b/tests/OrderSphere.IntegrationTests/PaymentResultToOrderFlowTests.cs index b280f5a3..b2a07490 100644 --- a/tests/OrderSphere.IntegrationTests/PaymentResultToOrderFlowTests.cs +++ b/tests/OrderSphere.IntegrationTests/PaymentResultToOrderFlowTests.cs @@ -14,8 +14,10 @@ using OrderSphere.BuildingBlocks.StronglyTypedIds; using OrderSphere.Ordering.Application.Abstractions; using OrderSphere.Ordering.Domain.Enums; +using OrderSphere.Ordering.Infrastructure.EventSourcing; using OrderSphere.Ordering.Infrastructure.Persistence; using OrderSphere.Ordering.Worker.Workers; +using static OrderSphere.Ordering.Worker.Workers.PaymentResultProcessor; using Xunit; using OrderItemDto = OrderSphere.BuildingBlocks.Contracts.Events.OrderItemDto; using ShippingAddressDto = OrderSphere.BuildingBlocks.Contracts.Events.ShippingAddressDto; @@ -211,6 +213,108 @@ public async Task Duplicate_PaymentProcessed_event_is_idempotent() unchanged.Status.Should().Be(OrderStatus.Created); } + [Fact] + public async Task Status_change_event_carries_the_order_owner_for_webhook_scoping() + { + var (orderId, correlationId) = await SeedOrderAsync(); + var evt = SucceededEvent(orderId, correlationId); + + await using var ctx = NewContext(); + await NewProcessor().ProcessPaymentResultAsync(evt, ctx, Substitute.For(), Catalog(), deliveryCount: 1, CancellationToken.None); + await ctx.SaveChangesAsync(); + + await using var verify = NewContext(); + var json = await verify.OutboxMessages + .Where(m => m.Type == nameof(OrderStatusChangedIntegrationEvent)) + .Select(m => m.Content) + .SingleAsync(); + JsonSerializer.Deserialize(json)!.CustomerId.Should().Be(CustomerGuid); + } + + [Fact] + public async Task Successful_payment_for_cancelled_order_is_refunded_and_order_stays_cancelled() + { + var (orderId, correlationId) = await SeedOrderAsync(); + await CancelOrderAsync(orderId); + var evt = SucceededEvent(orderId, correlationId); + var catalog = Catalog(); + + await using var ctx = NewContext(); + var outcome = await NewProcessor().ProcessPaymentResultAsync(evt, ctx, Substitute.For(), catalog, deliveryCount: 1, CancellationToken.None); + await ctx.SaveChangesAsync(); + + outcome.Should().Be(PaymentResultProcessor.PaymentResultOutcome.Processed); + await catalog.DidNotReceive().ConfirmReservationAsync(Arg.Any(), Arg.Any()); + + await using var verify = NewContext(); + (await verify.Orders.SingleAsync(o => o.Id == OrderId.From(orderId))).Status.Should().Be(OrderStatus.Cancelled); + + var types = await NewOutboxTypesAsync(verify); + types.Should().Contain(nameof(OrderConfirmationFailedIntegrationEvent)); + types.Should().NotContain(nameof(OrderPlacedIntegrationEvent)); + types.Should().NotContain(nameof(OrderStatusChangedIntegrationEvent), "the order was already cancelled"); + + var saga = await verify.OrderSagas.SingleAsync(s => s.CorrelationId == correlationId); + saga.State.Should().Be(SagaState.CompensationPending); + } + + [Fact] + public async Task Failed_payment_for_cancelled_order_only_closes_the_saga() + { + var (orderId, correlationId) = await SeedOrderAsync(); + await CancelOrderAsync(orderId); + var evt = FailedEvent(orderId, correlationId); + var inbox = Substitute.For(); + + await using var ctx = NewContext(); + var outcome = await NewProcessor().ProcessPaymentResultAsync(evt, ctx, inbox, Catalog(), deliveryCount: 1, CancellationToken.None); + await ctx.SaveChangesAsync(); + + outcome.Should().Be(PaymentResultOutcome.Processed); + await inbox.Received(1).MarkAsProcessedAsync(evt.Id, nameof(PaymentProcessedIntegrationEvent), Arg.Any()); + + await using var verify = NewContext(); + (await NewOutboxTypesAsync(verify)).Should().BeEmpty(); + (await verify.OrderSagas.SingleAsync(s => s.CorrelationId == correlationId)).State.Should().Be(SagaState.Cancelled); + } + + [Fact] + public async Task Second_success_for_paid_order_changes_nothing() + { + var (orderId, correlationId) = await SeedOrderAsync(); + await using (var first = NewContext()) + { + await NewProcessor().ProcessPaymentResultAsync(SucceededEvent(orderId, correlationId), first, Substitute.For(), Catalog(), 1, CancellationToken.None); + await first.SaveChangesAsync(); + } + + await using var ctx = NewContext(); + var catalog = Catalog(); + await NewProcessor().ProcessPaymentResultAsync(SucceededEvent(orderId, correlationId), ctx, Substitute.For(), catalog, 1, CancellationToken.None); + await ctx.SaveChangesAsync(); + + await catalog.DidNotReceive().ConfirmReservationAsync(Arg.Any(), Arg.Any()); + await using var verify = NewContext(); + (await NewOutboxTypesAsync(verify)).Should().HaveCount(3, "only the first result stages events"); + (await verify.Orders.SingleAsync(o => o.Id == OrderId.From(orderId))).Status.Should().Be(OrderStatus.Paid); + } + + private async Task CancelOrderAsync(Guid orderId) + { + await using var ctx = NewContext(); + var store = new OrderEventStore(ctx); + var order = await store.LoadAsync(OrderId.From(orderId)); + order!.Cancel().IsSuccess.Should().BeTrue(); + await store.AppendAsync(order); + await ctx.SaveChangesAsync(); + } + + private static Task> NewOutboxTypesAsync(OrderingDbContext ctx) => + ctx.OutboxMessages + .Where(m => m.Type != nameof(PaymentRequestedIntegrationEvent)) + .Select(m => m.Type) + .ToListAsync(); + [Fact] public async Task Missing_order_returns_OrderNotFound() { diff --git a/tests/OrderSphere.Payment.Tests/OrderConfirmationFailedProcessorTests.cs b/tests/OrderSphere.Payment.Tests/OrderConfirmationFailedProcessorTests.cs index 433120bc..d9b827dc 100644 --- a/tests/OrderSphere.Payment.Tests/OrderConfirmationFailedProcessorTests.cs +++ b/tests/OrderSphere.Payment.Tests/OrderConfirmationFailedProcessorTests.cs @@ -66,7 +66,7 @@ public async Task Captured_payment_is_refunded_through_provider_and_PaymentRefun await SeedCapturedPaymentAsync(context, orderId); var provider = Substitute.For(); - provider.RefundAsync(Arg.Any(), Arg.Any(), Arg.Any()) + provider.RefundAsync(Arg.Any(), Arg.Any(), Arg.Any(), Arg.Any()) .Returns(Result.Success()); var factory = Substitute.For(); factory.GetProvider(Method).Returns(provider); @@ -75,7 +75,7 @@ await NewProcessor().ProcessConfirmationFailureAsync( NewEvent(orderId), context, UnprocessedInbox(), factory, CancellationToken.None); await context.SaveChangesAsync(); - await provider.Received(1).RefundAsync("cap-1", 49.99m, Arg.Any()); + await provider.Received(1).RefundAsync("cap-1", 49.99m, "ocf", Arg.Any()); var record = await context.Payments.SingleAsync(p => p.OrderId == OrderId.From(orderId)); record.Status.Should().Be(PaymentStatus.Refunded); @@ -116,7 +116,7 @@ public async Task Provider_refund_failure_throws_so_the_message_is_retried() await SeedCapturedPaymentAsync(context, orderId); var provider = Substitute.For(); - provider.RefundAsync(Arg.Any(), Arg.Any(), Arg.Any()) + provider.RefundAsync(Arg.Any(), Arg.Any(), Arg.Any(), Arg.Any()) .Returns(Result.Failure(new Error("Payment.Refund", "Gateway timeout."))); var factory = Substitute.For(); factory.GetProvider(Method).Returns(provider); diff --git a/tests/OrderSphere.Payment.Tests/PaymentProcessorTests.cs b/tests/OrderSphere.Payment.Tests/PaymentProcessorTests.cs index 1aa687f4..e3b3b4d2 100644 --- a/tests/OrderSphere.Payment.Tests/PaymentProcessorTests.cs +++ b/tests/OrderSphere.Payment.Tests/PaymentProcessorTests.cs @@ -5,6 +5,7 @@ using Microsoft.Extensions.Logging.Abstractions; using Microsoft.Extensions.Options; using NSubstitute; +using NSubstitute.ExceptionExtensions; using OrderSphere.BuildingBlocks.Contracts.Events; using OrderSphere.BuildingBlocks.Primitives; using OrderSphere.BuildingBlocks.StronglyTypedIds; @@ -34,6 +35,7 @@ private static PaymentRequestedIntegrationEvent NewEvent(string method = Method) => new() { OrderId = Guid.NewGuid(), + CorrelationId = Guid.NewGuid(), Amount = 49.99m, Currency = "EUR", PaymentMethod = method, @@ -47,15 +49,23 @@ private static IPaymentProviderFactory FactoryReturning(IPaymentProvider? provid return factory; } + private static IPaymentProvider AuthorizingProvider(string authorizationId = "auth-1") + { + var provider = Substitute.For(); + provider.AuthorizeAsync(Arg.Any(), Arg.Any()) + .Returns(Result.Success(new PaymentProviderResult(authorizationId))); + provider.VoidAsync(Arg.Any(), Arg.Any()) + .Returns(Result.Success()); + return provider; + } + [Fact] - public async Task Authorize_and_capture_succeed_persists_captured_record_and_returns_true() + public async Task Authorize_and_capture_succeed_persists_captured_record() { await using var context = NewContext(); var evt = NewEvent(); - var provider = Substitute.For(); - provider.AuthorizeAsync(Arg.Any(), Arg.Any()) - .Returns(Result.Success(new PaymentProviderResult("auth-1"))); + var provider = AuthorizingProvider(); provider.CaptureAsync(Arg.Any(), Arg.Any(), Arg.Any()) .Returns(Result.Success(new PaymentProviderResult("cap-1"))); @@ -63,7 +73,7 @@ public async Task Authorize_and_capture_succeed_persists_captured_record_and_ret evt, context, FactoryReturning(provider), CancellationToken.None); await context.SaveChangesAsync(); // mirrors the single Save in OnMessageReceived - result.Should().BeTrue(); + result.Status.Should().Be(PaymentStatus.Captured); var record = await context.Payments.SingleAsync(p => p.OrderId == OrderId.From(evt.OrderId)); record.Status.Should().Be(PaymentStatus.Captured); @@ -71,7 +81,25 @@ public async Task Authorize_and_capture_succeed_persists_captured_record_and_ret } [Fact] - public async Task Unsupported_payment_method_persists_failed_record_and_returns_false() + public async Task Provider_request_carries_tenant_and_correlation_for_the_idempotency_scope() + { + await using var context = NewContext(); + var evt = NewEvent(); + var provider = AuthorizingProvider(); + provider.CaptureAsync(Arg.Any(), Arg.Any(), Arg.Any()) + .Returns(Result.Success(new PaymentProviderResult("auth-1"))); + + await NewProcessor().ProcessPaymentAsync(evt, context, FactoryReturning(provider), CancellationToken.None); + + await provider.Received(1).AuthorizeAsync( + Arg.Is(r => r.OrderId == evt.OrderId + && r.CorrelationId == evt.CorrelationId + && r.TenantId == evt.TenantId), + Arg.Any()); + } + + [Fact] + public async Task Unsupported_payment_method_persists_failed_record() { await using var context = NewContext(); var evt = NewEvent("bitcoin"); @@ -80,7 +108,7 @@ public async Task Unsupported_payment_method_persists_failed_record_and_returns_ evt, context, FactoryReturning(provider: null, method: "bitcoin"), CancellationToken.None); await context.SaveChangesAsync(); - result.Should().BeFalse(); + result.Status.Should().Be(PaymentStatus.Failed); var record = await context.Payments.SingleAsync(p => p.OrderId == OrderId.From(evt.OrderId)); record.Status.Should().Be(PaymentStatus.Failed); @@ -88,7 +116,7 @@ public async Task Unsupported_payment_method_persists_failed_record_and_returns_ } [Fact] - public async Task Authorization_failure_persists_failed_record_and_returns_false() + public async Task Authorization_failure_persists_failed_record_without_capture() { await using var context = NewContext(); var evt = NewEvent(); @@ -101,7 +129,7 @@ public async Task Authorization_failure_persists_failed_record_and_returns_false evt, context, FactoryReturning(provider), CancellationToken.None); await context.SaveChangesAsync(); - result.Should().BeFalse(); + result.Status.Should().Be(PaymentStatus.Failed); var record = await context.Payments.SingleAsync(p => p.OrderId == OrderId.From(evt.OrderId)); record.Status.Should().Be(PaymentStatus.Failed); @@ -110,14 +138,12 @@ public async Task Authorization_failure_persists_failed_record_and_returns_false } [Fact] - public async Task Capture_failure_persists_failed_record_and_returns_false() + public async Task Capture_failure_voids_the_authorization_and_keeps_its_reference() { await using var context = NewContext(); var evt = NewEvent(); - var provider = Substitute.For(); - provider.AuthorizeAsync(Arg.Any(), Arg.Any()) - .Returns(Result.Success(new PaymentProviderResult("auth-1"))); + var provider = AuthorizingProvider("auth-1"); provider.CaptureAsync(Arg.Any(), Arg.Any(), Arg.Any()) .Returns(Result.Failure(new Error("Payment.Capture", "Capture rejected."))); @@ -125,15 +151,55 @@ public async Task Capture_failure_persists_failed_record_and_returns_false() evt, context, FactoryReturning(provider), CancellationToken.None); await context.SaveChangesAsync(); - result.Should().BeFalse(); + result.Status.Should().Be(PaymentStatus.Failed); + await provider.Received(1).VoidAsync("auth-1", Arg.Any()); var record = await context.Payments.SingleAsync(p => p.OrderId == OrderId.From(evt.OrderId)); record.Status.Should().Be(PaymentStatus.Failed); record.FailureReason.Should().Be("Capture rejected."); + record.TransactionId.Should().Be("auth-1", "the webhook reconciles by the provider reference"); + } + + [Fact] + public async Task Failed_void_does_not_change_the_outcome() + { + await using var context = NewContext(); + var evt = NewEvent(); + + var provider = AuthorizingProvider(); + provider.CaptureAsync(Arg.Any(), Arg.Any(), Arg.Any()) + .Returns(Result.Failure(new Error("Payment.Capture", "Capture rejected."))); + provider.VoidAsync(Arg.Any(), Arg.Any()) + .Returns(Result.Failure(new Error("Payment.VoidFailed", "Cancel rejected."))); + + var result = await NewProcessor().ProcessPaymentAsync( + evt, context, FactoryReturning(provider), CancellationToken.None); + + result.Status.Should().Be(PaymentStatus.Failed); + } + + [Fact] + public async Task Transient_provider_exception_propagates_and_persists_nothing() + { + await using var context = NewContext(); + var evt = NewEvent(); + + var provider = AuthorizingProvider(); + provider.CaptureAsync(Arg.Any(), Arg.Any(), Arg.Any()) + .ThrowsAsync(new HttpRequestException("connection reset")); + + var act = () => NewProcessor().ProcessPaymentAsync( + evt, context, FactoryReturning(provider), CancellationToken.None); + + await act.Should().ThrowAsync(); + + // The worker abandons the message without saving; a fresh context sees no record. + context.ChangeTracker.Clear(); + (await context.Payments.CountAsync()).Should().Be(0); } [Fact] - public async Task Existing_captured_payment_is_idempotent_and_does_not_insert_duplicate() + public async Task Existing_captured_payment_is_returned_without_contacting_the_provider() { await using var context = NewContext(); var evt = NewEvent(); @@ -149,11 +215,30 @@ public async Task Existing_captured_payment_is_idempotent_and_does_not_insert_du var result = await NewProcessor().ProcessPaymentAsync( evt, context, FactoryReturning(provider), CancellationToken.None); - result.Should().BeTrue(); + result.Status.Should().Be(PaymentStatus.Captured); (await context.Payments.CountAsync(p => p.OrderId == OrderId.From(evt.OrderId))).Should().Be(1); await provider.DidNotReceive().AuthorizeAsync(Arg.Any(), Arg.Any()); } + [Fact] + public async Task PaymentProcessed_event_carries_the_transaction_reference() + { + await using var context = NewContext(); + var evt = NewEvent(); + var record = new PaymentRecord( + OrderId.From(evt.OrderId), evt.Amount, evt.Currency, evt.PaymentMethod, evt.CustomerEmail, evt.CorrelationId); + record.MarkAuthorized("auth-1"); + record.MarkFailed("Capture rejected."); + + PaymentProcessor.EnqueuePaymentProcessedOutboxMessage(context, evt, record); + await context.SaveChangesAsync(); + + var json = await context.OutboxMessages.Select(m => m.Content).SingleAsync(); + var processed = System.Text.Json.JsonSerializer.Deserialize(json)!; + processed.Succeeded.Should().BeFalse(); + processed.TransactionId.Should().Be("auth-1"); + } + [Fact] public async Task BypassProviders_marks_payment_as_captured_without_contacting_provider() { @@ -165,7 +250,7 @@ public async Task BypassProviders_marks_payment_as_captured_without_contacting_p evt, context, factory, CancellationToken.None); await context.SaveChangesAsync(); - result.Should().BeTrue(); + result.Status.Should().Be(PaymentStatus.Captured); var record = await context.Payments.SingleAsync(p => p.OrderId == OrderId.From(evt.OrderId)); record.Status.Should().Be(PaymentStatus.Captured); diff --git a/tests/OrderSphere.Payment.Tests/PaymentProviderTests.cs b/tests/OrderSphere.Payment.Tests/PaymentProviderTests.cs index d0e69247..f5cf8e3e 100644 --- a/tests/OrderSphere.Payment.Tests/PaymentProviderTests.cs +++ b/tests/OrderSphere.Payment.Tests/PaymentProviderTests.cs @@ -13,7 +13,7 @@ namespace OrderSphere.Payment.Tests; public sealed class PaymentProviderTests { private static readonly PaymentRequest Request = - new(Guid.NewGuid(), 49.99m, "EUR", "customer@example.com"); + new(Guid.NewGuid(), 49.99m, "EUR", "customer@example.com", Guid.Empty, Guid.NewGuid()); private static IPaymentProvider CreditCard() => new CreditCardPaymentProvider(NullLogger.Instance); @@ -72,7 +72,12 @@ public async Task Capture_ReturnsSuccess_EchoingTransactionId(IPaymentProvider p [Theory] [MemberData(nameof(AllProviders))] public async Task Refund_ReturnsSuccess(IPaymentProvider provider) - => (await provider.RefundAsync("txn-123", 49.99m)).IsSuccess.Should().BeTrue(); + => (await provider.RefundAsync("txn-123", 49.99m, "ref-1")).IsSuccess.Should().BeTrue(); + + [Theory] + [MemberData(nameof(AllProviders))] + public async Task Void_ReturnsSuccess(IPaymentProvider provider) + => (await provider.VoidAsync("txn-123")).IsSuccess.Should().BeTrue(); private static PaymentProviderFactory BuildFactory() => diff --git a/tests/OrderSphere.Payment.Tests/StripePaymentProviderTests.cs b/tests/OrderSphere.Payment.Tests/StripePaymentProviderTests.cs index f4db0856..5cf6a865 100644 --- a/tests/OrderSphere.Payment.Tests/StripePaymentProviderTests.cs +++ b/tests/OrderSphere.Payment.Tests/StripePaymentProviderTests.cs @@ -1,3 +1,4 @@ +using System.Net; using FluentAssertions; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; @@ -11,12 +12,24 @@ namespace OrderSphere.Payment.Tests; /// -/// B1 — the Stripe provider registers under the "CreditCard" method name so existing checkout -/// routes to it without a UI contract change, but only when an API key is configured. Without a -/// key the simulated provider stays active so local development needs no external credentials. +/// The Stripe provider registers under the "CreditCard" method name so existing checkout routes to +/// it without a UI contract change, but only when an API key is configured. Every Stripe call must +/// carry a deterministic idempotency key, and only declines/invalid requests may come back as a +/// Result failure — transient faults must propagate so Service Bus redelivers. /// public sealed class StripePaymentProviderTests { + private static readonly Guid OrderId = Guid.Parse("11111111-2222-3333-4444-555555555555"); + + private static PaymentRequest Request() => + new(OrderId, 49.99m, "EUR", "customer@example.com", Guid.Empty, Guid.NewGuid()); + + private static StripePaymentProvider Provider(FakeStripeClient client) => + new(client, NullLogger.Instance); + + private static StripeException StripeError(HttpStatusCode status, string type, string? code = null) => + new(status, new StripeError { Type = type, Code = code }, "stripe error"); + [Fact] public void MethodName_IsCreditCard() { @@ -27,6 +40,144 @@ public void MethodName_IsCreditCard() provider.MethodName.Should().Be("CreditCard"); } + [Fact] + public async Task Authorize_UsesOrderScopedIdempotencyKey_AndTagsTheIntent() + { + var client = new FakeStripeClient(); + var request = Request(); + + var result = await Provider(client).AuthorizeAsync(request); + + result.IsSuccess.Should().BeTrue(); + result.Value.TransactionId.Should().Be("pi_1"); + var call = client.Calls.Single(); + call.Path.Should().Be("/v1/payment_intents"); + call.RequestOptions!.IdempotencyKey.Should().Be($"os-pi-create-{OrderId:N}"); + var metadata = ((PaymentIntentCreateOptions)call.Options).Metadata; + metadata["orderId"].Should().Be(OrderId.ToString()); + metadata["correlationId"].Should().Be(request.CorrelationId.ToString()); + metadata.Should().ContainKey("tenantId"); + } + + [Fact] + public async Task Authorize_WhenIntentHoldsNoFunds_CancelsItAndFails() + { + var client = new FakeStripeClient + { + Respond = (_, path) => path.EndsWith("/cancel") + ? new PaymentIntent { Id = "pi_1", Status = "canceled" } + : new PaymentIntent { Id = "pi_1", Status = "requires_action" } + }; + + var result = await Provider(client).AuthorizeAsync(Request()); + + result.IsFailure.Should().BeTrue(); + client.Calls.Should().Contain(c => c.Path == "/v1/payment_intents/pi_1/cancel" + && c.RequestOptions!.IdempotencyKey == "os-pi-cancel-pi_1"); + } + + [Fact] + public async Task Authorize_CardDecline_ReturnsFailure() + { + var client = new FakeStripeClient + { + Respond = (_, _) => StripeError(HttpStatusCode.PaymentRequired, "card_error", "card_declined") + }; + + var result = await Provider(client).AuthorizeAsync(Request()); + + result.IsFailure.Should().BeTrue(); + } + + [Fact] + public async Task Capture_UsesIntentScopedIdempotencyKey() + { + var client = new FakeStripeClient { Respond = (_, _) => new PaymentIntent { Id = "pi_1", Status = "succeeded" } }; + + var result = await Provider(client).CaptureAsync("pi_1", 49.99m); + + result.IsSuccess.Should().BeTrue(); + client.Calls.Single().RequestOptions!.IdempotencyKey.Should().Be("os-pi-capture-pi_1"); + } + + [Theory] + [InlineData(HttpStatusCode.InternalServerError, "api_error")] + [InlineData(HttpStatusCode.TooManyRequests, "invalid_request_error")] + [InlineData(HttpStatusCode.Unauthorized, "invalid_request_error")] + [InlineData(HttpStatusCode.Conflict, "invalid_request_error")] + [InlineData(HttpStatusCode.BadRequest, "idempotency_error")] + public async Task Capture_NonPermanentStripeError_Propagates(HttpStatusCode status, string type) + { + var client = new FakeStripeClient { Respond = (_, _) => StripeError(status, type) }; + + var act = () => Provider(client).CaptureAsync("pi_1", 49.99m); + + await act.Should().ThrowAsync(); + } + + [Fact] + public async Task Capture_NetworkFailure_Propagates() + { + var client = new FakeStripeClient { Respond = (_, _) => new HttpRequestException("connection reset") }; + + var act = () => Provider(client).CaptureAsync("pi_1", 49.99m); + + await act.Should().ThrowAsync(); + } + + [Fact] + public async Task Capture_Decline_ReturnsFailure() + { + var client = new FakeStripeClient + { + Respond = (_, _) => StripeError(HttpStatusCode.PaymentRequired, "card_error", "card_declined") + }; + + var result = await Provider(client).CaptureAsync("pi_1", 49.99m); + + result.IsFailure.Should().BeTrue(); + } + + [Fact] + public async Task Capture_OfAlreadyCapturedIntent_ReturnsSuccess() + { + // An earlier attempt captured the intent but its idempotency key has expired. + var client = new FakeStripeClient + { + Respond = (method, _) => method == HttpMethod.Get + ? new PaymentIntent { Id = "pi_1", Status = "succeeded" } + : StripeError(HttpStatusCode.BadRequest, "invalid_request_error", "payment_intent_unexpected_state") + }; + + var result = await Provider(client).CaptureAsync("pi_1", 49.99m); + + result.IsSuccess.Should().BeTrue(); + } + + [Fact] + public async Task Refund_IsKeyedPerIntentAndReference() + { + var client = new FakeStripeClient { Respond = (_, _) => new Refund { Id = "re_1" } }; + + var result = await Provider(client).RefundAsync("pi_1", 10m, "ret42"); + + result.IsSuccess.Should().BeTrue(); + client.Calls.Single().RequestOptions!.IdempotencyKey.Should().Be("os-refund-pi_1-ret42"); + } + + [Fact] + public async Task Refund_OfAlreadyRefundedCharge_ReturnsSuccess() + { + var client = new FakeStripeClient + { + Respond = (_, _) => StripeError(HttpStatusCode.BadRequest, "invalid_request_error", "charge_already_refunded") + }; + + var result = await Provider(client).RefundAsync("pi_1", 10m, "ocf"); + + result.IsSuccess.Should().BeTrue(); + } + [Fact] public void WhenStripeApiKeyConfigured_CreditCardResolvesToStripe() { @@ -58,4 +209,39 @@ private static ServiceProvider BuildInfrastructure(Dictionary s services.AddPaymentInfrastructure(configuration); return services.BuildServiceProvider(); } + + /// + /// Records every request the Stripe services issue. returns the entity + /// to deserialize into, or an exception to throw. + /// + private sealed class FakeStripeClient : IStripeClient + { + public List<(HttpMethod Method, string Path, BaseOptions Options, RequestOptions? RequestOptions)> Calls { get; } = []; + + public Func Respond { get; init; } = + (_, _) => new PaymentIntent { Id = "pi_1", Status = "requires_capture" }; + + public string ApiBase => "https://api.stripe.com"; + public string ApiKey => "sk_test_fake"; + public string ClientId => ""; + public string ConnectBase => ""; + public string FilesBase => ""; + public string MeterEventsBase => ""; + + public Task RequestAsync(HttpMethod method, string path, BaseOptions options, + RequestOptions requestOptions, CancellationToken cancellationToken = default) + where T : IStripeEntity + { + Calls.Add((method, path, options, requestOptions)); + return Respond(method, path) switch + { + Exception ex => Task.FromException(ex), + var entity => Task.FromResult((T)entity) + }; + } + + public Task RequestStreamingAsync(HttpMethod method, string path, BaseOptions options, + RequestOptions requestOptions, CancellationToken cancellationToken = default) + => throw new NotSupportedException(); + } } diff --git a/tests/OrderSphere.Web.Tests/Services/LoginRedirectTests.cs b/tests/OrderSphere.Web.Tests/Services/LoginRedirectTests.cs new file mode 100644 index 00000000..6392d094 --- /dev/null +++ b/tests/OrderSphere.Web.Tests/Services/LoginRedirectTests.cs @@ -0,0 +1,34 @@ +using Bunit; +using Bunit.TestDoubles; +using Microsoft.AspNetCore.Components; +using Microsoft.Extensions.DependencyInjection; +using OrderSphere.Web.Services; + +namespace OrderSphere.Web.Tests.Services; + +/// +/// The BFF accepts only local paths as returnUrl; an absolute URI would be +/// replaced with the start page. The client must therefore send the current page as +/// a path relative to the origin. +/// +public sealed class LoginRedirectTests : BunitContext +{ + [Fact] + public void Go_SendsTheCurrentPageAsLocalPath() + { + var navigation = Services.GetRequiredService(); + navigation.NavigateTo("/products/shoe?size=42"); + + LoginRedirect.Go(navigation); + + var entry = navigation.History.First(); + entry.Uri.Should().EndWith("/bff/login?returnUrl=%2Fproducts%2Fshoe%3Fsize%3D42"); + entry.Options.ForceLoad.Should().BeTrue("the BFF login endpoint lives outside the Blazor router"); + } + + [Fact] + public void Url_EscapesTheReturnPath() + { + LoginRedirect.Url("/a b").Should().Be("/bff/login?returnUrl=%2Fa%20b"); + } +} diff --git a/tests/OrderSphere.Webhooks.Tests/Features/CreateSubscriptionCommandValidatorTests.cs b/tests/OrderSphere.Webhooks.Tests/Features/CreateSubscriptionCommandValidatorTests.cs index 7e8738bf..56b38592 100644 --- a/tests/OrderSphere.Webhooks.Tests/Features/CreateSubscriptionCommandValidatorTests.cs +++ b/tests/OrderSphere.Webhooks.Tests/Features/CreateSubscriptionCommandValidatorTests.cs @@ -40,6 +40,34 @@ public async Task Validate_HttpsUrl_Passes() result.IsValid.Should().BeTrue(); } + [Theory] + [InlineData("https://127.0.0.1/")] + [InlineData("https://localhost/hook")] + [InlineData("https://ordersphere-catalog/api/v1/products")] + [InlineData("https://169.254.169.254/latest/meta-data")] + [InlineData("https://10.0.0.5/")] + [InlineData("https://[::1]/")] + [InlineData("https://user:pw@example.com/hook")] + public async Task Validate_InternalTarget_Fails(string url) + { + var result = await _validator.ValidateAsync(ValidCommand(url: url)); + result.IsValid.Should().BeFalse(); + } + + [Fact] + public async Task Validate_OverlongUrl_Fails() + { + var result = await _validator.ValidateAsync(ValidCommand(url: "https://example.com/" + new string('a', 2048))); + result.IsValid.Should().BeFalse(); + } + + [Fact] + public async Task Validate_UndefinedEventType_Fails() + { + var result = await _validator.ValidateAsync(ValidCommand(events: [(WebhookEventType)999])); + result.IsValid.Should().BeFalse(); + } + [Fact] public async Task Validate_EmptyEvents_Fails() diff --git a/tests/OrderSphere.Webhooks.Tests/Features/UpdateSubscriptionCommandValidatorTests.cs b/tests/OrderSphere.Webhooks.Tests/Features/UpdateSubscriptionCommandValidatorTests.cs new file mode 100644 index 00000000..1d42a3d9 --- /dev/null +++ b/tests/OrderSphere.Webhooks.Tests/Features/UpdateSubscriptionCommandValidatorTests.cs @@ -0,0 +1,46 @@ +using OrderSphere.Webhooks.Application.Features.Subscriptions.UpdateSubscription; + +namespace OrderSphere.Webhooks.Tests.Features; + +public sealed class UpdateSubscriptionCommandValidatorTests +{ + private readonly UpdateSubscriptionCommandValidator _validator = new(); + + private static UpdateSubscriptionCommand Command( + string url = "https://example.com/hook", + string? secret = null, + WebhookEventType[]? events = null) => + new(Guid.NewGuid(), Guid.NewGuid(), url, secret, events ?? [WebhookEventType.OrderPlaced]); + + [Fact] + public async Task Validate_PublicHttpsUrl_Passes() + { + var result = await _validator.ValidateAsync(Command()); + result.IsValid.Should().BeTrue(); + } + + [Theory] + [InlineData("http://example.com/hook")] + [InlineData("https://127.0.0.1/")] + [InlineData("https://ordersphere-payment/")] + [InlineData("https://169.254.169.254/")] + public async Task Validate_NonHttpsOrInternalTarget_Fails(string url) + { + var result = await _validator.ValidateAsync(Command(url: url)); + result.IsValid.Should().BeFalse(); + } + + [Fact] + public async Task Validate_OverlongSecret_Fails() + { + var result = await _validator.ValidateAsync(Command(secret: new string('s', 257))); + result.IsValid.Should().BeFalse(); + } + + [Fact] + public async Task Validate_EmptyEvents_Fails() + { + var result = await _validator.ValidateAsync(Command(events: [])); + result.IsValid.Should().BeFalse(); + } +} diff --git a/tests/OrderSphere.Webhooks.Tests/OrderSphere.Webhooks.Tests.csproj b/tests/OrderSphere.Webhooks.Tests/OrderSphere.Webhooks.Tests.csproj index 986c1991..073f5d9b 100644 --- a/tests/OrderSphere.Webhooks.Tests/OrderSphere.Webhooks.Tests.csproj +++ b/tests/OrderSphere.Webhooks.Tests/OrderSphere.Webhooks.Tests.csproj @@ -29,6 +29,7 @@ + diff --git a/tests/OrderSphere.Webhooks.Tests/Security/WebhookTargetPolicyTests.cs b/tests/OrderSphere.Webhooks.Tests/Security/WebhookTargetPolicyTests.cs new file mode 100644 index 00000000..1ca02c37 --- /dev/null +++ b/tests/OrderSphere.Webhooks.Tests/Security/WebhookTargetPolicyTests.cs @@ -0,0 +1,82 @@ +using System.Net; +using OrderSphere.Webhooks.Application.Security; + +namespace OrderSphere.Webhooks.Tests.Security; + +public sealed class WebhookTargetPolicyTests +{ + [Theory] + [InlineData("https://example.com/hook")] + [InlineData("https://hooks.partner.example:8443/in?x=1")] + [InlineData("https://93.184.215.14/hook")] + [InlineData("https://[2606:4700:4700::1111]/hook")] + public void IsAllowedUrl_AcceptsPublicHttpsTargets(string url) + { + WebhookTargetPolicy.IsAllowedUrl(url).Should().BeTrue(); + } + + [Theory] + [InlineData(null)] + [InlineData("")] + [InlineData("/relative")] + [InlineData("http://example.com/hook")] + [InlineData("ftp://example.com/hook")] + [InlineData("https://user:pass@example.com/hook")] + [InlineData("https://localhost/hook")] + [InlineData("https://ordersphere-catalog/api")] + [InlineData("https://api.localhost/")] + [InlineData("https://printer.local/")] + [InlineData("https://metadata.google.internal/")] + [InlineData("https://127.0.0.1/")] + [InlineData("https://10.1.2.3/")] + [InlineData("https://172.16.0.1/")] + [InlineData("https://192.168.1.1/")] + [InlineData("https://169.254.169.254/latest/meta-data")] + [InlineData("https://100.64.0.1/")] + [InlineData("https://0.0.0.0/")] + [InlineData("https://[::1]/")] + [InlineData("https://[fd00::1]/")] + [InlineData("https://[fe80::1]/")] + [InlineData("https://[::ffff:127.0.0.1]/")] + public void IsAllowedUrl_RejectsNonHttpsAndInternalTargets(string? url) + { + WebhookTargetPolicy.IsAllowedUrl(url).Should().BeFalse(); + } + + [Theory] + [InlineData("127.0.0.1")] + [InlineData("127.255.255.254")] + [InlineData("10.0.0.1")] + [InlineData("172.31.255.255")] + [InlineData("192.168.0.1")] + [InlineData("169.254.169.254")] + [InlineData("100.127.255.255")] + [InlineData("0.0.0.0")] + [InlineData("224.0.0.1")] + [InlineData("255.255.255.255")] + [InlineData("168.63.129.16")] + [InlineData("::")] + [InlineData("::1")] + [InlineData("fc00::1")] + [InlineData("fe80::1")] + [InlineData("ff02::1")] + [InlineData("::ffff:10.0.0.1")] + [InlineData("64:ff9b::a9fe:a9fe")] + [InlineData("2002:0a00:0001::1")] + public void IsBlockedAddress_BlocksInternalRanges(string address) + { + WebhookTargetPolicy.IsBlockedAddress(IPAddress.Parse(address)).Should().BeTrue(); + } + + [Theory] + [InlineData("93.184.215.14")] + [InlineData("172.32.0.1")] + [InlineData("100.128.0.1")] + [InlineData("8.8.8.8")] + [InlineData("2606:4700:4700::1111")] + [InlineData("::ffff:8.8.8.8")] + public void IsBlockedAddress_AllowsPublicAddresses(string address) + { + WebhookTargetPolicy.IsBlockedAddress(IPAddress.Parse(address)).Should().BeFalse(); + } +} diff --git a/tests/OrderSphere.Webhooks.Tests/Worker/WebhookEventProcessorTests.cs b/tests/OrderSphere.Webhooks.Tests/Worker/WebhookEventProcessorTests.cs new file mode 100644 index 00000000..83b7b700 --- /dev/null +++ b/tests/OrderSphere.Webhooks.Tests/Worker/WebhookEventProcessorTests.cs @@ -0,0 +1,78 @@ +using System.Text.Json; +using Microsoft.EntityFrameworkCore; +using OrderSphere.BuildingBlocks.Contracts.Events; +using OrderSphere.Webhooks.Tests.Helpers; +using OrderSphere.Webhooks.Worker.Workers; + +namespace OrderSphere.Webhooks.Tests.Worker; + +/// +/// Subscriptions belong to a customer. An event must only produce deliveries for the +/// subscriptions of the customer it is about — the payload carries that customer's data. +/// +public sealed class WebhookEventProcessorTests +{ + private static readonly CustomerId CustomerA = CustomerId.New(); + private static readonly CustomerId CustomerB = CustomerId.New(); + + private static string StatusChangedBody(Guid? customerId) => + JsonSerializer.Serialize(new OrderStatusChangedIntegrationEvent + { + OrderId = Guid.NewGuid(), + PreviousStatus = "Pending", + NewStatus = "Confirmed", + CustomerEmail = "a@example.com", + CustomerId = customerId + }); + + private static WebhookSubscription Subscription(CustomerId owner, params WebhookEventType[] events) => + new(owner, "https://example.com/hook", "secret", events); + + [Fact] + public async Task StageDeliveries_OnlyMatchesSubscriptionsOfTheEventsCustomer() + { + await using var db = WebhooksDbContextFactory.Create(); + var ofA = Subscription(CustomerA, WebhookEventType.OrderStatusChanged); + var ofB = Subscription(CustomerB, WebhookEventType.OrderStatusChanged); + db.Subscriptions.AddRange(ofA, ofB); + await db.SaveChangesAsync(); + + var created = await WebhookEventProcessor.StageDeliveriesAsync( + db, WebhookEventType.OrderStatusChanged, Guid.NewGuid(), StatusChangedBody(CustomerA.Value), default); + await db.SaveChangesAsync(); + + created.Should().Be(1); + var deliveries = await db.Deliveries.ToListAsync(); + deliveries.Should().ContainSingle().Which.SubscriptionId.Should().Be(ofA.Id); + } + + [Fact] + public async Task StageDeliveries_EventWithoutCustomer_IsDeliveredToNobody() + { + await using var db = WebhooksDbContextFactory.Create(); + db.Subscriptions.Add(Subscription(CustomerA, WebhookEventType.OrderStatusChanged)); + await db.SaveChangesAsync(); + + var created = await WebhookEventProcessor.StageDeliveriesAsync( + db, WebhookEventType.OrderStatusChanged, Guid.NewGuid(), StatusChangedBody(null), default); + await db.SaveChangesAsync(); + + created.Should().Be(0); + (await db.Deliveries.CountAsync()).Should().Be(0); + } + + [Fact] + public async Task StageDeliveries_IgnoresInactiveAndOtherEventSubscriptions() + { + await using var db = WebhooksDbContextFactory.Create(); + var inactive = Subscription(CustomerA, WebhookEventType.OrderStatusChanged); + inactive.Deactivate(); + db.Subscriptions.AddRange(inactive, Subscription(CustomerA, WebhookEventType.OrderPlaced)); + await db.SaveChangesAsync(); + + var created = await WebhookEventProcessor.StageDeliveriesAsync( + db, WebhookEventType.OrderStatusChanged, Guid.NewGuid(), StatusChangedBody(CustomerA.Value), default); + + created.Should().Be(0); + } +} diff --git a/tests/OrderSphere.Webhooks.Tests/Worker/WebhookTargetConnectorTests.cs b/tests/OrderSphere.Webhooks.Tests/Worker/WebhookTargetConnectorTests.cs new file mode 100644 index 00000000..e9d1a1e3 --- /dev/null +++ b/tests/OrderSphere.Webhooks.Tests/Worker/WebhookTargetConnectorTests.cs @@ -0,0 +1,31 @@ +using OrderSphere.Webhooks.Worker.Delivery; + +namespace OrderSphere.Webhooks.Tests.Worker; + +/// +/// The connect-time check is the authoritative SSRF guard: it sees the address actually being +/// dialled, including names that pass the save-time URL check but resolve internally. +/// +public sealed class WebhookTargetConnectorTests +{ + private static HttpClient Client() => new(new SocketsHttpHandler + { + UseProxy = false, + ConnectCallback = WebhookTargetConnector.ConnectAsync, + }); + + [Theory] + [InlineData("https://127.0.0.1:9/hook")] + [InlineData("https://localhost:9/hook")] + [InlineData("https://[::1]:9/hook")] + [InlineData("https://169.254.169.254/latest/meta-data")] + public async Task Connect_ToBlockedAddress_FailsBeforeDialling(string url) + { + using var client = Client(); + + var act = () => client.PostAsync(url, new StringContent("{}")); + + (await act.Should().ThrowAsync()) + .Which.Message.Should().Contain("blocked address"); + } +}