using Dapr.Client; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.PubSub; using GB5Shared.EventLogPublish; using GB5Shared.GenerateAutoNumber; using GB5Shared.PubSub.OutBox; using GB5Shared.QueryExecutor; using GoodBooks.PAY.PAYBLL.Abstractions; using GoodBooks.PAY.PAYBLL.Factory; using GoodBooks.PAY.PAYBLL.Vault; using GoodBooks.PAY.PAYDAL.CustomCode.Cashback; using GoodBooks.PAY.PAYDAL.CustomCode.Gateway; using GoodBooks.PAY.PAYDAL.CustomCode.Loyalty; using GoodBooks.PAY.PAYDAL.CustomCode.Order; using GoodBooks.PAY.PAYDAL.CustomCode.Webhook; using GoodBooks.PAY.PAYDAL.DTO.Cashback; using GoodBooks.PAY.PAYDAL.DTO.Loyalty; using GoodBooks.PAY.PAYDAL.DTO.Order; using GoodBooks.PAY.PAYDAL.DTO.Webhook; using Microsoft.Extensions.Caching.Distributed; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.Logging; using static GB5Shared.GB5Constant.Constant; namespace GoodBooks.PAY.PAYBLL.Webhook { // ───────────────────────────────────────────────────────────────────────── // PayWebhookBLL // // 15-step webhook processing flow: // 1. INSERT webhook event row immediately (before any verification) // 2. Load gateway config → fetch webhook secret from Vault // 3-4. Parse & verify signature // 5. Update signature validity on event row // 6. Idempotency check // 7. Distributed lock (30-second advisory lock) // 8. Load TPAYORDER by gateway order ID // 9. Determine new order status // 10. AutoNumber for PayTransactionId // 11. BEGIN TX: insert TPAYTRANSACTION + update TPAYORDER status // + (on Paid) insert loyalty ledger earn row + cashback ledger row + Finance outbox // 12. COMMIT TX // 13. Update webhook ISPROCESSED=1 (no tx) // 14. Fire-and-forget Dapr publish "payorder.statuschanged" // 15. Log success // // RULE: ProcessWebhookAsync NEVER throws — catch at outermost level. // RULE: Always call UpdateProcessed so the webhook row reflects final state. // ───────────────────────────────────────────────────────────────────────── public class PayWebhookBLL : IPayWebhookBLL { private readonly IPayWebhookEventDAL _webhookDal; private readonly IPayGatewayConfigDAL _gwConfigDal; private readonly IPayOrderDAL _payOrderDal; private readonly IPayTransactionDAL _payTxnDal; private readonly ILoyaltyLedgerDAL _loyaltyLedgerDal; private readonly ICashbackLedgerDAL _cashbackLedgerDal; private readonly IPaymentGatewayFactory _gatewayFactory; private readonly IVaultService _vault; private readonly IOutBox _outBox; private readonly AutoNumber _autoNumber; private readonly IQueryExecutor _queryExecutor; private readonly IDistributedCache _distributedCache; private readonly DaprClient _daprClient; private readonly IConfiguration _config; private readonly ILogger _logger; public PayWebhookBLL( IPayWebhookEventDAL webhookDal, IPayGatewayConfigDAL gwConfigDal, IPayOrderDAL payOrderDal, IPayTransactionDAL payTxnDal, ILoyaltyLedgerDAL loyaltyLedgerDal, ICashbackLedgerDAL cashbackLedgerDal, IPaymentGatewayFactory gatewayFactory, IVaultService vault, IOutBox outBox, AutoNumber autoNumber, IQueryExecutor queryExecutor, IDistributedCache distributedCache, DaprClient daprClient, IConfiguration config, ILogger logger) { _webhookDal = webhookDal; _gwConfigDal = gwConfigDal; _payOrderDal = payOrderDal; _payTxnDal = payTxnDal; _loyaltyLedgerDal = loyaltyLedgerDal; _cashbackLedgerDal = cashbackLedgerDal; _gatewayFactory = gatewayFactory; _vault = vault; _outBox = outBox; _autoNumber = autoNumber; _queryExecutor = queryExecutor; _distributedCache = distributedCache; _daprClient = daprClient; _config = config; _logger = logger; } // ── BuildSystemLoginAsync ───────────────────────────────────────────── public Task BuildSystemLoginAsync(int tenantId, string gatewayCode, CancellationToken ct) { // Gateway calls don't carry a DB hint — use configured system database name. // The actual DB routing happens inside IQueryExecutor based on DatabaseName. var systemLogin = new LoginDTO { ClientId = tenantId, DatabaseName = _config["PayJobs:SystemDatabaseName"] ?? "GoodBooks_Main", UserId = -1 }; return Task.FromResult(systemLogin); } // ── GetWebhookEventList ─────────────────────────────────────────────── public async Task GetWebhookEventList(int pageOffset, int pageSize, LoginDTO login, CancellationToken ct) { try { return await _webhookDal.GetWebhookEventList(pageOffset, pageSize, login) .ConfigureAwait(false); } catch (Exception ex) { _logger.LogError(ex, "GetWebhookEventList failed for offset {Offset}/{PageSize}", pageOffset, pageSize); throw; } } // ── ProcessWebhookAsync — 15-step flow ─────────────────────────────── public async Task ProcessWebhookAsync( string gatewayCode, int tenantId, string rawPayload, string sigHeader, LoginDTO systemLogin, CancellationToken ct) { int webhookEventId = 0; try { // ── Step 1: INSERT webhook event immediately (before verification) ── var eventDto = new PayWebhookEventDTO { PayGatewayId = 0, // unknown until config loaded; updated below IdempotencyKey = string.Empty, // unknown until parsed; placeholder EventType = "UNVERIFIED", RawPayload = rawPayload, SignatureValid = 0, // 0 = Unverified IsProcessed = 0, CreatedById = -1, CreatedOn = DateTime.UtcNow, TenantId = tenantId }; webhookEventId = await _webhookDal.InsertWebhookEvent(eventDto, systemLogin) .ConfigureAwait(false); // ── Step 2: Load gateway config ──────────────────────────────── var config = await _gwConfigDal .GetPayGatewayConfigForWebhook(gatewayCode, tenantId, systemLogin) .ConfigureAwait(false); if (config is null) { _logger.LogWarning("No gateway config for {GatewayCode}/{TenantId} — cannot verify webhook", gatewayCode, tenantId); await _webhookDal.UpdateSignatureValid(webhookEventId, 2, systemLogin) .ConfigureAwait(false); await _webhookDal.UpdateProcessed(webhookEventId, 2, "No gateway config found", systemLogin) .ConfigureAwait(false); return; } // Step 2b: Fetch webhook secret from Vault string webhookSecret = await _vault .GetSecretAsync(config.WebhookSecretVaultKey, ct) .ConfigureAwait(false); // ── Steps 3-4: Parse and verify signature ─────────────────────── var gateway = _gatewayFactory.GetGateway(gatewayCode); GatewayWebhookEvent? webhookEvent = await gateway .ParseAndVerifyWebhookAsync(rawPayload, sigHeader, webhookSecret, ct) .ConfigureAwait(false); if (webhookEvent is null) { _logger.LogWarning("Webhook signature invalid for {GatewayCode}/{TenantId}", gatewayCode, tenantId); await _webhookDal.UpdateSignatureValid(webhookEventId, 2, systemLogin) .ConfigureAwait(false); await _webhookDal.UpdateProcessed(webhookEventId, 2, "Invalid signature", systemLogin) .ConfigureAwait(false); return; // ALWAYS return — never throw } // ── Step 5: Mark signature valid ──────────────────────────────── await _webhookDal.UpdateSignatureValid(webhookEventId, 1, systemLogin) .ConfigureAwait(false); // ── Step 6: Idempotency check ─────────────────────────────────── bool alreadyProcessed = await _webhookDal .CheckIdempotency(webhookEvent.IdempotencyKey, config.PayGatewayId, systemLogin) .ConfigureAwait(false); if (alreadyProcessed) { _logger.LogInformation("Duplicate webhook {Key} for {GatewayCode}", webhookEvent.IdempotencyKey, gatewayCode); await _webhookDal.UpdateProcessed(webhookEventId, 3, "Duplicate", systemLogin) .ConfigureAwait(false); return; } // ── Step 7: Distributed advisory lock ─────────────────────────── string lockKey = $"paywebhook:{gatewayCode}:{webhookEvent.IdempotencyKey}"; byte[]? existing = await _distributedCache.GetAsync(lockKey, ct).ConfigureAwait(false); if (existing is not null) { // Another instance is processing the same event — leave ISPROCESSED=0 for retry _logger.LogWarning("Concurrent webhook processing detected for {Key}", lockKey); return; } await _distributedCache.SetAsync(lockKey, [1], new DistributedCacheEntryOptions { AbsoluteExpirationRelativeToNow = TimeSpan.FromSeconds(30) }, ct).ConfigureAwait(false); // ── Step 8: Load TPAYORDER by gateway order ID ────────────────── var order = await _payOrderDal .GetPayOrderByGatewayOrderId(webhookEvent.GatewayOrderId, systemLogin) .ConfigureAwait(false); if (order is null) { _logger.LogWarning("No TPAYORDER found for GatewayOrderId {GwOrderId}", webhookEvent.GatewayOrderId); await _webhookDal.UpdateProcessed(webhookEventId, 2, "Order not found", systemLogin) .ConfigureAwait(false); return; } // ── Step 9: Determine new order status ────────────────────────── int newOrderStatus = webhookEvent.EventType.ToLowerInvariant() switch { "payment.captured" => 3, // Paid "payment.success" => 3, "charge.succeeded" => 3, "payment_intent.succeeded" => 3, "payment.failed" => 5, // Failed "charge.failed" => 5, "payment_intent.payment_failed" => 5, _ => 0 // Unknown — mark processed/ignored }; if (newOrderStatus == 0) { _logger.LogInformation("Unknown event type {EventType} — skipping order update", webhookEvent.EventType); await _webhookDal.UpdateProcessed(webhookEventId, 1, null, systemLogin) .ConfigureAwait(false); return; } // ── Step 10: AutoNumber for PayTransactionId ───────────────────── var txnAutoNum = await _autoNumber .GetAutoNumber(1, "PAYTRANSACTION", systemLogin) .ConfigureAwait(false); int payTransactionId = txnAutoNum.StartNumber; // ── Steps 11-12: DB transaction ───────────────────────────────── var tx = await _queryExecutor.BeginTransactionAsync(systemLogin).ConfigureAwait(false); try { // 11a: INSERT TPAYTRANSACTION var txnDto = new PayTransactionDTO { PayTransactionId = payTransactionId, PayOrderId = order.PayOrderId, TxnStatus = newOrderStatus == 3 ? (byte)2 : (byte)5, // 2=Captured, 5=Failed GatewayTxnId = webhookEvent.GatewayTxnId, TxnCurrency = order.OrderCurrency, TxnAmt = webhookEvent.Amount, TxnInitiatedOn = DateTime.UtcNow, TxnCompletedOn = newOrderStatus == 3 ? DateTime.UtcNow : (DateTime?)null, Version = 0, Status = 1, SortOrder = 9999, CreatedById = -1, CreatedOn = DateTime.UtcNow, ModifiedById = -1, ModifiedOn = DateTime.UtcNow, TenantId = tenantId }; await _payTxnDal.InsertPayTransaction(txnDto, systemLogin, tx).ConfigureAwait(false); // 11b: UPDATE TPAYORDER status await _payOrderDal.UpdatePayOrderStatus( order.PayOrderId, newOrderStatus, webhookEvent.GatewayOrderId, webhookEvent.Amount, systemLogin, tx).ConfigureAwait(false); // 11c: On successful payment — write subsidiary rows if (newOrderStatus == 3) { // i. INSERT TLOYALTYLEDGER earn row (append-only placeholder) var ledgerAutoNum = await _autoNumber .GetAutoNumber(1, "LOYALTYLEDGER", systemLogin).ConfigureAwait(false); var ledgerDto = new LoyaltyLedgerDTO { LoyaltyLedgerId = ledgerAutoNum.StartNumber, EnrollmentId = -1, // Placeholder — background job resolves LoyaltyProgramId = 0, CustomerId = order.CustomerId, EntryType = 1, // Earn Points = 0m, // Placeholder — background job calculates ReferenceOrderId = order.PayOrderId, Remarks = "Payment earn pending", CreatedById = -1, CreatedOn = DateTime.UtcNow, TenantId = tenantId }; await _loyaltyLedgerDal.InsertLedgerEntry(ledgerDto, systemLogin, tx) .ConfigureAwait(false); // ii. INSERT TCASHBACKLEDGER pending row var cashbackDto = new CashbackLedgerDTO { PayOrderId = order.PayOrderId, CashbackRuleId = 0, // Resolved by payout job CustomerId = order.CustomerId, CashbackAmt = 0m, // Calculated by payout job PayoutMode = 1, // Default CreditNote; job may adjust PayoutStatus = 1, // Pending ScheduledPayoutOn = DateTime.UtcNow.AddDays(1), RetryCount = 0, Version = 0, Status = 1, SortOrder = 9999, CreatedById = -1, CreatedOn = DateTime.UtcNow, ModifiedById = -1, ModifiedOn = DateTime.UtcNow, TenantId = tenantId }; await _cashbackLedgerDal.SaveCashbackLedger(cashbackDto, systemLogin, tx) .ConfigureAwait(false); // iii. Write Finance notification to outbox (same transaction) var outboxDto = new OutboxDTO { ObjectTypeId = EntityConstant.OBJECTTPAYORDER, ObjectId = order.PayOrderId, EventTypeId = EventTypeConstant.PAYORDERPAIDSENTTYPEID, CorrelationKey = $"payorder-{order.PayOrderId}", Payload = System.Text.Json.JsonSerializer.Serialize(new { PayOrderId = order.PayOrderId, Amount = webhookEvent.Amount, Currency = order.OrderCurrency, TenantId = tenantId }) }; await _outBox.PublishEventAsync(outboxDto, systemLogin, tx) .ConfigureAwait(false); } // Step 12: COMMIT await _queryExecutor.CommitAsync(tx).ConfigureAwait(false); } catch { await _queryExecutor.RollbackAsync(tx).ConfigureAwait(false); throw; } // ── Step 13: Mark webhook processed (no transaction) ───────────── await _webhookDal.UpdateProcessed(webhookEventId, 1, null, systemLogin) .ConfigureAwait(false); // ── Step 14: Fire-and-forget Dapr event ───────────────────────── try { await _daprClient.PublishEventAsync( _config["Dapr:PubSubName"] ?? "pubsub", "payorder.statuschanged", new { PayOrderId = order.PayOrderId, NewStatus = newOrderStatus, TenantId = tenantId }, ct).ConfigureAwait(false); } catch (Exception ex) { _logger.LogWarning(ex, "Dapr publish failed for PayOrder {PayOrderId} — non-fatal", order.PayOrderId); } // ── Step 15: Log success ───────────────────────────────────────── _logger.LogInformation( "Webhook processed {GatewayCode}/{TenantId} IdempotencyKey={Key} NewOrderStatus={Status}", gatewayCode, tenantId, webhookEvent.IdempotencyKey, newOrderStatus); } catch (Exception ex) { // Outermost safety net — BLL must never let an exception propagate to the HTTP layer. // Caller's HandleAsync always returns HTTP 200 regardless. _logger.LogError(ex, "Webhook processing FAILED {GatewayCode}/{TenantId}", gatewayCode, tenantId); if (webhookEventId > 0) { try { string errorMsg = ex.Message.Length > 500 ? ex.Message[..500] : ex.Message; await _webhookDal.UpdateProcessed(webhookEventId, 2, errorMsg, systemLogin) .ConfigureAwait(false); } catch { // Secondary failure — ignore; logging above is the record } } } } } }