Compare commits
8 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| d10759b26b | |||
| 96a9a0536d | |||
| 68eefff884 | |||
| 05a4aa775a | |||
| 235ae14693 | |||
| 72aa501a47 | |||
| 91642e4933 | |||
| 5dd324e79d |
+4
-12
@@ -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"))
|
||||||
@@ -65,10 +61,6 @@ dependencies {
|
|||||||
testImplementation("org.jetbrains.kotlinx:kotlinx-coroutines-test:1.10.2")
|
testImplementation("org.jetbrains.kotlinx:kotlinx-coroutines-test:1.10.2")
|
||||||
}
|
}
|
||||||
|
|
||||||
springBoot {
|
|
||||||
mainClass = "pw.binom.mirror.worker.WorkerApplication"
|
|
||||||
}
|
|
||||||
|
|
||||||
tasks.withType<Test> {
|
tasks.withType<Test> {
|
||||||
useJUnitPlatform()
|
useJUnitPlatform()
|
||||||
maxParallelForks = 1
|
maxParallelForks = 1
|
||||||
|
|||||||
+1
-1
@@ -17,7 +17,7 @@ image:
|
|||||||
db:
|
db:
|
||||||
host: null
|
host: null
|
||||||
port: 5432
|
port: 5432
|
||||||
name: glasses
|
name: media_mirror
|
||||||
user: null
|
user: null
|
||||||
password: null
|
password: null
|
||||||
maxConnections: 10
|
maxConnections: 10
|
||||||
|
|||||||
@@ -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"
|
try {
|
||||||
connection.instanceFollowRedirects = true
|
return@withContext downloadOnce(source)
|
||||||
connection.connectTimeout = 30_000
|
} catch (e: IOException) {
|
||||||
connection.readTimeout = 120_000
|
lastError = e
|
||||||
try {
|
logger.warn("Download attempt {}/3 failed: {}", attempt, e.message)
|
||||||
if (connection.responseCode != HttpURLConnection.HTTP_OK) {
|
if (attempt < 3) {
|
||||||
throw IllegalStateException("Invalid response code ${connection.responseCode}")
|
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,7 +2,7 @@ spring:
|
|||||||
application:
|
application:
|
||||||
name: media-mirror-worker
|
name: media-mirror-worker
|
||||||
datasource:
|
datasource:
|
||||||
url: jdbc:postgresql://192.168.76.106:5432/glasses
|
url: jdbc:postgresql://192.168.76.106:5432/media_mirror
|
||||||
username: postgres
|
username: postgres
|
||||||
password: postgres
|
password: postgres
|
||||||
flyway:
|
flyway:
|
||||||
@@ -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
|
||||||
|
|||||||
@@ -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",
|
||||||
|
|||||||
Reference in New Issue
Block a user