diff --git a/src/test/kotlin/pw/binom/mirror/worker/db/JobRepositoryTest.kt b/src/test/kotlin/pw/binom/mirror/worker/db/JobRepositoryTest.kt index 522d67b..e649c98 100644 --- a/src/test/kotlin/pw/binom/mirror/worker/db/JobRepositoryTest.kt +++ b/src/test/kotlin/pw/binom/mirror/worker/db/JobRepositoryTest.kt @@ -122,6 +122,81 @@ class JobRepositoryTest : AbstractPostgresTest() { assertEquals("boom", job.error) } + @Test + fun `markFailed increments retryCount`() { + val id = insertJob("item-1") + repository.takeNextJob() + + repository.markFailed(id, "boom") + + val job = repository.findByStatus("failed").single() + assertEquals(1, job.retryCount) + } + + @Test + fun `takeNextJob picks failed job with retryCount below max`() { + val id = insertJob(itemId = "item-1", status = "failed", retryCount = 0) + + val taken = repository.takeNextJob() + + assertNotNull(taken) + assertEquals(id, taken.id) + assertEquals("processing", taken.status) + } + + @Test + fun `takeNextJob does not pick failed job with retryCount at max`() { + insertJob(itemId = "item-1", status = "failed", retryCount = 3) + + val taken = repository.takeNextJob() + + assertNull(taken) + } + + @Test + fun `takeNextJob picks stale processing job`() { + val id = insertJob( + itemId = "item-1", + status = "processing", + takenAt = Instant.now().minus(31, ChronoUnit.SECONDS), + ) + + val taken = repository.takeNextJob() + + assertNotNull(taken) + assertEquals(id, taken.id) + assertEquals("processing", taken.status) + } + + @Test + fun `takeNextJob does not pick fresh processing job`() { + insertJob(itemId = "item-1", status = "processing", takenAt = Instant.now()) + + val taken = repository.takeNextJob() + + assertNull(taken) + } + + @Test + fun `takeNextJob priority new over failed over stale`() { + val newId = insertJob(itemId = "item-new") + val failedId = insertJob(itemId = "item-failed", status = "failed", retryCount = 0) + insertJob( + itemId = "item-stale", + status = "processing", + takenAt = Instant.now().minus(31, ChronoUnit.SECONDS), + ) + + val first = repository.takeNextJob() + assertEquals(newId, first?.id) + + val second = repository.takeNextJob() + assertEquals(failedId, second?.id) + + val third = repository.takeNextJob() + assertEquals("item-stale", third?.itemId) + } + @Test fun `stale processing job returns to new`() { val id = insertJob( @@ -177,18 +252,20 @@ class JobRepositoryTest : AbstractPostgresTest() { status: String = "new", sourceUrl: String = "http://localhost/video.mp4", takenAt: Instant? = null, + retryCount: Int = 0, ): UUID { val id = UUID.randomUUID() JdbcTemplate(dataSource).update( """ - INSERT INTO media_mirror.jobs (id, item_id, source_url, status, taken_at) - VALUES (?, ?, ?, ?, ?) + INSERT INTO media_mirror.jobs (id, item_id, source_url, status, taken_at, retry_count) + VALUES (?, ?, ?, ?, ?, ?) """.trimIndent(), id, itemId, sourceUrl, status, takenAt?.let { java.sql.Timestamp.from(it) }, + retryCount, ) return id }