From 3596cee4fd2a2aafcc0ebf71caa855ea8dad3a31 Mon Sep 17 00:00:00 2001 From: Moritz Waldau Date: Wed, 23 Sep 2026 21:58:27 +0200 Subject: [PATCH 1/6] fix(webhooks): scope deliveries to the event's customer and block internal targets Subscriptions were matched by event type only, so every subscriber of OrderStatusChanged received all customers' events. Events now carry an optional CustomerId (additive contract change) and only that customer's subscriptions receive a delivery; events without it go to nobody. Webhook URLs are restricted to public HTTPS hosts when saved, and every resolved address is checked again at connect time (covers DNS rebinding and service-discovery rewrites). Redirects are not followed, no proxy is used, and a failed delivery stores only the status code, not the response body. Co-Authored-By: Claude Opus 5.5 --- .../OrderStatusChangedIntegrationEvent.cs | 6 + .../CreateSubscriptionCommandValidator.cs | 13 +- .../UpdateSubscriptionCommandValidator.cs | 13 +- .../Security/WebhookTargetPolicy.cs | 124 ++++++++++++++++++ .../Delivery/WebhookTargetConnector.cs | 40 ++++++ .../OrderSphere.Webhooks.Worker.csproj | 4 + .../OrderSphere.Webhooks.Worker/Program.cs | 10 ++ .../Workers/WebhookDeliveryProcessor.cs | 7 +- .../Workers/WebhookEventProcessor.cs | 76 +++++++---- ...CreateSubscriptionCommandValidatorTests.cs | 28 ++++ ...UpdateSubscriptionCommandValidatorTests.cs | 46 +++++++ .../OrderSphere.Webhooks.Tests.csproj | 1 + .../Security/WebhookTargetPolicyTests.cs | 82 ++++++++++++ .../Worker/WebhookEventProcessorTests.cs | 78 +++++++++++ .../Worker/WebhookTargetConnectorTests.cs | 31 +++++ 15 files changed, 528 insertions(+), 31 deletions(-) create mode 100644 src/Services/Webhooks/OrderSphere.Webhooks.Application/Security/WebhookTargetPolicy.cs create mode 100644 src/Services/Webhooks/OrderSphere.Webhooks.Worker/Delivery/WebhookTargetConnector.cs create mode 100644 tests/OrderSphere.Webhooks.Tests/Features/UpdateSubscriptionCommandValidatorTests.cs create mode 100644 tests/OrderSphere.Webhooks.Tests/Security/WebhookTargetPolicyTests.cs create mode 100644 tests/OrderSphere.Webhooks.Tests/Worker/WebhookEventProcessorTests.cs create mode 100644 tests/OrderSphere.Webhooks.Tests/Worker/WebhookTargetConnectorTests.cs 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/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.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"); + } +} From 27d8617a124db303ba4812196290b7c216bd7ac3 Mon Sep 17 00:00:00 2001 From: Moritz Waldau Date: Wed, 23 Sep 2026 21:58:40 +0200 Subject: [PATCH 2/6] fix(ordering): guard order transitions and stop re-confirming cancelled orders Order.Confirm/MarkShipped/MarkDelivered/Cancel return Result and reject invalid transitions; Confirm is only allowed from Created. Rehydration stays unguarded so existing streams still load. PaymentResultProcessor checks the order status before confirming the stock reservation: a captured payment for a cancelled order is refunded instead of turning the order Paid, a failed payment for a cancelled order only closes the saga, and a failure reported for a paid order is logged for review instead of cancelling it. OrderStatusChanged events carry the CustomerId. Co-Authored-By: Claude Opus 5.5 --- .../Order/Admin/CancelOrderCommandHandler.cs | 8 +- .../Admin/UpdateOrderStatusCommandHandler.cs | 27 ++-- .../Entities/Order.cs | 34 +++-- .../Workers/PaymentResultProcessor.cs | 119 ++++++++++++++---- .../Aggregates/OrderTests.cs | 64 ++++++++-- .../PaymentResultToOrderFlowTests.cs | 104 +++++++++++++++ 6 files changed, 290 insertions(+), 66 deletions(-) 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/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.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() { From f1b44ecae3eacd34d8de7e6d6d1f07d938b12ec9 Mon Sep 17 00:00:00 2001 From: Moritz Waldau Date: Wed, 23 Sep 2026 21:58:54 +0200 Subject: [PATCH 3/6] fix(payment): make the Stripe payment path idempotent and reconcilable Every Stripe call carries a deterministic idempotency key (order- or intent-scoped; refunds per refund reason), so a Service Bus redelivery replays the original request instead of creating a second PaymentIntent. Only declines and invalid requests become Result failures; transient faults propagate and are retried by redelivery. The SDK timeout drops to 30 s so one run stays inside the lock-renewal window. A failed capture voids the authorization and keeps the intent id on the record. PaymentRecord transitions return Result and reject invalid moves. The Stripe webhook finds the record by orderId metadata or intent id, answers 503 while the worker has not stored it yet, rejects a missing signature with 400 instead of 500, applies only valid transitions, and logs contradictions as EventId 5001. Adds a non-unique index on payments.TransactionId. Co-Authored-By: Claude Opus 5.5 --- docs/operations.md | 26 +++ .../Endpoints/StripeWebhookEndpoints.cs | 132 ++++++++--- .../Logging/PaymentApiLog.cs | 27 +++ .../Entities/PaymentRecord.cs | 41 +++- .../Errors/PaymentErrors.cs | 2 + .../DependencyInjection.cs | 9 +- .../PaymentRecordConfiguration.cs | 2 + ...4_AddPaymentTransactionIdIndex.Designer.cs | 217 ++++++++++++++++++ ...0923194654_AddPaymentTransactionIdIndex.cs | 27 +++ .../PaymentDbContextModelSnapshot.cs | 2 + .../Providers/CreditCardPaymentProvider.cs | 9 +- .../Providers/IPaymentProvider.cs | 22 +- .../Providers/InvoicePaymentProvider.cs | 9 +- .../Providers/PayPalPaymentProvider.cs | 9 +- .../Providers/StripePaymentProvider.cs | 141 ++++++++++-- .../OrderConfirmationFailedProcessor.cs | 3 +- .../Workers/PaymentProcessor.cs | 73 ++++-- .../Workers/RefundRequestedProcessor.cs | 3 +- .../Aggregates/PaymentRecordTests.cs | 85 +++++++ .../Api/StripeWebhookTests.cs | 149 ++++++++++++ .../OrderConfirmationFailedProcessorTests.cs | 6 +- .../PaymentProcessorTests.cs | 119 ++++++++-- .../PaymentProviderTests.cs | 9 +- .../StripePaymentProviderTests.cs | 192 +++++++++++++++- 24 files changed, 1202 insertions(+), 112 deletions(-) create mode 100644 src/Services/Payment/OrderSphere.Payment.Api/Logging/PaymentApiLog.cs create mode 100644 src/Services/Payment/OrderSphere.Payment.Infrastructure/Migrations/20260923194654_AddPaymentTransactionIdIndex.Designer.cs create mode 100644 src/Services/Payment/OrderSphere.Payment.Infrastructure/Migrations/20260923194654_AddPaymentTransactionIdIndex.cs create mode 100644 tests/OrderSphere.IntegrationTests/Api/StripeWebhookTests.cs 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/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/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.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(); + } } From cdeebae229ef099cfcaa02e8a5a22f1666cf4505 Mon Sep 17 00:00:00 2001 From: Moritz Waldau Date: Wed, 23 Sep 2026 21:59:10 +0200 Subject: [PATCH 4/6] fix(bff): accept only local returnUrl on /bff/login The login endpoint redirected to any returnUrl after sign-in, an open redirect. It now accepts local paths only (same rule as ASP.NET Core's IsLocalUrl) and falls back to "/". The client sends the current page as a relative path. Co-Authored-By: Claude Opus 5.5 --- .../OrderSphere.Web/Services/LoginRedirect.cs | 8 +++-- .../OrderSphere.Bff/Auth/LocalReturnUrl.cs | 21 +++++++++++ src/Gateways/OrderSphere.Bff/Program.cs | 9 ++--- .../LocalReturnUrlTests.cs | 36 +++++++++++++++++++ .../Services/LoginRedirectTests.cs | 34 ++++++++++++++++++ 5 files changed, 100 insertions(+), 8 deletions(-) create mode 100644 src/Gateways/OrderSphere.Bff/Auth/LocalReturnUrl.cs create mode 100644 tests/OrderSphere.Bff.Tests/LocalReturnUrlTests.cs create mode 100644 tests/OrderSphere.Web.Tests/Services/LoginRedirectTests.cs 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.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/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.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"); + } +} From 50b9ecf4ba1456eee6b38f383e2be77668f46406 Mon Sep 17 00:00:00 2001 From: Moritz Waldau Date: Wed, 23 Sep 2026 21:59:11 +0200 Subject: [PATCH 5/6] fix(gateway): route the Stripe webhook and admin coupons Stripe could not reach the webhook: the BFF required a session and CSRF token on /api/**, and the gateway required a JWT on /api/v1/payments/**. The BFF now exposes an anonymous POST /webhooks/stripe (outside /api) that forwards to an anonymous gateway route; the Payment API authenticates it by Stripe signature. Adds the missing /api/v1/admin/coupons gateway route (admin UI got 404). Co-Authored-By: Claude Opus 5.5 --- .../OrderSphere.ApiGateway/appsettings.json | 16 ++++++ .../appsettings.Development.json | 12 +++++ src/Gateways/OrderSphere.Bff/appsettings.json | 12 +++++ .../StripeWebhookRouteTests.cs | 49 +++++++++++++++++++ 4 files changed, 89 insertions(+) create mode 100644 tests/OrderSphere.Bff.Tests/StripeWebhookRouteTests.cs 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/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/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); + } +} From 01faac06d434e1e0f009ab032b3cea4251f355fe Mon Sep 17 00:00:00 2001 From: Moritz Waldau Date: Wed, 23 Sep 2026 21:59:11 +0200 Subject: [PATCH 6/6] docs: changelog for audit phase 1 fixes Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 26 ++++++++++++++++++++++++++ 1 file changed, 26 insertions(+) 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).