Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 68eefff884 | |||
| 05a4aa775a |
@@ -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 = "",
|
||||||
|
|||||||
@@ -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()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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()
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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",
|
||||||
|
|||||||
Reference in New Issue
Block a user