6 Commits
3 ... main

Author SHA1 Message Date
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
Hermes Agent 235ae14693 fix: FileInfo.Stream.Data — data-потоки (tmcd/субтитры) не ломают ffprobe-парсинг
Build Media Mirror Worker / Build and publish (release) Successful in 1m35s
- подкласс Stream.Data с @SerialName(data), Subtitle.startTime nullable
- JobExecutor не трогался: filterIsInstance игнорирует Data
- тесты: JSON-фикстура с data-потоком, реальный .mov с tmcd-треком
2026-08-11 00:00:39 +03:00
Hermes Agent 72aa501a47 fix: V1__init.sql приведён к api-версии (uq_jobs_active_item), jackson-exclude из webmvc
Build Media Mirror Worker / Build and publish (release) Successful in 1m36s
2026-08-10 23:43:11 +03:00
17 changed files with 245 additions and 37 deletions
+4 -8
View File
@@ -24,13 +24,6 @@ kotlin {
configurations { configurations {
named("runtimeClasspath") { named("runtimeClasspath") {
exclude(group = "tools.jackson.core")
exclude(group = "tools.jackson.databind")
exclude(group = "com.fasterxml.jackson.core")
exclude(group = "com.fasterxml.jackson.databind")
exclude(group = "com.fasterxml.jackson.datatype")
exclude(group = "com.fasterxml.jackson.dataformat")
exclude(group = "com.fasterxml.jackson.module")
} }
} }
@@ -45,7 +38,9 @@ configurations.all {
dependencies { dependencies {
implementation(platform("org.springframework.boot:spring-boot-dependencies:4.1.0")) implementation(platform("org.springframework.boot:spring-boot-dependencies:4.1.0"))
implementation("org.springframework.boot:spring-boot-starter-webmvc") implementation("org.springframework.boot:spring-boot-starter-webmvc") {
exclude(group = "org.springframework.boot", module = "spring-boot-starter-jackson")
}
implementation("org.springframework.boot:spring-boot-starter-jdbc") implementation("org.springframework.boot:spring-boot-starter-jdbc")
implementation("org.springframework.boot:spring-boot-starter-flyway") implementation("org.springframework.boot:spring-boot-starter-flyway")
implementation("org.flywaydb:flyway-database-postgresql") implementation("org.flywaydb:flyway-database-postgresql")
@@ -55,6 +50,7 @@ dependencies {
implementation("org.jetbrains.kotlinx:kotlinx-serialization-json:1.11.0") implementation("org.jetbrains.kotlinx:kotlinx-serialization-json:1.11.0")
implementation(platform("software.amazon.awssdk:bom:2.46.21")) implementation(platform("software.amazon.awssdk:bom:2.46.21"))
implementation("software.amazon.awssdk:s3") implementation("software.amazon.awssdk:s3")
implementation("com.squareup.okhttp3:okhttp:4.12.0")
testImplementation("org.springframework.boot:spring-boot-starter-test") testImplementation("org.springframework.boot:spring-boot-starter-test")
testImplementation(kotlin("test-junit5")) testImplementation(kotlin("test-junit5"))
@@ -8,12 +8,14 @@ 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",
val ffprobePath: String = "ffprobe", val ffprobePath: String = "ffprobe",
val localFileStorage: Path = Path.of("/tmp/media-mirror-worker"), val localFileStorage: Path = Path.of("/tmp/media-mirror-worker"),
val ffmpegThreadCount: Int? = null, val ffmpegThreadCount: Int? = null,
val videoCodec: String = "h264_nvenc",
) { ) {
data class Jellyfin( data class Jellyfin(
val url: String = "", val url: String = "",
@@ -82,7 +82,9 @@ class FfmpegService(
} }
if (width != null || height != null) { if (width != null || height != null) {
args += "-vf" 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 (!removeAudio) {
if (audioCodec != null) { if (audioCodec != null) {
@@ -37,7 +37,7 @@ data class FileInfo(
override val duration: Duration? = null, override val duration: Duration? = null,
@SerialName("start_time") @SerialName("start_time")
@Serializable(DurationAsSecondsString::class) @Serializable(DurationAsSecondsString::class)
override val startTime: Duration?, override val startTime: Duration? = null,
override val tags: Map<String, String> = emptyMap(), override val tags: Map<String, String> = emptyMap(),
) : Stream ) : Stream
@@ -95,6 +95,18 @@ data class FileInfo(
get() = tags["language"] ?: tags["lang"] get() = tags["language"] ?: tags["lang"]
} }
@Serializable
@SerialName("data")
data class Data(
override val index: Int,
@Serializable(DurationAsSecondsString::class)
override val duration: Duration? = null,
@SerialName("start_time")
@Serializable(DurationAsSecondsString::class)
override val startTime: Duration? = null,
override val tags: Map<String, String> = emptyMap(),
) : Stream
@Serializable @Serializable
data class Format( data class Format(
@SerialName("format_name") @SerialName("format_name")
@@ -28,6 +28,7 @@ class InputService(
logger.info("Download cancelled: {}", source) logger.info("Download cancelled: {}", source)
throw e throw e
} catch (e: Throwable) { } catch (e: Throwable) {
logger.warn("Download failed: {}", e.message, e)
throw IllegalStateException("Can't download video from $source", e) throw IllegalStateException("Can't download video from $source", e)
} }
logger.info("Downloaded success in {} ms", (System.nanoTime() - startTime) / 1_000_000) logger.info("Downloaded success in {} ms", (System.nanoTime() - startTime) / 1_000_000)
@@ -1,15 +1,18 @@
package pw.binom.mirror.worker.convert package pw.binom.mirror.worker.convert
import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.delay
import kotlinx.coroutines.withContext import kotlinx.coroutines.withContext
import okhttp3.OkHttpClient
import okhttp3.Request
import org.slf4j.Logger import org.slf4j.Logger
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.stereotype.Component import org.springframework.stereotype.Component
import java.net.HttpURLConnection import java.io.IOException
import java.net.URL
import java.nio.file.Files import java.nio.file.Files
import java.nio.file.Path import java.nio.file.Path
import java.nio.file.StandardCopyOption import java.nio.file.StandardCopyOption
import java.util.concurrent.TimeUnit
@Component @Component
class JellyfinInput( class JellyfinInput(
@@ -17,24 +20,46 @@ class JellyfinInput(
) : InputService.InputImplementation { ) : InputService.InputImplementation {
private val logger: Logger = LoggerFactory.getLogger(JellyfinInput::class.java) 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 fun isSupport(source: InputSource): Boolean = source is InputSource.Jellyfin
override suspend fun download(source: InputSource): Path = withContext(Dispatchers.IO) { override suspend fun download(source: InputSource): Path = withContext(Dispatchers.IO) {
source as InputSource.Jellyfin source as InputSource.Jellyfin
val url = URL(source.streamUrl) var lastError: IOException? = null
val connection = url.openConnection() as HttpURLConnection for (attempt in 1..3) {
connection.requestMethod = "GET"
connection.instanceFollowRedirects = true
connection.connectTimeout = 30_000
connection.readTimeout = 120_000
try { try {
if (connection.responseCode != HttpURLConnection.HTTP_OK) { return@withContext downloadOnce(source)
throw IllegalStateException("Invalid response code ${connection.responseCode}") } 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() val resultFile = localStorageService.genFile()
logger.info("Downloading file in to ${resultFile}...") logger.info("Downloading file in to {}...", resultFile)
try { val path = try {
connection.inputStream.use { input -> val body = response.body ?: throw IllegalStateException("Empty response body")
body.byteStream().use { input ->
Files.copy(input, resultFile, StandardCopyOption.REPLACE_EXISTING) Files.copy(input, resultFile, StandardCopyOption.REPLACE_EXISTING)
} }
logger.info("File success downloaded!") logger.info("File success downloaded!")
@@ -43,8 +68,7 @@ class JellyfinInput(
runCatching { Files.deleteIfExists(resultFile) } runCatching { Files.deleteIfExists(resultFile) }
throw e throw e
} }
} finally { return path
connection.disconnect()
} }
} }
} }
@@ -100,7 +100,7 @@ class JobExecutor(
output = videoFile, output = videoFile,
width = newWidth, width = newWidth,
height = newHeight, height = newHeight,
videoCodec = VIDEO_CODEC, videoCodec = props.videoCodec,
removeAudio = true, removeAudio = true,
removeSubtitles = true, removeSubtitles = true,
removeMetadata = true, removeMetadata = true,
@@ -154,7 +154,6 @@ class JobExecutor(
} }
companion object { companion object {
private const val VIDEO_CODEC = "libvpx-vp9"
private const val MAX_HEIGHT = 480 private const val MAX_HEIGHT = 480
private const val VIDEO_CONTAINER = "mkv" private const val VIDEO_CONTAINER = "mkv"
private const val AUDIO_CODEC = "libvorbis" private const val AUDIO_CODEC = "libvorbis"
@@ -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()
} }
} }
+2
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:}
@@ -24,6 +25,7 @@ app:
ffmpegPath: ffmpeg ffmpegPath: ffmpeg
ffprobePath: ffprobe ffprobePath: ffprobe
localFileStorage: /tmp/media-mirror-worker localFileStorage: /tmp/media-mirror-worker
videoCodec: h264_nvenc
server: server:
port: 8081 port: 8081
+10 -9
View File
@@ -1,18 +1,19 @@
CREATE SCHEMA IF NOT EXISTS media_mirror; CREATE SCHEMA IF NOT EXISTS media_mirror;
CREATE TABLE media_mirror.jobs (
CREATE TABLE IF NOT EXISTS media_mirror.jobs (
id uuid PRIMARY KEY DEFAULT gen_random_uuid(), id uuid PRIMARY KEY DEFAULT gen_random_uuid(),
item_id text NOT NULL, item_id text NOT NULL, -- itemId из Jellyfin
source_url text NOT NULL, source_url text NOT NULL, -- URL исходника (jellyfin stream)
source_type text NOT NULL DEFAULT 'jellyfin', source_type text NOT NULL DEFAULT 'jellyfin',
status text NOT NULL DEFAULT 'new', status text NOT NULL DEFAULT 'new', -- new | processing | done | failed | cancelled
progress int NOT NULL DEFAULT 0, progress int NOT NULL DEFAULT 0, -- 0..100
video_key text, video_key text, -- S3-ключ видео (после done)
audio_keys jsonb, audio_keys jsonb, -- [{key, title, language, index}]
error text, error text,
created_at timestamptz NOT NULL DEFAULT now(), created_at timestamptz NOT NULL DEFAULT now(),
updated_at timestamptz NOT NULL DEFAULT now(), updated_at timestamptz NOT NULL DEFAULT now(),
taken_at timestamptz taken_at timestamptz
); );
CREATE INDEX IF NOT EXISTS idx_jobs_status ON media_mirror.jobs (status, created_at); CREATE INDEX IF NOT EXISTS idx_jobs_status ON media_mirror.jobs (status, created_at);
CREATE UNIQUE INDEX IF NOT EXISTS uq_jobs_active_item
ON media_mirror.jobs (item_id)
WHERE status NOT IN ('failed', 'cancelled');
@@ -54,6 +54,7 @@ class JobPollerTest : AbstractMinioTest() {
registry.add("app.s3.prefix") { "mirror" } registry.add("app.s3.prefix") { "mirror" }
registry.add("app.ffmpegPath") { "/usr/bin/ffmpeg" } registry.add("app.ffmpegPath") { "/usr/bin/ffmpeg" }
registry.add("app.ffprobePath") { "/usr/bin/ffprobe" } registry.add("app.ffprobePath") { "/usr/bin/ffprobe" }
registry.add("app.videoCodec") { "libx264" }
registry.add("app.localFileStorage") { registry.add("app.localFileStorage") {
Files.createTempDirectory("mmw-worker-test").toString() Files.createTempDirectory("mmw-worker-test").toString()
} }
@@ -33,6 +33,23 @@ object TestMedia {
return out return out
} }
fun generateSourceWithDataStream(dir: Path, durationSeconds: Int = 2): Path {
val base = generateSource(dir, durationSeconds)
val out = dir.resolve("source-with-data.mov")
val command = listOf(
FFMPEG,
"-y",
"-i", base.toString(),
"-map", "0:v",
"-map", "0:a",
"-c", "copy",
"-timecode", "00:00:00:00",
out.toString(),
)
run(command)
return out
}
fun generateBrokenFile(dir: Path): Path { fun generateBrokenFile(dir: Path): Path {
Files.createDirectories(dir) Files.createDirectories(dir)
val out = dir.resolve("broken.bin") val out = dir.resolve("broken.bin")
@@ -43,6 +43,16 @@ class FfprobeServiceTest {
assertTrue(audio.channels >= 1) assertTrue(audio.channels >= 1)
} }
@Test
fun `getInfo ignores data streams and keeps audio and video`() {
val source = TestMedia.generateSourceWithDataStream(tempDir)
val info = service().getInfo(source)
assertTrue(info.streams.any { it is FileInfo.Data })
assertEquals(1, info.videoStreams.size)
assertEquals(2, info.audioStreams.size)
}
@Test @Test
fun `getInfo fails on broken file`() { fun `getInfo fails on broken file`() {
val broken = TestMedia.generateBrokenFile(tempDir) val broken = TestMedia.generateBrokenFile(tempDir)
@@ -4,6 +4,7 @@ import kotlinx.serialization.json.Json
import kotlin.time.Duration.Companion.seconds import kotlin.time.Duration.Companion.seconds
import kotlin.test.Test import kotlin.test.Test
import kotlin.test.assertEquals import kotlin.test.assertEquals
import kotlin.test.assertIs
import kotlin.test.assertNotNull import kotlin.test.assertNotNull
import kotlin.test.assertTrue import kotlin.test.assertTrue
@@ -92,6 +93,91 @@ class FileInfoTest {
assertEquals(3, info.videoStreams[0].index) assertEquals(3, info.videoStreams[0].index)
} }
@Test
fun `parses ffprobe json with data stream`() {
val fixture = """
{
"streams": [
{
"index": 0,
"codec_name": "h264",
"codec_type": "video",
"width": 1920,
"height": 1080,
"duration": "10.000000",
"start_time": "0.000000",
"tags": {"title": "Main"}
},
{
"index": 1,
"codec_name": "aac",
"codec_type": "audio",
"channels": 2,
"channel_layout": "stereo",
"duration": "10.000000",
"start_time": "0.000000",
"tags": {"language": "eng", "title": "English"}
},
{
"index": 2,
"codec_type": "data",
"codec_tag_string": "tmcd",
"tags": {"language": "eng", "handler_name": "SubtitleHandler"}
},
{
"index": 3,
"codec_name": "subrip",
"codec_type": "subtitle",
"start_time": "0.000000",
"tags": {"language": "rus"}
}
],
"format": {
"format_name": "matroska,webm",
"format_long_name": "Matroska / WebM",
"duration": "10.000000",
"size": "1000",
"bit_rate": "2000",
"tags": {"encoder": "x"}
}
}
""".trimIndent()
val info = json.decodeFromString(FileInfo.serializer(), fixture)
assertEquals(4, info.streams.size)
assertEquals(1, info.videoStreams.size)
assertEquals(1, info.audioStreams.size)
assertIs<FileInfo.VideoStream>(info.streams[0])
assertIs<FileInfo.AudioStream>(info.streams[1])
assertIs<FileInfo.Data>(info.streams[2])
assertIs<FileInfo.Subtitle>(info.streams[3])
assertEquals(1920, info.videoStreams.single().width)
assertEquals(2, info.audioStreams.single().channels)
assertEquals(10.seconds, info.format.duration)
}
@Test
fun `data stream round trip`() {
val info = FileInfo(
streams = listOf(
FileInfo.VideoStream(index = 0, codecName = "vp9", width = 1920, height = 1080, startTime = 0.seconds),
FileInfo.Data(index = 1, tags = mapOf("language" to "eng")),
),
format = FileInfo.Format(
formatName = "matroska",
formatLongName = "M",
duration = 5.seconds,
size = 1,
bitRate = 1,
),
)
val decoded = json.decodeFromString(FileInfo.serializer(), json.encodeToString(FileInfo.serializer(), info))
assertEquals(1, decoded.videoStreams.size)
assertEquals(0, decoded.audioStreams.size)
assertIs<FileInfo.Data>(decoded.streams[1])
}
@Test @Test
fun `audio stream title and language fallbacks`() { fun `audio stream title and language fallbacks`() {
val a1 = FileInfo.AudioStream( val a1 = FileInfo.AudioStream(
@@ -23,6 +23,7 @@ class JobExecutorTest : AbstractMinioTest() {
ffmpegPath = "/usr/bin/ffmpeg", ffmpegPath = "/usr/bin/ffmpeg",
ffprobePath = "/usr/bin/ffprobe", ffprobePath = "/usr/bin/ffprobe",
localFileStorage = tempDir, localFileStorage = tempDir,
videoCodec = "libx264",
) )
private fun executor(): JobExecutor { private fun executor(): JobExecutor {
@@ -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",