2 Commits
9 ... 10

Author SHA1 Message Date
Hermes Agent b9c248e359 test: retry-механизм (6 тестов)
Build Media Mirror Worker / Build and publish (release) Successful in 1m20s
- markFailed increments retryCount
- takeNextJob picks failed job (retryCount < max)
- takeNextJob does not pick failed job (retryCount >= max)
- takeNextJob picks stale processing job (30с)
- takeNextJob does not pick fresh processing job
- takeNextJob priority: new > failed > stale
2026-08-21 12:30:50 +03:00
Hermes Agent 12233e9bfa feat: retry-механизм (max 3 попытки)
- V2__add_retry_count.sql: retry_count int NOT NULL DEFAULT 0
- takeNextJob: new + failed (retry_count < 3) + stale processing (30с)
- markFailed: retry_count = retry_count + 1
- JobProcessor: логирование retry (X/3)
- Heartbeat: не трогать (уже работает, 10с)
2026-08-21 12:17:42 +03:00
5 changed files with 100 additions and 5 deletions
@@ -26,6 +26,7 @@ data class Job(
val createdAt: Instant,
val updatedAt: Instant,
val takenAt: Instant? = null,
val retryCount: Int = 0,
) {
companion object {
const val STATUS_NEW = "new"
@@ -33,5 +34,6 @@ data class Job(
const val STATUS_DONE = "done"
const val STATUS_FAILED = "failed"
const val STATUS_CANCELLED = "cancelled"
const val MAX_RETRIES = 3
}
}
@@ -71,7 +71,7 @@ open class JobRepository(
jdbcTemplate.update(
"""
UPDATE media_mirror.jobs
SET status = 'failed', error = ?, updated_at = now()
SET status = 'failed', error = ?, retry_count = retry_count + 1, updated_at = now()
WHERE id = ?
""".trimIndent(),
error.take(4000),
@@ -110,8 +110,16 @@ open class JobRepository(
val jobs = jdbcTemplate.query(
"""
SELECT * FROM media_mirror.jobs
WHERE status = 'new'
ORDER BY created_at
WHERE (
status = 'new'
OR (status = 'failed' AND retry_count < 3)
OR (status = 'processing' AND taken_at < now() - interval '30 seconds')
)
ORDER BY
CASE WHEN status = 'new' THEN 0
WHEN status = 'failed' THEN 1
ELSE 2 END,
created_at
LIMIT 1
FOR UPDATE SKIP LOCKED
""".trimIndent(),
@@ -152,6 +160,7 @@ open class JobRepository(
createdAt = rs.getTimestamp("created_at").toInstant(),
updatedAt = rs.getTimestamp("updated_at").toInstant(),
takenAt = rs.getTimestamp("taken_at")?.toInstant(),
retryCount = rs.getInt("retry_count"),
)
}
}
@@ -58,6 +58,12 @@ class JobProcessor(
throw e
} catch (e: Throwable) {
logger.warn("Job {} finished with error: {}", job.id, e.message)
val currentRetries = jobRepository.findById(job.id)?.retryCount ?: 0
if (currentRetries >= Job.MAX_RETRIES) {
logger.warn("Job {} exceeded max retries ({}), permanently failed", job.id, currentRetries)
} else {
logger.info("Job {} failed (retry {}/{}), will be retried", job.id, currentRetries, Job.MAX_RETRIES)
}
jobRepository.markFailed(job.id, e.message ?: e.toString())
} finally {
heartbeatJob.cancel()
@@ -0,0 +1 @@
ALTER TABLE media_mirror.jobs ADD COLUMN IF NOT EXISTS retry_count int NOT NULL DEFAULT 0;
@@ -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
}