feat: media-mirror-worker — воркер конвертации зеркал (очередь Postgres, ядро из media-transformer)

- очередь: jobs в схеме media_mirror, SELECT ... FOR UPDATE SKIP LOCKED, stale-возврат
- ядро перенесено из media-transformer (FfmpegService, FfprobeService, JellyfinInput, S3Input, OutputService...)
- Spring Boot 4.1.0, Kotlin 2.4.10, kotlinx-serialization 1.11.0, NATS выпилен
- тесты: 45 (SKIP LOCKED параллельный, конец-в-конец до done, конвертация 480p+2 аудио), LINE 93%
This commit is contained in:
Hermes Agent
2026-08-10 17:29:42 +03:00
parent dd62744e02
commit 5ad865715e
46 changed files with 2851 additions and 0 deletions
+3
View File
@@ -0,0 +1,3 @@
.gradle/
build/
.kotlin/
+12
View File
@@ -0,0 +1,12 @@
FROM eclipse-temurin:21-jre
RUN apt-get update \
&& apt-get install -y --no-install-recommends ffmpeg \
&& rm -rf /var/lib/apt/lists/*
WORKDIR /app
COPY build/libs/*.jar app.jar
EXPOSE 8081
ENTRYPOINT ["java", "-jar", "/app/app.jar"]
+110
View File
@@ -0,0 +1,110 @@
import org.jetbrains.kotlin.gradle.dsl.JvmTarget
plugins {
kotlin("jvm") version "2.4.10"
kotlin("plugin.serialization") version "2.4.10"
id("org.springframework.boot") version "4.1.0"
jacoco
}
group = "pw.binom.mirror.worker"
version = "1.0.0-SNAPSHOT"
java {
sourceCompatibility = JavaVersion.VERSION_21
targetCompatibility = JavaVersion.VERSION_21
}
kotlin {
compilerOptions {
jvmTarget = JvmTarget.JVM_21
freeCompilerArgs.add("-Xjsr305=strict")
}
}
configurations {
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")
}
}
configurations.all {
resolutionStrategy {
force("org.testcontainers:testcontainers:1.21.4")
force("org.testcontainers:junit-jupiter:1.21.4")
force("org.testcontainers:postgresql:1.21.4")
force("org.testcontainers:minio:1.21.4")
}
}
dependencies {
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-jdbc")
implementation("org.springframework.boot:spring-boot-starter-flyway")
implementation("org.flywaydb:flyway-database-postgresql")
implementation("org.postgresql:postgresql")
implementation("org.jetbrains.kotlin:kotlin-reflect")
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-core:1.10.2")
implementation("org.jetbrains.kotlinx:kotlinx-serialization-json:1.11.0")
implementation(platform("software.amazon.awssdk:bom:2.46.21"))
implementation("software.amazon.awssdk:s3")
testImplementation("org.springframework.boot:spring-boot-starter-test")
testImplementation(kotlin("test-junit5"))
testImplementation("org.testcontainers:testcontainers:1.21.4")
testImplementation("org.testcontainers:junit-jupiter:1.21.4")
testImplementation("org.testcontainers:postgresql:1.21.4")
testImplementation("org.testcontainers:minio:1.21.4")
testImplementation("org.jetbrains.kotlinx:kotlinx-coroutines-test:1.10.2")
}
springBoot {
mainClass = "pw.binom.mirror.worker.WorkerApplication"
}
tasks.withType<Test> {
useJUnitPlatform()
maxParallelForks = 1
testLogging {
showStandardStreams = true
events("passed", "failed", "skipped")
}
}
tasks.withType<Jar> {
exclude("META-INF/*.SF", "META-INF/*.DSA", "META-INF/*.RSA")
}
tasks.jacocoTestReport {
dependsOn(tasks.test)
reports {
xml.required.set(true)
html.required.set(true)
}
}
tasks.jacocoTestCoverageVerification {
dependsOn(tasks.jacocoTestReport)
violationRules {
rule {
element = "PACKAGE"
includes = listOf("pw.binom.mirror.worker.convert")
limit {
counter = "LINE"
value = "COVEREDRATIO"
minimum = "0.80".toBigDecimal()
}
}
}
}
tasks.check {
dependsOn(tasks.jacocoTestReport, tasks.jacocoTestCoverageVerification)
}
+2
View File
@@ -0,0 +1,2 @@
org.gradle.jvmargs=-Xmx2g -Dfile.encoding=UTF-8
kotlin.code.style=official
Binary file not shown.
+6
View File
@@ -0,0 +1,6 @@
#Sat Apr 15 17:57:35 CST 2023
distributionBase=GRADLE_USER_HOME
distributionUrl=https\://services.gradle.org/distributions/gradle-8.14-bin.zip
distributionPath=wrapper/dists
zipStorePath=wrapper/dists
zipStoreBase=GRADLE_USER_HOME
Vendored Executable
+185
View File
@@ -0,0 +1,185 @@
#!/usr/bin/env sh
#
# Copyright 2015 the original author or authors.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# https://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
#
##############################################################################
##
## Gradle start up script for UN*X
##
##############################################################################
# Attempt to set APP_HOME
# Resolve links: $0 may be a link
PRG="$0"
# Need this for relative symlinks.
while [ -h "$PRG" ] ; do
ls=`ls -ld "$PRG"`
link=`expr "$ls" : '.*-> \(.*\)$'`
if expr "$link" : '/.*' > /dev/null; then
PRG="$link"
else
PRG=`dirname "$PRG"`"/$link"
fi
done
SAVED="`pwd`"
cd "`dirname \"$PRG\"`/" >/dev/null
APP_HOME="`pwd -P`"
cd "$SAVED" >/dev/null
APP_NAME="Gradle"
APP_BASE_NAME=`basename "$0"`
# Add default JVM options here. You can also use JAVA_OPTS and GRADLE_OPTS to pass JVM options to this script.
DEFAULT_JVM_OPTS='"-Xmx64m" "-Xms64m"'
# Use the maximum available, or set MAX_FD != -1 to use that value.
MAX_FD="maximum"
warn () {
echo "$*"
}
die () {
echo
echo "$*"
echo
exit 1
}
# OS specific support (must be 'true' or 'false').
cygwin=false
msys=false
darwin=false
nonstop=false
case "`uname`" in
CYGWIN* )
cygwin=true
;;
Darwin* )
darwin=true
;;
MINGW* )
msys=true
;;
NONSTOP* )
nonstop=true
;;
esac
CLASSPATH=$APP_HOME/gradle/wrapper/gradle-wrapper.jar
# Determine the Java command to use to start the JVM.
if [ -n "$JAVA_HOME" ] ; then
if [ -x "$JAVA_HOME/jre/sh/java" ] ; then
# IBM's JDK on AIX uses strange locations for the executables
JAVACMD="$JAVA_HOME/jre/sh/java"
else
JAVACMD="$JAVA_HOME/bin/java"
fi
if [ ! -x "$JAVACMD" ] ; then
die "ERROR: JAVA_HOME is set to an invalid directory: $JAVA_HOME
Please set the JAVA_HOME variable in your environment to match the
location of your Java installation."
fi
else
JAVACMD="java"
which java >/dev/null 2>&1 || die "ERROR: JAVA_HOME is not set and no 'java' command could be found in your PATH.
Please set the JAVA_HOME variable in your environment to match the
location of your Java installation."
fi
# Increase the maximum file descriptors if we can.
if [ "$cygwin" = "false" -a "$darwin" = "false" -a "$nonstop" = "false" ] ; then
MAX_FD_LIMIT=`ulimit -H -n`
if [ $? -eq 0 ] ; then
if [ "$MAX_FD" = "maximum" -o "$MAX_FD" = "max" ] ; then
MAX_FD="$MAX_FD_LIMIT"
fi
ulimit -n $MAX_FD
if [ $? -ne 0 ] ; then
warn "Could not set maximum file descriptor limit: $MAX_FD"
fi
else
warn "Could not query maximum file descriptor limit: $MAX_FD_LIMIT"
fi
fi
# For Darwin, add options to specify how the application appears in the dock
if $darwin; then
GRADLE_OPTS="$GRADLE_OPTS \"-Xdock:name=$APP_NAME\" \"-Xdock:icon=$APP_HOME/media/gradle.icns\""
fi
# For Cygwin or MSYS, switch paths to Windows format before running java
if [ "$cygwin" = "true" -o "$msys" = "true" ] ; then
APP_HOME=`cygpath --path --mixed "$APP_HOME"`
CLASSPATH=`cygpath --path --mixed "$CLASSPATH"`
JAVACMD=`cygpath --unix "$JAVACMD"`
# We build the pattern for arguments to be converted via cygpath
ROOTDIRSRAW=`find -L / -maxdepth 1 -mindepth 1 -type d 2>/dev/null`
SEP=""
for dir in $ROOTDIRSRAW ; do
ROOTDIRS="$ROOTDIRS$SEP$dir"
SEP="|"
done
OURCYGPATTERN="(^($ROOTDIRS))"
# Add a user-defined pattern to the cygpath arguments
if [ "$GRADLE_CYGPATTERN" != "" ] ; then
OURCYGPATTERN="$OURCYGPATTERN|($GRADLE_CYGPATTERN)"
fi
# Now convert the arguments - kludge to limit ourselves to /bin/sh
i=0
for arg in "$@" ; do
CHECK=`echo "$arg"|egrep -c "$OURCYGPATTERN" -`
CHECK2=`echo "$arg"|egrep -c "^-"` ### Determine if an option
if [ $CHECK -ne 0 ] && [ $CHECK2 -eq 0 ] ; then ### Added a condition
eval `echo args$i`=`cygpath --path --ignore --mixed "$arg"`
else
eval `echo args$i`="\"$arg\""
fi
i=`expr $i + 1`
done
case $i in
0) set -- ;;
1) set -- "$args0" ;;
2) set -- "$args0" "$args1" ;;
3) set -- "$args0" "$args1" "$args2" ;;
4) set -- "$args0" "$args1" "$args2" "$args3" ;;
5) set -- "$args0" "$args1" "$args2" "$args3" "$args4" ;;
6) set -- "$args0" "$args1" "$args2" "$args3" "$args4" "$args5" ;;
7) set -- "$args0" "$args1" "$args2" "$args3" "$args4" "$args5" "$args6" ;;
8) set -- "$args0" "$args1" "$args2" "$args3" "$args4" "$args5" "$args6" "$args7" ;;
9) set -- "$args0" "$args1" "$args2" "$args3" "$args4" "$args5" "$args6" "$args7" "$args8" ;;
esac
fi
# Escape application args
save () {
for i do printf %s\\n "$i" | sed "s/'/'\\\\''/g;1s/^/'/;\$s/\$/' \\\\/" ; done
echo " "
}
APP_ARGS=`save "$@"`
# Collect all arguments for the java command, following the shell quoting and substitution rules
eval set -- $DEFAULT_JVM_OPTS $JAVA_OPTS $GRADLE_OPTS "\"-Dorg.gradle.appname=$APP_BASE_NAME\"" -classpath "\"$CLASSPATH\"" org.gradle.wrapper.GradleWrapperMain "$APP_ARGS"
exec "$JAVACMD" "$@"
Vendored
+89
View File
@@ -0,0 +1,89 @@
@rem
@rem Copyright 2015 the original author or authors.
@rem
@rem Licensed under the Apache License, Version 2.0 (the "License");
@rem you may not use this file except in compliance with the License.
@rem You may obtain a copy of the License at
@rem
@rem https://www.apache.org/licenses/LICENSE-2.0
@rem
@rem Unless required by applicable law or agreed to in writing, software
@rem distributed under the License is distributed on an "AS IS" BASIS,
@rem WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
@rem See the License for the specific language governing permissions and
@rem limitations under the License.
@rem
@if "%DEBUG%" == "" @echo off
@rem ##########################################################################
@rem
@rem Gradle startup script for Windows
@rem
@rem ##########################################################################
@rem Set local scope for the variables with windows NT shell
if "%OS%"=="Windows_NT" setlocal
set DIRNAME=%~dp0
if "%DIRNAME%" == "" set DIRNAME=.
set APP_BASE_NAME=%~n0
set APP_HOME=%DIRNAME%
@rem Resolve any "." and ".." in APP_HOME to make it shorter.
for %%i in ("%APP_HOME%") do set APP_HOME=%%~fi
@rem Add default JVM options here. You can also use JAVA_OPTS and GRADLE_OPTS to pass JVM options to this script.
set DEFAULT_JVM_OPTS="-Xmx64m" "-Xms64m"
@rem Find java.exe
if defined JAVA_HOME goto findJavaFromJavaHome
set JAVA_EXE=java.exe
%JAVA_EXE% -version >NUL 2>&1
if "%ERRORLEVEL%" == "0" goto execute
echo.
echo ERROR: JAVA_HOME is not set and no 'java' command could be found in your PATH.
echo.
echo Please set the JAVA_HOME variable in your environment to match the
echo location of your Java installation.
goto fail
:findJavaFromJavaHome
set JAVA_HOME=%JAVA_HOME:"=%
set JAVA_EXE=%JAVA_HOME%/bin/java.exe
if exist "%JAVA_EXE%" goto execute
echo.
echo ERROR: JAVA_HOME is set to an invalid directory: %JAVA_HOME%
echo.
echo Please set the JAVA_HOME variable in your environment to match the
echo location of your Java installation.
goto fail
:execute
@rem Setup the command line
set CLASSPATH=%APP_HOME%\gradle\wrapper\gradle-wrapper.jar
@rem Execute Gradle
"%JAVA_EXE%" %DEFAULT_JVM_OPTS% %JAVA_OPTS% %GRADLE_OPTS% "-Dorg.gradle.appname=%APP_BASE_NAME%" -classpath "%CLASSPATH%" org.gradle.wrapper.GradleWrapperMain %*
:end
@rem End local scope for the variables with windows NT shell
if "%ERRORLEVEL%"=="0" goto mainEnd
:fail
rem Set variable GRADLE_EXIT_CONSOLE if you need the _script_ return code instead of
rem the _cmd.exe /c_ return code!
if not "" == "%GRADLE_EXIT_CONSOLE%" exit 1
exit /b 1
:mainEnd
if "%OS%"=="Windows_NT" endlocal
:omega
+15
View File
@@ -0,0 +1,15 @@
pluginManagement {
repositories {
mavenCentral()
gradlePluginPortal()
}
}
dependencyResolutionManagement {
repositoriesMode.set(RepositoriesMode.FAIL_ON_PROJECT_REPOS)
repositories {
mavenCentral()
}
}
rootProject.name = "media-mirror-worker"
@@ -0,0 +1,11 @@
package pw.binom.mirror.worker
import org.springframework.boot.autoconfigure.SpringBootApplication
import org.springframework.boot.runApplication
@SpringBootApplication
class WorkerApplication
fun main(args: Array<String>) {
runApplication<WorkerApplication>(*args)
}
@@ -0,0 +1,23 @@
package pw.binom.mirror.worker.config
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.serialization.json.Json
import org.springframework.boot.context.properties.EnableConfigurationProperties
import org.springframework.context.annotation.Bean
import org.springframework.context.annotation.Configuration
@Configuration
@EnableConfigurationProperties(AppProperties::class)
open class AppConfig {
@Bean
open fun json(): Json = Json {
ignoreUnknownKeys = true
encodeDefaults = false
}
@Bean
open fun jobScope(): CoroutineScope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
}
@@ -0,0 +1,31 @@
package pw.binom.mirror.worker.config
import org.springframework.boot.context.properties.ConfigurationProperties
import java.nio.file.Path
import java.time.Duration
@ConfigurationProperties(prefix = "app")
data class AppProperties(
val pollInterval: Duration = Duration.ofSeconds(2),
val staleTimeout: Duration = Duration.ofMinutes(30),
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,
) {
data class Jellyfin(
val url: String = "",
val apiKey: String = "",
)
data class S3(
val url: String = "",
val accessKey: String = "",
val secretKey: String = "",
val bucket: String = "media",
val region: String = "us-east-1",
val prefix: String = "mirror",
)
}
@@ -0,0 +1,24 @@
package pw.binom.mirror.worker.convert
import kotlinx.serialization.KSerializer
import kotlinx.serialization.descriptors.PrimitiveKind
import kotlinx.serialization.descriptors.PrimitiveSerialDescriptor
import kotlinx.serialization.descriptors.SerialDescriptor
import kotlinx.serialization.encoding.Decoder
import kotlinx.serialization.encoding.Encoder
import kotlin.time.Duration
import kotlin.time.Duration.Companion.seconds
object DurationAsSecondsString : KSerializer<Duration> {
override val descriptor: SerialDescriptor =
PrimitiveSerialDescriptor("DurationAsSecondsString", PrimitiveKind.STRING)
override fun serialize(encoder: Encoder, value: Duration) {
encoder.encodeString((value.inWholeMicroseconds / 1_000_000.0).toString())
}
override fun deserialize(decoder: Decoder): Duration {
val raw = decoder.decodeString().trim()
return raw.toDoubleOrNull()?.seconds ?: Duration.ZERO
}
}
@@ -0,0 +1,246 @@
package pw.binom.mirror.worker.convert
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.launch
import kotlinx.coroutines.suspendCancellableCoroutine
import org.springframework.stereotype.Service
import pw.binom.mirror.worker.config.AppProperties
import java.io.IOException
import java.nio.file.Files
import java.nio.file.Path
import kotlin.coroutines.resume
import kotlin.coroutines.resumeWithException
import kotlin.time.Duration
import kotlin.time.Duration.Companion.microseconds
@Service
class FfmpegService(
private val props: AppProperties,
private val jobScope: CoroutineScope,
) {
data class Progress(
val speed: Float,
val outTime: Duration,
val totalSize: Long,
val fps: Float?,
val frame: Long?,
)
suspend fun extractAudio(
input: Path,
output: Path,
audioIndex: Int,
channelCount: Int? = null,
codec: String? = null,
progress: (suspend (Progress) -> Unit)? = null,
) {
require(audioIndex >= 0) { "Invalid audioIndex: index should be greater or equals zero" }
require(channelCount == null || channelCount > 0) { "Invalid channelCount value: $channelCount" }
val args = ArrayList<String>()
args += "-map"
args += "0:a:$audioIndex"
if (channelCount != null) {
args += "-ac"
args += channelCount.toString()
}
if (codec != null) {
args += "-c:a"
args += codec
if (codec == "vorbis") {
args += "-strict"
args += "-2"
}
}
try {
run(input = input, output = output, args = args, progress = progress)
} catch (e: Exception) {
runCatching { Files.deleteIfExists(output) }
throw e
}
}
suspend fun transform(
input: Path,
output: Path,
width: Int? = null,
height: Int? = null,
videoCodec: String? = null,
audioCodec: String? = null,
channels: Int? = null,
removeAudio: Boolean = false,
removeSubtitles: Boolean = false,
removeMetadata: Boolean = false,
removeChapters: Boolean = false,
progress: (suspend (Progress) -> Unit)? = null,
) {
require(videoCodec != "copy") { "videoCodec can't be equals \"copy\"" }
val args = ArrayList<String>()
if (videoCodec != null) {
args += "-c:v"
args += videoCodec
}
if (width != null || height != null) {
args += "-vf"
args += "scale=${width ?: "-1"}:${height ?: "-1"}"
}
if (!removeAudio) {
if (audioCodec != null) {
args += "-c:a"
args += audioCodec
if (audioCodec == "vorbis") {
args += "-strict"
args += "-2"
}
}
if (channels != null) {
args += "-ac"
args += channels.toString()
}
} else {
args += "-an"
}
if (removeSubtitles) {
args += "-sn"
}
if (removeMetadata) {
args += "-map_metadata"
args += "-1"
}
if (removeChapters) {
args += "-map_chapters"
args += "-1"
}
try {
run(input = input, output = output, args = args, progress = progress)
} catch (e: Exception) {
runCatching { Files.deleteIfExists(output) }
throw e
}
}
private suspend fun run(
args: List<String>,
input: Path,
output: Path,
progress: (suspend (Progress) -> Unit)?,
) {
output.parent?.let { Files.createDirectories(it) }
val command = ArrayList<String>()
command += props.ffmpegPath
command += "-y"
command += "-i"
command += input.toString()
command += "-progress"
command += "pipe:1"
command += "-stats_period"
command += "1"
props.ffmpegThreadCount?.let {
command += "-threads"
command += it.toString()
}
command += args
command += output.toString()
val process = try {
ProcessBuilder(command).start()
} catch (e: IOException) {
throw IllegalStateException("Can't start ffmpeg ${props.ffmpegPath}", e)
}
val stdout = process.inputStream.bufferedReader()
val stderr = process.errorStream.bufferedReader()
val stderrText = StringBuilder()
suspendCancellableCoroutine<Unit> { cont ->
cont.invokeOnCancellation {
process.destroyForcibly()
}
val progressChannel = Channel<Progress>(Channel.UNLIMITED)
val progressJob = if (progress != null) {
jobScope.launch {
for (item in progressChannel) {
progress(item)
}
}
} else {
null
}
val readerThread = Thread {
runCatching {
val kv = HashMap<String, String>()
stdout.forEachLine { line ->
val idx = line.indexOf('=')
if (idx > 0) {
val key = line.substring(0, idx)
val value = line.substring(idx + 1)
kv[key] = value
if (key == "progress") {
val outTime = kv["out_time_us"]
?.takeIf { it != "N/A" }
?.toLongOrNull()
?.microseconds ?: Duration.ZERO
val speed = kv["speed"]
?.takeIf { it != "N/A" }
?.removeSuffix("x")
?.toFloatOrNull() ?: 0f
val fps = kv["fps"]
?.takeIf { it != "N/A" }
?.toFloatOrNull()
val frame = kv["frame"]
?.takeIf { it != "N/A" }
?.toLongOrNull()
val totalSize = kv["total_size"]
?.takeIf { it != "N/A" }
?.toLongOrNull() ?: 0L
progressChannel.trySend(
Progress(
speed = speed,
outTime = outTime,
totalSize = totalSize,
fps = fps,
frame = frame,
),
)
kv.clear()
}
}
}
}
progressChannel.close()
}
val errorThread = Thread {
runCatching {
stderr.forEachLine { line ->
stderrText.append(line).append('\n')
}
}
}
val waiterThread = Thread {
val exitCode = runCatching { process.waitFor() }.getOrDefault(-1)
runCatching { readerThread.join(1000) }
runCatching { errorThread.join(1000) }
progressJob?.cancel()
if (!cont.isCancelled) {
if (exitCode != 0) {
cont.resumeWithException(
IllegalStateException(
"ffmpeg process exited with $exitCode\nerror:\n$stderrText",
),
)
} else {
cont.resume(Unit)
}
}
}
readerThread.start()
errorThread.start()
waiterThread.start()
}
}
}
@@ -0,0 +1,42 @@
package pw.binom.mirror.worker.convert
import kotlinx.serialization.json.Json
import org.springframework.stereotype.Service
import pw.binom.mirror.worker.config.AppProperties
import java.io.IOException
import java.nio.file.Path
@Service
class FfprobeService(
private val props: AppProperties,
) {
private val json = Json {
ignoreUnknownKeys = true
}
fun getInfo(file: Path): FileInfo {
val command = listOf(
props.ffprobePath,
"-v",
"quiet",
"-print_format",
"json",
"-show_format",
"-show_streams",
file.toString(),
)
val process = try {
ProcessBuilder(command)
.redirectErrorStream(true)
.start()
} catch (e: IOException) {
throw IllegalStateException("Can't execute ffprobe ${props.ffprobePath}", e)
}
val output = process.inputStream.readBytes().decodeToString()
val exitCode = process.waitFor()
if (exitCode != 0) {
throw IllegalStateException("Can't execute ffprobe: $output")
}
return json.decodeFromString(FileInfo.serializer(), output)
}
}
@@ -0,0 +1,113 @@
package pw.binom.mirror.worker.convert
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
import kotlinx.serialization.builtins.LongAsStringSerializer
import kotlinx.serialization.json.JsonClassDiscriminator
import kotlin.time.Duration
@Serializable
data class FileInfo(
val streams: List<Stream>,
val format: Format,
) {
val videoStreams: List<VideoStream>
get() = streams.filterIsInstance<VideoStream>()
.filter { it.codecName != "png" }
.filter { it.codecName != "mjpeg" }
.filter { it.tags["filename"] != "cover.jpg" }
val audioStreams: List<AudioStream>
get() = streams.filterIsInstance<AudioStream>()
@Serializable
@JsonClassDiscriminator("codec_type")
sealed interface Stream {
val index: Int
val duration: Duration?
val startTime: Duration?
val tags: Map<String, String>
}
@Serializable
@SerialName("subtitle")
data class Subtitle(
override val index: Int,
@Serializable(DurationAsSecondsString::class)
override val duration: Duration? = null,
@SerialName("start_time")
@Serializable(DurationAsSecondsString::class)
override val startTime: Duration?,
override val tags: Map<String, String> = emptyMap(),
) : Stream
@Serializable
@SerialName("video")
data class VideoStream(
override val index: Int,
@SerialName("codec_name")
val codecName: String,
val width: Int,
val height: Int,
@SerialName("bit_rate")
@Serializable(LongAsStringSerializer::class)
val bitRate: Long? = null,
@Serializable(DurationAsSecondsString::class)
override val duration: Duration? = null,
@SerialName("start_time")
@Serializable(DurationAsSecondsString::class)
override val startTime: Duration,
override val tags: Map<String, String> = emptyMap(),
) : Stream
@Serializable
@SerialName("attachment")
data class Attachment(
override val index: Int,
@Serializable(DurationAsSecondsString::class)
override val duration: Duration?,
@Serializable(DurationAsSecondsString::class)
@SerialName("start_time")
override val startTime: Duration? = null,
override val tags: Map<String, String>,
) : Stream
@Serializable
@SerialName("audio")
data class AudioStream(
override val index: Int,
val channels: Int,
@SerialName("channel_layout")
val channelLayout: String? = null,
@SerialName("bit_rate")
@Serializable(LongAsStringSerializer::class)
val bitRate: Long? = null,
@Serializable(DurationAsSecondsString::class)
override val duration: Duration? = null,
@SerialName("start_time")
@Serializable(DurationAsSecondsString::class)
override val startTime: Duration,
override val tags: Map<String, String> = emptyMap(),
) : Stream {
val title: String
get() = tags["title"] ?: "#$index"
val language: String?
get() = tags["language"] ?: tags["lang"]
}
@Serializable
data class Format(
@SerialName("format_name")
val formatName: String,
@SerialName("format_long_name")
val formatLongName: String,
@Serializable(DurationAsSecondsString::class)
val duration: Duration,
@Serializable(LongAsStringSerializer::class)
val size: Long,
@SerialName("bit_rate")
@Serializable(LongAsStringSerializer::class)
val bitRate: Long,
val tags: Map<String, String> = emptyMap(),
)
}
@@ -0,0 +1,36 @@
package pw.binom.mirror.worker.convert
import kotlinx.coroutines.CancellationException
import org.slf4j.Logger
import org.slf4j.LoggerFactory
import org.springframework.stereotype.Service
import java.nio.file.Path
@Service
class InputService(
private val inputs: List<InputImplementation>,
) {
private val logger: Logger = LoggerFactory.getLogger(InputService::class.java)
interface InputImplementation {
fun isSupport(source: InputSource): Boolean
suspend fun download(source: InputSource): Path
}
suspend fun downloadIncome(source: InputSource): Path {
val implementation = inputs.find { it.isSupport(source) }
?: throw IllegalStateException("Can't find valid implementation for ${source::class.simpleName}")
logger.info("Downloading {}", source)
val startTime = System.nanoTime()
val resultFile = try {
implementation.download(source)
} catch (e: CancellationException) {
logger.info("Download cancelled: {}", source)
throw e
} catch (e: Throwable) {
throw IllegalStateException("Can't download video from $source", e)
}
logger.info("Downloaded success in {} ms", (System.nanoTime() - startTime) / 1_000_000)
return resultFile
}
}
@@ -0,0 +1,7 @@
package pw.binom.mirror.worker.convert
sealed interface InputSource {
data class Jellyfin(val streamUrl: String) : InputSource
data class S3(val destination: S3Destination, val key: String) : InputSource
}
@@ -0,0 +1,50 @@
package pw.binom.mirror.worker.convert
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
import org.slf4j.Logger
import org.slf4j.LoggerFactory
import org.springframework.stereotype.Component
import java.net.HttpURLConnection
import java.net.URL
import java.nio.file.Files
import java.nio.file.Path
import java.nio.file.StandardCopyOption
@Component
class JellyfinInput(
private val localStorageService: LocalStorageService,
) : InputService.InputImplementation {
private val logger: Logger = LoggerFactory.getLogger(JellyfinInput::class.java)
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
try {
if (connection.responseCode != HttpURLConnection.HTTP_OK) {
throw IllegalStateException("Invalid response code ${connection.responseCode}")
}
val resultFile = localStorageService.genFile()
logger.info("Downloading file in to ${resultFile}...")
try {
connection.inputStream.use { input ->
Files.copy(input, resultFile, StandardCopyOption.REPLACE_EXISTING)
}
logger.info("File success downloaded!")
resultFile
} catch (e: Throwable) {
runCatching { Files.deleteIfExists(resultFile) }
throw e
}
} finally {
connection.disconnect()
}
}
}
@@ -0,0 +1,163 @@
package pw.binom.mirror.worker.convert
import kotlinx.coroutines.CancellationException
import org.slf4j.Logger
import org.slf4j.LoggerFactory
import org.springframework.stereotype.Service
import pw.binom.mirror.worker.config.AppProperties
import pw.binom.mirror.worker.db.AudioKey
import pw.binom.mirror.worker.db.Job
import java.nio.file.Files
import java.nio.file.Path
import kotlin.time.Duration
import kotlin.time.Duration.Companion.seconds
@Service
class JobExecutor(
private val props: AppProperties,
private val inputService: InputService,
private val outputService: OutputService,
private val localStorageService: LocalStorageService,
private val ffmpegService: FfmpegService,
private val ffprobeService: FfprobeService,
) {
private val logger: Logger = LoggerFactory.getLogger(JobExecutor::class.java)
data class ConversionResult(
val videoKey: String,
val audio: List<AudioKey>,
val duration: Duration,
)
data class Progress(
val audioProgress: List<Duration>,
val videoProgress: Duration,
val totalDuration: Duration,
)
suspend fun execute(
job: Job,
destination: S3Destination,
progress: (suspend (Progress) -> Unit)? = null,
): ConversionResult {
logger.info("Start processing job {}", job.id)
val sourceFile = inputService.downloadIncome(InputSource.Jellyfin(streamUrl(job)))
val tmpFiles = ArrayList<Path>()
try {
val info = ffprobeService.getInfo(sourceFile)
if (info.videoStreams.size > 1) {
throw IllegalStateException("Video has more than one video streams")
}
if (info.videoStreams.isEmpty()) {
throw IllegalStateException("Video has no any video streams")
}
val videoStream = info.videoStreams.single()
val audioProgress = info.audioStreams.map { 0.seconds }.toMutableList()
val audioKeys = ArrayList<AudioKey>()
val audioFiles = ArrayList<Path>()
info.audioStreams.forEachIndexed { index, stream ->
val file = localStorageService.genFile(AUDIO_CONTAINER)
tmpFiles.add(file)
ffmpegService.extractAudio(
input = sourceFile,
output = file,
audioIndex = index,
codec = AUDIO_CODEC,
progress = { p ->
audioProgress[index] = p.outTime
progress?.invoke(
Progress(
audioProgress = audioProgress.toList(),
videoProgress = 0.seconds,
totalDuration = info.format.duration,
),
)
},
)
val audioName = "${job.itemId}/audio-$index.$AUDIO_CONTAINER"
audioKeys.add(
AudioKey(
key = destination.fullKey(audioName),
title = stream.title,
language = stream.language,
index = index,
),
)
audioFiles.add(file)
}
val videoFile = localStorageService.genFile(VIDEO_CONTAINER)
tmpFiles.add(videoFile)
val (newWidth, newHeight) = ViewSizeCalculate.calc(
width = videoStream.width,
height = videoStream.height,
maxWidth = null,
maxHeight = MAX_HEIGHT,
)
ffmpegService.transform(
input = sourceFile,
output = videoFile,
width = newWidth,
height = newHeight,
videoCodec = VIDEO_CODEC,
removeAudio = true,
removeSubtitles = true,
removeMetadata = true,
removeChapters = true,
progress = { p ->
progress?.invoke(
Progress(
audioProgress = audioProgress.toList(),
videoProgress = p.outTime,
totalDuration = info.format.duration,
),
)
},
)
val videoName = "${job.itemId}/video.$VIDEO_CONTAINER"
val videoKey = destination.fullKey(videoName)
logger.info("Converting finished. Start upload data")
logger.info("Audio Files: {}", audioKeys)
logger.info("Video File: {}", videoKey)
audioFiles.zip(audioKeys).forEach { (file, audioKey) ->
val audioName = "${job.itemId}/audio-${audioKey.index}.$AUDIO_CONTAINER"
logger.info("Upload as name {}", audioName)
outputService.pushResult(destination, file, audioName)
}
logger.info("Upload as name {}", videoName)
outputService.pushResult(destination, videoFile, videoName)
logger.info("Upload finished!")
return ConversionResult(
videoKey = videoKey,
audio = audioKeys,
duration = info.format.duration,
)
} catch (e: CancellationException) {
throw e
} catch (e: Throwable) {
logger.warn("Job {} finished with error: {}", job.id, e.message)
throw e
} finally {
tmpFiles.forEach { runCatching { Files.deleteIfExists(it) } }
runCatching { Files.deleteIfExists(sourceFile) }
}
}
private fun streamUrl(job: Job): String {
return job.sourceUrl.ifBlank {
props.jellyfin.url.trimEnd('/') +
"/Videos/${job.itemId}/stream?api_key=${props.jellyfin.apiKey}&static=true"
}
}
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"
private const val AUDIO_CONTAINER = "ogg"
}
}
@@ -0,0 +1,24 @@
package pw.binom.mirror.worker.convert
import org.springframework.stereotype.Service
import pw.binom.mirror.worker.config.AppProperties
import java.nio.file.Files
import java.nio.file.Path
import java.nio.file.Paths
import java.util.UUID
@Service
class LocalStorageService(
private val props: AppProperties,
) {
private val storageDirectory: Path = Paths.get(props.localFileStorage.toString()).toAbsolutePath()
fun genFile(ext: String? = null): Path {
Files.createDirectories(storageDirectory)
var name = UUID.randomUUID().toString()
if (ext != null) {
name = "$name.$ext"
}
return storageDirectory.resolve(name)
}
}
@@ -0,0 +1,31 @@
package pw.binom.mirror.worker.convert
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
import org.springframework.stereotype.Service
import software.amazon.awssdk.services.s3.model.PutObjectRequest
import java.nio.file.Files
import java.nio.file.Path
@Service
class OutputService {
suspend fun pushResult(destination: S3Destination, file: Path, name: String = file.fileName.toString()) {
check(Files.isRegularFile(file)) { "File $file is not found" }
withContext(Dispatchers.IO) {
val key = destination.fullKey(name)
val client = S3Clients.client(destination)
try {
client.putObject(
PutObjectRequest.builder()
.bucket(destination.bucket)
.key(key)
.build(),
file,
)
} finally {
client.close()
}
}
}
}
@@ -0,0 +1,26 @@
package pw.binom.mirror.worker.convert
import software.amazon.awssdk.auth.credentials.AwsBasicCredentials
import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider
import software.amazon.awssdk.regions.Region
import software.amazon.awssdk.services.s3.S3Client
import software.amazon.awssdk.services.s3.S3Configuration
import java.net.URI
object S3Clients {
fun client(config: S3Destination): S3Client =
S3Client.builder()
.endpointOverride(URI.create(config.url))
.region(Region.of(config.region))
.credentialsProvider(
StaticCredentialsProvider.create(
AwsBasicCredentials.create(config.accessKey, config.secretKey),
),
)
.serviceConfiguration(
S3Configuration.builder()
.pathStyleAccessEnabled(true)
.build(),
)
.build()
}
@@ -0,0 +1,13 @@
package pw.binom.mirror.worker.convert
data class S3Destination(
val url: String,
val accessKey: String,
val secretKey: String,
val bucket: String,
val region: String,
val prefix: String = "",
) {
fun fullKey(name: String): String =
if (prefix.isEmpty()) name else "${prefix.removeSuffix("/")}/$name"
}
@@ -0,0 +1,47 @@
package pw.binom.mirror.worker.convert
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
import org.springframework.stereotype.Component
import software.amazon.awssdk.core.exception.SdkServiceException
import software.amazon.awssdk.services.s3.model.GetObjectRequest
import java.nio.file.Files
import java.nio.file.Path
import java.nio.file.StandardCopyOption
@Component
class S3Input(
private val localStorageService: LocalStorageService,
) : InputService.InputImplementation {
override fun isSupport(source: InputSource): Boolean = source is InputSource.S3
override suspend fun download(source: InputSource): Path = withContext(Dispatchers.IO) {
source as InputSource.S3
val client = S3Clients.client(source.destination)
try {
val resultFile = localStorageService.genFile()
try {
client.getObject(
GetObjectRequest.builder()
.bucket(source.destination.bucket)
.key(source.key)
.build(),
).use { obj ->
Files.copy(obj, resultFile, StandardCopyOption.REPLACE_EXISTING)
}
resultFile
} catch (e: Throwable) {
runCatching { Files.deleteIfExists(resultFile) }
if (e is SdkServiceException && e.statusCode() == 404) {
throw IllegalStateException(
"File ${source.key} not found in ${source.destination.url}/${source.destination.bucket}",
e,
)
}
throw e
}
} finally {
client.close()
}
}
}
@@ -0,0 +1,20 @@
package pw.binom.mirror.worker.convert
object ViewSizeCalculate {
fun calc(width: Int, height: Int, maxWidth: Int?, maxHeight: Int?): Pair<Int, Int> {
val maxWidth = maxWidth ?: width
val maxHeight = maxHeight ?: height
if (width <= maxWidth && height <= maxHeight) {
return width to height
}
val widthScale = maxWidth.toDouble() / width
val heightScale = maxHeight.toDouble() / height
val scale = minOf(widthScale, heightScale)
val newWidth = (width * scale).toInt()
val newHeight = (height * scale).toInt()
return newWidth to newHeight
}
}
@@ -0,0 +1,37 @@
package pw.binom.mirror.worker.db
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable
import java.time.Instant
import java.util.UUID
@Serializable
data class AudioKey(
val key: String,
val title: String,
val language: String? = null,
val index: Int,
)
data class Job(
val id: UUID,
val itemId: String,
val sourceUrl: String,
val sourceType: String,
val status: String,
val progress: Int,
val videoKey: String? = null,
val audioKeys: List<AudioKey> = emptyList(),
val error: String? = null,
val createdAt: Instant,
val updatedAt: Instant,
val takenAt: Instant? = null,
) {
companion object {
const val STATUS_NEW = "new"
const val STATUS_PROCESSING = "processing"
const val STATUS_DONE = "done"
const val STATUS_FAILED = "failed"
const val STATUS_CANCELLED = "cancelled"
}
}
@@ -0,0 +1,146 @@
package pw.binom.mirror.worker.db
import kotlinx.serialization.builtins.ListSerializer
import kotlinx.serialization.json.Json
import org.springframework.jdbc.core.JdbcTemplate
import org.springframework.stereotype.Repository
import org.springframework.transaction.support.TransactionTemplate
import java.sql.ResultSet
import java.time.Duration
import java.time.Instant
import java.util.UUID
@Repository
open class JobRepository(
private val jdbcTemplate: JdbcTemplate,
private val transactionTemplate: TransactionTemplate,
private val json: Json,
) {
private val audioKeysSerializer = ListSerializer(AudioKey.serializer())
open fun takeNextJob(): Job? {
return transactionTemplate.execute {
val job = queryNewJob() ?: return@execute null
updateStatus(job.id, Job.STATUS_PROCESSING)
job.copy(
status = Job.STATUS_PROCESSING,
takenAt = Instant.now(),
)
}
}
open fun updateProgress(id: UUID, progress: Int) {
jdbcTemplate.update(
"""
UPDATE media_mirror.jobs
SET progress = ?, updated_at = now()
WHERE id = ?
""".trimIndent(),
progress,
id,
)
}
open fun markDone(id: UUID, videoKey: String, audioKeys: List<AudioKey>) {
val audioJson = json.encodeToString(audioKeysSerializer, audioKeys)
jdbcTemplate.update(
"""
UPDATE media_mirror.jobs
SET status = 'done', progress = 100, video_key = ?, audio_keys = CAST(? AS jsonb), error = NULL, updated_at = now()
WHERE id = ?
""".trimIndent(),
videoKey,
audioJson,
id,
)
}
open fun markFailed(id: UUID, error: String) {
jdbcTemplate.update(
"""
UPDATE media_mirror.jobs
SET status = 'failed', error = ?, updated_at = now()
WHERE id = ?
""".trimIndent(),
error.take(4000),
id,
)
}
open fun returnStaleToNew(staleTimeout: Duration): Int {
return jdbcTemplate.update(
"""
UPDATE media_mirror.jobs
SET status = 'new', taken_at = NULL, updated_at = now()
WHERE status = 'processing' AND taken_at < now() - (? || ' seconds')::interval
""".trimIndent(),
staleTimeout.seconds,
)
}
open fun findByStatus(status: String): List<Job> {
return jdbcTemplate.query(
"SELECT * FROM media_mirror.jobs WHERE status = ? ORDER BY created_at",
{ rs, _ -> rowToJob(rs) },
status,
)
}
open fun findById(id: UUID): Job? {
return jdbcTemplate.query(
"SELECT * FROM media_mirror.jobs WHERE id = ?",
{ rs, _ -> rowToJob(rs) },
id,
).firstOrNull()
}
private fun queryNewJob(): Job? {
val jobs = jdbcTemplate.query(
"""
SELECT * FROM media_mirror.jobs
WHERE status = 'new'
ORDER BY created_at
LIMIT 1
FOR UPDATE SKIP LOCKED
""".trimIndent(),
{ rs, _ -> rowToJob(rs) },
)
return jobs.firstOrNull()
}
private fun updateStatus(id: UUID, status: String) {
jdbcTemplate.update(
"""
UPDATE media_mirror.jobs
SET status = ?, taken_at = now(), updated_at = now()
WHERE id = ?
""".trimIndent(),
status,
id,
)
}
private fun rowToJob(rs: ResultSet): Job {
val audioKeysRaw = rs.getString("audio_keys")
val audioKeys = if (audioKeysRaw == null) {
emptyList()
} else {
runCatching { json.decodeFromString(audioKeysSerializer, audioKeysRaw) }.getOrElse { emptyList() }
}
return Job(
id = rs.getObject("id", UUID::class.java),
itemId = rs.getString("item_id"),
sourceUrl = rs.getString("source_url"),
sourceType = rs.getString("source_type"),
status = rs.getString("status"),
progress = rs.getInt("progress"),
videoKey = rs.getString("video_key"),
audioKeys = audioKeys,
error = rs.getString("error"),
createdAt = rs.getTimestamp("created_at").toInstant(),
updatedAt = rs.getTimestamp("updated_at").toInstant(),
takenAt = rs.getTimestamp("taken_at")?.toInstant(),
)
}
}
@@ -0,0 +1,67 @@
package pw.binom.mirror.worker.worker
import kotlinx.coroutines.CancellationException
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.delay
import kotlinx.coroutines.isActive
import kotlinx.coroutines.launch
import kotlinx.coroutines.withContext
import org.slf4j.Logger
import org.slf4j.LoggerFactory
import org.springframework.boot.ApplicationArguments
import org.springframework.boot.ApplicationRunner
import org.springframework.stereotype.Service
import pw.binom.mirror.worker.config.AppProperties
import pw.binom.mirror.worker.db.JobRepository
import kotlin.coroutines.coroutineContext
@Service
class JobPoller(
private val jobRepository: JobRepository,
private val jobProcessor: JobProcessor,
private val props: AppProperties,
private val jobScope: CoroutineScope,
) : ApplicationRunner {
private val logger: Logger = LoggerFactory.getLogger(JobPoller::class.java)
override fun run(args: ApplicationArguments) {
jobScope.launch { pollLoop() }
jobScope.launch { staleLoop() }
logger.info("Job poller started with interval {}", props.pollInterval)
}
private suspend fun pollLoop() {
while (coroutineContext.isActive) {
try {
val job = withContext(Dispatchers.IO) { jobRepository.takeNextJob() }
if (job != null) {
logger.info("Taken job {}", job.id)
jobProcessor.process(job)
}
} catch (e: CancellationException) {
throw e
} catch (e: Exception) {
logger.warn("Poll loop error", e)
} finally {
delay(props.pollInterval.toMillis())
}
}
}
private suspend fun staleLoop() {
while (coroutineContext.isActive) {
delay(props.pollInterval.toMillis())
try {
val count = jobRepository.returnStaleToNew(props.staleTimeout)
if (count > 0) {
logger.info("Returned {} stale jobs to new", count)
}
} catch (e: CancellationException) {
throw e
} catch (e: Exception) {
logger.warn("Stale check error", e)
}
}
}
}
@@ -0,0 +1,62 @@
package pw.binom.mirror.worker.worker
import kotlinx.coroutines.CancellationException
import org.slf4j.Logger
import org.slf4j.LoggerFactory
import org.springframework.stereotype.Service
import pw.binom.mirror.worker.config.AppProperties
import pw.binom.mirror.worker.convert.JobExecutor
import pw.binom.mirror.worker.convert.S3Destination
import pw.binom.mirror.worker.db.Job
import pw.binom.mirror.worker.db.JobRepository
import kotlin.time.Duration.Companion.seconds
@Service
class JobProcessor(
private val jobExecutor: JobExecutor,
private val jobRepository: JobRepository,
private val props: AppProperties,
) {
private val logger: Logger = LoggerFactory.getLogger(JobProcessor::class.java)
suspend fun process(job: Job) {
logger.info("Processing job {}", job.id)
val destination = S3Destination(
url = props.s3.url,
accessKey = props.s3.accessKey,
secretKey = props.s3.secretKey,
bucket = props.s3.bucket,
region = props.s3.region,
prefix = props.s3.prefix,
)
try {
val result = jobExecutor.execute(
job = job,
destination = destination,
progress = { p ->
jobRepository.updateProgress(job.id, calcPercent(p))
},
)
jobRepository.markDone(job.id, result.videoKey, result.audio)
logger.info("Job {} done, videoKey={}", job.id, result.videoKey)
} catch (e: CancellationException) {
logger.info("Job {} cancelled", job.id)
throw e
} catch (e: Throwable) {
logger.warn("Job {} finished with error: {}", job.id, e.message)
jobRepository.markFailed(job.id, e.message ?: e.toString())
}
}
private fun calcPercent(progress: JobExecutor.Progress): Int {
if (!progress.totalDuration.isPositive()) {
return 0
}
val maxProgress = maxOf(
progress.videoProgress.inWholeMilliseconds,
(progress.audioProgress.maxOrNull() ?: 0.seconds).inWholeMilliseconds,
)
val percent = maxProgress * 100 / progress.totalDuration.inWholeMilliseconds
return percent.coerceIn(0, 99).toInt()
}
}
+33
View File
@@ -0,0 +1,33 @@
spring:
application:
name: media-mirror-worker
datasource:
url: jdbc:postgresql://192.168.76.106:5432/glasses
username: postgres
password: postgres
flyway:
schemas: media_mirror
app:
poll-interval: 2s
stale-timeout: 30m
jellyfin:
url: https://jellyfin.binom.pw/
apiKey: ${JELLYFIN_API_KEY:}
s3:
url: https://s3.binom.pw
accessKey: ${S3_ACCESS_KEY:}
secretKey: ${S3_SECRET_KEY:}
bucket: media
region: us-east-1
prefix: mirror
ffmpegPath: ffmpeg
ffprobePath: ffprobe
localFileStorage: /tmp/media-mirror-worker
server:
port: 8081
logging:
level:
root: INFO
@@ -0,0 +1,18 @@
CREATE SCHEMA IF NOT EXISTS media_mirror;
CREATE TABLE IF NOT EXISTS media_mirror.jobs (
id uuid PRIMARY KEY DEFAULT gen_random_uuid(),
item_id text NOT NULL,
source_url text NOT NULL,
source_type text NOT NULL DEFAULT 'jellyfin',
status text NOT NULL DEFAULT 'new',
progress int NOT NULL DEFAULT 0,
video_key text,
audio_keys jsonb,
error text,
created_at timestamptz NOT NULL DEFAULT now(),
updated_at timestamptz NOT NULL DEFAULT now(),
taken_at timestamptz
);
CREATE INDEX IF NOT EXISTS idx_jobs_status ON media_mirror.jobs (status, created_at);
@@ -0,0 +1,60 @@
package pw.binom.mirror.worker
import org.junit.jupiter.api.io.TempDir
import org.testcontainers.containers.MinIOContainer
import org.testcontainers.junit.jupiter.Container
import org.testcontainers.junit.jupiter.Testcontainers
import pw.binom.mirror.worker.convert.S3Destination
import software.amazon.awssdk.auth.credentials.AwsBasicCredentials
import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider
import software.amazon.awssdk.regions.Region
import software.amazon.awssdk.services.s3.S3Client
import software.amazon.awssdk.services.s3.S3Configuration
import software.amazon.awssdk.services.s3.model.CreateBucketRequest
import java.net.URI
import java.nio.file.Path
@Testcontainers
abstract class AbstractMinioTest {
@TempDir
lateinit var tempDir: Path
companion object {
@JvmStatic
@Container
val minio: MinIOContainer = MinIOContainer("minio/minio:latest")
.withUserName("minioadmin")
.withPassword("minioadmin")
}
protected fun s3Destination(bucket: String = "media", prefix: String = "mirror"): S3Destination =
S3Destination(
url = minio.s3URL,
accessKey = minio.userName,
secretKey = minio.password,
bucket = bucket,
region = "us-east-1",
prefix = prefix,
)
protected fun s3Client(): S3Client =
S3Client.builder()
.endpointOverride(URI.create(minio.s3URL))
.region(Region.of("us-east-1"))
.credentialsProvider(
StaticCredentialsProvider.create(
AwsBasicCredentials.create(minio.userName, minio.password),
),
)
.serviceConfiguration(S3Configuration.builder().pathStyleAccessEnabled(true).build())
.build()
protected fun createBucket(bucket: String) {
s3Client().use { client ->
runCatching {
client.createBucket(CreateBucketRequest.builder().bucket(bucket).build())
}
}
}
}
@@ -0,0 +1,18 @@
package pw.binom.mirror.worker
import org.testcontainers.containers.PostgreSQLContainer
import org.testcontainers.junit.jupiter.Container
import org.testcontainers.junit.jupiter.Testcontainers
@Testcontainers
abstract class AbstractPostgresTest {
companion object {
@JvmStatic
@Container
val postgres: PostgreSQLContainer<*> = PostgreSQLContainer("postgres:16-alpine")
.withDatabaseName("glasses")
.withUsername("postgres")
.withPassword("postgres")
}
}
@@ -0,0 +1,127 @@
package pw.binom.mirror.worker
import com.sun.net.httpserver.HttpServer
import org.springframework.beans.factory.annotation.Autowired
import org.springframework.boot.test.context.SpringBootTest
import org.springframework.jdbc.core.JdbcTemplate
import org.springframework.test.context.DynamicPropertyRegistry
import org.springframework.test.context.DynamicPropertySource
import org.testcontainers.containers.PostgreSQLContainer
import org.testcontainers.junit.jupiter.Container
import org.testcontainers.junit.jupiter.Testcontainers
import pw.binom.mirror.worker.db.Job
import pw.binom.mirror.worker.db.JobRepository
import java.net.InetSocketAddress
import java.nio.file.Files
import java.nio.file.Path
import java.util.UUID
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
@Testcontainers
@SpringBootTest
class JobPollerTest : AbstractMinioTest() {
@Autowired
lateinit var jobRepository: JobRepository
@Autowired
lateinit var jdbcTemplate: JdbcTemplate
companion object {
@JvmStatic
@Container
val postgres: PostgreSQLContainer<*> = PostgreSQLContainer("postgres:16-alpine")
.withDatabaseName("glasses")
.withUsername("postgres")
.withPassword("postgres")
@JvmStatic
@DynamicPropertySource
fun properties(registry: DynamicPropertyRegistry) {
registry.add("spring.datasource.url") { postgres.jdbcUrl }
registry.add("spring.datasource.username") { postgres.username }
registry.add("spring.datasource.password") { postgres.password }
registry.add("app.poll-interval") { "200ms" }
registry.add("app.jellyfin.url") { "http://jellyfin.test" }
registry.add("app.jellyfin.apiKey") { "secret" }
registry.add("app.s3.url") { AbstractMinioTest.minio.s3URL }
registry.add("app.s3.accessKey") { AbstractMinioTest.minio.userName }
registry.add("app.s3.secretKey") { AbstractMinioTest.minio.password }
registry.add("app.s3.bucket") { "media" }
registry.add("app.s3.region") { "us-east-1" }
registry.add("app.s3.prefix") { "mirror" }
registry.add("app.ffmpegPath") { "/usr/bin/ffmpeg" }
registry.add("app.ffprobePath") { "/usr/bin/ffprobe" }
registry.add("app.localFileStorage") {
Files.createTempDirectory("mmw-worker-test").toString()
}
}
}
@Test
fun `poller picks up job and processes it to done`() {
createBucket("media")
val source = TestMedia.generateSource(tempDir)
val server = HttpServer.create(InetSocketAddress(0), 0)
server.createContext("/video.mp4") { exchange ->
exchange.sendResponseHeaders(200, Files.size(source))
exchange.responseBody.use { out -> Files.newInputStream(source).use { it.copyTo(out) } }
}
server.start()
try {
val id = insertJob(itemId = "item-1", sourceUrl = "http://127.0.0.1:${server.address.port}/video.mp4")
val done = await(120_000) {
jobRepository.findById(id)?.takeIf { it.status == Job.STATUS_DONE }
} ?: error("Timed out waiting for job to become done")
assertEquals("mirror/item-1/video.mkv", done.videoKey)
assertEquals(100, done.progress)
assertEquals(2, done.audioKeys.size)
assertEquals("mirror/item-1/audio-0.ogg", done.audioKeys[0].key)
assertEquals("mirror/item-1/audio-1.ogg", done.audioKeys[1].key)
} finally {
server.stop(0)
}
s3Client().use { client ->
val keys = client.listObjectsV2 { it.bucket("media") }.contents().map { it.key() }
assertTrue(keys.contains("mirror/item-1/video.mkv"))
assertTrue(keys.contains("mirror/item-1/audio-0.ogg"))
assertTrue(keys.contains("mirror/item-1/audio-1.ogg"))
}
}
@Test
fun `poller marks job as failed on broken source`() {
val id = insertJob(itemId = "broken", sourceUrl = "http://127.0.0.1:1/video.mp4")
val failed = await(120_000) {
jobRepository.findById(id)?.takeIf { it.status == Job.STATUS_FAILED }
} ?: error("Timed out waiting for job to become failed")
assertTrue(failed.error?.isNotBlank() == true)
}
private fun insertJob(itemId: String, sourceUrl: String): UUID {
val id = UUID.randomUUID()
jdbcTemplate.update(
"""
INSERT INTO media_mirror.jobs (id, item_id, source_url, status)
VALUES (?, ?, ?, 'new')
""".trimIndent(),
id,
itemId,
sourceUrl,
)
return id
}
private fun <T> await(timeoutMs: Long, condition: () -> T?): T? {
val deadline = System.currentTimeMillis() + timeoutMs
while (System.currentTimeMillis() < deadline) {
condition()?.let { return it }
Thread.sleep(100)
}
return null
}
}
@@ -0,0 +1,53 @@
package pw.binom.mirror.worker
import java.nio.file.Files
import java.nio.file.Path
object TestMedia {
private const val FFMPEG = "/usr/bin/ffmpeg"
fun generateSource(dir: Path, durationSeconds: Int = 2): Path {
Files.createDirectories(dir)
val out = dir.resolve("source.mp4")
val command = listOf(
FFMPEG,
"-y",
"-f", "lavfi",
"-i", "testsrc=duration=$durationSeconds:size=160x120:rate=10",
"-f", "lavfi",
"-i", "sine=frequency=440:duration=$durationSeconds",
"-f", "lavfi",
"-i", "sine=frequency=880:duration=$durationSeconds",
"-map", "0:v",
"-map", "1:a",
"-map", "2:a",
"-c:v", "libx264",
"-pix_fmt", "yuv420p",
"-c:a", "aac",
"-t", "$durationSeconds",
"-movflags", "+faststart",
out.toString(),
)
run(command)
return out
}
fun generateBrokenFile(dir: Path): Path {
Files.createDirectories(dir)
val out = dir.resolve("broken.bin")
Files.write(out, byteArrayOf(0x00, 0x01, 0x02, 0x03, 0x04))
return out
}
private fun run(command: List<String>) {
val process = ProcessBuilder(command)
.redirectErrorStream(true)
.start()
val output = process.inputStream.readBytes().decodeToString()
val exitCode = process.waitFor()
check(exitCode == 0) {
"Command failed with exit $exitCode:\n$command\n$output"
}
}
}
@@ -0,0 +1,28 @@
package pw.binom.mirror.worker.convert
import kotlinx.serialization.json.Json
import kotlin.time.Duration
import kotlin.time.Duration.Companion.seconds
import kotlin.test.Test
import kotlin.test.assertEquals
class DurationAsSecondsStringTest {
private val json = Json
@Test
fun `serializes duration to seconds string`() {
val encoded = json.encodeToString(DurationAsSecondsString, 12.5.seconds)
assertEquals("\"12.5\"", encoded)
}
@Test
fun `deserializes seconds string to duration`() {
assertEquals(12.5.seconds, json.decodeFromString(DurationAsSecondsString, "\"12.5\""))
}
@Test
fun `deserializes invalid value to zero`() {
assertEquals(Duration.ZERO, json.decodeFromString(DurationAsSecondsString, "\"N/A\""))
}
}
@@ -0,0 +1,136 @@
package pw.binom.mirror.worker.convert
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.runBlocking
import org.junit.jupiter.api.io.TempDir
import pw.binom.mirror.worker.TestMedia
import pw.binom.mirror.worker.config.AppProperties
import java.nio.file.Files
import java.nio.file.Path
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertTrue
class FfmpegServiceTest {
@TempDir
lateinit var tempDir: Path
private fun service() = FfmpegService(
props = AppProperties(ffmpegPath = "/usr/bin/ffmpeg", localFileStorage = tempDir),
jobScope = CoroutineScope(Dispatchers.Default),
)
@Test
fun `extractAudio extracts single audio track`() = runBlocking {
val source = TestMedia.generateSource(tempDir)
val output = tempDir.resolve("audio-0.m4a")
service().extractAudio(
input = source,
output = output,
audioIndex = 0,
codec = "aac",
channelCount = 2,
)
assertTrue(Files.size(output) > 0)
}
@Test
fun `extractAudio supports vorbis codec`() = runBlocking {
val source = TestMedia.generateSource(tempDir)
val output = tempDir.resolve("audio-0.ogg")
service().extractAudio(
input = source,
output = output,
audioIndex = 0,
codec = "vorbis",
channelCount = 2,
)
assertTrue(Files.size(output) > 0)
}
@Test
fun `extractAudio deletes output on failure`() = runBlocking {
val source = TestMedia.generateBrokenFile(tempDir)
val output = tempDir.resolve("bad.ogg")
runCatching {
service().extractAudio(
input = source,
output = output,
audioIndex = 0,
)
}
assertFalse(Files.exists(output))
}
@Test
fun `transform removes audio for video-only output`() = runBlocking {
val source = TestMedia.generateSource(tempDir)
val output = tempDir.resolve("video.mp4")
service().transform(
input = source,
output = output,
width = 160,
height = 120,
videoCodec = "libx264",
removeAudio = true,
removeSubtitles = true,
removeMetadata = true,
removeChapters = true,
)
assertTrue(Files.size(output) > 0)
val info = FfprobeService(AppProperties(ffprobePath = "/usr/bin/ffprobe")).getInfo(output)
assertEquals(1, info.videoStreams.size)
assertEquals(0, info.audioStreams.size)
}
@Test
fun `transform keeps audio for single channel output`() = runBlocking {
val source = TestMedia.generateSource(tempDir)
val output = tempDir.resolve("single.mp4")
service().transform(
input = source,
output = output,
width = 160,
height = 120,
videoCodec = "libx264",
audioCodec = "aac",
channels = 2,
)
assertTrue(Files.size(output) > 0)
val info = FfprobeService(AppProperties(ffprobePath = "/usr/bin/ffprobe")).getInfo(output)
assertEquals(1, info.videoStreams.size)
assertEquals(1, info.audioStreams.size)
}
@Test
fun `transform reports progress`() = runBlocking {
val source = TestMedia.generateSource(tempDir)
val output = tempDir.resolve("progress.mp4")
val progresses = mutableListOf<FfmpegService.Progress>()
service().transform(
input = source,
output = output,
videoCodec = "libx264",
removeAudio = true,
progress = { progresses += it },
)
assertTrue(progresses.isNotEmpty(), "Expected progress callbacks to be reported")
}
@Test
fun `transform rejects copy codec`() = runBlocking {
val source = TestMedia.generateSource(tempDir)
val output = tempDir.resolve("copy.mp4")
val result = runCatching {
service().transform(
input = source,
output = output,
videoCodec = "copy",
)
}
assertTrue(result.isFailure)
}
}
@@ -0,0 +1,52 @@
package pw.binom.mirror.worker.convert
import org.junit.jupiter.api.io.TempDir
import pw.binom.mirror.worker.TestMedia
import pw.binom.mirror.worker.config.AppProperties
import java.nio.file.Path
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNotNull
import kotlin.test.assertTrue
class FfprobeServiceTest {
@TempDir
lateinit var tempDir: Path
private fun service() = FfprobeService(AppProperties(ffprobePath = "/usr/bin/ffprobe"))
@Test
fun `getInfo returns video and audio streams`() {
val source = TestMedia.generateSource(tempDir)
val info = service().getInfo(source)
assertEquals(1, info.videoStreams.size)
assertEquals(2, info.audioStreams.size)
val video = info.videoStreams.single()
assertEquals(160, video.width)
assertEquals(120, video.height)
assertTrue(video.duration!!.isPositive())
assertTrue(info.format.duration.isPositive())
assertNotNull(info.format.formatName)
}
@Test
fun `getInfo parses audio stream metadata`() {
val source = TestMedia.generateSource(tempDir)
val info = service().getInfo(source)
val audio = info.audioStreams[0]
assertTrue(audio.index >= 0)
assertTrue(audio.channels >= 1)
}
@Test
fun `getInfo fails on broken file`() {
val broken = TestMedia.generateBrokenFile(tempDir)
val result = runCatching { service().getInfo(broken) }
assertTrue(result.isFailure)
}
}
@@ -0,0 +1,115 @@
package pw.binom.mirror.worker.convert
import kotlinx.serialization.json.Json
import kotlin.time.Duration.Companion.seconds
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertNotNull
import kotlin.test.assertTrue
class FileInfoTest {
private val json = Json { ignoreUnknownKeys = true }
@Test
fun `round trip with subtitle and attachment streams`() {
val info = FileInfo(
streams = listOf(
FileInfo.VideoStream(
index = 0,
codecName = "vp9",
width = 1920,
height = 1080,
bitRate = 1000,
duration = 10.seconds,
startTime = 0.seconds,
tags = mapOf("title" to "Main"),
),
FileInfo.AudioStream(
index = 1,
channels = 2,
channelLayout = "stereo",
bitRate = 500,
duration = 10.seconds,
startTime = 0.seconds,
tags = mapOf("title" to "English", "language" to "eng"),
),
FileInfo.Subtitle(
index = 2,
duration = 10.seconds,
startTime = 0.seconds,
tags = mapOf("language" to "rus"),
),
FileInfo.Attachment(
index = 3,
duration = null,
startTime = null,
tags = mapOf("filename" to "font.ttf"),
),
),
format = FileInfo.Format(
formatName = "matroska",
formatLongName = "Matroska",
duration = 10.seconds,
size = 1000,
bitRate = 2000,
tags = mapOf("encoder" to "x"),
),
)
val decoded = json.decodeFromString(FileInfo.serializer(), json.encodeToString(FileInfo.serializer(), info))
assertEquals(1, decoded.videoStreams.size)
assertEquals(1, decoded.audioStreams.size)
assertNotNull(decoded.format.formatName)
assertTrue(decoded.format.duration.isPositive())
}
@Test
fun `videoStreams filters png mjpeg and cover jpg`() {
val info = FileInfo(
streams = listOf(
FileInfo.VideoStream(index = 0, codecName = "png", width = 100, height = 100, startTime = 0.seconds),
FileInfo.VideoStream(index = 1, codecName = "mjpeg", width = 100, height = 100, startTime = 0.seconds),
FileInfo.VideoStream(
index = 2,
codecName = "vp9",
width = 1920,
height = 1080,
tags = mapOf("filename" to "cover.jpg"),
startTime = 0.seconds,
),
FileInfo.VideoStream(index = 3, codecName = "vp9", width = 1920, height = 1080, startTime = 0.seconds),
),
format = FileInfo.Format(
formatName = "matroska",
formatLongName = "M",
duration = 5.seconds,
size = 1,
bitRate = 1,
),
)
assertEquals(1, info.videoStreams.size)
assertEquals(3, info.videoStreams[0].index)
}
@Test
fun `audio stream title and language fallbacks`() {
val a1 = FileInfo.AudioStream(
index = 0,
channels = 2,
startTime = 0.seconds,
tags = mapOf("language" to "eng"),
)
assertEquals("#0", a1.title)
assertEquals("eng", a1.language)
val a2 = FileInfo.AudioStream(
index = 1,
channels = 2,
startTime = 0.seconds,
tags = emptyMap(),
)
assertEquals("#1", a2.title)
assertEquals(null, a2.language)
}
}
@@ -0,0 +1,84 @@
package pw.binom.mirror.worker.convert
import com.sun.net.httpserver.HttpServer
import kotlinx.coroutines.runBlocking
import org.junit.jupiter.api.Test
import pw.binom.mirror.worker.AbstractMinioTest
import pw.binom.mirror.worker.TestMedia
import pw.binom.mirror.worker.config.AppProperties
import software.amazon.awssdk.services.s3.model.PutObjectRequest
import java.net.InetSocketAddress
import java.nio.file.Files
import kotlin.test.assertContentEquals
import kotlin.test.assertEquals
import kotlin.test.assertTrue
class InputServiceTest : AbstractMinioTest() {
private val localStorage: LocalStorageService by lazy {
LocalStorageService(AppProperties(localFileStorage = tempDir))
}
@Test
fun `downloadIncome downloads s3 source`() = runBlocking {
createBucket("media")
val sourceFile = TestMedia.generateSource(tempDir)
val destination = s3Destination("media")
s3Client().use { client ->
client.putObject(
PutObjectRequest.builder().bucket("media").key("in/video.mp4").build(),
sourceFile,
)
}
val inputService = InputService(listOf(S3Input(localStorage)))
val result = inputService.downloadIncome(
InputSource.S3(destination = destination, key = "in/video.mp4"),
)
assertTrue(Files.exists(result))
assertEquals(Files.size(sourceFile), Files.size(result))
}
@Test
fun `downloadIncome fails for missing s3 key`() = runBlocking {
createBucket("media")
val inputService = InputService(listOf(S3Input(localStorage)))
val result = runCatching {
inputService.downloadIncome(
InputSource.S3(destination = s3Destination("media"), key = "missing.mp4"),
)
}
assertTrue(result.isFailure)
}
@Test
fun `jellyfin input downloads stream url`() = runBlocking {
val server = HttpServer.create(InetSocketAddress(0), 0)
val content = "jellyfin-video-content".toByteArray()
server.createContext("/video.mp4") { exchange ->
assertEquals("/video.mp4", exchange.requestURI.path)
exchange.sendResponseHeaders(200, content.size.toLong())
exchange.responseBody.use { it.write(content) }
}
server.start()
try {
val inputService = InputService(listOf(JellyfinInput(localStorage)))
val result = inputService.downloadIncome(
InputSource.Jellyfin(streamUrl = "http://127.0.0.1:${server.address.port}/video.mp4"),
)
assertContentEquals(content, Files.readAllBytes(result))
} finally {
server.stop(0)
}
}
@Test
fun `downloadIncome fails when no implementation supports source`() = runBlocking {
val inputService = InputService(emptyList())
val result = runCatching {
inputService.downloadIncome(
InputSource.Jellyfin(streamUrl = "http://localhost/x.mp4"),
)
}
assertTrue(result.isFailure)
}
}
@@ -0,0 +1,124 @@
package pw.binom.mirror.worker.convert
import com.sun.net.httpserver.HttpServer
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.runBlocking
import org.junit.jupiter.api.Test
import pw.binom.mirror.worker.AbstractMinioTest
import pw.binom.mirror.worker.TestMedia
import pw.binom.mirror.worker.config.AppProperties
import pw.binom.mirror.worker.db.Job
import java.net.InetSocketAddress
import java.nio.file.Files
import java.time.Instant
import java.util.UUID
import kotlin.test.assertEquals
import kotlin.test.assertTrue
class JobExecutorTest : AbstractMinioTest() {
private fun props() = AppProperties(
jellyfin = AppProperties.Jellyfin(url = "http://jellyfin.test", apiKey = "secret"),
ffmpegPath = "/usr/bin/ffmpeg",
ffprobePath = "/usr/bin/ffprobe",
localFileStorage = tempDir,
)
private fun executor(): JobExecutor {
val p = props()
return JobExecutor(
props = p,
inputService = InputService(listOf(JellyfinInput(LocalStorageService(p)))),
outputService = OutputService(),
localStorageService = LocalStorageService(p),
ffmpegService = FfmpegService(p, CoroutineScope(Dispatchers.Default)),
ffprobeService = FfprobeService(p),
)
}
@Test
fun `separate channels produces video and audio files in s3`() {
createBucket("media")
val source = TestMedia.generateSource(tempDir)
val server = HttpServer.create(InetSocketAddress(0), 0)
server.createContext("/video.mp4") { exchange ->
exchange.sendResponseHeaders(200, Files.size(source))
exchange.responseBody.use { out -> Files.newInputStream(source).use { it.copyTo(out) } }
}
server.start()
try {
val job = job(itemId = "item-1", sourceUrl = "http://127.0.0.1:${server.address.port}/video.mp4")
val result = runBlocking {
executor().execute(job = job, destination = s3Destination("media"))
}
assertEquals("mirror/item-1/video.mkv", result.videoKey)
assertEquals(2, result.audio.size)
assertEquals("mirror/item-1/audio-0.ogg", result.audio[0].key)
assertEquals("mirror/item-1/audio-1.ogg", result.audio[1].key)
assertTrue(result.duration.isPositive())
s3Client().use { client ->
val keys = client.listObjectsV2 { it.bucket("media") }.contents().map { it.key() }
assertTrue(keys.contains("mirror/item-1/video.mkv"))
assertTrue(keys.contains("mirror/item-1/audio-0.ogg"))
assertTrue(keys.contains("mirror/item-1/audio-1.ogg"))
}
} finally {
server.stop(0)
}
}
@Test
fun `downloaded video has no audio and fits 480p height`() {
createBucket("media")
val source = TestMedia.generateSource(tempDir)
val server = HttpServer.create(InetSocketAddress(0), 0)
server.createContext("/video.mp4") { exchange ->
exchange.sendResponseHeaders(200, Files.size(source))
exchange.responseBody.use { out -> Files.newInputStream(source).use { it.copyTo(out) } }
}
server.start()
try {
val job = job(itemId = "item-2", sourceUrl = "http://127.0.0.1:${server.address.port}/video.mp4")
runBlocking {
executor().execute(job = job, destination = s3Destination("media"))
}
val probe = FfprobeService(props())
s3Client().use { client ->
client.getObject { it.bucket("media").key("mirror/item-2/video.mkv") }.use { obj ->
val videoFile = tempDir.resolve("video-download.mkv")
Files.copy(obj, videoFile)
val info = probe.getInfo(videoFile)
assertEquals(1, info.videoStreams.size)
assertTrue(info.audioStreams.isEmpty(), "video should have no audio streams")
assertTrue(info.videoStreams.single().height <= 480)
}
}
} finally {
server.stop(0)
}
}
@Test
fun `broken source fails with exception`() {
val job = job(itemId = "broken", sourceUrl = "http://127.0.0.1:1/video.mp4")
val result = runCatching {
runBlocking { executor().execute(job = job, destination = s3Destination("media")) }
}
assertTrue(result.isFailure)
}
private fun job(itemId: String, sourceUrl: String): Job = Job(
id = UUID.randomUUID(),
itemId = itemId,
sourceUrl = sourceUrl,
sourceType = "jellyfin",
status = "new",
progress = 0,
createdAt = Instant.now(),
updatedAt = Instant.now(),
)
}
@@ -0,0 +1,38 @@
package pw.binom.mirror.worker.convert
import org.junit.jupiter.api.io.TempDir
import pw.binom.mirror.worker.config.AppProperties
import java.nio.file.Files
import java.nio.file.Path
import kotlin.test.Test
import kotlin.test.assertFalse
import kotlin.test.assertTrue
class LocalStorageServiceTest {
@TempDir
lateinit var tempDir: Path
@Test
fun `genFile without extension`() {
val service = LocalStorageService(AppProperties(localFileStorage = tempDir))
val file = service.genFile()
assertTrue(file.fileName.toString().isNotBlank())
assertFalse(Files.exists(file))
}
@Test
fun `genFile with extension`() {
val service = LocalStorageService(AppProperties(localFileStorage = tempDir))
val file = service.genFile("mp4")
assertTrue(file.fileName.toString().endsWith(".mp4"))
}
@Test
fun `genFile creates storage directory`() {
val nested = tempDir.resolve("nested").resolve("dir")
val service = LocalStorageService(AppProperties(localFileStorage = nested))
service.genFile()
assertTrue(nested.toFile().isDirectory)
}
}
@@ -0,0 +1,107 @@
package pw.binom.mirror.worker.convert
import kotlinx.coroutines.runBlocking
import org.junit.jupiter.api.Test
import pw.binom.mirror.worker.AbstractMinioTest
import software.amazon.awssdk.services.s3.model.GetObjectRequest
import java.nio.file.Files
import kotlin.test.assertEquals
import kotlin.test.assertTrue
class OutputServiceTest : AbstractMinioTest() {
@Test
fun `pushResult uploads file under prefix`() = runBlocking {
createBucket("media")
val destination = s3Destination("media", prefix = "out")
val file = tempDir.resolve("result.mp4")
Files.write(file, ByteArray(1024) { 1 })
OutputService().pushResult(
destination = destination,
file = file,
name = "video.mp4",
)
s3Client().use { client ->
val obj = client.getObject(
GetObjectRequest.builder().bucket("media").key("out/video.mp4").build(),
)
val bytes = obj.readAllBytes()
assertEquals(Files.size(file), bytes.size.toLong())
}
}
@Test
fun `pushResult uploads without prefix when empty`() = runBlocking {
createBucket("media")
val destination = s3Destination("media", prefix = "")
val file = tempDir.resolve("plain.mp4")
Files.write(file, byteArrayOf(1, 2, 3))
OutputService().pushResult(
destination = destination,
file = file,
name = "plain.mp4",
)
s3Client().use { client ->
val obj = client.getObject(
GetObjectRequest.builder().bucket("media").key("plain.mp4").build(),
)
assertEquals(3L, obj.readAllBytes().size.toLong())
}
}
@Test
fun `pushResult trims trailing slash from prefix`() = runBlocking {
createBucket("media")
val destination = s3Destination("media", prefix = "out/")
val file = tempDir.resolve("x.mp4")
Files.write(file, byteArrayOf(9))
OutputService().pushResult(
destination = destination,
file = file,
name = "x.mp4",
)
s3Client().use { client ->
val keys = client.listObjectsV2 { it.bucket("media") }.contents().map { it.key() }
assertTrue(keys.contains("out/x.mp4"))
}
}
@Test
fun `pushResult defaults name to file name`() = runBlocking {
createBucket("media")
val destination = s3Destination("media", prefix = "out")
val file = tempDir.resolve("default.mp4")
Files.write(file, byteArrayOf(5))
OutputService().pushResult(
destination = destination,
file = file,
)
s3Client().use { client ->
val keys = client.listObjectsV2 { it.bucket("media") }.contents().map { it.key() }
assertTrue(keys.contains("out/default.mp4"))
}
}
@Test
fun `pushResult fails for missing file`() = runBlocking {
createBucket("media")
val destination = s3Destination("media", prefix = "")
val missing = tempDir.resolve("missing.mp4")
val result = runCatching {
OutputService().pushResult(
destination = destination,
file = missing,
name = "missing.mp4",
)
}
assertTrue(result.isFailure)
}
}
@@ -0,0 +1,48 @@
package pw.binom.mirror.worker.convert
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
class ViewSizeCalculateTest {
@Test
fun `no limits keeps original size`() {
assertEquals(1920 to 1080, ViewSizeCalculate.calc(1920, 1080, null, null))
}
@Test
fun `scales down preserving aspect ratio`() {
val (w, h) = ViewSizeCalculate.calc(1920, 1080, 1280, 720)
assertTrue(w <= 1280)
assertTrue(h <= 720)
assertEquals(1280, w)
assertEquals(720, h)
}
@Test
fun `width limit only`() {
val (w, h) = ViewSizeCalculate.calc(1920, 1080, 640, null)
assertEquals(640, w)
assertEquals(360, h)
}
@Test
fun `height limit only`() {
val (w, h) = ViewSizeCalculate.calc(1920, 1080, null, 540)
assertEquals(960, w)
assertEquals(540, h)
}
@Test
fun `smaller than limits stays the same`() {
assertEquals(640 to 360, ViewSizeCalculate.calc(640, 360, 1920, 1080))
}
@Test
fun `portrait scaling`() {
val (w, h) = ViewSizeCalculate.calc(1080, 1920, 720, 1280)
assertEquals(720, w)
assertEquals(1280, h)
}
}
@@ -0,0 +1,169 @@
package pw.binom.mirror.worker.db
import com.zaxxer.hikari.HikariConfig
import com.zaxxer.hikari.HikariDataSource
import kotlinx.serialization.json.Json
import org.flywaydb.core.Flyway
import org.junit.jupiter.api.AfterEach
import org.junit.jupiter.api.BeforeEach
import org.junit.jupiter.api.Test
import org.springframework.jdbc.core.JdbcTemplate
import org.springframework.jdbc.datasource.DataSourceTransactionManager
import org.springframework.transaction.support.TransactionTemplate
import pw.binom.mirror.worker.AbstractPostgresTest
import java.time.Duration
import java.time.Instant
import java.time.temporal.ChronoUnit
import java.util.UUID
import java.util.concurrent.ConcurrentLinkedQueue
import java.util.concurrent.CountDownLatch
import kotlin.test.assertEquals
import kotlin.test.assertNotNull
import kotlin.test.assertNull
import kotlin.test.assertTrue
class JobRepositoryTest : AbstractPostgresTest() {
private lateinit var dataSource: HikariDataSource
private lateinit var repository: JobRepository
@BeforeEach
fun setup() {
dataSource = HikariDataSource(
HikariConfig().apply {
jdbcUrl = postgres.jdbcUrl
username = postgres.username
password = postgres.password
maximumPoolSize = 5
},
)
Flyway.configure()
.dataSource(dataSource)
.schemas("media_mirror")
.load()
.migrate()
JdbcTemplate(dataSource).execute("TRUNCATE TABLE media_mirror.jobs")
repository = JobRepository(
jdbcTemplate = JdbcTemplate(dataSource),
transactionTemplate = TransactionTemplate(DataSourceTransactionManager(dataSource)),
json = Json { ignoreUnknownKeys = true },
)
}
@AfterEach
fun teardown() {
dataSource.close()
}
@Test
fun `takeNextJob marks job as processing with takenAt`() {
val id = insertJob("item-1")
val taken = repository.takeNextJob()
assertNotNull(taken)
assertEquals(id, taken.id)
assertEquals("processing", taken.status)
assertNotNull(taken.takenAt)
assertEquals(1, repository.findByStatus("processing").size)
assertTrue(repository.findByStatus("new").isEmpty())
}
@Test
fun `skip locked gives a job to only one of two concurrent workers`() {
insertJob("item-1")
val startGate = CountDownLatch(2)
val taken = ConcurrentLinkedQueue<UUID>()
val attempts = java.util.concurrent.atomic.AtomicInteger()
val threads = (1..2).map {
Thread {
startGate.countDown()
startGate.await()
attempts.incrementAndGet()
repository.takeNextJob()?.id?.let { taken.add(it) }
}
}
threads.forEach { it.start() }
threads.forEach { it.join(30_000) }
assertEquals(2, attempts.get())
assertEquals(1, taken.size)
assertEquals(1, repository.findByStatus("processing").size)
assertTrue(repository.findByStatus("new").isEmpty())
}
@Test
fun `markDone stores video key and audio keys`() {
val id = insertJob("item-1")
repository.takeNextJob()
val audioKeys = listOf(
AudioKey(key = "mirror/item-1/audio-0.ogg", title = "Track 1", language = "eng", index = 0),
AudioKey(key = "mirror/item-1/audio-1.ogg", title = "#1", language = null, index = 1),
)
repository.markDone(id, "mirror/item-1/video.mkv", audioKeys)
val job = repository.findByStatus("done").single()
assertEquals(id, job.id)
assertEquals("mirror/item-1/video.mkv", job.videoKey)
assertEquals(audioKeys, job.audioKeys)
assertEquals(100, job.progress)
}
@Test
fun `markFailed stores error message`() {
val id = insertJob("item-1")
repository.takeNextJob()
repository.markFailed(id, "boom")
val job = repository.findByStatus("failed").single()
assertEquals("boom", job.error)
}
@Test
fun `stale processing job returns to new`() {
val id = insertJob(
itemId = "item-1",
status = "processing",
takenAt = Instant.now().minus(31, ChronoUnit.MINUTES),
)
val count = repository.returnStaleToNew(Duration.ofMinutes(30))
assertEquals(1, count)
val job = repository.findByStatus("new").single()
assertEquals(id, job.id)
assertNull(job.takenAt)
}
@Test
fun `fresh processing job is not stale`() {
insertJob(itemId = "item-1", status = "processing", takenAt = Instant.now())
assertEquals(0, repository.returnStaleToNew(Duration.ofMinutes(30)))
assertEquals(1, repository.findByStatus("processing").size)
}
private fun insertJob(
itemId: String,
status: String = "new",
sourceUrl: String = "http://localhost/video.mp4",
takenAt: Instant? = null,
): UUID {
val id = UUID.randomUUID()
JdbcTemplate(dataSource).update(
"""
INSERT INTO media_mirror.jobs (id, item_id, source_url, status, taken_at)
VALUES (?, ?, ?, ?, ?)
""".trimIndent(),
id,
itemId,
sourceUrl,
status,
takenAt?.let { java.sql.Timestamp.from(it) },
)
return id
}
}