From 68eefff884f4b77722cc78e58f8a68398744f7e5 Mon Sep 17 00:00:00 2001 From: Hermes Agent Date: Wed, 19 Aug 2026 05:16:23 +0300 Subject: [PATCH] =?UTF-8?q?fix:=20heartbeat=20=D0=B2=D0=BE=D1=80=D0=BA?= =?UTF-8?q?=D0=B5=D1=80=D0=B0=20=E2=80=94=20=D0=B6=D0=B8=D0=B2=D0=BE=D0=B9?= =?UTF-8?q?=20=D0=B4=D0=B6=D0=BE=D0=B1=20=D0=BD=D0=B5=20=D0=BF=D0=B5=D1=80?= =?UTF-8?q?=D0=B5=D1=85=D0=B2=D0=B0=D1=82=D1=8B=D0=B2=D0=B0=D0=B5=D1=82?= =?UTF-8?q?=D1=81=D1=8F=20=D0=BA=D0=B0=D0=BA=20stale?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Проблема: returnStaleToNew (stale-timeout 30м) смотрит только taken_at, который ставится один раз при взятии джоба. Долгий джоб (скачивание 7.2 ГБ 11 мин + конвертация 24 мин > 30 мин) возвращался в очередь и перехватывался вторым воркером → дубль работы, падение на скачивании, статус failed при успешно работающем первом воркере. Решение: heartbeat — пока джоб выполняется, воркер каждые 10 сек обновляет taken_at (UPDATE ... WHERE status='processing'). Мёртвый воркер по-прежнему перехватывается через stale-timeout; живой — никогда не считается зависшим. - JobRepository.heartbeat(id): UPDATE taken_at/updated_at с guard processing - JobProcessor: параллельная корутина heartbeat (свой CoroutineScope+Dispatchers.IO), отмена в finally (done/failed/cancelled) - AppProperties.heartbeatInterval = 10s, application.yaml heartbeat-interval: 10s - JobRepositoryTest: heartbeat обновляет processing-джоб, не трогает non-processing --- .../mirror/worker/config/AppProperties.kt | 1 + .../binom/mirror/worker/db/JobRepository.kt | 11 ++++++++ .../mirror/worker/worker/JobProcessor.kt | 17 ++++++++++++ src/main/resources/application.yaml | 1 + .../mirror/worker/db/JobRepositoryTest.kt | 26 +++++++++++++++++++ 5 files changed, 56 insertions(+) diff --git a/src/main/kotlin/pw/binom/mirror/worker/config/AppProperties.kt b/src/main/kotlin/pw/binom/mirror/worker/config/AppProperties.kt index 54d6af9..3aff786 100644 --- a/src/main/kotlin/pw/binom/mirror/worker/config/AppProperties.kt +++ b/src/main/kotlin/pw/binom/mirror/worker/config/AppProperties.kt @@ -8,6 +8,7 @@ import java.time.Duration data class AppProperties( val pollInterval: Duration = Duration.ofSeconds(2), val staleTimeout: Duration = Duration.ofMinutes(30), + val heartbeatInterval: Duration = Duration.ofSeconds(10), val jellyfin: Jellyfin = Jellyfin(), val s3: S3 = S3(), val ffmpegPath: String = "ffmpeg", diff --git a/src/main/kotlin/pw/binom/mirror/worker/db/JobRepository.kt b/src/main/kotlin/pw/binom/mirror/worker/db/JobRepository.kt index 11926c3..4df7465 100644 --- a/src/main/kotlin/pw/binom/mirror/worker/db/JobRepository.kt +++ b/src/main/kotlin/pw/binom/mirror/worker/db/JobRepository.kt @@ -42,6 +42,17 @@ open class JobRepository( ) } + open fun heartbeat(id: UUID) { + jdbcTemplate.update( + """ + UPDATE media_mirror.jobs + SET taken_at = now(), updated_at = now() + WHERE id = ? AND status = 'processing' + """.trimIndent(), + id, + ) + } + open fun markDone(id: UUID, videoKey: String, audioKeys: List) { val audioJson = json.encodeToString(audioKeysSerializer, audioKeys) jdbcTemplate.update( diff --git a/src/main/kotlin/pw/binom/mirror/worker/worker/JobProcessor.kt b/src/main/kotlin/pw/binom/mirror/worker/worker/JobProcessor.kt index 80726b4..42b1bc0 100644 --- a/src/main/kotlin/pw/binom/mirror/worker/worker/JobProcessor.kt +++ b/src/main/kotlin/pw/binom/mirror/worker/worker/JobProcessor.kt @@ -1,6 +1,13 @@ package pw.binom.mirror.worker.worker import kotlinx.coroutines.CancellationException +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.cancel +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.delay +import kotlinx.coroutines.isActive +import kotlinx.coroutines.launch import org.slf4j.Logger import org.slf4j.LoggerFactory import org.springframework.stereotype.Service @@ -29,6 +36,13 @@ class JobProcessor( region = props.s3.region, prefix = props.s3.prefix, ) + val heartbeatScope = CoroutineScope(SupervisorJob() + Dispatchers.IO) + val heartbeatJob = heartbeatScope.launch { + while (isActive) { + delay(props.heartbeatInterval.toMillis()) + runCatching { jobRepository.heartbeat(job.id) } + } + } try { val result = jobExecutor.execute( job = job, @@ -45,6 +59,9 @@ class JobProcessor( } catch (e: Throwable) { logger.warn("Job {} finished with error: {}", job.id, e.message) jobRepository.markFailed(job.id, e.message ?: e.toString()) + } finally { + heartbeatJob.cancel() + heartbeatScope.cancel() } } diff --git a/src/main/resources/application.yaml b/src/main/resources/application.yaml index f106462..4e26201 100644 --- a/src/main/resources/application.yaml +++ b/src/main/resources/application.yaml @@ -11,6 +11,7 @@ spring: app: poll-interval: 2s stale-timeout: 30m + heartbeat-interval: 10s jellyfin: url: https://jellyfin.binom.pw/ apiKey: ${JELLYFIN_API_KEY:} 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 6340ba6..522d67b 100644 --- a/src/test/kotlin/pw/binom/mirror/worker/db/JobRepositoryTest.kt +++ b/src/test/kotlin/pw/binom/mirror/worker/db/JobRepositoryTest.kt @@ -146,6 +146,32 @@ class JobRepositoryTest : AbstractPostgresTest() { assertEquals(1, repository.findByStatus("processing").size) } + @Test + fun `heartbeat refreshes takenAt of a processing job`() { + val id = insertJob( + itemId = "item-1", + status = "processing", + takenAt = Instant.now().minus(1, ChronoUnit.HOURS), + ) + + repository.heartbeat(id) + + val job = repository.findById(id)!! + assertEquals("processing", job.status) + assertTrue(job.takenAt!!.isAfter(Instant.now().minus(5, ChronoUnit.SECONDS))) + } + + @Test + fun `heartbeat does not touch a non-processing job`() { + val id = insertJob(itemId = "item-1", status = "new") + + repository.heartbeat(id) + + val job = repository.findById(id)!! + assertEquals("new", job.status) + assertNull(job.takenAt) + } + private fun insertJob( itemId: String, status: String = "new",