package com.test.payment.service; import com.test.payment.dto.CallbackAckDto; import com.test.payment.dto.PaymentRequest; import com.test.payment.dto.PaymentResultDto; import com.test.payment.dto.TransactionStatusDto; import com.test.payment.models.AirtelPaymentCallback; import com.test.payment.models.AirtelPaymentResponse; import com.test.payment.models.MpesaPaymentCallback; import com.test.payment.models.MpesaPaymentResponse; import com.test.payment.models.MtnPaymentCallback; import com.test.payment.models.MtnPaymentResponse; import com.test.payment.models.PaymentInitiation; import com.test.payment.models.PaymentProviderType; import com.test.payment.models.Status; import com.test.payment.models.Transaction; import com.test.payment.repository.AirtelPaymentCallbackRepository; import com.test.payment.repository.AirtelPaymentResponseRepository; import com.test.payment.repository.MpesaPaymentCallbackRepository; import com.test.payment.repository.MpesaPaymentResponseRepository; import com.test.payment.repository.MtnPaymentCallbackRepository; import com.test.payment.repository.MtnPaymentResponseRepository; import com.test.payment.repository.PaymentInitiationRepository; import com.test.payment.repository.audit.AirtelRequestRepository; import com.test.payment.repository.audit.MpesaRequestRepository; import com.test.payment.repository.audit.MtnRequestRepository; import com.test.payment.repository.TransactionRepository; import com.test.payment.service.PaymentLifecycleService.CallbackData; import com.test.payment.service.PaymentLifecycleService.ProviderResponseData; import com.test.payment.service.PaymentLifecycleService.QueryOutcome; import com.test.payment.service.PaymentLifecycleService.StoredResponse; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.dao.DataIntegrityViolationException; import org.springframework.http.HttpStatus; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import org.springframework.web.server.ResponseStatusException; import java.math.BigDecimal; import java.time.LocalDateTime; import java.util.List; import java.util.Optional; /** * The blocking, transactional half of {@link PaymentLifecycleService}. Every method * here is one unit of work against JPA; the facade is what puts them on a * bounded-elastic thread so the WebFlux event loop is never blocked. * *

It lives in its own bean deliberately: {@code @Transactional} is applied by a * Spring AOP proxy, which self-invocation from the facade would bypass — the same * trap the Resilience4j annotations have in the client classes. * *

Responses and callbacks live in per-operator tables that share no supertype, so * this class dispatches on {@link com.test.payment.models.Operator} and hands the * rest of the lifecycle a provider-neutral {@link StoredResponse} view. Note the * provider equality checks: one table holds every market for that operator, so a * MPESA_TZ reference must not resolve against a MPESA_KE payment. */ @Service @RequiredArgsConstructor @Slf4j public class PaymentLifecycleStore { private final PaymentInitiationRepository initiationRepository; private final TransactionRepository transactionRepository; private final StatusCatalog statuses; private final MpesaPaymentResponseRepository mpesaResponses; private final MpesaPaymentCallbackRepository mpesaCallbacks; private final AirtelPaymentResponseRepository airtelResponses; private final AirtelPaymentCallbackRepository airtelCallbacks; private final MtnPaymentResponseRepository mtnResponses; private final MtnPaymentCallbackRepository mtnCallbacks; // the outbound-call audit tables, so a transaction can point at the call it made private final MpesaRequestRepository mpesaRequests; private final AirtelRequestRepository airtelRequests; private final MtnRequestRepository mtnRequests; /** * What a status check needs before deciding whether to hit the provider: the * stored response, whether the payment is still open, and the state to serve * when it is not. */ public record StatusCheckContext(StoredResponse response, boolean pending, TransactionStatusDto currentState) { } /** The callback fields the status DTO reads back. */ public record StoredCallback(String resultCode, String resultDesc, String receiptNumber) { } /** * Opens a payment attempt: the initiation and its consolidated transaction row are * created together. The transaction exists from the request onwards rather than * appearing only at resolution, so it can be linked to the provider's response and * callback as each arrives, and a payment is never invisible in TRANSACTIONS. */ @Transactional public PaymentInitiation saveInitiation(PaymentProviderType provider, PaymentRequest request) { PaymentInitiation initiation = initiationRepository.save(PaymentInitiation.builder() .Provider(provider) .PhoneNumber(request.getPhoneNumber()) .Amount(BigDecimal.valueOf(request.getAmount())) .AccountReference(request.getAccountReference()) .TransactionDesc(request.getTransactionDesc()) .Status(statuses.pending()) .CreatedAt(LocalDateTime.now()) .build()); transactionRepository.save(Transaction.builder() .Initiation(initiation) .Provider(provider) .PhoneNumber(initiation.getPhoneNumber()) .Amount(initiation.getAmount()) .AccountReference(initiation.getAccountReference()) .Status(initiation.getStatus()) .CreatedAt(LocalDateTime.now()) .build()); return initiation; } @Transactional public PaymentResultDto persistResponse(Long initiationId, ProviderResponseData data) { PaymentInitiation initiation = requireInitiation(initiationId); Status newStatus = data.accepted() ? statuses.pending() : statuses.failed(); // one response per initiation — an existing row wins, a concurrent insert falls back to it StoredResponse saved = findResponseByInitiation(initiation.getProvider(), initiationId) .orElseGet(() -> insertResponse(initiation, data)); PaymentInitiation updated = updateStatus(initiation, newStatus); // linked either way: a rejection resolves the payment, an acceptance just // records which response row belongs to it recordTransaction(updated, saved, null, data.responseDescription(), null, null, data.accepted() ? null : "REJECTION"); return PaymentResultDto.builder() .initiationId(updated.getId()) .provider(updated.getProvider()) .status(updated.statusName()) .providerReference(saved.providerReference()) .secondaryReference(saved.secondaryReference()) .responseCode(saved.responseCode()) .responseDescription(saved.responseDescription()) .customerMessage(saved.customerMessage()) .build(); } @Transactional public void markFailed(Long initiationId, String reason) { PaymentInitiation updated = updateStatus(requireInitiation(initiationId), statuses.failed()); recordTransaction(updated, null, null, truncate(reason), null, null, "ERROR"); } @Transactional public CallbackAckDto applyCallback(PaymentProviderType provider, CallbackData data, String rawPayload) { Optional match = findResponseByReference(provider, data.providerReference()); if (match.isEmpty()) { log.warn("[{}] callback for unknown reference {}", provider, data.providerReference()); return CallbackAckDto.accepted("Unknown reference"); } StoredResponse response = match.get(); if (findCallbackByInitiation(provider, response.initiationId()).isPresent()) { log.info("[{}] duplicate callback for {} ignored", provider, data.providerReference()); return CallbackAckDto.accepted("Duplicate callback ignored"); } return saveCallback(response, data, rawPayload); } /** * Loads what a status check needs. Never calls the provider — the HTTP query is * the facade's job, so no transaction is held open across the network. */ @Transactional(readOnly = true) public StatusCheckContext loadForStatusCheck(PaymentProviderType provider, String providerReference) { StoredResponse response = findResponseByReference(provider, providerReference) .orElseThrow(() -> new ResponseStatusException(HttpStatus.NOT_FOUND, "No " + provider + " transaction found for reference " + providerReference)); PaymentInitiation initiation = requireInitiation(response.initiationId()); boolean pending = statuses.isPending(initiation.statusName()); return new StatusCheckContext(response, pending, pending ? null : buildStatusDto(initiation, response)); } /** Applies the outcome of a live provider status query. */ @Transactional public TransactionStatusDto applyQueryOutcome(Long initiationId, QueryOutcome outcome) { PaymentInitiation initiation = requireInitiation(initiationId); StoredResponse response = findResponseByInitiation(initiation.getProvider(), initiationId).orElse(null); if (statuses.isPending(outcome.newStatus())) { TransactionStatusDto dto = buildStatusDto(initiation, response); if (outcome.resultDesc() != null) { dto.setResultDesc(outcome.resultDesc()); } return dto; } PaymentInitiation updated = updateStatus(initiation, outcome.newStatus()); recordTransaction(updated, response, outcome.resultCode(), outcome.resultDesc(), outcome.receiptNumber(), null, "QUERY"); TransactionStatusDto dto = buildStatusDto(updated, response); dto.setResultCode(outcome.resultCode()); dto.setResultDesc(outcome.resultDesc()); if (outcome.receiptNumber() != null) { dto.setReceiptNumber(outcome.receiptNumber()); } return dto; } /** * First half of a reconciliation: resolves the provider reference to re-query, * or terminally fails the initiation when there is nothing to query with. * * @return the provider reference, or empty when the initiation was failed outright */ @Transactional public Optional beginReconcile(Long initiationId) { PaymentInitiation initiation = requireInitiation(initiationId); Optional response = findResponseByInitiation(initiation.getProvider(), initiationId); if (response.isEmpty()) { log.warn("[{}] initiation {} never received a provider response — marking FAILED", initiation.getProvider(), initiationId); failTerminal(initiation, "No provider response received"); return Optional.empty(); } if (response.get().providerReference() == null) { failTerminal(initiation, "No provider reference on response"); return Optional.empty(); } return Optional.of(response.get().providerReference()); } @Transactional(readOnly = true) public List listTransactions(PaymentProviderType provider) { return provider == null ? transactionRepository.findAllWithAssociations() : transactionRepository.findByProvider(provider); } @Transactional(readOnly = true) public List findPendingOlderThan(LocalDateTime cutoff) { return initiationRepository.findByStatusNameAndCreatedAtBefore(statuses.pending().getName(), cutoff); } // --- per-operator dispatch ------------------------------------------------- private StoredResponse insertResponse(PaymentInitiation initiation, ProviderResponseData data) { PaymentProviderType provider = initiation.getProvider(); try { return switch (provider.operator()) { case MPESA -> view(mpesaResponses.saveAndFlush(MpesaPaymentResponse.builder() .Initiation(initiation).Provider(provider) .ProviderReference(data.providerReference()) .SecondaryReference(data.secondaryReference()) .ResponseCode(data.responseCode()) .ResponseDescription(data.responseDescription()) .CustomerMessage(data.customerMessage()) .CreatedAt(LocalDateTime.now()).build())); case AIRTEL -> view(airtelResponses.saveAndFlush(AirtelPaymentResponse.builder() .Initiation(initiation).Provider(provider) .ProviderReference(data.providerReference()) .SecondaryReference(data.secondaryReference()) .ResponseCode(data.responseCode()) .ResponseDescription(data.responseDescription()) .CustomerMessage(data.customerMessage()) .CreatedAt(LocalDateTime.now()).build())); case MTN -> view(mtnResponses.saveAndFlush(MtnPaymentResponse.builder() .Initiation(initiation).Provider(provider) .ProviderReference(data.providerReference()) .SecondaryReference(data.secondaryReference()) .ResponseCode(data.responseCode()) .ResponseDescription(data.responseDescription()) .CustomerMessage(data.customerMessage()) .CreatedAt(LocalDateTime.now()).build())); }; } catch (DataIntegrityViolationException ex) { return findResponseByInitiation(provider, initiation.getId()).orElseThrow(() -> ex); } } private Optional findResponseByInitiation(PaymentProviderType provider, Long initiationId) { return switch (provider.operator()) { case MPESA -> mpesaResponses.findByInitiationId(initiationId).map(this::view); case AIRTEL -> airtelResponses.findByInitiationId(initiationId).map(this::view); case MTN -> mtnResponses.findByInitiationId(initiationId).map(this::view); }; } /** The provider filter matters: one table holds every market for that operator. */ private Optional findResponseByReference(PaymentProviderType provider, String providerReference) { if (providerReference == null) { return Optional.empty(); } Optional found = switch (provider.operator()) { case MPESA -> mpesaResponses.findByProviderReference(providerReference).map(this::view); case AIRTEL -> airtelResponses.findByProviderReference(providerReference).map(this::view); case MTN -> mtnResponses.findByProviderReference(providerReference).map(this::view); }; return found.filter(response -> provider == response.provider()); } private Optional findCallbackByInitiation(PaymentProviderType provider, Long initiationId) { return switch (provider.operator()) { case MPESA -> mpesaCallbacks.findByInitiationId(initiationId) .map(c -> new StoredCallback(c.getResultCode(), c.getResultDesc(), c.getReceiptNumber())); case AIRTEL -> airtelCallbacks.findByInitiationId(initiationId) .map(c -> new StoredCallback(c.getResultCode(), c.getResultDesc(), c.getReceiptNumber())); case MTN -> mtnCallbacks.findByInitiationId(initiationId) .map(c -> new StoredCallback(c.getResultCode(), c.getResultDesc(), c.getReceiptNumber())); }; } private CallbackAckDto saveCallback(StoredResponse response, CallbackData data, String rawPayload) { PaymentProviderType provider = response.provider(); PaymentInitiation initiation = requireInitiation(response.initiationId()); LocalDateTime now = LocalDateTime.now(); try { switch (provider.operator()) { case MPESA -> mpesaCallbacks.saveAndFlush(MpesaPaymentCallback.builder() .Initiation(initiation).Provider(provider) .ProviderReference(data.providerReference()).ResultCode(data.resultCode()) .ResultDesc(truncate(data.resultDesc())).ReceiptNumber(data.receiptNumber()) .Amount(data.amount()).PhoneNumber(data.phoneNumber()) .TransactionDate(data.transactionDate()).RawPayload(rawPayload) .CreatedAt(now).build()); case AIRTEL -> airtelCallbacks.saveAndFlush(AirtelPaymentCallback.builder() .Initiation(initiation).Provider(provider) .ProviderReference(data.providerReference()).ResultCode(data.resultCode()) .ResultDesc(truncate(data.resultDesc())).ReceiptNumber(data.receiptNumber()) .Amount(data.amount()).PhoneNumber(data.phoneNumber()) .TransactionDate(data.transactionDate()).RawPayload(rawPayload) .CreatedAt(now).build()); case MTN -> mtnCallbacks.saveAndFlush(MtnPaymentCallback.builder() .Initiation(initiation).Provider(provider) .ProviderReference(data.providerReference()).ResultCode(data.resultCode()) .ResultDesc(truncate(data.resultDesc())).ReceiptNumber(data.receiptNumber()) .Amount(data.amount()).PhoneNumber(data.phoneNumber()) .TransactionDate(data.transactionDate()).RawPayload(rawPayload) .CreatedAt(now).build()); } } catch (DataIntegrityViolationException ex) { log.info("[{}] concurrent duplicate callback for {} ignored", provider, data.providerReference()); return CallbackAckDto.accepted("Duplicate callback ignored"); } // A receipt means the money actually moved, which is Paid rather than a bare Success. Status newStatus = data.success() ? (data.receiptNumber() != null ? statuses.paid() : statuses.success()) : statuses.failed(); PaymentInitiation updated = updateStatus(initiation, newStatus); Transaction tx = recordTransaction(updated, response, data.resultCode(), data.resultDesc(), data.receiptNumber(), data.transactionDate(), "CALLBACK"); log.info("[{}] callback processed for initiation {} — status {}", tx.getProvider(), tx.initiationId(), tx.statusName()); return CallbackAckDto.accepted("Callback processed"); } private StoredResponse view(MpesaPaymentResponse r) { return new StoredResponse(r.initiationId(), r.getProvider(), r.getProviderReference(), r.getSecondaryReference(), r.getResponseCode(), r.getResponseDescription(), r.getCustomerMessage()); } private StoredResponse view(AirtelPaymentResponse r) { return new StoredResponse(r.initiationId(), r.getProvider(), r.getProviderReference(), r.getSecondaryReference(), r.getResponseCode(), r.getResponseDescription(), r.getCustomerMessage()); } private StoredResponse view(MtnPaymentResponse r) { return new StoredResponse(r.initiationId(), r.getProvider(), r.getProviderReference(), r.getSecondaryReference(), r.getResponseCode(), r.getResponseDescription(), r.getCustomerMessage()); } // --- shared lifecycle ------------------------------------------------------ private Transaction failTerminal(PaymentInitiation initiation, String reason) { PaymentInitiation updated = updateStatus(initiation, statuses.failed()); return recordTransaction(updated, null, null, reason, null, null, "RECONCILIATION"); } private TransactionStatusDto buildStatusDto(PaymentInitiation initiation, StoredResponse response) { // result details come from the callback when we have one, otherwise from the // consolidated transaction row (e.g. when a status query resolved the payment) Optional cb = findCallbackByInitiation(initiation.getProvider(), initiation.getId()); Optional tx = transactionRepository.findByInitiationId(initiation.getId()); return TransactionStatusDto.builder() .initiationId(initiation.getId()) .provider(initiation.getProvider()) .providerReference(response == null ? null : response.providerReference()) .secondaryReference(response == null ? null : response.secondaryReference()) .status(initiation.statusName()) .phoneNumber(initiation.getPhoneNumber()) .amount(initiation.getAmount()) .accountReference(initiation.getAccountReference()) .resultCode(cb.map(StoredCallback::resultCode) .or(() -> tx.map(Transaction::getResultCode)).orElse(null)) .resultDesc(cb.map(StoredCallback::resultDesc) .or(() -> tx.map(Transaction::getResultDesc)).orElse(null)) .receiptNumber(cb.map(StoredCallback::receiptNumber) .or(() -> tx.map(Transaction::getReceiptNumber)).orElse(null)) .createdAt(initiation.getCreatedAt()) .updatedAt(initiation.getUpdatedAt()) .build(); } private PaymentInitiation updateStatus(PaymentInitiation initiation, Status status) { initiation.setStatus(status); initiation.setUpdatedAt(LocalDateTime.now()); return initiationRepository.save(initiation); } /** * Upserts the consolidated TRANSACTIONS row for an initiation that reached a * terminal state. Keyed by initiation (UNIQUE) so it can never duplicate; * a later, richer resolution (e.g. a callback after a query) updates the row. */ private Transaction recordTransaction(PaymentInitiation initiation, StoredResponse response, String resultCode, String resultDesc, String receiptNumber, String transactionDate, String resolvedBy) { Transaction tx = transactionRepository.findByInitiationId(initiation.getId()) .orElseGet(() -> Transaction.builder() .Initiation(initiation) .Provider(initiation.getProvider()) .PhoneNumber(initiation.getPhoneNumber()) .Amount(initiation.getAmount()) .AccountReference(initiation.getAccountReference()) .CreatedAt(LocalDateTime.now()) .build()); if (response != null) { tx.setProviderReference(response.providerReference()); tx.setSecondaryReference(response.secondaryReference()); } tx.setStatus(initiation.getStatus()); if (resultCode != null) { tx.setResultCode(resultCode); } if (resultDesc != null) { tx.setResultDesc(resultDesc); } if (receiptNumber != null) { tx.setReceiptNumber(receiptNumber); } if (transactionDate != null) { tx.setTransactionDate(transactionDate); } if (resolvedBy != null) { tx.setResolvedBy(resolvedBy); } if (tx.getId() != null) { tx.setUpdatedAt(LocalDateTime.now()); } attachRequest(tx, initiation.getProvider(), initiation.getId()); attachCallback(tx, initiation.getProvider(), initiation.getId()); Transaction saved = transactionRepository.save(tx); log.info("[{}] transaction {} recorded for initiation {} — status {} (via {})", initiation.getProvider(), saved.getId(), initiation.getId(), saved.statusName(), resolvedBy); return saved; } /** * Links the call that started this payment — the earliest audit request for the * initiation, since later rows are status queries. Best effort: audit rows are * written asynchronously, so on the rare occasion the row has not landed yet the * link is simply picked up by the next update. */ private void attachRequest(Transaction tx, PaymentProviderType provider, Long initiationId) { String reference = String.valueOf(initiationId); switch (provider.operator()) { case MPESA -> earliest(mpesaRequests.findByInitiationId(reference)).ifPresent(tx::setMpesaRequest); case AIRTEL -> earliest(airtelRequests.findByInitiationId(reference)).ifPresent(tx::setAirtelRequest); case MTN -> earliest(mtnRequests.findByInitiationId(reference)).ifPresent(tx::setMtnRequest); } } private Optional earliest(List rows) { return rows.isEmpty() ? Optional.empty() : Optional.of(rows.get(0)); } /** Same for the callback, which only exists once the operator has reported back. */ private void attachCallback(Transaction tx, PaymentProviderType provider, Long initiationId) { switch (provider.operator()) { case MPESA -> mpesaCallbacks.findByInitiationId(initiationId).ifPresent(tx::setMpesaCallback); case AIRTEL -> airtelCallbacks.findByInitiationId(initiationId).ifPresent(tx::setAirtelCallback); case MTN -> mtnCallbacks.findByInitiationId(initiationId).ifPresent(tx::setMtnCallback); } } private PaymentInitiation requireInitiation(Long initiationId) { return initiationRepository.findById(initiationId) .orElseThrow(() -> new IllegalStateException("Initiation " + initiationId + " no longer exists")); } private String truncate(String value) { return value == null || value.length() <= 255 ? value : value.substring(0, 255); } }