6 Commits
5 ... 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
Hermes Agent d10759b26b fix: pad кадра до кратности 16 (H.264 макроблоки) — scale+pad=ceil(iw/16)*16, картинка без искажений
Build Media Mirror Worker / Build and publish (release) Successful in 1m20s
Эмуляторный декодер (c2.goldfish) отказывался декодировать 893x480 (не кратно 16). Теперь после scale любой ширины pad доводит кадр до кратного 16 чёрными полями, пропорции не меняются.
2026-08-19 11:50:16 +03:00
Hermes Agent 96a9a0536d fix: устойчивое скачивание из Jellyfin — OkHttp + retry вместо HttpURLConnection
Build Media Mirror Worker / Build and publish (release) Successful in 1m25s
Прод-инцидент (3 джоба failed подряд): HttpURLConnection рвал скачивание
7.2 ГБ на ~5-й минуте ('Can't download video from Jellyfin'), при этом curl
с той же машины качал стабильно (HTTP/1.1 и HTTP/2). Сеть и Jellyfin чисты.

- JellyfinInput: OkHttpClient (readTimeout=0, followRedirects), byteStream→Files.copy,
  retry до 3 попыток с паузой 5с на IOException, warn-лог каждой попытки
- InputService: лог с полным stacktrace (реальная причина обрыва теперь видна)
2026-08-19 06:00:18 +03:00
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
Hermes Agent 05a4aa775a feat: видео-кодек h264_nvenc (GPU) вместо libvpx-vp9 (CPU), кодек в конфиг
Build Media Mirror Worker / Build and publish (release) Successful in 1m35s
- AppProperties.videoCodec (дефолт h264_nvenc), application.yaml
- тесты: libx264 (без GPU), JobExecutorTest/JobPollerTest
- аудио libvorbis/ogg и mkv-контейнер не тронуты
2026-08-11 00:34:01 +03:00
14 changed files with 206 additions and 24 deletions
+1
View File
@@ -50,6 +50,7 @@ dependencies {
implementation("org.jetbrains.kotlinx:kotlinx-serialization-json:1.11.0")
implementation(platform("software.amazon.awssdk:bom:2.46.21"))
implementation("software.amazon.awssdk:s3")
implementation("com.squareup.okhttp3:okhttp:4.12.0")
testImplementation("org.springframework.boot:spring-boot-starter-test")
testImplementation(kotlin("test-junit5"))
@@ -8,12 +8,14 @@ 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",
val ffprobePath: String = "ffprobe",
val localFileStorage: Path = Path.of("/tmp/media-mirror-worker"),
val ffmpegThreadCount: Int? = null,
val videoCodec: String = "h264_nvenc",
) {
data class Jellyfin(
val url: String = "",
@@ -82,7 +82,9 @@ class FfmpegService(
}
if (width != null || height != null) {
args += "-vf"
args += "scale=${width ?: "-1"}:${height ?: "-1"}"
// Высота прибита (480), ширина — по пропорциям; pad доводит кадр до кратности 16
// (H.264 макроблоки 16x16): чёрные поля по краям, картинка не искажается.
args += "scale=${width ?: "-1"}:${height ?: "-1"},pad=ceil(iw/16)*16:ceil(ih/16)*16:(ow-iw)/2:(oh-ih)/2"
}
if (!removeAudio) {
if (audioCodec != null) {
@@ -28,6 +28,7 @@ class InputService(
logger.info("Download cancelled: {}", source)
throw e
} catch (e: Throwable) {
logger.warn("Download failed: {}", e.message, e)
throw IllegalStateException("Can't download video from $source", e)
}
logger.info("Downloaded success in {} ms", (System.nanoTime() - startTime) / 1_000_000)
@@ -1,15 +1,18 @@
package pw.binom.mirror.worker.convert
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.delay
import kotlinx.coroutines.withContext
import okhttp3.OkHttpClient
import okhttp3.Request
import org.slf4j.Logger
import org.slf4j.LoggerFactory
import org.springframework.stereotype.Component
import java.net.HttpURLConnection
import java.net.URL
import java.io.IOException
import java.nio.file.Files
import java.nio.file.Path
import java.nio.file.StandardCopyOption
import java.util.concurrent.TimeUnit
@Component
class JellyfinInput(
@@ -17,24 +20,46 @@ class JellyfinInput(
) : InputService.InputImplementation {
private val logger: Logger = LoggerFactory.getLogger(JellyfinInput::class.java)
private val client: OkHttpClient by lazy {
OkHttpClient.Builder()
.connectTimeout(30, TimeUnit.SECONDS)
.readTimeout(0, TimeUnit.SECONDS)
.followRedirects(true)
.followSslRedirects(true)
.build()
}
override fun isSupport(source: InputSource): Boolean = source is InputSource.Jellyfin
override suspend fun download(source: InputSource): Path = withContext(Dispatchers.IO) {
source as InputSource.Jellyfin
val url = URL(source.streamUrl)
val connection = url.openConnection() as HttpURLConnection
connection.requestMethod = "GET"
connection.instanceFollowRedirects = true
connection.connectTimeout = 30_000
connection.readTimeout = 120_000
var lastError: IOException? = null
for (attempt in 1..3) {
try {
if (connection.responseCode != HttpURLConnection.HTTP_OK) {
throw IllegalStateException("Invalid response code ${connection.responseCode}")
return@withContext downloadOnce(source)
} catch (e: IOException) {
lastError = e
logger.warn("Download attempt {}/3 failed: {}", attempt, e.message)
if (attempt < 3) {
delay(5_000)
}
}
}
throw lastError!!
}
private fun downloadOnce(source: InputSource.Jellyfin): Path {
val request = Request.Builder().url(source.streamUrl).get().build()
val response = client.newCall(request).execute()
response.use {
if (!response.isSuccessful) {
throw IllegalStateException("Invalid response code ${response.code}")
}
val resultFile = localStorageService.genFile()
logger.info("Downloading file in to ${resultFile}...")
try {
connection.inputStream.use { input ->
logger.info("Downloading file in to {}...", resultFile)
val path = try {
val body = response.body ?: throw IllegalStateException("Empty response body")
body.byteStream().use { input ->
Files.copy(input, resultFile, StandardCopyOption.REPLACE_EXISTING)
}
logger.info("File success downloaded!")
@@ -43,8 +68,7 @@ class JellyfinInput(
runCatching { Files.deleteIfExists(resultFile) }
throw e
}
} finally {
connection.disconnect()
return path
}
}
}
@@ -100,7 +100,7 @@ class JobExecutor(
output = videoFile,
width = newWidth,
height = newHeight,
videoCodec = VIDEO_CODEC,
videoCodec = props.videoCodec,
removeAudio = true,
removeSubtitles = true,
removeMetadata = true,
@@ -154,7 +154,6 @@ class JobExecutor(
}
companion object {
private const val VIDEO_CODEC = "libvpx-vp9"
private const val MAX_HEIGHT = 480
private const val VIDEO_CONTAINER = "mkv"
private const val AUDIO_CODEC = "libvorbis"
@@ -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
}
}
@@ -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>) {
val audioJson = json.encodeToString(audioKeysSerializer, audioKeys)
jdbcTemplate.update(
@@ -60,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),
@@ -99,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(),
@@ -141,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"),
)
}
}
@@ -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,
@@ -44,7 +58,16 @@ 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()
heartbeatScope.cancel()
}
}
+2
View File
@@ -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:}
@@ -24,6 +25,7 @@ app:
ffmpegPath: ffmpeg
ffprobePath: ffprobe
localFileStorage: /tmp/media-mirror-worker
videoCodec: h264_nvenc
server:
port: 8081
@@ -0,0 +1 @@
ALTER TABLE media_mirror.jobs ADD COLUMN IF NOT EXISTS retry_count int NOT NULL DEFAULT 0;
@@ -54,6 +54,7 @@ class JobPollerTest : AbstractMinioTest() {
registry.add("app.s3.prefix") { "mirror" }
registry.add("app.ffmpegPath") { "/usr/bin/ffmpeg" }
registry.add("app.ffprobePath") { "/usr/bin/ffprobe" }
registry.add("app.videoCodec") { "libx264" }
registry.add("app.localFileStorage") {
Files.createTempDirectory("mmw-worker-test").toString()
}
@@ -23,6 +23,7 @@ class JobExecutorTest : AbstractMinioTest() {
ffmpegPath = "/usr/bin/ffmpeg",
ffprobePath = "/usr/bin/ffprobe",
localFileStorage = tempDir,
videoCodec = "libx264",
)
private fun executor(): JobExecutor {
@@ -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(
@@ -146,23 +221,51 @@ 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",
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
}