diff --git a/.gradle/8.14.3/checksums/checksums.lock b/.gradle/8.14.3/checksums/checksums.lock index 083f4f1..9fd90df 100644 Binary files a/.gradle/8.14.3/checksums/checksums.lock and b/.gradle/8.14.3/checksums/checksums.lock differ diff --git a/.gradle/8.14.3/checksums/md5-checksums.bin b/.gradle/8.14.3/checksums/md5-checksums.bin new file mode 100644 index 0000000..85c3bc5 Binary files /dev/null and b/.gradle/8.14.3/checksums/md5-checksums.bin differ diff --git a/.gradle/8.14.3/checksums/sha1-checksums.bin b/.gradle/8.14.3/checksums/sha1-checksums.bin new file mode 100644 index 0000000..028991a Binary files /dev/null and b/.gradle/8.14.3/checksums/sha1-checksums.bin differ diff --git a/.gradle/8.14.3/executionHistory/executionHistory.bin b/.gradle/8.14.3/executionHistory/executionHistory.bin new file mode 100644 index 0000000..54e6c4a Binary files /dev/null and b/.gradle/8.14.3/executionHistory/executionHistory.bin differ diff --git a/.gradle/8.14.3/executionHistory/executionHistory.lock b/.gradle/8.14.3/executionHistory/executionHistory.lock index 73c4251..927f37c 100644 Binary files a/.gradle/8.14.3/executionHistory/executionHistory.lock and b/.gradle/8.14.3/executionHistory/executionHistory.lock differ diff --git a/.gradle/8.14.3/fileHashes/fileHashes.bin b/.gradle/8.14.3/fileHashes/fileHashes.bin index 04d0d25..d312e43 100644 Binary files a/.gradle/8.14.3/fileHashes/fileHashes.bin and b/.gradle/8.14.3/fileHashes/fileHashes.bin differ diff --git a/.gradle/8.14.3/fileHashes/fileHashes.lock b/.gradle/8.14.3/fileHashes/fileHashes.lock index 668f751..a6230ae 100644 Binary files a/.gradle/8.14.3/fileHashes/fileHashes.lock and b/.gradle/8.14.3/fileHashes/fileHashes.lock differ diff --git a/.gradle/8.14.3/fileHashes/resourceHashesCache.bin b/.gradle/8.14.3/fileHashes/resourceHashesCache.bin new file mode 100644 index 0000000..31128db Binary files /dev/null and b/.gradle/8.14.3/fileHashes/resourceHashesCache.bin differ diff --git a/.gradle/buildOutputCleanup/buildOutputCleanup.lock b/.gradle/buildOutputCleanup/buildOutputCleanup.lock index 62a9154..8a3c9b3 100644 Binary files a/.gradle/buildOutputCleanup/buildOutputCleanup.lock and b/.gradle/buildOutputCleanup/buildOutputCleanup.lock differ diff --git a/.gradle/buildOutputCleanup/outputFiles.bin b/.gradle/buildOutputCleanup/outputFiles.bin new file mode 100644 index 0000000..bb3069c Binary files /dev/null and b/.gradle/buildOutputCleanup/outputFiles.bin differ diff --git a/.gradle/file-system.probe b/.gradle/file-system.probe new file mode 100644 index 0000000..0d90605 Binary files /dev/null and b/.gradle/file-system.probe differ diff --git a/build.gradle b/build.gradle index a474213..3dc7c69 100644 --- a/build.gradle +++ b/build.gradle @@ -21,22 +21,26 @@ repositories { dependencies { implementation("org.springframework.boot:spring-boot-starter-webflux") - implementation("org.springframework.boot:spring-boot-starter-data-jpa") implementation("org.springframework.boot:spring-boot-starter-data-r2dbc")/* implementation("io.r2dbc:r2dbc-postgresql")*/ - implementation("org.springframework.boot:spring-boot-starter-data-redis-reactive") - implementation("io.lettuce:lettuce-core:6.2.2.RELEASE") + /*implementation("org.springframework.boot:spring-boot-starter-data-redis-reactive") + */implementation("io.lettuce:lettuce-core:6.2.2.RELEASE") - implementation 'com.h2database:h2' +// Choose ONE driver depending on your DB: + runtimeOnly("io.r2dbc:r2dbc-h2") // for in-memory - implementation("io.github.resilience4j:resilience4j-reactor:2.0.2") implementation("io.github.resilience4j:resilience4j-ratelimiter:2.0.2") implementation("io.github.resilience4j:resilience4j-circuitbreaker:2.0.2") implementation("io.github.resilience4j:resilience4j-retry:2.0.2") - implementation("io.github.resilience4j:resilience4j-spring-boot3:2.0.2") implementation("org.springframework.boot:spring-boot-starter") + implementation("org.springdoc:springdoc-openapi-starter-webflux-ui:2.6.0") + + implementation 'io.github.resilience4j:resilience4j-spring-boot2:1.7.1' + implementation 'io.github.resilience4j:resilience4j-reactor:1.7.1' + + implementation 'org.projectlombok:lombok:1.18.32' // Use the latest stable version annotationProcessor 'org.projectlombok:lombok:1.18.32' // For annotation processing diff --git a/build/reports/problems/problems-report.html b/build/reports/problems/problems-report.html index f6cb495..200197d 100644 --- a/build/reports/problems/problems-report.html +++ b/build/reports/problems/problems-report.html @@ -650,7 +650,7 @@ code + .copy-button { diff --git a/build/resources/main/application.yml b/build/resources/main/application.yml index c09bf79..8d45436 100644 --- a/build/resources/main/application.yml +++ b/build/resources/main/application.yml @@ -19,8 +19,11 @@ resilience4j: retry-exceptions: - org.springframework.web.reactive.function.client.WebClientRequestException - java.io.IOException + - com.example.mpesa.exceptions.MpesaBusyException ignore-exceptions: - com.test.payment.exceptions.MpesaPermanentException + - java.lang.IllegalArgumentException + circuitbreaker: instances: @@ -30,11 +33,38 @@ resilience4j: failure-rate-threshold: 50 wait-duration-in-open-state: 10s -spring: - datasource: - url: jdbc:h2:mem:testdb;DB_CLOSE_DELAY=-1;DB_CLOSE_ON_EXIT=FALSE - driverClassName: org.h2.Driver - username: sa - password: password - platform: h2 + + + +spring: + r2dbc: + url: r2dbc:h2:mem:///mpesa_db;DB_CLOSE_DELAY=-1;DB_CLOSE_ON_EXIT=FALSE + username: sa + password: + sql: + init: + mode: always + schema-locations: classpath:schema.sql + main: + web-application-type: reactive + +logging: + level: + org.springframework.data.r2dbc: DEBUG + +springdoc: + swagger-ui: + path: /swagger-ui.html + operationsSorter: method + tagsSorter: alpha + api-docs: + path: /v3/api-docs + packages-to-scan: com.test.payment.controller + +mpesa: + base-url: https://sandbox.safaricom.co.ke + consumer-key: k6e7LtBNeVX7V8MPqB7P83FsZio8cRZD + consumer-secret: cGwiWzhDGopC3dho + business-short-code: 174379 + pass-key: bfb279f9aa9bdbcf158e97dd71a467cd2e0c893059b10f78e6b72ada1ed2c919 \ No newline at end of file diff --git a/build/resources/main/schema.sql b/build/resources/main/schema.sql new file mode 100644 index 0000000..9930216 --- /dev/null +++ b/build/resources/main/schema.sql @@ -0,0 +1,8 @@ +CREATE TABLE IF NOT EXISTS transactions ( + id SERIAL PRIMARY KEY, + mpesa_reference VARCHAR(255), + checkout_request_id VARCHAR(255), + status VARCHAR(50), + amount DECIMAL(10,2), + created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP +); diff --git a/build/tmp/compileJava/compileTransaction/stash-dir/MpesaController.class.uniqueId1 b/build/tmp/compileJava/compileTransaction/stash-dir/MpesaController.class.uniqueId1 new file mode 100644 index 0000000..521c44f Binary files /dev/null and b/build/tmp/compileJava/compileTransaction/stash-dir/MpesaController.class.uniqueId1 differ diff --git a/build/tmp/compileJava/compileTransaction/stash-dir/MpesaService.class.uniqueId0 b/build/tmp/compileJava/compileTransaction/stash-dir/MpesaService.class.uniqueId0 new file mode 100644 index 0000000..a0c6481 Binary files /dev/null and b/build/tmp/compileJava/compileTransaction/stash-dir/MpesaService.class.uniqueId0 differ diff --git a/build/tmp/compileJava/previous-compilation-data.bin b/build/tmp/compileJava/previous-compilation-data.bin index 3bd3e61..19a15ca 100644 Binary files a/build/tmp/compileJava/previous-compilation-data.bin and b/build/tmp/compileJava/previous-compilation-data.bin differ diff --git a/src/main/java/com/test/payment/PaymentApplication.java b/src/main/java/com/test/payment/PaymentApplication.java index f02a856..49aacad 100644 --- a/src/main/java/com/test/payment/PaymentApplication.java +++ b/src/main/java/com/test/payment/PaymentApplication.java @@ -2,10 +2,12 @@ package com.test.payment; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.data.r2dbc.repository.config.EnableR2dbcRepositories; import org.springframework.scheduling.annotation.EnableScheduling; @SpringBootApplication @EnableScheduling +@EnableR2dbcRepositories(basePackages = "com.test.payment.repository") public class PaymentApplication { public static void main(String[] args) { diff --git a/src/main/java/com/test/payment/configurations/RedisConfig.java b/src/main/java/com/test/payment/configurations/RedisConfig.java index 9c7a9d9..0bed43b 100644 --- a/src/main/java/com/test/payment/configurations/RedisConfig.java +++ b/src/main/java/com/test/payment/configurations/RedisConfig.java @@ -2,13 +2,13 @@ package com.test.payment.configurations; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; -import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory; -import org.springframework.data.redis.core.RedisTemplate; +//import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory; +//import org.springframework.data.redis.core.RedisTemplate; -@Configuration +//@Configuration public class RedisConfig { - @Bean + /* @Bean public LettuceConnectionFactory redisConnectionFactory() { return new LettuceConnectionFactory(); } @@ -18,5 +18,5 @@ public class RedisConfig { RedisTemplate template = new RedisTemplate<>(); template.setConnectionFactory(factory); return template; - } + }*/ } diff --git a/src/main/java/com/test/payment/controller/MpesaController.java b/src/main/java/com/test/payment/controller/MpesaController.java index a7bda73..4287e86 100644 --- a/src/main/java/com/test/payment/controller/MpesaController.java +++ b/src/main/java/com/test/payment/controller/MpesaController.java @@ -1,10 +1,10 @@ package com.test.payment.controller; -import com.test.payment.models.*; import com.test.payment.models.MpesaResponse; import com.test.payment.models.PaymentRequest; import com.test.payment.service.MpesaService; +import com.test.payment.service.MpesaServiceaa; import lombok.RequiredArgsConstructor; import org.springframework.http.ResponseEntity; import org.springframework.web.bind.annotation.*; @@ -22,7 +22,7 @@ public class MpesaController { return mpesaService.initiatePayment(request) .map(ResponseEntity::ok) .onErrorResume(ex -> Mono.just(ResponseEntity.badRequest() - .body(new MpesaResponse("FAILED", ex.getMessage())))); + .body(new MpesaResponse("FAILED", ex.getMessage(),"","","")))); } } diff --git a/src/main/java/com/test/payment/dto/MpesaRequestDto.java b/src/main/java/com/test/payment/dto/MpesaRequestDto.java new file mode 100644 index 0000000..22999c3 --- /dev/null +++ b/src/main/java/com/test/payment/dto/MpesaRequestDto.java @@ -0,0 +1,44 @@ +package com.test.payment.dto; + +import com.fasterxml.jackson.annotation.JsonProperty; +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; + +@Data +@AllArgsConstructor +@NoArgsConstructor +public class MpesaRequestDto { + @JsonProperty("BusinessShortCode") + private Long businessShortCode; + + @JsonProperty("Password") + private String password; + + @JsonProperty("Timestamp") + private String timestamp; + + @JsonProperty("TransactionType") + private String transactionType; + + @JsonProperty("Amount") + private Integer amount; + + @JsonProperty("PartyA") + private Long partyA; + + @JsonProperty("PartyB") + private Long partyB; + + @JsonProperty("PhoneNumber") + private Long phoneNumber; + + @JsonProperty("CallBackURL") + private String callBackURL; + + @JsonProperty("AccountReference") + private String accountReference; + + @JsonProperty("TransactionDesc") + private String transactionDesc; +} diff --git a/src/main/java/com/test/payment/dto/MpesaTokenResponse.java b/src/main/java/com/test/payment/dto/MpesaTokenResponse.java new file mode 100644 index 0000000..79a90f7 --- /dev/null +++ b/src/main/java/com/test/payment/dto/MpesaTokenResponse.java @@ -0,0 +1,14 @@ +package com.test.payment.dto; + +import com.fasterxml.jackson.annotation.JsonProperty; +import lombok.Data; + +@Data +public class MpesaTokenResponse { + + @JsonProperty("access_token") + private String accessToken; + + @JsonProperty("expires_in") + private String expiresIn; +} diff --git a/src/main/java/com/test/payment/exceptions/MpesaBusyException.java b/src/main/java/com/test/payment/exceptions/MpesaBusyException.java new file mode 100644 index 0000000..8f366ca --- /dev/null +++ b/src/main/java/com/test/payment/exceptions/MpesaBusyException.java @@ -0,0 +1,8 @@ +package com.test.payment.exceptions; + + +public class MpesaBusyException extends RuntimeException { + public MpesaBusyException(String msg) { + super(msg); + } +} diff --git a/src/main/java/com/test/payment/jobs/MpesaTransactionJob.java b/src/main/java/com/test/payment/jobs/MpesaTransactionJob.java index b94fa60..6dfeda0 100644 --- a/src/main/java/com/test/payment/jobs/MpesaTransactionJob.java +++ b/src/main/java/com/test/payment/jobs/MpesaTransactionJob.java @@ -16,7 +16,7 @@ public class MpesaTransactionJob { @Scheduled(fixedDelay = 60000) public void pullTransactions() { - mpesaRepository.findAllTransactions() + mpesaRepository.findAll() .doOnNext(tx -> log.info("Checking transaction: {}", tx)) .subscribe(); } diff --git a/src/main/java/com/test/payment/models/MpesaResponse.java b/src/main/java/com/test/payment/models/MpesaResponse.java index 60d80ed..4bda54f 100644 --- a/src/main/java/com/test/payment/models/MpesaResponse.java +++ b/src/main/java/com/test/payment/models/MpesaResponse.java @@ -1,5 +1,6 @@ package com.test.payment.models; +import com.fasterxml.jackson.annotation.JsonProperty; import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; @@ -8,6 +9,17 @@ import lombok.NoArgsConstructor; @AllArgsConstructor @NoArgsConstructor public class MpesaResponse { - private String status; - private String message; + + @JsonProperty("MerchantRequestID") + private String merchantRequestID; + + @JsonProperty("ResponseCode") + private String responseCode; + @JsonProperty("CheckoutRequestID") + private String checkoutRequestId; + @JsonProperty("CustomerMessage") + private String customerMessage; + @JsonProperty("ResponseDescription") + private String responseDescription; + } diff --git a/src/main/java/com/test/payment/models/PaymentRequest.java b/src/main/java/com/test/payment/models/PaymentRequest.java index 888f464..7000d4a 100644 --- a/src/main/java/com/test/payment/models/PaymentRequest.java +++ b/src/main/java/com/test/payment/models/PaymentRequest.java @@ -5,8 +5,8 @@ import lombok.Data; @Data public class PaymentRequest { - private String phoneNumber; - private double amount; + private long phoneNumber; + private int amount; private String accountReference; private String transactionDesc; } diff --git a/src/main/java/com/test/payment/models/Transaction.java b/src/main/java/com/test/payment/models/Transaction.java index 7c7c3d9..076a26b 100644 --- a/src/main/java/com/test/payment/models/Transaction.java +++ b/src/main/java/com/test/payment/models/Transaction.java @@ -1,21 +1,23 @@ package com.test.payment.models; -import jakarta.persistence.Entity; -import jakarta.persistence.Id; + import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; +import org.springframework.data.annotation.Id; +import org.springframework.data.relational.core.mapping.*; @Data @AllArgsConstructor @NoArgsConstructor -@Entity +@Table("transactions") public class Transaction { @Id private String id; - private String phoneNumber; - private double amount; + private long phoneNumber; + private long amount; private String status; + private String checkoutRequestId; } diff --git a/src/main/java/com/test/payment/repository/MpesaRepository.java b/src/main/java/com/test/payment/repository/MpesaRepository.java index 18298ae..289fa77 100644 --- a/src/main/java/com/test/payment/repository/MpesaRepository.java +++ b/src/main/java/com/test/payment/repository/MpesaRepository.java @@ -2,23 +2,20 @@ package com.test.payment.repository; import com.test.payment.models.Transaction; -import org.springframework.data.redis.core.RedisTemplate; +//import org.springframework.data.redis.core.RedisTemplate; +import org.springframework.data.repository.reactive.ReactiveCrudRepository; import org.springframework.stereotype.Repository; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.util.List; +import java.util.Locale; + @Repository -public class MpesaRepository { +public interface MpesaRepository extends ReactiveCrudRepository { - private final RedisTemplate redisTemplate; - - public MpesaRepository(RedisTemplate redisTemplate) { - this.redisTemplate = redisTemplate; - } - - public Mono saveTransaction(Transaction transaction) { + /* public Mono saveTransaction(Transaction transaction) { return Mono.fromRunnable(() -> redisTemplate.opsForHash().put("mpesa:transactions", transaction.getId(), transaction.getStatus()) ).then(); @@ -27,5 +24,16 @@ public class MpesaRepository { public Flux findAllTransactions() { List values = redisTemplate.opsForHash().values("mpesa:transactions"); return Flux.fromIterable(values).cast(Transaction.class); + }*/ + + Flux findByStatus(String status); + Mono findByCheckoutRequestId(String checkoutRequestId); + + /* private final RedisTemplate redisTemplate; + + public MpesaRepository(RedisTemplate redisTemplate) { + this.redisTemplate = redisTemplate; } + + */ } diff --git a/src/main/java/com/test/payment/service/MpesaService.java b/src/main/java/com/test/payment/service/MpesaService.java index ba3d913..9657783 100644 --- a/src/main/java/com/test/payment/service/MpesaService.java +++ b/src/main/java/com/test/payment/service/MpesaService.java @@ -1,21 +1,27 @@ package com.test.payment.service; -import com.test.payment.exceptions.MpesaPermanentException; +import com.test.payment.dto.MpesaRequestDto; +import com.test.payment.dto.MpesaTokenResponse; +import com.test.payment.exceptions.MpesaBusyException; import com.test.payment.exceptions.MpesaTransientException; -import com.test.payment.models.*; +import com.test.payment.models.MpesaResponse; +import com.test.payment.models.PaymentRequest; import com.test.payment.repository.MpesaRepository; -import io.github.resilience4j.circuitbreaker.*; -import io.github.resilience4j.ratelimiter.*; -import io.github.resilience4j.retry.*; +import com.test.payment.utils.MpesaUtils; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.core.env.Environment; +import org.springframework.http.HttpStatusCode; import org.springframework.stereotype.Service; import org.springframework.web.reactive.function.client.WebClient; import reactor.core.publisher.Mono; +import io.github.resilience4j.circuitbreaker.annotation.CircuitBreaker; +import io.github.resilience4j.ratelimiter.annotation.RateLimiter; +import io.github.resilience4j.retry.annotation.Retry; -import java.util.UUID; -import java.util.function.Supplier; +import java.time.Duration; +import java.util.Base64; @Service @RequiredArgsConstructor @@ -24,42 +30,175 @@ public class MpesaService { private final WebClient mpesaWebClient; private final MpesaRepository mpesaRepository; - private final MpesaTokenService tokenService; - private final RateLimiter rateLimiter; - private final Retry retry; - private final CircuitBreaker circuitBreaker; + private final Environment environment; - public Mono initiatePayment(PaymentRequest request) { - Supplier> decoratedSupplier = - CircuitBreaker.decorateSupplier(circuitBreaker, - RateLimiter.decorateSupplier(rateLimiter, - Retry.decorateSupplier(retry, () -> callMpesa(request)) - ) - ); - - return Mono.defer(decoratedSupplier) - .flatMap(response -> { - Transaction tx = new Transaction(UUID.randomUUID().toString(), - request.getPhoneNumber(), request.getAmount(), response.getStatus()); - return mpesaRepository.saveTransaction(tx).thenReturn(response); - }) - .doOnError(e -> log.error("M-Pesa API failed: {}", e.getMessage())); + public Mono getToken() { + return Mono.just("hVG5ybD47dHUMQ6RfWby2ZpIs1Ul"); } - private Mono callMpesa(PaymentRequest request) { - return tokenService.getToken() + @CircuitBreaker(name = "mpesaCircuitBreaker", fallbackMethod = "mpesaFallback") + @RateLimiter(name = "mpesaLimiter") + @Retry(name = "mpesaRetry") + public Mono initiatePayment(PaymentRequest request) { + String businessShortCode = environment.getProperty("mpesa.business-short-code"); + String passkey = environment.getProperty("mpesa.pass-key"); + String callback = "https://mydomain.com/path"; + + MpesaUtils.MpesaAuthData authData = MpesaUtils.generateAuthData(businessShortCode, passkey); + + MpesaRequestDto mpesaRequestDto = new MpesaRequestDto( + Long.valueOf(businessShortCode), + authData.getPassword(), + authData.getTimestamp(), + "CustomerPayBillOnline", + request.getAmount(), + request.getPhoneNumber(), + Long.valueOf(businessShortCode), + request.getPhoneNumber(), + callback, + request.getAccountReference(), + request.getTransactionDesc() + ); + + // Use Mono.defer so each subscription is independent (multi-user safe) + return Mono.defer(() -> + getToken() + .flatMap(token -> callMpesa(mpesaRequestDto)) + ) + // Reactive retry only on MpesaBusyException, with configurable backoff + .retryWhen( + reactor.util.retry.Retry.backoff(3, Duration.ofSeconds(10)) + .filter(ex -> ex instanceof MpesaBusyException) + .onRetryExhaustedThrow((spec, signal) -> signal.failure()) + ) + .doOnNext(resp -> log.info("M-Pesa STK Response: {}", resp)) + .doOnError(e -> log.error("M-Pesa call failed: {}", e.getMessage())); + } + +/* @CircuitBreaker(name = "mpesaCircuitBreaker", fallbackMethod = "mpesaFallback") + @RateLimiter(name = "mpesaLimiter") + @Retry(name = "mpesaRetry") + public Mono initiatePayment(PaymentRequest request) { + String businessShortCode = environment.getProperty("mpesa.business-short-code"); + String passkey = environment.getProperty("mpesa.pass-key"); + String callback = "https://mydomain.com/path"; + + MpesaUtils.MpesaAuthData authData = MpesaUtils.generateAuthData(businessShortCode, passkey); + + + MpesaRequestDto mpesaRequestDto = new MpesaRequestDto( + Long.valueOf(businessShortCode), + authData.getPassword(), + authData.getTimestamp(), + "CustomerPayBillOnline", + request.getAmount(), + request.getPhoneNumber(), + Long.valueOf(businessShortCode), + request.getPhoneNumber(), + callback, + request.getAccountReference(), + request.getTransactionDesc() + ); + *//*return Mono.delay(Duration.ofSeconds(1)) + .then(callMpesa(mpesaRequestDto)) + .doOnNext(resp -> log.info("M-Pesa STK Response: {}", resp)) + .doOnError(e -> log.error("M-Pesa API failed: {}", e.getMessage()));*//* + + *//* return callMpesa(mpesaRequestDto) + .doOnNext(resp -> log.info("M-Pesa STK Response: {}", resp)) + .doOnError(e -> { + if (e.toString().contains("System is busy")) { + log.warn("M-Pesa system is busy"); + throe Mono.error(new MpesaBusyException("System busy")); + } else { + log.error("M-Pesa API failed: {}", e.getMessage()); + } + });*//* + + return getToken() + .flatMap(token -> callMpesa(mpesaRequestDto)) + .retryWhen( + reactor.util.retry.Retry.backoff(3, Duration.ofSeconds(10)) + .filter(ex -> ex instanceof MpesaBusyException) + .onRetryExhaustedThrow((retryBackoffSpec, retrySignal) -> + retrySignal.failure() + ) + ) .doOnError(e -> log.error("M-Pesa call failed: {}", e.getMessage())); + }*/ + + /*@Retry(name = "mpesaRetry", fallbackMethod = "mpesaFallback") + @RateLimiter(name = "mpesaLimiter") + @CircuitBreaker(name = "mpesaCB", fallbackMethod = "mpesaFallback") + public Mono initiatePayment(PaymentRequest request) { + + String businessShortCode = environment.getProperty("mpesa.business-short-code"); + String passkey = environment.getProperty("mpesa.pass-key"); + String callback = "https://mydomain.com/path"; + + MpesaUtils.MpesaAuthData authData = MpesaUtils.generateAuthData(businessShortCode, passkey); + + + MpesaRequestDto mpesaRequestDto = new MpesaRequestDto( + Long.valueOf(businessShortCode), + authData.getPassword(), + authData.getTimestamp(), + "CustomerPayBillOnline", + request.getAmount(), + request.getPhoneNumber(), + Long.valueOf(businessShortCode), + request.getPhoneNumber(), + callback, + request.getAccountReference(), + request.getTransactionDesc() + ); + + return mpesaWebClient.post() + .uri("/mpesa/stkpush/v1/processrequest") + .bodyValue(mpesaRequestDto) + .retrieve() + .bodyToMono(MpesaResponse.class) + .flatMap(resp -> { + if (resp.toString().contains("System is busy")) { + return Mono.error(new MpesaBusyException("System busy")); + } + return Mono.just(resp); + }); + }*/ + + + + private Mono callMpesa(MpesaRequestDto request) { + return getToken() .flatMap(token -> mpesaWebClient.post() .uri("/mpesa/stkpush/v1/processrequest") .header("Authorization", "Bearer " + token) .bodyValue(request) .retrieve() + .onStatus(HttpStatusCode::isError, clientResponse -> + clientResponse.bodyToMono(String.class) + .flatMap(errorBody -> { + log.error("M-Pesa returned {} with body: {}", clientResponse.statusCode(), errorBody); + if (errorBody.contains("System is busy")) { + return Mono.error(new MpesaBusyException("System busy")); + } + return Mono.error(new MpesaTransientException(errorBody, null)); + }) + ) .bodyToMono(MpesaResponse.class) - .timeout(java.time.Duration.ofSeconds(20)) - .onErrorResume(ex -> { - log.error("Error calling M-Pesa: {}", ex.getMessage()); - return Mono.error(new MpesaTransientException("Temporary M-Pesa issue", ex)); - }) ); } + + // 🧯 Fallback if circuit is open or all retries fail + private Mono mpesaFallback(PaymentRequest request, Throwable ex) { + log.error("⚠️ Mpesa fallback triggered: {}", ex.getMessage()); + return Mono.just(new MpesaResponse( + "500", + "Fallback triggered due to service unavailability", + null, + "", + null + )); + } } + diff --git a/src/main/java/com/test/payment/service/MpesaServiceaa.java b/src/main/java/com/test/payment/service/MpesaServiceaa.java new file mode 100644 index 0000000..c8d9aa9 --- /dev/null +++ b/src/main/java/com/test/payment/service/MpesaServiceaa.java @@ -0,0 +1,190 @@ +package com.test.payment.service; + + +import com.test.payment.dto.MpesaRequestDto; +import com.test.payment.exceptions.MpesaTransientException; +import com.test.payment.models.*; +import com.test.payment.repository.MpesaRepository; +import com.test.payment.utils.MpesaUtils; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.core.env.Environment; +import org.springframework.http.HttpStatusCode; +import org.springframework.stereotype.Service; +import org.springframework.web.reactive.function.client.WebClient; +import reactor.core.publisher.Mono; + +/*@Service +@RequiredArgsConstructor +@Slf4j*/ +public class MpesaServiceaa { +/* + private final WebClient mpesaWebClient; + private final MpesaRepository mpesaRepository; + private final RateLimiter rateLimiter; + private final Retry retry; + private final CircuitBreaker circuitBreaker; + private final Environment environment; + + public Mono getToken() { +*//* // 1️⃣ Check Redis cache first + return redisTemplate.opsForValue().get("mpesa:token") + .flatMap(cachedToken -> { + if (cachedToken != null) { + return Mono.just(cachedToken); + } + return fetchNewToken(); + }) + // if no cached token, fetch and cache new one + .switchIfEmpty(fetchNewToken());*//* + return Mono.fromSupplier(() -> "Q3HCot6dNLTpUtpitkr5Tatsa4KB"); + // return fetchNewToken(); + } + + private Mono fetchNewToken() { + + String baseUrl = environment.getProperty("mpesa.base-url"); + String consumerKey = environment.getProperty("mpesa.consumer-key"); + String consumerSecret = environment.getProperty("mpesa.consumer-secret"); + + String credentials = consumerKey + ":" + consumerSecret; + String encodedCredentials = Base64.getEncoder().encodeToString(credentials.getBytes()); + + return mpesaWebClient.get() + .uri(baseUrl + "/oauth/v1/generate?grant_type=client_credentials") + .header("Authorization", "Basic " + encodedCredentials) + .retrieve() + .bodyToMono(MpesaTokenResponse.class) + .map(MpesaTokenResponse::getAccessToken); + *//*.flatMap(token -> + redisTemplate.opsForValue() + .set("mpesa:token", token, Duration.ofMinutes(50)) + .thenReturn(token) + );*//* + } + + public Mono initiatePayment(PaymentRequest request) { + + String businessShortCode = environment.getProperty("mpesa.business-short-code"); + + String passkey = environment.getProperty("mpesa.pass-key"); + + String callback = "https://mydomain.com/path"; + + MpesaUtils.MpesaAuthData authData = MpesaUtils.generateAuthData(businessShortCode, passkey); + + MpesaRequestDto mpesaRequestDto = new MpesaRequestDto(Long.valueOf(businessShortCode),authData.getPassword(),authData.getTimestamp(),"CustomerPayBillOnline",request.getAmount(),request.getPhoneNumber(), + Long.valueOf(businessShortCode),request.getPhoneNumber(),callback,request.getAccountReference(), request.getTransactionDesc()); + + Supplier> decoratedSupplier = + CircuitBreaker.decorateSupplier(circuitBreaker, + RateLimiter.decorateSupplier(rateLimiter, + Retry.decorateSupplier(retry, () -> callMpesa(mpesaRequestDto)) + ) + ); + + return Mono.defer(decoratedSupplier) + *//* .flatMap(response -> { + *//**//* Transaction tx = new Transaction(UUID.randomUUID().toString(), + request.getPhoneNumber(), request.getAmount(), response.getResponseCode(), response.getCheckoutRequestId()); + return mpesaRepository.save(tx).thenReturn(response);*//**//* + return response; + })*//* + .doOnError(e -> log.error("M-Pesa API failed: {}", e.getMessage())); + } + + private Mono callMpesa(MpesaRequestDto request) { + + getToken() + .map(token -> { + System.out.println("Token: " + token); + return token; + }) + .subscribe(); + + return getToken() + .flatMap(token -> + mpesaWebClient.post() + .uri("/mpesa/stkpush/v1/processrequest") + .header("Authorization", "Bearer " + token) + .bodyValue(request) + .retrieve() + // Intercept 4xx/5xx and extract actual body + .onStatus(HttpStatusCode::isError, clientResponse -> + clientResponse.bodyToMono(String.class) + .flatMap(errorBody -> { + log.error("M-Pesa returned {} with body: {}", clientResponse.statusCode(), errorBody); + // Return a Mono.error so onErrorResume below can handle it + return Mono.error(new RuntimeException(errorBody)); + }) + ) + .bodyToMono(MpesaResponse.class) + // Catch the above RuntimeException and return the error body as normal data + .onErrorResume(RuntimeException.class, ex -> { + log.error("Returning raw M-Pesa 500 error body: {}", ex.getMessage()); + return Mono.error(new MpesaTransientException(ex.getMessage(), ex)); // This is the actual 500 error body from M-Pesa + }) + .doOnNext(response -> log.info("M-Pesa STK Response: {}", response))); + + }*/ + + /* private final WebClient mpesaWebClient; + private final MpesaRepository mpesaRepository; + private final io.github.resilience4j.ratelimiter.RateLimiter rateLimiter; + private final io.github.resilience4j.retry.Retry retry; + private final io.github.resilience4j.circuitbreaker.CircuitBreaker circuitBreaker; + private final Environment environment; + + public Mono getToken() { + return Mono.fromSupplier(() -> "Q3HCot6dNLTpUtpitkr5Tatsa4KB"); + } + + public Mono initiatePayment(PaymentRequest request) { + String businessShortCode = environment.getProperty("mpesa.business-short-code"); + String passkey = environment.getProperty("mpesa.pass-key"); + String callback = "https://mydomain.com/path"; + + MpesaUtils.MpesaAuthData authData = MpesaUtils.generateAuthData(businessShortCode, passkey); + + MpesaRequestDto mpesaRequestDto = new MpesaRequestDto( + Long.valueOf(businessShortCode), + authData.getPassword(), + authData.getTimestamp(), + "CustomerPayBillOnline", + request.getAmount(), + request.getPhoneNumber(), + Long.valueOf(businessShortCode), + request.getPhoneNumber(), + callback, + request.getAccountReference(), + request.getTransactionDesc() + ); + + // Reactive chaining with operators (not Supplier) + return callMpesa(mpesaRequestDto) + .transformDeferred(ReactorRateLimiterOperator.of(rateLimiter)) + .transformDeferred(ReactorRetryOperator.of(retry)) + .transformDeferred(ReactorCircuitBreakerOperator.of(circuitBreaker)) + .doOnNext(resp -> log.info("M-Pesa STK Response: {}", resp)) + .doOnError(e -> log.error("M-Pesa API failed: {}", e.getMessage())); + } + + private Mono callMpesa(MpesaRequestDto request) { + return getToken() + .flatMap(token -> + mpesaWebClient.post() + .uri("/mpesa/stkpush/v1/processrequest") + .header("Authorization", "Bearer " + token) + .bodyValue(request) + .retrieve() + .onStatus(HttpStatusCode::isError, clientResponse -> + clientResponse.bodyToMono(String.class) + .flatMap(errorBody -> { + log.error("M-Pesa returned {} with body: {}", clientResponse.statusCode(), errorBody); + return Mono.error(new MpesaTransientException(errorBody, null)); + }) + ) + .bodyToMono(MpesaResponse.class) + ); + }*/ +} diff --git a/src/main/java/com/test/payment/service/MpesaTokenService.java b/src/main/java/com/test/payment/service/MpesaTokenService.java index 2d10bc5..f39537d 100644 --- a/src/main/java/com/test/payment/service/MpesaTokenService.java +++ b/src/main/java/com/test/payment/service/MpesaTokenService.java @@ -2,7 +2,7 @@ package com.test.payment.service; import lombok.RequiredArgsConstructor; -import org.springframework.data.redis.core.RedisTemplate; +//import org.springframework.data.redis.core.RedisTemplate; import org.springframework.stereotype.Service; import reactor.core.publisher.Mono; @@ -12,16 +12,16 @@ import java.time.Duration; @RequiredArgsConstructor public class MpesaTokenService { - private final RedisTemplate redisTemplate; + // private final RedisTemplate redisTemplate; public Mono getToken() { - String cachedToken = redisTemplate.opsForValue().get("mpesa:token"); + /* String cachedToken = redisTemplate.opsForValue().get("mpesa:token"); if (cachedToken != null) { return Mono.just(cachedToken); } // Simulate token fetch from M-Pesa auth endpoint String newToken = "access_token_" + System.currentTimeMillis(); - redisTemplate.opsForValue().set("mpesa:token", newToken, Duration.ofMinutes(50)); - return Mono.just(newToken); + redisTemplate.opsForValue().set("mpesa:token", newToken, Duration.ofMinutes(50));*/ + return Mono.just("newToken"); } } diff --git a/src/main/java/com/test/payment/utils/MpesaUtils.java b/src/main/java/com/test/payment/utils/MpesaUtils.java new file mode 100644 index 0000000..6cedc3a --- /dev/null +++ b/src/main/java/com/test/payment/utils/MpesaUtils.java @@ -0,0 +1,60 @@ +package com.test.payment.utils; + + + +import java.nio.charset.StandardCharsets; +import java.text.SimpleDateFormat; +import java.util.Base64; +import java.util.Date; +import java.util.TimeZone; + +public class MpesaUtils { + + /** + * Generates a timestamp in the format yyyyMMddHHmmss + * Example: 20251010162455 + */ + public static String generateTimestamp() { + SimpleDateFormat sdf = new SimpleDateFormat("yyyyMMddHHmmss"); + // Set timezone to Africa/Nairobi + sdf.setTimeZone(TimeZone.getTimeZone("Africa/Nairobi")); + return sdf.format(new Date()); + } + + /** + * Generates a Base64-encoded password for M-Pesa STK Push. + * Format: Base64(BusinessShortCode + Passkey + Timestamp) + */ + public static String generatePassword(String businessShortCode, String passkey, String timestamp) { + String dataToEncode = businessShortCode + passkey + timestamp; + return Base64.getEncoder().encodeToString(dataToEncode.getBytes(StandardCharsets.UTF_8)); + } + + /** + * Helper method to generate both password and timestamp together. + */ + public static MpesaAuthData generateAuthData(String businessShortCode, String passkey) { + String timestamp = generateTimestamp(); + String password = generatePassword(businessShortCode, passkey, timestamp); + return new MpesaAuthData(password, timestamp); + } + + // Inner class to hold both password and timestamp + public static class MpesaAuthData { + private final String password; + private final String timestamp; + + public MpesaAuthData(String password, String timestamp) { + this.password = password; + this.timestamp = timestamp; + } + + public String getPassword() { + return password; + } + + public String getTimestamp() { + return timestamp; + } + } +} diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index c09bf79..8d45436 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -19,8 +19,11 @@ resilience4j: retry-exceptions: - org.springframework.web.reactive.function.client.WebClientRequestException - java.io.IOException + - com.example.mpesa.exceptions.MpesaBusyException ignore-exceptions: - com.test.payment.exceptions.MpesaPermanentException + - java.lang.IllegalArgumentException + circuitbreaker: instances: @@ -30,11 +33,38 @@ resilience4j: failure-rate-threshold: 50 wait-duration-in-open-state: 10s -spring: - datasource: - url: jdbc:h2:mem:testdb;DB_CLOSE_DELAY=-1;DB_CLOSE_ON_EXIT=FALSE - driverClassName: org.h2.Driver - username: sa - password: password - platform: h2 + + + +spring: + r2dbc: + url: r2dbc:h2:mem:///mpesa_db;DB_CLOSE_DELAY=-1;DB_CLOSE_ON_EXIT=FALSE + username: sa + password: + sql: + init: + mode: always + schema-locations: classpath:schema.sql + main: + web-application-type: reactive + +logging: + level: + org.springframework.data.r2dbc: DEBUG + +springdoc: + swagger-ui: + path: /swagger-ui.html + operationsSorter: method + tagsSorter: alpha + api-docs: + path: /v3/api-docs + packages-to-scan: com.test.payment.controller + +mpesa: + base-url: https://sandbox.safaricom.co.ke + consumer-key: k6e7LtBNeVX7V8MPqB7P83FsZio8cRZD + consumer-secret: cGwiWzhDGopC3dho + business-short-code: 174379 + pass-key: bfb279f9aa9bdbcf158e97dd71a467cd2e0c893059b10f78e6b72ada1ed2c919 \ No newline at end of file diff --git a/src/main/resources/schema.sql b/src/main/resources/schema.sql new file mode 100644 index 0000000..9930216 --- /dev/null +++ b/src/main/resources/schema.sql @@ -0,0 +1,8 @@ +CREATE TABLE IF NOT EXISTS transactions ( + id SERIAL PRIMARY KEY, + mpesa_reference VARCHAR(255), + checkout_request_id VARCHAR(255), + status VARCHAR(50), + amount DECIMAL(10,2), + created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP +);