diff --git a/accountant/accountant-app/src/main/kotlin/co/nilin/opex/accountant/app/scheduler/FinancialActionsArchiveJob.kt b/accountant/accountant-app/src/main/kotlin/co/nilin/opex/accountant/app/scheduler/FinancialActionsArchiveJob.kt index f0d26804b..a6acea9e8 100644 --- a/accountant/accountant-app/src/main/kotlin/co/nilin/opex/accountant/app/scheduler/FinancialActionsArchiveJob.kt +++ b/accountant/accountant-app/src/main/kotlin/co/nilin/opex/accountant/app/scheduler/FinancialActionsArchiveJob.kt @@ -25,15 +25,26 @@ class FinancialActionsArchiveJob( @Value("\${app.fi-action.archive.batch-size:1000}") private var batchSize: Int = 1000 + @Value("\${app.fi-action.archive.max-batches-per-run:20}") + private var maxBatchesPerRun: Int = 20 + @Scheduled(fixedDelayString = "\${app.fi-action.archive.fixed-delay-ms:300000}", initialDelay = 60000) fun archiveProcessedActions() { - if (!enabled || batchSize <= 0 || retentionDays <= 0) return + if (!enabled || batchSize <= 0 || retentionDays <= 0 || maxBatchesPerRun <= 0) return runBlocking { val before = LocalDateTime.now().minusDays(retentionDays) - val archived = financialActionPersister.archiveProcessedActions(before, batchSize) - if (archived > 0) { - log.info("Archived $archived processed financial actions older than $before") + var totalArchived = 0 + var shouldContinue = true + repeat(maxBatchesPerRun) { + if (!shouldContinue) return@repeat + val archived = financialActionPersister.archiveProcessedActions(before, batchSize) + totalArchived += archived + if (archived < batchSize) shouldContinue = false + } + + if (totalArchived > 0) { + log.info("Archived $totalArchived processed financial actions older than $before") } } } diff --git a/accountant/accountant-ports/accountant-persister-postgres/src/main/kotlin/co/nilin/opex/accountant/ports/postgres/dao/FinancialActionRepository.kt b/accountant/accountant-ports/accountant-persister-postgres/src/main/kotlin/co/nilin/opex/accountant/ports/postgres/dao/FinancialActionRepository.kt index 16dbeddf4..86b7fb960 100644 --- a/accountant/accountant-ports/accountant-persister-postgres/src/main/kotlin/co/nilin/opex/accountant/ports/postgres/dao/FinancialActionRepository.kt +++ b/accountant/accountant-ports/accountant-persister-postgres/src/main/kotlin/co/nilin/opex/accountant/ports/postgres/dao/FinancialActionRepository.kt @@ -9,7 +9,6 @@ import org.springframework.data.repository.query.Param import org.springframework.data.repository.reactive.ReactiveCrudRepository import org.springframework.stereotype.Repository import reactor.core.publisher.Mono -import java.math.BigDecimal import java.time.LocalDateTime @Repository @@ -22,13 +21,23 @@ interface FinancialActionRepository : ReactiveCrudRepository - @Query("select count(1) from fi_actions fi where fi.sender = :uuid and fi.symbol = :symbol and fi.event_type = :eventType and fi.status != :status") - fun countByUuidAndSymbolAndEventTypeAndStatusNot( + @Query( + """ + select exists( + select 1 + from fi_actions fi + where fi.sender = :uuid + and fi.symbol = :symbol + and fi.event_type = :eventType + and fi.status <> 'PROCESSED' + ) + """ + ) + fun existsUnprocessedBySenderAndSymbolAndEventType( @Param("uuid") uuid: String, @Param("symbol") symbol: String, - @Param("eventType") eventType: String, - @Param("status") financialActionStatus: FinancialActionStatus - ): Mono + @Param("eventType") eventType: String + ): Mono @Query("select * from fi_actions fi where status != :status") fun findByStatusNot(@Param("status") status: String, paging: Pageable): Flow @@ -69,6 +78,11 @@ interface FinancialActionRepository : ReactiveCrudRepository 'PROCESSED' + ) order by create_date limit :limit ), diff --git a/accountant/accountant-ports/accountant-persister-postgres/src/main/kotlin/co/nilin/opex/accountant/ports/postgres/impl/FinancialActionLoaderImpl.kt b/accountant/accountant-ports/accountant-persister-postgres/src/main/kotlin/co/nilin/opex/accountant/ports/postgres/impl/FinancialActionLoaderImpl.kt index 1fb65e088..c4e281c2c 100644 --- a/accountant/accountant-ports/accountant-persister-postgres/src/main/kotlin/co/nilin/opex/accountant/ports/postgres/impl/FinancialActionLoaderImpl.kt +++ b/accountant/accountant-ports/accountant-persister-postgres/src/main/kotlin/co/nilin/opex/accountant/ports/postgres/impl/FinancialActionLoaderImpl.kt @@ -15,7 +15,6 @@ import kotlinx.coroutines.reactive.awaitFirstOrElse import org.springframework.data.domain.PageRequest import org.springframework.data.domain.Sort import org.springframework.stereotype.Component -import java.math.BigDecimal import java.time.LocalDateTime @Component @@ -48,12 +47,11 @@ class FinancialActionLoaderImpl( } override suspend fun countUnprocessed(userUuid: String, symbol: String, eventType: String): Long { - return financialActionRepository.countByUuidAndSymbolAndEventTypeAndStatusNot( + return if (financialActionRepository.existsUnprocessedBySenderAndSymbolAndEventType( userUuid, symbol, - eventType, - FinancialActionStatus.PROCESSED - ).awaitFirstOrElse { BigDecimal.ZERO }.toLong() + eventType + ).awaitFirstOrElse { false }) 1L else 0L } override suspend fun loadFinancialAction(id: Long?): FinancialAction? { diff --git a/accountant/accountant-ports/accountant-persister-postgres/src/main/kotlin/co/nilin/opex/accountant/ports/postgres/impl/FinancialActionPersisterImpl.kt b/accountant/accountant-ports/accountant-persister-postgres/src/main/kotlin/co/nilin/opex/accountant/ports/postgres/impl/FinancialActionPersisterImpl.kt index 12292c2e6..e5412a6e7 100644 --- a/accountant/accountant-ports/accountant-persister-postgres/src/main/kotlin/co/nilin/opex/accountant/ports/postgres/impl/FinancialActionPersisterImpl.kt +++ b/accountant/accountant-ports/accountant-persister-postgres/src/main/kotlin/co/nilin/opex/accountant/ports/postgres/impl/FinancialActionPersisterImpl.kt @@ -103,7 +103,7 @@ class FinancialActionPersisterImpl( faRetryRepository.scheduleNext( id!!, retries + 1, - LocalDateTime.now().plusSeconds(retries * delayMultiplier * delaySeconds), + LocalDateTime.now().plusSeconds((retries + 1L) * delayMultiplier * delaySeconds), giveUp ).awaitSingleOrNull() diff --git a/accountant/accountant-ports/accountant-persister-postgres/src/main/resources/schema.sql b/accountant/accountant-ports/accountant-persister-postgres/src/main/resources/schema.sql index 8b7b1b772..994e76d3d 100644 --- a/accountant/accountant-ports/accountant-persister-postgres/src/main/resources/schema.sql +++ b/accountant/accountant-ports/accountant-persister-postgres/src/main/resources/schema.sql @@ -51,6 +51,15 @@ CREATE INDEX IF NOT EXISTS idx_fi_actions_status ON fi_actions (status); CREATE INDEX IF NOT EXISTS idx_fi_actions_pointer ON fi_actions (pointer); CREATE INDEX IF NOT EXISTS idx_fi_actions_status_create_date ON fi_actions (status, create_date); CREATE INDEX IF NOT EXISTS idx_fi_actions_parent_status ON fi_actions (parent_id, status); +CREATE INDEX IF NOT EXISTS idx_fi_actions_unprocessed_lookup + ON fi_actions (sender, symbol, event_type) + WHERE status <> 'PROCESSED'; +CREATE INDEX IF NOT EXISTS idx_fi_actions_archive_candidates + ON fi_actions (create_date, id) + WHERE status = 'PROCESSED'; +CREATE INDEX IF NOT EXISTS idx_fi_actions_unprocessed_children_by_parent + ON fi_actions (parent_id) + WHERE status <> 'PROCESSED'; ALTER TABLE fi_actions ADD COLUMN IF NOT EXISTS category_name VARCHAR(36); diff --git a/accountant/accountant-ports/accountant-persister-postgres/src/test/kotlin/co/nilin/opex/accountant/ports/postgres/FAPersisterImplTest.kt b/accountant/accountant-ports/accountant-persister-postgres/src/test/kotlin/co/nilin/opex/accountant/ports/postgres/FAPersisterImplTest.kt index 38c794ffd..ea63bfebf 100644 --- a/accountant/accountant-ports/accountant-persister-postgres/src/test/kotlin/co/nilin/opex/accountant/ports/postgres/FAPersisterImplTest.kt +++ b/accountant/accountant-ports/accountant-persister-postgres/src/test/kotlin/co/nilin/opex/accountant/ports/postgres/FAPersisterImplTest.kt @@ -9,10 +9,14 @@ import co.nilin.opex.accountant.ports.postgres.model.FinancialActionModel import io.mockk.coEvery import io.mockk.coVerify import io.mockk.mockk +import io.mockk.slot +import co.nilin.opex.accountant.ports.postgres.model.FinancialActionRetryModel import kotlinx.coroutines.runBlocking +import org.junit.jupiter.api.Assertions.assertTrue import org.junit.jupiter.api.Test import reactor.core.publisher.Flux import reactor.core.publisher.Mono +import java.time.LocalDateTime @Suppress("ReactiveStreamsUnusedPublisher") class FAPersisterImplTest { @@ -48,4 +52,29 @@ class FAPersisterImplTest { } } + @Test + fun givenRetryableAction_whenUpdateWithError_thenScheduleUsesBackoffDelay(): Unit = runBlocking { + val retryModel = FinancialActionRetryModel( + faId = Valid.fa.id!!, + nextRunTime = LocalDateTime.now(), + retries = 0, + isResolved = false, + hasGivenUp = false, + id = 10 + ) + val nextRunSlot = slot() + + coEvery { faRetryRepository.findByFaId(Valid.fa.id!!) } returns Mono.just(retryModel) + coEvery { faRetryRepository.scheduleNext(eq(10), eq(1), capture(nextRunSlot), eq(false)) } returns Mono.empty() + coEvery { financialActionRepository.updateStatus(eq(Valid.fa.id!!), eq(FinancialActionStatus.RETRYING)) } returns Mono.empty() + coEvery { faErrorRepository.save(any()) } returns Mono.empty() + + val before = LocalDateTime.now() + faPersister.updateWithError(Valid.fa, "ERR", "message", null) + + coVerify(exactly = 1) { faRetryRepository.scheduleNext(eq(10), eq(1), any(), eq(false)) } + assertTrue(nextRunSlot.isCaptured) + assertTrue(nextRunSlot.captured.isAfter(before.plusSeconds(10))) + } + } \ No newline at end of file diff --git a/wallet/wallet-app/pom.xml b/wallet/wallet-app/pom.xml index 9e1e23288..13dd14714 100644 --- a/wallet/wallet-app/pom.xml +++ b/wallet/wallet-app/pom.xml @@ -251,6 +251,11 @@ 5.4.0 test + + com.zaxxer + HikariCP + test + diff --git a/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/exc/ConcurrentBalanceChangException.kt b/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/exc/ConcurrentBalanceChangException.kt index edf52fad0..c8ab6b465 100644 --- a/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/exc/ConcurrentBalanceChangException.kt +++ b/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/exc/ConcurrentBalanceChangException.kt @@ -1,3 +1,3 @@ package co.nilin.opex.wallet.core.exc -class ConcurrentBalanceChangException(override val message: String?) : Exception() \ No newline at end of file +class ConcurrentBalanceChangException(override val message: String?) : RuntimeException() \ No newline at end of file diff --git a/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/model/PersistedTransaction.kt b/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/model/PersistedTransaction.kt new file mode 100644 index 000000000..439062dec --- /dev/null +++ b/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/model/PersistedTransaction.kt @@ -0,0 +1,6 @@ +package co.nilin.opex.wallet.core.model + +data class PersistedTransaction( + val id: Long, + val transaction: Transaction +) diff --git a/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/service/TransferManagerImpl.kt b/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/service/TransferManagerImpl.kt index b03fb97ed..8fc6cd9dc 100644 --- a/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/service/TransferManagerImpl.kt +++ b/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/service/TransferManagerImpl.kt @@ -7,6 +7,7 @@ import co.nilin.opex.wallet.core.inout.TransferResultDetailed import co.nilin.opex.wallet.core.model.* import co.nilin.opex.wallet.core.spi.* import org.slf4j.LoggerFactory +import org.springframework.dao.DuplicateKeyException import org.springframework.stereotype.Component import org.springframework.transaction.annotation.Transactional import java.time.LocalDateTime @@ -25,6 +26,8 @@ class TransferManagerImpl( @Transactional override suspend fun transfer(transferCommand: TransferCommand): TransferResultDetailed { + resolveIdempotentTransfer(transferCommand)?.let { return it } + //pre transfer hook (dispatch pre transfer event) val srcWallet = transferCommand.sourceWallet val srcWalletOwner = srcWallet.owner @@ -53,20 +56,26 @@ class TransferManagerImpl( if (!walletManager.isDepositAllowed(destWallet, amountToTransfer)) throw OpexError.DepositLimitExceeded.exception() + val tx = try { + transactionManager.save( + Transaction( + srcWallet, + destWallet, + transferCommand.amount.amount, + amountToTransfer, + transferCommand.description, + transferCommand.transferRef, + transferCommand.transferCategory, + LocalDateTime.now() + ) + ) + } catch (e: DuplicateKeyException) { + resolveIdempotentTransfer(transferCommand)?.let { return it } + throw e + } + walletManager.decreaseBalance(srcWallet, transferCommand.amount.amount) walletManager.increaseBalance(destWallet, amountToTransfer) - val tx = transactionManager.save( - Transaction( - srcWallet, - destWallet, - transferCommand.amount.amount, - amountToTransfer, - transferCommand.description, - transferCommand.transferRef, - transferCommand.transferCategory, - LocalDateTime.now() - ) - ) //TODO make tx long by default createUserTX(transferCommand, tx) @@ -93,6 +102,50 @@ class TransferManagerImpl( ) } + private suspend fun resolveIdempotentTransfer(transferCommand: TransferCommand): TransferResultDetailed? { + val transferRef = transferCommand.transferRef ?: return null + val persistedTransaction = transactionManager.findTransactionByTransferRef(transferRef) ?: return null + val existingTxId = persistedTransaction.id + val existingTransaction = persistedTransaction.transaction + + if (!matchesIdempotentTransfer(transferCommand, existingTransaction)) { + throw OpexError.BadRequest.exception("transferRef=$transferRef already exists with different parameters") + } + + logger.info("Idempotent transfer hit for transferRef={}", transferRef) + return buildIdempotentResult(existingTransaction, existingTxId) + } + + private fun matchesIdempotentTransfer(transferCommand: TransferCommand, existingTransaction: Transaction): Boolean { + return transferCommand.sourceWallet.id == existingTransaction.sourceWallet.id && + transferCommand.destWallet.id == existingTransaction.destWallet.id && + transferCommand.amount == Amount(existingTransaction.sourceWallet.currency, existingTransaction.sourceAmount) && + transferCommand.destAmount == Amount(existingTransaction.destWallet.currency, existingTransaction.destAmount) && + transferCommand.transferCategory == existingTransaction.transferCategory && + transferCommand.description == existingTransaction.description + } + + private fun buildIdempotentResult(existingTransaction: Transaction, existingTxId: Long): TransferResultDetailed { + val srcWallet = existingTransaction.sourceWallet + val destWallet = existingTransaction.destWallet + return TransferResultDetailed( + TransferResult( + Date().time, + srcWallet.owner.uuid, + srcWallet.type, + srcWallet.balance, + srcWallet.balance, + Amount(srcWallet.currency, existingTransaction.sourceAmount), + destWallet.owner.uuid, + destWallet.type, + Amount(destWallet.currency, existingTransaction.destAmount), + srcWallet.id, + destWallet.id, + ), + existingTxId.toString() + ) + } + private suspend fun createUserTX(command: TransferCommand, txId: Long) { val currency = command.amount.currency.symbol val amount = command.amount.amount diff --git a/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/spi/TransactionManager.kt b/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/spi/TransactionManager.kt index 9f646d823..d8fff39c0 100644 --- a/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/spi/TransactionManager.kt +++ b/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/spi/TransactionManager.kt @@ -6,6 +6,7 @@ import java.time.LocalDateTime interface TransactionManager { suspend fun save(transaction: Transaction): Long + suspend fun findTransactionByTransferRef(transferRef: String): PersistedTransaction? suspend fun findDepositTransactions( uuid: String, diff --git a/wallet/wallet-core/src/test/kotlin/co/nilin/opex/wallet/core/service/TransferManagerImplTest.kt b/wallet/wallet-core/src/test/kotlin/co/nilin/opex/wallet/core/service/TransferManagerImplTest.kt index 0fc406bfa..08320d369 100644 --- a/wallet/wallet-core/src/test/kotlin/co/nilin/opex/wallet/core/service/TransferManagerImplTest.kt +++ b/wallet/wallet-core/src/test/kotlin/co/nilin/opex/wallet/core/service/TransferManagerImplTest.kt @@ -1,15 +1,20 @@ package co.nilin.opex.wallet.core.service +import co.nilin.opex.common.OpexError import co.nilin.opex.wallet.core.model.Amount +import co.nilin.opex.wallet.core.model.Transaction import co.nilin.opex.wallet.core.service.sample.VALID import co.nilin.opex.wallet.core.spi.* import io.mockk.MockKException import io.mockk.coEvery +import io.mockk.coVerify import io.mockk.mockk import kotlinx.coroutines.runBlocking import org.assertj.core.api.Assertions.assertThat import org.assertj.core.api.Assertions.assertThatThrownBy +import org.junit.jupiter.api.Assertions import org.junit.jupiter.api.Test +import java.math.BigDecimal private class TransferManagerImplTest { private val walletOwnerManager: WalletOwnerManager = mockk() @@ -169,4 +174,68 @@ private class TransferManagerImplTest { } }.isNotInstanceOf(MockKException::class.java) } + + @Test + fun givenExistingTransferRef_whenTransfer_thenReturnIdempotentSuccessWithoutBalanceChanges(): Unit = runBlocking { + val command = VALID.TRANSFER_COMMAND.copy(transferRef = "accountant:fiActions:abc") + coEvery { transactionManager.findTransactionByTransferRef(eq(command.transferRef!!)) } returns co.nilin.opex.wallet.core.model.PersistedTransaction( + 100L, + Transaction( + VALID.SOURCE_WALLET, + VALID.DEST_WALLET, + command.amount.amount, + command.destAmount.amount, + command.description, + command.transferRef, + command.transferCategory, + java.time.LocalDateTime.now() + ) + ) + + val result = transferManager.transfer(command) + + assertThat(result.tx).isEqualTo("100") + assertThat(result.transferResult.sourceUuid).isEqualTo(command.sourceWallet.owner.uuid) + assertThat(result.transferResult.destUuid).isEqualTo(command.destWallet.owner.uuid) + + coVerify(exactly = 0) { walletManager.decreaseBalance(any(), any()) } + coVerify(exactly = 0) { walletManager.increaseBalance(any(), any()) } + coVerify(exactly = 0) { transactionManager.save(any()) } + coVerify(exactly = 0) { walletListener.onDeposit(any(), any(), any(), any(), any()) } + coVerify(exactly = 0) { walletListener.onWithdraw(any(), any(), any(), any()) } + } + + @Test + fun givenExistingTransferRefWithDifferentParams_whenTransfer_thenThrowBadRequest(): Unit = runBlocking { + val command = VALID.TRANSFER_COMMAND.copy( + transferRef = "accountant:fiActions:abc", + destWallet = VALID.DEST_WALLET.copy(id = 999L), + amount = Amount(VALID.CURRENCY, BigDecimal("0.75")), + destAmount = Amount(VALID.CURRENCY, BigDecimal("0.75")) + ) + coEvery { transactionManager.findTransactionByTransferRef(eq(command.transferRef!!)) } returns co.nilin.opex.wallet.core.model.PersistedTransaction( + 100L, + Transaction( + VALID.SOURCE_WALLET, + VALID.DEST_WALLET, + VALID.TRANSFER_COMMAND.amount.amount, + VALID.TRANSFER_COMMAND.destAmount.amount, + VALID.TRANSFER_COMMAND.description, + VALID.TRANSFER_COMMAND.transferRef, + VALID.TRANSFER_COMMAND.transferCategory, + java.time.LocalDateTime.now() + ) + ) + + val ex = Assertions.assertThrows(co.nilin.opex.utility.error.data.OpexException::class.java) { + runBlocking { + transferManager.transfer(command) + } + } + + assertThat(ex.error).isEqualTo(OpexError.BadRequest) + coVerify(exactly = 0) { walletManager.decreaseBalance(any(), any()) } + coVerify(exactly = 0) { walletManager.increaseBalance(any(), any()) } + coVerify(exactly = 0) { transactionManager.save(any()) } + } } diff --git a/wallet/wallet-ports/wallet-persister-postgres/src/main/kotlin/co/nilin/opex/wallet/ports/postgres/dao/TransactionRepository.kt b/wallet/wallet-ports/wallet-persister-postgres/src/main/kotlin/co/nilin/opex/wallet/ports/postgres/dao/TransactionRepository.kt index f47ca1312..0eba0f408 100644 --- a/wallet/wallet-ports/wallet-persister-postgres/src/main/kotlin/co/nilin/opex/wallet/ports/postgres/dao/TransactionRepository.kt +++ b/wallet/wallet-ports/wallet-persister-postgres/src/main/kotlin/co/nilin/opex/wallet/ports/postgres/dao/TransactionRepository.kt @@ -16,6 +16,9 @@ import java.time.LocalDateTime @Repository interface TransactionRepository : ReactiveCrudRepository { + @Query("select * from transaction where transfer_ref = :transferRef limit 1") + fun findByTransferRef(transferRef: String): Mono + @Query( """ SELECT count(1) cnt, COALESCE(sum(source_amount), 0) total diff --git a/wallet/wallet-ports/wallet-persister-postgres/src/main/kotlin/co/nilin/opex/wallet/ports/postgres/impl/TransactionManagerImpl.kt b/wallet/wallet-ports/wallet-persister-postgres/src/main/kotlin/co/nilin/opex/wallet/ports/postgres/impl/TransactionManagerImpl.kt index 7e7ec45e0..d60f33884 100644 --- a/wallet/wallet-ports/wallet-persister-postgres/src/main/kotlin/co/nilin/opex/wallet/ports/postgres/impl/TransactionManagerImpl.kt +++ b/wallet/wallet-ports/wallet-persister-postgres/src/main/kotlin/co/nilin/opex/wallet/ports/postgres/impl/TransactionManagerImpl.kt @@ -2,12 +2,14 @@ package co.nilin.opex.wallet.ports.postgres.impl import co.nilin.opex.wallet.core.model.* import co.nilin.opex.wallet.core.spi.TransactionManager +import co.nilin.opex.wallet.core.spi.WalletManager import co.nilin.opex.wallet.ports.postgres.dao.CurrencyRepositoryV2 import co.nilin.opex.wallet.ports.postgres.dao.TransactionRepository import co.nilin.opex.wallet.ports.postgres.model.TransactionModel import com.fasterxml.jackson.databind.ObjectMapper import kotlinx.coroutines.reactive.awaitFirstOrElse import kotlinx.coroutines.reactive.awaitSingle +import kotlinx.coroutines.reactor.awaitSingleOrNull import org.slf4j.LoggerFactory import org.springframework.stereotype.Service import java.time.LocalDateTime @@ -17,6 +19,7 @@ import java.time.ZoneId class TransactionManagerImpl( private val transactionRepository: TransactionRepository, private val currencyRepositoryV2: CurrencyRepositoryV2, + private val walletManager: WalletManager, private val objectMapper: ObjectMapper ) : TransactionManager { private val logger = LoggerFactory.getLogger(TransactionManagerImpl::class.java) @@ -36,6 +39,26 @@ class TransactionManagerImpl( ).awaitSingle().id!! } + override suspend fun findTransactionByTransferRef(transferRef: String): PersistedTransaction? { + val transaction = transactionRepository.findByTransferRef(transferRef).awaitSingleOrNull() ?: return null + val sourceWallet = walletManager.findWalletById(transaction.sourceWallet) ?: return null + val destWallet = walletManager.findWalletById(transaction.destWallet) ?: return null + + return PersistedTransaction( + transaction.id!!, + Transaction( + sourceWallet, + destWallet, + transaction.sourceAmount, + transaction.destAmount, + transaction.description, + transaction.transferRef, + transaction.transferCategory, + transaction.transactionDate + ) + ) + } + override suspend fun findDepositTransactions( uuid: String, @@ -148,6 +171,3 @@ class TransactionManagerImpl( .collectList().awaitFirstOrElse { emptyList() } } } - - -