feat(home): 크리에이터 랭킹 fallback refresh job을 실행한다
This commit is contained in:
@@ -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 markProcessing(jobId: Long) {
|
||||
transactionTemplate.executeWithoutResult {
|
||||
jobPort.markProcessing(jobId, LocalDateTime.now())
|
||||
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()) }
|
||||
}
|
||||
|
||||
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> = CreatorRankingSnapshotJobStatus.values().toList()
|
||||
): List<CreatorRankingSnapshotJobRecord> {
|
||||
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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
}
|
||||
|
||||
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)
|
||||
Mockito.`when`(redissonClient.getLock(Mockito.anyString())).thenReturn(lock)
|
||||
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<CreatorRankingSnapshotJobRecord>()
|
||||
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<CreatorRankingSnapshotJobStatus>
|
||||
): List<CreatorRankingSnapshotJobRecord> {
|
||||
return jobs.filter {
|
||||
it.aggregationStartAtUtc == aggregationStartAtUtc &&
|
||||
it.aggregationEndAtUtc == aggregationEndAtUtc &&
|
||||
it.status in statuses
|
||||
}
|
||||
}
|
||||
) = jobs.filter { it.status in statuses }
|
||||
|
||||
override fun markProcessing(jobId: Long, processingStartedAt: LocalDateTime): CreatorRankingSnapshotJobRecord? {
|
||||
return update(jobId) {
|
||||
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 markDone(jobId: Long, processedAt: LocalDateTime): CreatorRankingSnapshotJobRecord? {
|
||||
return update(jobId) {
|
||||
override fun markDone(jobId: Long, processedAt: LocalDateTime) = update(jobId) {
|
||||
it.copy(
|
||||
status = CreatorRankingSnapshotJobStatus.DONE,
|
||||
processedAt = processedAt,
|
||||
lastError = null
|
||||
processedAt = processedAt
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
override fun markFailed(jobId: Long, processedAt: LocalDateTime, lastError: String?): CreatorRankingSnapshotJobRecord? {
|
||||
return update(jobId) {
|
||||
override fun markFailed(
|
||||
jobId: Long,
|
||||
processedAt: LocalDateTime,
|
||||
lastError: String?
|
||||
) = update(jobId) {
|
||||
it.copy(
|
||||
status = CreatorRankingSnapshotJobStatus.FAILED,
|
||||
processedAt = processedAt,
|
||||
lastError = lastError
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
override fun markPending(jobId: Long): CreatorRankingSnapshotJobRecord? {
|
||||
return update(jobId) {
|
||||
override fun markPending(jobId: Long) = update(jobId) {
|
||||
it.copy(
|
||||
status = CreatorRankingSnapshotJobStatus.PENDING,
|
||||
lastError = null,
|
||||
processingStartedAt = null,
|
||||
processedAt = 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]
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user