1 Commits
6 ... 7

Author SHA1 Message Date
Hermes Agent 68eefff884 fix: heartbeat воркера — живой джоб не перехватывается как stale
Build Media Mirror Worker / Build and publish (release) Successful in 1m25s
Проблема: 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
2026-08-19 05:16:23 +03:00
5 changed files with 56 additions and 0 deletions
@@ -8,6 +8,7 @@ import java.time.Duration
data class AppProperties( data class AppProperties(
val pollInterval: Duration = Duration.ofSeconds(2), val pollInterval: Duration = Duration.ofSeconds(2),
val staleTimeout: Duration = Duration.ofMinutes(30), val staleTimeout: Duration = Duration.ofMinutes(30),
val heartbeatInterval: Duration = Duration.ofSeconds(10),
val jellyfin: Jellyfin = Jellyfin(), val jellyfin: Jellyfin = Jellyfin(),
val s3: S3 = S3(), val s3: S3 = S3(),
val ffmpegPath: String = "ffmpeg", val ffmpegPath: String = "ffmpeg",
@@ -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<AudioKey>) { open fun markDone(id: UUID, videoKey: String, audioKeys: List<AudioKey>) {
val audioJson = json.encodeToString(audioKeysSerializer, audioKeys) val audioJson = json.encodeToString(audioKeysSerializer, audioKeys)
jdbcTemplate.update( jdbcTemplate.update(
@@ -1,6 +1,13 @@
package pw.binom.mirror.worker.worker package pw.binom.mirror.worker.worker
import kotlinx.coroutines.CancellationException 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.Logger
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.stereotype.Service import org.springframework.stereotype.Service
@@ -29,6 +36,13 @@ class JobProcessor(
region = props.s3.region, region = props.s3.region,
prefix = props.s3.prefix, 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 { try {
val result = jobExecutor.execute( val result = jobExecutor.execute(
job = job, job = job,
@@ -45,6 +59,9 @@ class JobProcessor(
} catch (e: Throwable) { } catch (e: Throwable) {
logger.warn("Job {} finished with error: {}", job.id, e.message) logger.warn("Job {} finished with error: {}", job.id, e.message)
jobRepository.markFailed(job.id, e.message ?: e.toString()) jobRepository.markFailed(job.id, e.message ?: e.toString())
} finally {
heartbeatJob.cancel()
heartbeatScope.cancel()
} }
} }
+1
View File
@@ -11,6 +11,7 @@ spring:
app: app:
poll-interval: 2s poll-interval: 2s
stale-timeout: 30m stale-timeout: 30m
heartbeat-interval: 10s
jellyfin: jellyfin:
url: https://jellyfin.binom.pw/ url: https://jellyfin.binom.pw/
apiKey: ${JELLYFIN_API_KEY:} apiKey: ${JELLYFIN_API_KEY:}
@@ -146,6 +146,32 @@ class JobRepositoryTest : AbstractPostgresTest() {
assertEquals(1, repository.findByStatus("processing").size) 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( private fun insertJob(
itemId: String, itemId: String,
status: String = "new", status: String = "new",