diff --git a/services/account-service/src/main/java/com/ntropy/account/client/LocalFinancialAccountCommandClient.java b/services/account-service/src/main/java/com/ntropy/account/client/LocalFinancialAccountCommandClient.java index a3296e28..7b56c335 100644 --- a/services/account-service/src/main/java/com/ntropy/account/client/LocalFinancialAccountCommandClient.java +++ b/services/account-service/src/main/java/com/ntropy/account/client/LocalFinancialAccountCommandClient.java @@ -10,6 +10,7 @@ import com.ntropy.account.domain.InstitutionKeys; import com.ntropy.account.domain.PersonalBank; import com.ntropy.account.domain.entity.CodefConnection; +import com.ntropy.account.event.AccountTransactionClassificationEventPublisher; import com.ntropy.account.exception.AccountErrorCode; import com.ntropy.account.mapper.AccountLifecycleMapper; import com.ntropy.account.mapper.AccountMapper; @@ -20,7 +21,6 @@ import com.ntropy.account.service.VirtualFinancialDataService; import com.ntropy.account.service.VirtualFinancialDataService.GenerationSummary; import com.ntropy.common.client.FinancialAccountCommandClient; -import com.ntropy.common.client.TransactionClassificationCommandClient; import com.ntropy.common.dto.account.AccountRegistrationCommand; import com.ntropy.common.dto.account.AccountRegistrationSummary; import com.ntropy.common.dto.account.BankSummary; @@ -45,7 +45,7 @@ public class LocalFinancialAccountCommandClient implements FinancialAccountComma private final AccountLifecycleMapper accountLifecycleMapper; private final AccountMapper accountMapper; private final CodefConnectionMapper codefConnectionMapper; - private final TransactionClassificationCommandClient transactionClassificationCommandClient; + private final AccountTransactionClassificationEventPublisher classificationEventPublisher; @Override public List findSupportedBanks() { @@ -82,7 +82,7 @@ public AccountRegistrationSummary registerAccount(Long userId, AccountRegistrati if ("VIRTUAL".equals(connectionType)) { GenerationSummary summary = virtualAccountRegenerationService.regenerateForUser(userId, bank); - classifyTransactionsSafely(userId); + requestTransactionClassificationSafely(userId); return new AccountRegistrationSummary(connectionType, bank.getOrganizationCode(), summary.accounts()); } @@ -97,16 +97,18 @@ public AccountRegistrationSummary registerAccount(Long userId, AccountRegistrati ); ensureVirtualDatasetSafely(userId, bank); int accountCount = accountCollectionService.collect(userId, bank, birthDate, startDate, endDate).size(); - classifyTransactionsSafely(userId); + requestTransactionClassificationSafely(userId); return new AccountRegistrationSummary(connectionType, bank.getOrganizationCode(), accountCount); } - private void classifyTransactionsSafely(Long userId) { + private void requestTransactionClassificationSafely(Long userId) { try { - int processed = transactionClassificationCommandClient.classifyUnanalyzedTransactions(userId); - log.info("계좌 연동 후 소비 분류 완료: userId={}, processed={}", userId, processed); + classificationEventPublisher.publishAfterCommit(userId); } catch (RuntimeException e) { - log.warn("계좌 연동 후 소비 분류 실패: userId={}", userId, e); + log.warn( + "계좌 연동 후 소비 분류 작업 등록 실패: userId={}, errorKind={}", + userId, e.getClass().getSimpleName() + ); } } diff --git a/services/account-service/src/main/java/com/ntropy/account/config/AccountClassificationAsyncConfig.java b/services/account-service/src/main/java/com/ntropy/account/config/AccountClassificationAsyncConfig.java new file mode 100644 index 00000000..02a7fccb --- /dev/null +++ b/services/account-service/src/main/java/com/ntropy/account/config/AccountClassificationAsyncConfig.java @@ -0,0 +1,44 @@ +package com.ntropy.account.config; + +import java.util.concurrent.ThreadPoolExecutor; + +import org.springframework.beans.factory.annotation.Value; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.PropertySource; +import org.springframework.context.annotation.PropertySources; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; + +/** 계좌 연동 요청과 소비 분류 작업이 서로의 스레드 풀을 점유하지 않도록 분리합니다. */ +@Configuration +@PropertySources({ + @PropertySource( + value = "classpath:account-classification.properties", + ignoreResourceNotFound = true + ), + @PropertySource( + value = "file:${NTROPY_CONFIG_DIR:./config}/account-classification.properties", + ignoreResourceNotFound = true + ) +}) +public class AccountClassificationAsyncConfig { + + @Bean("accountClassificationJobExecutor") + public ThreadPoolTaskExecutor accountClassificationJobExecutor( + @Value("${account.classification.async.core-pool-size:1}") int configuredCorePoolSize, + @Value("${account.classification.async.max-pool-size:2}") int configuredMaxPoolSize, + @Value("${account.classification.async.queue-capacity:100}") int configuredQueueCapacity + ) { + int corePoolSize = Math.max(1, Math.min(configuredCorePoolSize, 8)); + int maxPoolSize = Math.max(corePoolSize, Math.min(configuredMaxPoolSize, 8)); + int queueCapacity = Math.max(1, Math.min(configuredQueueCapacity, 1000)); + + ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); + executor.setCorePoolSize(corePoolSize); + executor.setMaxPoolSize(maxPoolSize); + executor.setQueueCapacity(queueCapacity); + executor.setThreadNamePrefix("account-classification-job-"); + executor.setRejectedExecutionHandler(new ThreadPoolExecutor.AbortPolicy()); + return executor; + } +} diff --git a/services/account-service/src/main/java/com/ntropy/account/event/AccountTransactionClassificationEventPublisher.java b/services/account-service/src/main/java/com/ntropy/account/event/AccountTransactionClassificationEventPublisher.java new file mode 100644 index 00000000..debcc444 --- /dev/null +++ b/services/account-service/src/main/java/com/ntropy/account/event/AccountTransactionClassificationEventPublisher.java @@ -0,0 +1,37 @@ +package com.ntropy.account.event; + +import org.springframework.context.ApplicationEventPublisher; +import org.springframework.stereotype.Component; +import org.springframework.transaction.support.TransactionSynchronization; +import org.springframework.transaction.support.TransactionSynchronizationManager; + +import lombok.RequiredArgsConstructor; + +/** 거래 저장 트랜잭션이 있으면 커밋된 뒤에만 소비 분류 이벤트를 발행합니다. */ +@Component +@RequiredArgsConstructor +public class AccountTransactionClassificationEventPublisher { + + private final ApplicationEventPublisher applicationEventPublisher; + + public void publishAfterCommit(Long userId) { + AccountTransactionsCollectedEvent event = + new AccountTransactionsCollectedEvent(userId); + + if (TransactionSynchronizationManager.isActualTransactionActive() + && TransactionSynchronizationManager.isSynchronizationActive()) { + TransactionSynchronizationManager.registerSynchronization( + new TransactionSynchronization() { + @Override + public void afterCommit() { + applicationEventPublisher.publishEvent(event); + } + } + ); + return; + } + + /* 현재 계좌 조합 경로처럼 하위 저장 메서드가 이미 커밋된 경우입니다. */ + applicationEventPublisher.publishEvent(event); + } +} diff --git a/services/account-service/src/main/java/com/ntropy/account/event/AccountTransactionClassificationListener.java b/services/account-service/src/main/java/com/ntropy/account/event/AccountTransactionClassificationListener.java new file mode 100644 index 00000000..fe64a63d --- /dev/null +++ b/services/account-service/src/main/java/com/ntropy/account/event/AccountTransactionClassificationListener.java @@ -0,0 +1,74 @@ +package com.ntropy.account.event; + +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; + +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.event.EventListener; +import org.springframework.core.task.TaskExecutor; +import org.springframework.stereotype.Component; + +import com.ntropy.common.client.TransactionClassificationCommandClient; + +import lombok.extern.slf4j.Slf4j; + +/** 계좌 연동 응답과 분리된 전용 executor에서 미분류 거래를 처리합니다. */ +@Component +@Slf4j +public class AccountTransactionClassificationListener { + + private final TransactionClassificationCommandClient classificationCommandClient; + private final TaskExecutor classificationJobExecutor; + private final Set runningUserIds = ConcurrentHashMap.newKeySet(); + + public AccountTransactionClassificationListener( + TransactionClassificationCommandClient classificationCommandClient, + @Qualifier("accountClassificationJobExecutor") TaskExecutor classificationJobExecutor + ) { + this.classificationCommandClient = classificationCommandClient; + this.classificationJobExecutor = classificationJobExecutor; + } + + @EventListener + public void onTransactionsCollected(AccountTransactionsCollectedEvent event) { + Long userId = event.userId(); + if (!runningUserIds.add(userId)) { + log.info("[비동기 소비 분류] 중복 작업 건너뜀. scope=userId={}", userId); + return; + } + + try { + classificationJobExecutor.execute(() -> classify(userId)); + } catch (RuntimeException e) { + runningUserIds.remove(userId); + log.warn( + "[비동기 소비 분류] 작업 등록 실패. scope=userId={}, errorKind={}", + userId, e.getClass().getSimpleName() + ); + } + } + + private void classify(Long userId) { + long startedAt = System.nanoTime(); + log.info("[비동기 소비 분류] 작업 시작. scope=userId={}", userId); + + try { + int processed = classificationCommandClient.classifyUnanalyzedTransactions(userId); + log.info( + "[비동기 소비 분류] 작업 완료. scope=userId={}, totalProcessed={}, elapsedMs={}", + userId, processed, elapsedMillis(startedAt) + ); + } catch (RuntimeException e) { + log.warn( + "[비동기 소비 분류] 작업 실패. scope=userId={}, elapsedMs={}, errorKind={}", + userId, elapsedMillis(startedAt), e.getClass().getSimpleName() + ); + } finally { + runningUserIds.remove(userId); + } + } + + private static long elapsedMillis(long startedAt) { + return (System.nanoTime() - startedAt) / 1_000_000L; + } +} diff --git a/services/account-service/src/main/java/com/ntropy/account/event/AccountTransactionsCollectedEvent.java b/services/account-service/src/main/java/com/ntropy/account/event/AccountTransactionsCollectedEvent.java new file mode 100644 index 00000000..f3941bc6 --- /dev/null +++ b/services/account-service/src/main/java/com/ntropy/account/event/AccountTransactionsCollectedEvent.java @@ -0,0 +1,5 @@ +package com.ntropy.account.event; + +/** 계좌 연동으로 거래 저장을 마친 뒤 소비 분류를 요청하는 내부 이벤트입니다. */ +public record AccountTransactionsCollectedEvent(Long userId) { +} diff --git a/services/account-service/src/main/resources/account-classification.properties.example b/services/account-service/src/main/resources/account-classification.properties.example new file mode 100644 index 00000000..ab7c96b9 --- /dev/null +++ b/services/account-service/src/main/resources/account-classification.properties.example @@ -0,0 +1,4 @@ +# 계좌 연동 응답과 분리해 실행하는 사용자 단위 소비 분류 작업 executor 설정 +account.classification.async.core-pool-size=1 +account.classification.async.max-pool-size=2 +account.classification.async.queue-capacity=100 diff --git a/services/account-service/src/test/java/com/ntropy/account/client/LocalFinancialAccountCommandClientTest.java b/services/account-service/src/test/java/com/ntropy/account/client/LocalFinancialAccountCommandClientTest.java index 7991a656..0efa0442 100644 --- a/services/account-service/src/test/java/com/ntropy/account/client/LocalFinancialAccountCommandClientTest.java +++ b/services/account-service/src/test/java/com/ntropy/account/client/LocalFinancialAccountCommandClientTest.java @@ -10,12 +10,15 @@ import java.util.List; import org.junit.jupiter.api.Test; +import org.springframework.context.ApplicationEventPublisher; import com.ntropy.account.domain.ConnectionProvider; import com.ntropy.account.domain.PersonalBank; import com.ntropy.account.domain.entity.Account; import com.ntropy.account.domain.entity.CodefConnection; import com.ntropy.account.exception.AccountErrorCode; +import com.ntropy.account.event.AccountTransactionClassificationEventPublisher; +import com.ntropy.account.event.AccountTransactionsCollectedEvent; import com.ntropy.account.mapper.AccountLifecycleMapper; import com.ntropy.account.mapper.AccountMapper; import com.ntropy.account.mapper.CodefConnectionMapper; @@ -25,7 +28,6 @@ import com.ntropy.account.service.VirtualFinancialDataService; import com.ntropy.account.service.VirtualFinancialDataService.GenerationSummary; import com.ntropy.common.dto.account.AccountRegistrationCommand; -import com.ntropy.common.client.TransactionClassificationCommandClient; import com.ntropy.common.exception.ServiceException; class LocalFinancialAccountCommandClientTest { @@ -137,33 +139,36 @@ collectionService, new StubVirtualAccountRegenerationService(), } @Test - void classifiesUsersTransactionsAfterCodefCollectionCompletes() { - StubTransactionClassificationCommandClient classificationClient = - new StubTransactionClassificationCommandClient(); + void requestsUsersTransactionClassificationAfterCodefCollectionCompletes() { + StubApplicationEventPublisher eventPublisher = new StubApplicationEventPublisher(); LocalFinancialAccountCommandClient client = newClient( new StubPersonalBankAccountService(), new StubAccountCollectionService(), new StubVirtualAccountRegenerationService(), new StubVirtualFinancialDataService(), new StubAccountLifecycleMapper(1, 1), new StubAccountMapper(false), - new StubCodefConnectionMapper(), classificationClient + new StubCodefConnectionMapper(), + new AccountTransactionClassificationEventPublisher(eventPublisher) ); client.registerAccount( 42L, new AccountRegistrationCommand("CODEF", "0088", "bank-id", "bank-password", null) ); - assertEquals(List.of(42L), classificationClient.userIds); + assertEquals( + List.of(new AccountTransactionsCollectedEvent(42L)), + eventPublisher.events + ); } @Test - void keepsAccountRegistrationSuccessfulWhenClassificationFails() { - StubTransactionClassificationCommandClient classificationClient = - new StubTransactionClassificationCommandClient(); - classificationClient.failure = new IllegalStateException("분류 실패"); + void keepsAccountRegistrationSuccessfulWhenClassificationRequestFails() { + StubApplicationEventPublisher eventPublisher = new StubApplicationEventPublisher(); + eventPublisher.failure = new IllegalStateException("작업 등록 실패"); LocalFinancialAccountCommandClient client = newClient( new StubPersonalBankAccountService(), new StubAccountCollectionService(), new StubVirtualAccountRegenerationService(), new StubVirtualFinancialDataService(), new StubAccountLifecycleMapper(1, 1), new StubAccountMapper(false), - new StubCodefConnectionMapper(), classificationClient + new StubCodefConnectionMapper(), + new AccountTransactionClassificationEventPublisher(eventPublisher) ); var result = client.registerAccount( @@ -171,7 +176,6 @@ void keepsAccountRegistrationSuccessfulWhenClassificationFails() { ); assertEquals("VIRTUAL", result.connectionType()); - assertEquals(List.of(42L), classificationClient.userIds); } @Test @@ -461,7 +465,9 @@ private static LocalFinancialAccountCommandClient newClient( return newClient( personalBankAccountService, collectionService, regenerationService, virtualFinancialDataService, lifecycleMapper, accountMapper, connectionMapper, - new StubTransactionClassificationCommandClient() + new AccountTransactionClassificationEventPublisher( + new StubApplicationEventPublisher() + ) ); } @@ -473,26 +479,24 @@ private static LocalFinancialAccountCommandClient newClient( AccountLifecycleMapper lifecycleMapper, AccountMapper accountMapper, CodefConnectionMapper connectionMapper, - TransactionClassificationCommandClient classificationClient + AccountTransactionClassificationEventPublisher classificationEventPublisher ) { return new LocalFinancialAccountCommandClient( personalBankAccountService, collectionService, regenerationService, virtualFinancialDataService, - lifecycleMapper, accountMapper, connectionMapper, classificationClient + lifecycleMapper, accountMapper, connectionMapper, classificationEventPublisher ); } - private static class StubTransactionClassificationCommandClient - implements TransactionClassificationCommandClient { - private final List userIds = new ArrayList<>(); + private static class StubApplicationEventPublisher implements ApplicationEventPublisher { + private final List events = new ArrayList<>(); private RuntimeException failure; @Override - public int classifyUnanalyzedTransactions(Long userId) { - userIds.add(userId); + public void publishEvent(Object event) { if (failure != null) { throw failure; } - return 0; + events.add(event); } } diff --git a/services/account-service/src/test/java/com/ntropy/account/config/AccountClassificationAsyncConfigTest.java b/services/account-service/src/test/java/com/ntropy/account/config/AccountClassificationAsyncConfigTest.java new file mode 100644 index 00000000..0e7f557a --- /dev/null +++ b/services/account-service/src/test/java/com/ntropy/account/config/AccountClassificationAsyncConfigTest.java @@ -0,0 +1,45 @@ +package com.ntropy.account.config; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; + +import java.util.concurrent.ThreadPoolExecutor; + +import org.junit.jupiter.api.Test; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; + +class AccountClassificationAsyncConfigTest { + + private final AccountClassificationAsyncConfig config = + new AccountClassificationAsyncConfig(); + + @Test + void usesConfiguredPoolAndQueueSizes() { + ThreadPoolTaskExecutor executor = + config.accountClassificationJobExecutor(2, 4, 25); + executor.initialize(); + + assertEquals(2, executor.getCorePoolSize()); + assertEquals(4, executor.getMaxPoolSize()); + assertEquals(25, executor.getQueueCapacity()); + assertInstanceOf( + ThreadPoolExecutor.AbortPolicy.class, + executor.getThreadPoolExecutor().getRejectedExecutionHandler() + ); + + executor.shutdown(); + } + + @Test + void clampsUnsafeConfigurationValues() { + ThreadPoolTaskExecutor executor = + config.accountClassificationJobExecutor(0, 99, 0); + executor.initialize(); + + assertEquals(1, executor.getCorePoolSize()); + assertEquals(8, executor.getMaxPoolSize()); + assertEquals(1, executor.getQueueCapacity()); + + executor.shutdown(); + } +} diff --git a/services/account-service/src/test/java/com/ntropy/account/event/AccountTransactionClassificationEventPublisherTest.java b/services/account-service/src/test/java/com/ntropy/account/event/AccountTransactionClassificationEventPublisherTest.java new file mode 100644 index 00000000..d92be439 --- /dev/null +++ b/services/account-service/src/test/java/com/ntropy/account/event/AccountTransactionClassificationEventPublisherTest.java @@ -0,0 +1,88 @@ +package com.ntropy.account.event; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import java.sql.Connection; +import java.util.ArrayList; +import java.util.List; + +import javax.sql.DataSource; + +import org.junit.jupiter.api.Test; +import org.springframework.context.ApplicationEventPublisher; +import org.springframework.jdbc.datasource.DataSourceTransactionManager; +import org.springframework.transaction.support.TransactionTemplate; + +class AccountTransactionClassificationEventPublisherTest { + + @Test + void publishesImmediatelyWhenNoTransactionIsActive() { + RecordingApplicationEventPublisher events = + new RecordingApplicationEventPublisher(); + AccountTransactionClassificationEventPublisher publisher = + new AccountTransactionClassificationEventPublisher(events); + + publisher.publishAfterCommit(42L); + + assertEquals( + List.of(new AccountTransactionsCollectedEvent(42L)), + events.publishedEvents + ); + } + + @Test + void publishesOnlyAfterTransactionCommits() throws Exception { + RecordingApplicationEventPublisher events = + new RecordingApplicationEventPublisher(); + AccountTransactionClassificationEventPublisher publisher = + new AccountTransactionClassificationEventPublisher(events); + TransactionTemplate transaction = transactionTemplate(); + + transaction.executeWithoutResult(status -> { + publisher.publishAfterCommit(42L); + assertTrue(events.publishedEvents.isEmpty()); + }); + + assertEquals( + List.of(new AccountTransactionsCollectedEvent(42L)), + events.publishedEvents + ); + } + + @Test + void doesNotPublishWhenTransactionRollsBack() throws Exception { + RecordingApplicationEventPublisher events = + new RecordingApplicationEventPublisher(); + AccountTransactionClassificationEventPublisher publisher = + new AccountTransactionClassificationEventPublisher(events); + TransactionTemplate transaction = transactionTemplate(); + + transaction.executeWithoutResult(status -> { + publisher.publishAfterCommit(42L); + status.setRollbackOnly(); + }); + + assertTrue(events.publishedEvents.isEmpty()); + } + + private static TransactionTemplate transactionTemplate() throws Exception { + DataSource dataSource = mock(DataSource.class); + Connection connection = mock(Connection.class); + when(dataSource.getConnection()).thenReturn(connection); + when(connection.getAutoCommit()).thenReturn(true); + return new TransactionTemplate(new DataSourceTransactionManager(dataSource)); + } + + private static class RecordingApplicationEventPublisher + implements ApplicationEventPublisher { + private final List publishedEvents = new ArrayList<>(); + + @Override + public void publishEvent(Object event) { + publishedEvents.add(event); + } + } +} diff --git a/services/account-service/src/test/java/com/ntropy/account/event/AccountTransactionClassificationListenerTest.java b/services/account-service/src/test/java/com/ntropy/account/event/AccountTransactionClassificationListenerTest.java new file mode 100644 index 00000000..0d24a904 --- /dev/null +++ b/services/account-service/src/test/java/com/ntropy/account/event/AccountTransactionClassificationListenerTest.java @@ -0,0 +1,119 @@ +package com.ntropy.account.event; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; +import java.util.ArrayList; +import java.util.List; + +import org.junit.jupiter.api.Test; +import org.springframework.core.task.TaskExecutor; +import org.springframework.core.task.TaskRejectedException; + +import com.ntropy.common.client.TransactionClassificationCommandClient; + +class AccountTransactionClassificationListenerTest { + + @Test + void schedulesClassificationWithoutRunningItOnPublisherThread() { + RecordingExecutor executor = new RecordingExecutor(); + RecordingClassificationClient classificationClient = new RecordingClassificationClient(); + AccountTransactionClassificationListener listener = + new AccountTransactionClassificationListener(classificationClient, executor); + + listener.onTransactionsCollected(new AccountTransactionsCollectedEvent(42L)); + + assertTrue(classificationClient.userIds.isEmpty()); + assertEquals(1, executor.tasks.size()); + + executor.runNext(); + + assertEquals(List.of(42L), classificationClient.userIds); + } + + @Test + void ignoresDuplicateJobForSameUserUntilRunningJobFinishes() { + RecordingExecutor executor = new RecordingExecutor(); + RecordingClassificationClient classificationClient = new RecordingClassificationClient(); + AccountTransactionClassificationListener listener = + new AccountTransactionClassificationListener(classificationClient, executor); + AccountTransactionsCollectedEvent event = new AccountTransactionsCollectedEvent(42L); + + listener.onTransactionsCollected(event); + listener.onTransactionsCollected(event); + + assertEquals(1, executor.tasks.size()); + + executor.runNext(); + listener.onTransactionsCollected(event); + + assertEquals(1, executor.tasks.size()); + executor.runNext(); + assertEquals(List.of(42L, 42L), classificationClient.userIds); + } + + @Test + void keepsListenerSuccessfulWhenBackgroundClassificationFails() { + RecordingClassificationClient classificationClient = new RecordingClassificationClient(); + classificationClient.failure = new IllegalStateException("분류 실패"); + AccountTransactionClassificationListener listener = + new AccountTransactionClassificationListener(classificationClient, Runnable::run); + + assertDoesNotThrow( + () -> listener.onTransactionsCollected( + new AccountTransactionsCollectedEvent(42L) + ) + ); + assertEquals(List.of(42L), classificationClient.userIds); + } + + @Test + void releasesDuplicateGuardWhenExecutorRejectsJob() { + RecordingClassificationClient classificationClient = new RecordingClassificationClient(); + TaskExecutor rejectingExecutor = task -> { + throw new TaskRejectedException("포화"); + }; + AccountTransactionClassificationListener listener = + new AccountTransactionClassificationListener(classificationClient, rejectingExecutor); + + assertDoesNotThrow( + () -> listener.onTransactionsCollected( + new AccountTransactionsCollectedEvent(42L) + ) + ); + assertDoesNotThrow( + () -> listener.onTransactionsCollected( + new AccountTransactionsCollectedEvent(42L) + ) + ); + } + + private static class RecordingExecutor implements TaskExecutor { + private final List tasks = new ArrayList<>(); + + @Override + public void execute(Runnable task) { + tasks.add(task); + } + + void runNext() { + tasks.remove(0).run(); + } + } + + private static class RecordingClassificationClient + implements TransactionClassificationCommandClient { + private final List userIds = new ArrayList<>(); + private RuntimeException failure; + + @Override + public int classifyUnanalyzedTransactions(Long userId) { + userIds.add(userId); + if (failure != null) { + throw failure; + } + return 1; + } + } + +}