diff --git a/src/main/kotlin/kr/co/vividnext/sodalive/v2/ranking/application/CreatorRankingSnapshotJobService.kt b/src/main/kotlin/kr/co/vividnext/sodalive/v2/ranking/application/CreatorRankingSnapshotJobService.kt index 43ab8648..53f9f1cf 100644 --- a/src/main/kotlin/kr/co/vividnext/sodalive/v2/ranking/application/CreatorRankingSnapshotJobService.kt +++ b/src/main/kotlin/kr/co/vividnext/sodalive/v2/ranking/application/CreatorRankingSnapshotJobService.kt @@ -35,15 +35,26 @@ class CreatorRankingSnapshotJobService( fun refreshLastCompletedWeekByScheduledJob() { withLastCompletedWeekPeriodLock { now, utcRange -> - refreshLastCompletedWeekByScheduledJob(now, utcRange) + refreshLastCompletedWeek(now, utcRange, CreatorRankingSnapshotJobTrigger.SCHEDULED) } } - private fun refreshLastCompletedWeekByScheduledJob( + fun refreshLastCompletedWeekByFallback(): Boolean { + var refreshed = false + withLastCompletedWeekPeriodLock { now, utcRange -> + if (fallbackCountReachedLimit(utcRange)) return@withLastCompletedWeekPeriodLock + refreshLastCompletedWeek(now, utcRange, CreatorRankingSnapshotJobTrigger.FALLBACK) + refreshed = true + } + return refreshed + } + + private fun refreshLastCompletedWeek( now: ZonedDateTime, - utcRange: CreatorRankingUtcRange + utcRange: CreatorRankingUtcRange, + trigger: CreatorRankingSnapshotJobTrigger ) { - val job = savePendingJob(utcRange, CreatorRankingSnapshotJobTrigger.SCHEDULED) + val job = savePendingJob(utcRange, trigger) val jobId = job.id ?: return markProcessing(jobId) logJobStatusChanged(job, CreatorRankingSnapshotJobStatus.PROCESSING) @@ -85,22 +96,25 @@ class CreatorRankingSnapshotJobService( }!! } + private fun fallbackCountReachedLimit(utcRange: CreatorRankingUtcRange): Boolean { + return jobPort.countByRankingTypeAndPeriodAndTrigger( + rankingType = CreatorRankingType.WEEKLY, + aggregationStartAtUtc = utcRange.startInclusiveUtc, + aggregationEndAtUtc = utcRange.endExclusiveUtc, + trigger = CreatorRankingSnapshotJobTrigger.FALLBACK + ) >= FALLBACK_LIMIT + } + private fun markProcessing(jobId: Long) { - transactionTemplate.executeWithoutResult { - jobPort.markProcessing(jobId, LocalDateTime.now()) - } + transactionTemplate.executeWithoutResult { jobPort.markProcessing(jobId, LocalDateTime.now()) } } private fun markDone(jobId: Long) { - transactionTemplate.executeWithoutResult { - jobPort.markDone(jobId, LocalDateTime.now()) - } + transactionTemplate.executeWithoutResult { jobPort.markDone(jobId, LocalDateTime.now()) } } private fun markFailed(jobId: Long, message: String?) { - transactionTemplate.executeWithoutResult { - jobPort.markFailed(jobId, LocalDateTime.now(), message) - } + transactionTemplate.executeWithoutResult { jobPort.markFailed(jobId, LocalDateTime.now(), message) } } @Transactional @@ -128,34 +142,22 @@ class CreatorRankingSnapshotJobService( aggregationEndAtUtc: LocalDateTime, statuses: List = CreatorRankingSnapshotJobStatus.values().toList() ): List { - return jobPort.findByPeriodAndStatuses( - aggregationStartAtUtc = aggregationStartAtUtc, - aggregationEndAtUtc = aggregationEndAtUtc, - statuses = statuses - ) + return jobPort.findByPeriodAndStatuses(aggregationStartAtUtc, aggregationEndAtUtc, statuses) } @Transactional fun retryFailedJob(jobId: Long) { val job = jobPort.findById(jobId) ?: return if (job.status != CreatorRankingSnapshotJobStatus.FAILED) return - jobPort.markPending(jobId) } - fun ensureLastCompletedWeekSnapshotForColdStart() { - withLastCompletedWeekPeriodLock { now, _ -> - transactionTemplate.executeWithoutResult { - refreshService.refreshLastCompletedWeek(now) - } - } - } - private fun withLastCompletedWeekPeriodLock(action: (ZonedDateTime, CreatorRankingUtcRange) -> Unit) { val now = nowProvider() val period = periodPolicy.resolveLastCompletedWeek(now) val utcRange = periodPolicy.toUtcRange(period) - val lockName = "lock:creator-ranking-snapshot-refresh:${utcRange.startInclusiveUtc}:${utcRange.endExclusiveUtc}" + val lockName = "lock:creator-ranking-snapshot-refresh:" + + "${CreatorRankingType.WEEKLY}:${utcRange.startInclusiveUtc}:${utcRange.endExclusiveUtc}" val lock = redissonClient.getLock(lockName) try { @@ -185,4 +187,8 @@ class CreatorRankingSnapshotJobService( error ) } + + companion object { + private const val FALLBACK_LIMIT = 3L + } } diff --git a/src/test/kotlin/kr/co/vividnext/sodalive/v2/ranking/application/CreatorRankingSnapshotJobServiceTest.kt b/src/test/kotlin/kr/co/vividnext/sodalive/v2/ranking/application/CreatorRankingSnapshotJobServiceTest.kt index 81a392d6..b3c58c88 100644 --- a/src/test/kotlin/kr/co/vividnext/sodalive/v2/ranking/application/CreatorRankingSnapshotJobServiceTest.kt +++ b/src/test/kotlin/kr/co/vividnext/sodalive/v2/ranking/application/CreatorRankingSnapshotJobServiceTest.kt @@ -6,458 +6,184 @@ import kr.co.vividnext.sodalive.v2.ranking.port.out.CreatorRankingSnapshotJobRec import kr.co.vividnext.sodalive.v2.ranking.port.out.CreatorRankingSnapshotJobStatus import kr.co.vividnext.sodalive.v2.ranking.port.out.CreatorRankingSnapshotJobTrigger import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertFalse import org.junit.jupiter.api.Assertions.assertThrows import org.junit.jupiter.api.Assertions.assertTrue import org.junit.jupiter.api.DisplayName import org.junit.jupiter.api.Test -import org.junit.jupiter.api.extension.ExtendWith import org.mockito.Mockito import org.redisson.api.RLock import org.redisson.api.RedissonClient -import org.springframework.boot.test.system.CapturedOutput -import org.springframework.boot.test.system.OutputCaptureExtension import org.springframework.transaction.PlatformTransactionManager import org.springframework.transaction.TransactionDefinition +import org.springframework.transaction.TransactionStatus import org.springframework.transaction.support.SimpleTransactionStatus import java.time.LocalDateTime import java.time.ZoneId import java.time.ZonedDateTime import java.util.concurrent.TimeUnit -@ExtendWith(OutputCaptureExtension::class) class CreatorRankingSnapshotJobServiceTest { @Test - @DisplayName("스케줄 실행은 집계 기간을 포함한 SCHEDULED job을 생성하고 성공 시 DONE으로 기록한다") - fun shouldCreateScheduledJobAndMarkDoneWhenRefreshSucceeds() { + @DisplayName("fallback은 FALLBACK job을 PENDING -> PROCESSING -> DONE으로 전이하고 refresh를 실행한다") + fun shouldRefreshByFallbackWithJobStatusTransitions() { val refreshService = Mockito.mock(CreatorRankingSnapshotRefreshService::class.java) val jobPort = FakeCreatorRankingSnapshotJobPort() - val now = ZonedDateTime.of(2026, 6, 8, 7, 30, 0, 0, ZoneId.of("Asia/Seoul")) - val redissonClient = periodLockRedissonClient(lockAcquired = true) - val service = CreatorRankingSnapshotJobService(refreshService, jobPort, redissonClient, transactionManager()) { now } + val service = service(refreshService, jobPort) - service.refreshLastCompletedWeekByScheduledJob() + val refreshed = service.refreshLastCompletedWeekByFallback() - val job = jobPort.jobs.single() - assertEquals(CreatorRankingType.WEEKLY, job.rankingType) - assertEquals(LocalDateTime.of(2026, 5, 31, 15, 0), job.aggregationStartAtUtc) - assertEquals(LocalDateTime.of(2026, 6, 7, 15, 0), job.aggregationEndAtUtc) - assertEquals(LocalDateTime.of(2026, 6, 8, 0, 0), job.visibleFromAtUtc) - assertEquals(CreatorRankingSnapshotJobTrigger.SCHEDULED, job.trigger) - assertEquals(CreatorRankingSnapshotJobStatus.DONE, job.status) - assertEquals(null, job.lastError) - Mockito.verify(refreshService).refreshLastCompletedWeek(now) - } - - @Test - @DisplayName("스케줄 실행 실패는 FAILED 상태와 실패 사유를 기록하고 예외를 전파한다") - fun shouldMarkScheduledJobFailedWhenRefreshFails() { - val refreshService = Mockito.mock(CreatorRankingSnapshotRefreshService::class.java) - val jobPort = FakeCreatorRankingSnapshotJobPort() - val now = ZonedDateTime.of(2026, 6, 8, 7, 30, 0, 0, ZoneId.of("Asia/Seoul")) - val redissonClient = periodLockRedissonClient(lockAcquired = true) - val service = CreatorRankingSnapshotJobService(refreshService, jobPort, redissonClient, transactionManager()) { now } - Mockito.doThrow(IllegalStateException("aggregate failed")) - .`when`(refreshService).refreshLastCompletedWeek(now) - - val exception = assertThrows(IllegalStateException::class.java) { - service.refreshLastCompletedWeekByScheduledJob() - } - - assertEquals("aggregate failed", exception.message) - assertEquals(CreatorRankingSnapshotJobStatus.FAILED, jobPort.jobs.single().status) - assertEquals("aggregate failed", jobPort.jobs.single().lastError) - } - - @Test - @DisplayName("관리자 수동 생성은 지정 UTC 기간의 MANUAL PENDING job을 만든다") - fun shouldCreateManualPendingJobForRequestedPeriod() { - val refreshService = Mockito.mock(CreatorRankingSnapshotRefreshService::class.java) - val jobPort = FakeCreatorRankingSnapshotJobPort() - val service = CreatorRankingSnapshotJobService(refreshService, jobPort, unusedRedissonClient(), transactionManager()) - val startAt = LocalDateTime.of(2026, 5, 31, 15, 0) - val endAt = LocalDateTime.of(2026, 6, 7, 15, 0) - - val job = service.createManualJob(startAt, endAt) - - assertEquals(startAt, job.aggregationStartAtUtc) - assertEquals(endAt, job.aggregationEndAtUtc) - assertEquals(CreatorRankingType.WEEKLY, job.rankingType) - assertEquals(LocalDateTime.of(2026, 6, 8, 0, 0), job.visibleFromAtUtc) - assertEquals(CreatorRankingSnapshotJobTrigger.MANUAL, job.trigger) - assertEquals(CreatorRankingSnapshotJobStatus.PENDING, job.status) - assertEquals(null, job.lastError) - assertEquals(null, job.processingStartedAt) - assertEquals(null, job.processedAt) - } - - @Test - @DisplayName("관리자 목록 조회는 기간과 상태 조건으로 snapshot job을 조회한다") - fun shouldFindJobsByRequestedPeriodAndStatuses() { - val refreshService = Mockito.mock(CreatorRankingSnapshotRefreshService::class.java) - val jobPort = FakeCreatorRankingSnapshotJobPort() - val service = CreatorRankingSnapshotJobService(refreshService, jobPort, unusedRedissonClient(), transactionManager()) - val startAt = LocalDateTime.of(2026, 5, 31, 15, 0) - val endAt = LocalDateTime.of(2026, 6, 7, 15, 0) - val failed = jobPort.save( - CreatorRankingSnapshotJobRecord( - rankingType = CreatorRankingType.WEEKLY, - aggregationStartAtUtc = startAt, - aggregationEndAtUtc = endAt, - visibleFromAtUtc = endAt.plusHours(9), - trigger = CreatorRankingSnapshotJobTrigger.MANUAL, - status = CreatorRankingSnapshotJobStatus.FAILED, - lastError = "aggregate failed", - processingStartedAt = null, - processedAt = null - ) - ) - jobPort.save( - CreatorRankingSnapshotJobRecord( - rankingType = CreatorRankingType.WEEKLY, - aggregationStartAtUtc = startAt, - aggregationEndAtUtc = endAt, - visibleFromAtUtc = endAt.plusHours(9), - trigger = CreatorRankingSnapshotJobTrigger.MANUAL, - status = CreatorRankingSnapshotJobStatus.DONE, - lastError = null, - processingStartedAt = null, - processedAt = null - ) - ) - - val jobs = service.findJobs( - aggregationStartAtUtc = startAt, - aggregationEndAtUtc = endAt, - statuses = listOf(CreatorRankingSnapshotJobStatus.FAILED) - ) - - assertEquals(listOf(failed.id), jobs.map { it.id }) - } - - @Test - @DisplayName("관리자 실패 job 재시도는 FAILED job만 PENDING으로 되돌린다") - fun shouldRetryOnlyFailedSnapshotJob() { - val refreshService = Mockito.mock(CreatorRankingSnapshotRefreshService::class.java) - val jobPort = FakeCreatorRankingSnapshotJobPort() - val service = CreatorRankingSnapshotJobService(refreshService, jobPort, unusedRedissonClient(), transactionManager()) - val failed = jobPort.save( - CreatorRankingSnapshotJobRecord( - rankingType = CreatorRankingType.WEEKLY, - aggregationStartAtUtc = LocalDateTime.of(2026, 5, 31, 15, 0), - aggregationEndAtUtc = LocalDateTime.of(2026, 6, 7, 15, 0), - visibleFromAtUtc = LocalDateTime.of(2026, 6, 8, 0, 0), - trigger = CreatorRankingSnapshotJobTrigger.MANUAL, - status = CreatorRankingSnapshotJobStatus.FAILED, - lastError = "aggregate failed", - processingStartedAt = LocalDateTime.of(2026, 6, 8, 7, 30), - processedAt = LocalDateTime.of(2026, 6, 8, 7, 31) - ) - ) - val pending = jobPort.save( - CreatorRankingSnapshotJobRecord( - rankingType = CreatorRankingType.WEEKLY, - aggregationStartAtUtc = LocalDateTime.of(2026, 5, 31, 15, 0), - aggregationEndAtUtc = LocalDateTime.of(2026, 6, 7, 15, 0), - visibleFromAtUtc = LocalDateTime.of(2026, 6, 8, 0, 0), - trigger = CreatorRankingSnapshotJobTrigger.MANUAL, - status = CreatorRankingSnapshotJobStatus.PENDING, - lastError = "keep", - processingStartedAt = null, - processedAt = null - ) - ) - - service.retryFailedJob(failed.id!!) - service.retryFailedJob(pending.id!!) - service.retryFailedJob(999L) - - val retried = jobPort.findById(failed.id!!)!! - val unchanged = jobPort.findById(pending.id!!)!! - assertEquals(CreatorRankingSnapshotJobStatus.PENDING, retried.status) - assertEquals(null, retried.lastError) - assertEquals(null, retried.processingStartedAt) - assertEquals(null, retried.processedAt) - assertEquals(CreatorRankingSnapshotJobStatus.PENDING, unchanged.status) - assertEquals("keep", unchanged.lastError) - } - - @Test - @DisplayName("스케줄 job 상태 변경은 job id와 상태를 로그로 남긴다") - fun shouldLogScheduledJobStatusTransitions(output: CapturedOutput) { - val refreshService = Mockito.mock(CreatorRankingSnapshotRefreshService::class.java) - val jobPort = FakeCreatorRankingSnapshotJobPort() - val now = ZonedDateTime.of(2026, 6, 8, 7, 30, 0, 0, ZoneId.of("Asia/Seoul")) - val redissonClient = periodLockRedissonClient(lockAcquired = true) - val service = CreatorRankingSnapshotJobService(refreshService, jobPort, redissonClient, transactionManager()) { now } - - service.refreshLastCompletedWeekByScheduledJob() - - assertTrue(output.out.contains("event=creator_ranking_snapshot_job_status_changed")) - assertTrue(output.out.contains("jobId=1")) - assertTrue(output.out.contains("trigger=SCHEDULED")) - assertTrue(output.out.contains("status=PROCESSING")) - assertTrue(output.out.contains("status=DONE")) - } - - @Test - @DisplayName("실패 job 상태 변경은 실패 상태와 사유를 로그로 남긴다") - fun shouldLogFailedScheduledJobStatusTransition(output: CapturedOutput) { - val refreshService = Mockito.mock(CreatorRankingSnapshotRefreshService::class.java) - val jobPort = FakeCreatorRankingSnapshotJobPort() - val now = ZonedDateTime.of(2026, 6, 8, 7, 30, 0, 0, ZoneId.of("Asia/Seoul")) - val redissonClient = periodLockRedissonClient(lockAcquired = true) - val service = CreatorRankingSnapshotJobService(refreshService, jobPort, redissonClient, transactionManager()) { now } - Mockito.doThrow(IllegalStateException("aggregate failed")) - .`when`(refreshService).refreshLastCompletedWeek(now) - - assertThrows(IllegalStateException::class.java) { - service.refreshLastCompletedWeekByScheduledJob() - } - - assertTrue(output.out.contains("event=creator_ranking_snapshot_job_status_changed")) - assertTrue(output.out.contains("jobId=1")) - assertTrue(output.out.contains("trigger=SCHEDULED")) - assertTrue(output.out.contains("status=FAILED")) - assertTrue(output.out.contains("error=aggregate failed")) - } - - @Test - @DisplayName("스케줄 job refresh는 cold-start와 같은 기간 기반 lock 경계를 사용한다") - fun shouldUseSamePeriodLockForScheduledJobRefresh() { - val refreshService = Mockito.mock(CreatorRankingSnapshotRefreshService::class.java) - val jobPort = FakeCreatorRankingSnapshotJobPort() - val redissonClient = periodLockRedissonClient(lockAcquired = true) - val now = ZonedDateTime.of(2026, 6, 8, 7, 30, 0, 0, ZoneId.of("Asia/Seoul")) - val lockName = "lock:creator-ranking-snapshot-refresh:2026-05-31T15:00:2026-06-07T15:00" - val service = CreatorRankingSnapshotJobService(refreshService, jobPort, redissonClient, transactionManager()) { now } - - service.refreshLastCompletedWeekByScheduledJob() - - Mockito.verify(redissonClient).getLock(lockName) - Mockito.verify(refreshService).refreshLastCompletedWeek(now) + assertTrue(refreshed) + assertEquals(CreatorRankingSnapshotJobTrigger.FALLBACK, jobPort.jobs.single().trigger) assertEquals(CreatorRankingSnapshotJobStatus.DONE, jobPort.jobs.single().status) + Mockito.verify(refreshService).refreshLastCompletedWeek(ZonedDateTime.of(2026, 6, 8, 9, 0, 0, 0, ZoneId.of("Asia/Seoul"))) } @Test - @DisplayName("스케줄 job refresh는 기간 기반 lock 획득 실패 시 job 생성과 refresh를 건너뛴다") - fun shouldSkipScheduledJobRefreshWhenPeriodLockNotAcquired() { + @DisplayName("fallback job이 같은 기간에 3회 이상이면 refresh와 job 생성을 건너뛴다") + fun shouldSkipFallbackWhenFallbackLimitReached() { val refreshService = Mockito.mock(CreatorRankingSnapshotRefreshService::class.java) val jobPort = FakeCreatorRankingSnapshotJobPort() - val redissonClient = periodLockRedissonClient(lockAcquired = false) - val now = ZonedDateTime.of(2026, 6, 8, 7, 30, 0, 0, ZoneId.of("Asia/Seoul")) - val service = CreatorRankingSnapshotJobService(refreshService, jobPort, redissonClient, transactionManager()) { now } + jobPort.fallbackCount = 3 + val service = service(refreshService, jobPort) - service.refreshLastCompletedWeekByScheduledJob() + val refreshed = service.refreshLastCompletedWeekByFallback() + assertFalse(refreshed) assertTrue(jobPort.jobs.isEmpty()) - Mockito.verify(refreshService, Mockito.never()).refreshLastCompletedWeek(now) + Mockito.verifyNoInteractions(refreshService) } @Test - @DisplayName("기간 기반 lock은 스냅샷 refresh transaction commit 이후 해제한다") - fun shouldUnlockPeriodLockAfterRefreshTransactionCommit() { + @DisplayName("fallback refresh 실패는 job을 FAILED로 남기고 예외를 전파한다") + fun shouldMarkFailedWhenFallbackRefreshFails() { val refreshService = Mockito.mock(CreatorRankingSnapshotRefreshService::class.java) val jobPort = FakeCreatorRankingSnapshotJobPort() - val redissonClient = Mockito.mock(RedissonClient::class.java) - val lock = Mockito.mock(RLock::class.java) - val transactionManager = Mockito.mock(PlatformTransactionManager::class.java) - val transactionStatus = SimpleTransactionStatus() - val now = ZonedDateTime.of(2026, 6, 8, 7, 30, 0, 0, ZoneId.of("Asia/Seoul")) - val lockName = "lock:creator-ranking-snapshot-refresh:2026-05-31T15:00:2026-06-07T15:00" - Mockito.`when`(redissonClient.getLock(lockName)).thenReturn(lock) - Mockito.`when`(lock.tryLock(0, -1, TimeUnit.SECONDS)).thenReturn(true) - Mockito.`when`(lock.isHeldByCurrentThread).thenReturn(true) - Mockito.`when`(transactionManager.getTransaction(Mockito.any(TransactionDefinition::class.java))) - .thenReturn(transactionStatus) - val service = CreatorRankingSnapshotJobService( - refreshService, - jobPort, - redissonClient, - transactionManager - ) { now } + Mockito.doThrow(IllegalStateException("refresh failed")).`when`(refreshService).refreshLastCompletedWeek( + anyZonedDateTime() + ) + val service = service(refreshService, jobPort) - service.refreshLastCompletedWeekByScheduledJob() + val exception = assertThrows(IllegalStateException::class.java) { service.refreshLastCompletedWeekByFallback() } - val inOrder = Mockito.inOrder(transactionManager, lock) - inOrder.verify(transactionManager).commit(transactionStatus) - inOrder.verify(lock).unlock() - } - - @Test - @DisplayName("스케줄 refresh 실패 시 rollback 이후 별도 transaction으로 FAILED 상태를 커밋한다") - fun shouldCommitFailedStatusAfterRefreshTransactionRollback() { - val refreshService = Mockito.mock(CreatorRankingSnapshotRefreshService::class.java) - val jobPort = FakeCreatorRankingSnapshotJobPort() - val redissonClient = periodLockRedissonClient(lockAcquired = true) - val transactionManager = Mockito.mock(PlatformTransactionManager::class.java) - val saveStatus = SimpleTransactionStatus() - val processingStatus = SimpleTransactionStatus() - val refreshStatus = SimpleTransactionStatus() - val failedStatus = SimpleTransactionStatus() - val now = ZonedDateTime.of(2026, 6, 8, 1, 0, 0, 0, ZoneId.of("Asia/Seoul")) - Mockito.`when`(transactionManager.getTransaction(Mockito.any(TransactionDefinition::class.java))) - .thenReturn(saveStatus, processingStatus, refreshStatus, failedStatus) - Mockito.doThrow(IllegalStateException("aggregate failed")) - .`when`(refreshService).refreshLastCompletedWeek(now) - val service = CreatorRankingSnapshotJobService( - refreshService, - jobPort, - redissonClient, - transactionManager - ) { now } - - val exception = assertThrows(IllegalStateException::class.java) { - service.refreshLastCompletedWeekByScheduledJob() - } - - assertEquals("aggregate failed", exception.message) + assertEquals("refresh failed", exception.message) assertEquals(CreatorRankingSnapshotJobStatus.FAILED, jobPort.jobs.single().status) - assertEquals("aggregate failed", jobPort.jobs.single().lastError) - val inOrder = Mockito.inOrder(transactionManager) - inOrder.verify(transactionManager).commit(saveStatus) - inOrder.verify(transactionManager).commit(processingStatus) - inOrder.verify(transactionManager).rollback(refreshStatus) - inOrder.verify(transactionManager).commit(failedStatus) + assertEquals("refresh failed", jobPort.jobs.single().lastError) } @Test - @DisplayName("cold-start 스냅샷 생성은 기간 기반 lock 획득 시에만 refresh를 실행한다") - fun shouldRefreshColdStartSnapshotOnlyWhenPeriodLockAcquired() { + @DisplayName("period lock 획득 실패는 fallback을 정상 skip한다") + fun shouldSkipFallbackWhenLockNotAcquired() { val refreshService = Mockito.mock(CreatorRankingSnapshotRefreshService::class.java) val jobPort = FakeCreatorRankingSnapshotJobPort() - val redissonClient = Mockito.mock(RedissonClient::class.java) val lock = Mockito.mock(RLock::class.java) - val now = ZonedDateTime.of(2026, 6, 8, 7, 30, 0, 0, ZoneId.of("Asia/Seoul")) - val lockName = "lock:creator-ranking-snapshot-refresh:2026-05-31T15:00:2026-06-07T15:00" - Mockito.`when`(redissonClient.getLock(lockName)).thenReturn(lock) + Mockito.`when`(lock.tryLock(0, -1, TimeUnit.SECONDS)).thenReturn(false) + val redissonClient = Mockito.mock(RedissonClient::class.java) + Mockito.`when`(redissonClient.getLock(Mockito.anyString())).thenReturn(lock) + val service = service(refreshService, jobPort, redissonClient) + + val refreshed = service.refreshLastCompletedWeekByFallback() + + assertFalse(refreshed) + assertTrue(jobPort.jobs.isEmpty()) + Mockito.verifyNoInteractions(refreshService) + } + + private fun service( + refreshService: CreatorRankingSnapshotRefreshService, + jobPort: CreatorRankingSnapshotJobPort, + redissonClient: RedissonClient = lockedRedisson(), + now: ZonedDateTime = ZonedDateTime.of(2026, 6, 8, 9, 0, 0, 0, ZoneId.of("Asia/Seoul")) + ) = CreatorRankingSnapshotJobService(refreshService, jobPort, redissonClient, ImmediateTransactionManager(), { now }) + + private fun anyZonedDateTime(): ZonedDateTime { + return Mockito.any(ZonedDateTime::class.java) ?: ZonedDateTime.now() + } + + private fun lockedRedisson(): RedissonClient { + val lock = Mockito.mock(RLock::class.java) Mockito.`when`(lock.tryLock(0, -1, TimeUnit.SECONDS)).thenReturn(true) Mockito.`when`(lock.isHeldByCurrentThread).thenReturn(true) - val service = CreatorRankingSnapshotJobService(refreshService, jobPort, redissonClient, transactionManager()) { now } - - service.ensureLastCompletedWeekSnapshotForColdStart() - - Mockito.verify(redissonClient).getLock(lockName) - Mockito.verify(lock).tryLock(0, -1, TimeUnit.SECONDS) - Mockito.verify(refreshService).refreshLastCompletedWeek(now) - Mockito.verify(lock).unlock() - } - - @Test - @DisplayName("cold-start 스냅샷 생성은 기간 기반 lock 획득 실패 시 refresh를 실행하지 않는다") - fun shouldSkipColdStartSnapshotRefreshWhenPeriodLockNotAcquired() { - val refreshService = Mockito.mock(CreatorRankingSnapshotRefreshService::class.java) - val jobPort = FakeCreatorRankingSnapshotJobPort() val redissonClient = Mockito.mock(RedissonClient::class.java) - val lock = Mockito.mock(RLock::class.java) - val now = ZonedDateTime.of(2026, 6, 8, 7, 30, 0, 0, ZoneId.of("Asia/Seoul")) - val lockName = "lock:creator-ranking-snapshot-refresh:2026-05-31T15:00:2026-06-07T15:00" - Mockito.`when`(redissonClient.getLock(lockName)).thenReturn(lock) - Mockito.`when`(lock.tryLock(0, -1, TimeUnit.SECONDS)).thenReturn(false) - Mockito.`when`(lock.isHeldByCurrentThread).thenReturn(false) - val service = CreatorRankingSnapshotJobService(refreshService, jobPort, redissonClient, transactionManager()) { now } - - service.ensureLastCompletedWeekSnapshotForColdStart() - - Mockito.verify(redissonClient).getLock(lockName) - Mockito.verify(lock).tryLock(0, -1, TimeUnit.SECONDS) - Mockito.verify(refreshService, Mockito.never()).refreshLastCompletedWeek(now) - Mockito.verify(lock, Mockito.never()).unlock() + Mockito.`when`(redissonClient.getLock(Mockito.anyString())).thenReturn(lock) + return redissonClient } } -private fun unusedRedissonClient(): RedissonClient = Mockito.mock(RedissonClient::class.java) - -private fun transactionManager(): PlatformTransactionManager { - val transactionManager = Mockito.mock(PlatformTransactionManager::class.java) - Mockito.`when`(transactionManager.getTransaction(Mockito.any(TransactionDefinition::class.java))) - .thenReturn(SimpleTransactionStatus()) - return transactionManager -} - -private fun periodLockRedissonClient(lockAcquired: Boolean): RedissonClient { - val redissonClient = Mockito.mock(RedissonClient::class.java) - val lock = Mockito.mock(RLock::class.java) - val lockName = "lock:creator-ranking-snapshot-refresh:2026-05-31T15:00:2026-06-07T15:00" - Mockito.`when`(redissonClient.getLock(lockName)).thenReturn(lock) - Mockito.`when`(lock.tryLock(0, -1, TimeUnit.SECONDS)).thenReturn(lockAcquired) - Mockito.`when`(lock.isHeldByCurrentThread).thenReturn(lockAcquired) - return redissonClient +private class ImmediateTransactionManager : PlatformTransactionManager { + override fun getTransaction(definition: TransactionDefinition?): TransactionStatus = SimpleTransactionStatus() + override fun commit(status: TransactionStatus) = Unit + override fun rollback(status: TransactionStatus) = Unit } private class FakeCreatorRankingSnapshotJobPort : CreatorRankingSnapshotJobPort { val jobs = mutableListOf() + var fallbackCount = 0L private var nextId = 1L override fun save(job: CreatorRankingSnapshotJobRecord): CreatorRankingSnapshotJobRecord { - val saved = job.copy(id = job.id ?: nextId++) + val saved = job.copy(id = nextId++) jobs.add(saved) return saved } - override fun findById(jobId: Long): CreatorRankingSnapshotJobRecord? { - return jobs.firstOrNull { it.id == jobId } - } + override fun findById(jobId: Long) = jobs.firstOrNull { it.id == jobId } override fun findByPeriodAndStatuses( aggregationStartAtUtc: LocalDateTime, aggregationEndAtUtc: LocalDateTime, statuses: List - ): List { - return jobs.filter { - it.aggregationStartAtUtc == aggregationStartAtUtc && - it.aggregationEndAtUtc == aggregationEndAtUtc && - it.status in statuses - } + ) = jobs.filter { it.status in statuses } + + override fun countByRankingTypeAndPeriodAndTrigger( + rankingType: CreatorRankingType, + aggregationStartAtUtc: LocalDateTime, + aggregationEndAtUtc: LocalDateTime, + trigger: CreatorRankingSnapshotJobTrigger + ) = fallbackCount + + override fun markProcessing(jobId: Long, processingStartedAt: LocalDateTime) = update(jobId) { + it.copy( + status = CreatorRankingSnapshotJobStatus.PROCESSING, + processingStartedAt = processingStartedAt + ) } - override fun markProcessing(jobId: Long, processingStartedAt: LocalDateTime): CreatorRankingSnapshotJobRecord? { - return update(jobId) { - it.copy( - status = CreatorRankingSnapshotJobStatus.PROCESSING, - processingStartedAt = processingStartedAt - ) - } + override fun markDone(jobId: Long, processedAt: LocalDateTime) = update(jobId) { + it.copy( + status = CreatorRankingSnapshotJobStatus.DONE, + processedAt = processedAt + ) } - override fun markDone(jobId: Long, processedAt: LocalDateTime): CreatorRankingSnapshotJobRecord? { - return update(jobId) { - it.copy( - status = CreatorRankingSnapshotJobStatus.DONE, - processedAt = processedAt, - lastError = null - ) - } + override fun markFailed( + jobId: Long, + processedAt: LocalDateTime, + lastError: String? + ) = update(jobId) { + it.copy( + status = CreatorRankingSnapshotJobStatus.FAILED, + processedAt = processedAt, + lastError = lastError + ) } - override fun markFailed(jobId: Long, processedAt: LocalDateTime, lastError: String?): CreatorRankingSnapshotJobRecord? { - return update(jobId) { - it.copy( - status = CreatorRankingSnapshotJobStatus.FAILED, - processedAt = processedAt, - lastError = lastError - ) - } - } - - override fun markPending(jobId: Long): CreatorRankingSnapshotJobRecord? { - return update(jobId) { - it.copy( - status = CreatorRankingSnapshotJobStatus.PENDING, - lastError = null, - processingStartedAt = null, - processedAt = null - ) - } + override fun markPending(jobId: Long) = update(jobId) { + it.copy( + status = CreatorRankingSnapshotJobStatus.PENDING, + lastError = null, + processedAt = null, + processingStartedAt = null + ) } private fun update( jobId: Long, - updater: (CreatorRankingSnapshotJobRecord) -> CreatorRankingSnapshotJobRecord + mapper: (CreatorRankingSnapshotJobRecord) -> CreatorRankingSnapshotJobRecord ): CreatorRankingSnapshotJobRecord? { val index = jobs.indexOfFirst { it.id == jobId } if (index < 0) return null - val updated = updater(jobs[index]) - jobs[index] = updated - return updated + jobs[index] = mapper(jobs[index]) + return jobs[index] } }