diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..9f5f768 --- /dev/null +++ b/.gitignore @@ -0,0 +1,3 @@ +.gradle/ +build/ +.kotlin/ diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..7d0b70e --- /dev/null +++ b/Dockerfile @@ -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"] diff --git a/build.gradle.kts b/build.gradle.kts new file mode 100644 index 0000000..a09c900 --- /dev/null +++ b/build.gradle.kts @@ -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 { + useJUnitPlatform() + maxParallelForks = 1 + testLogging { + showStandardStreams = true + events("passed", "failed", "skipped") + } +} + +tasks.withType { + 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) +} diff --git a/gradle.properties b/gradle.properties new file mode 100644 index 0000000..265ded2 --- /dev/null +++ b/gradle.properties @@ -0,0 +1,2 @@ +org.gradle.jvmargs=-Xmx2g -Dfile.encoding=UTF-8 +kotlin.code.style=official diff --git a/gradle/wrapper/gradle-wrapper.jar b/gradle/wrapper/gradle-wrapper.jar new file mode 100644 index 0000000..e708b1c Binary files /dev/null and b/gradle/wrapper/gradle-wrapper.jar differ diff --git a/gradle/wrapper/gradle-wrapper.properties b/gradle/wrapper/gradle-wrapper.properties new file mode 100644 index 0000000..c1f93d8 --- /dev/null +++ b/gradle/wrapper/gradle-wrapper.properties @@ -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 diff --git a/gradlew b/gradlew new file mode 100755 index 0000000..4f906e0 --- /dev/null +++ b/gradlew @@ -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" "$@" diff --git a/gradlew.bat b/gradlew.bat new file mode 100644 index 0000000..107acd3 --- /dev/null +++ b/gradlew.bat @@ -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 diff --git a/settings.gradle.kts b/settings.gradle.kts new file mode 100644 index 0000000..60171e6 --- /dev/null +++ b/settings.gradle.kts @@ -0,0 +1,15 @@ +pluginManagement { + repositories { + mavenCentral() + gradlePluginPortal() + } +} + +dependencyResolutionManagement { + repositoriesMode.set(RepositoriesMode.FAIL_ON_PROJECT_REPOS) + repositories { + mavenCentral() + } +} + +rootProject.name = "media-mirror-worker" diff --git a/src/main/kotlin/pw/binom/mirror/worker/WorkerApplication.kt b/src/main/kotlin/pw/binom/mirror/worker/WorkerApplication.kt new file mode 100644 index 0000000..e6246a7 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/WorkerApplication.kt @@ -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) { + runApplication(*args) +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/config/AppConfig.kt b/src/main/kotlin/pw/binom/mirror/worker/config/AppConfig.kt new file mode 100644 index 0000000..8b7c3b6 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/config/AppConfig.kt @@ -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) +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/config/AppProperties.kt b/src/main/kotlin/pw/binom/mirror/worker/config/AppProperties.kt new file mode 100644 index 0000000..5e72997 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/config/AppProperties.kt @@ -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", + ) +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/convert/DurationAsSecondsString.kt b/src/main/kotlin/pw/binom/mirror/worker/convert/DurationAsSecondsString.kt new file mode 100644 index 0000000..272ac6c --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/convert/DurationAsSecondsString.kt @@ -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 { + 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 + } +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/convert/FfmpegService.kt b/src/main/kotlin/pw/binom/mirror/worker/convert/FfmpegService.kt new file mode 100644 index 0000000..991bbf4 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/convert/FfmpegService.kt @@ -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() + 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() + 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, + input: Path, + output: Path, + progress: (suspend (Progress) -> Unit)?, + ) { + output.parent?.let { Files.createDirectories(it) } + val command = ArrayList() + 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 { cont -> + cont.invokeOnCancellation { + process.destroyForcibly() + } + + val progressChannel = Channel(Channel.UNLIMITED) + val progressJob = if (progress != null) { + jobScope.launch { + for (item in progressChannel) { + progress(item) + } + } + } else { + null + } + + val readerThread = Thread { + runCatching { + val kv = HashMap() + 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() + } + } +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/convert/FfprobeService.kt b/src/main/kotlin/pw/binom/mirror/worker/convert/FfprobeService.kt new file mode 100644 index 0000000..4e493d6 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/convert/FfprobeService.kt @@ -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) + } +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/convert/FileInfo.kt b/src/main/kotlin/pw/binom/mirror/worker/convert/FileInfo.kt new file mode 100644 index 0000000..7796897 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/convert/FileInfo.kt @@ -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, + val format: Format, +) { + val videoStreams: List + get() = streams.filterIsInstance() + .filter { it.codecName != "png" } + .filter { it.codecName != "mjpeg" } + .filter { it.tags["filename"] != "cover.jpg" } + + val audioStreams: List + get() = streams.filterIsInstance() + + @Serializable + @JsonClassDiscriminator("codec_type") + sealed interface Stream { + val index: Int + val duration: Duration? + val startTime: Duration? + val tags: Map + } + + @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 = 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 = 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, + ) : 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 = 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 = emptyMap(), + ) +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/convert/InputService.kt b/src/main/kotlin/pw/binom/mirror/worker/convert/InputService.kt new file mode 100644 index 0000000..03f0017 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/convert/InputService.kt @@ -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, +) { + 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 + } +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/convert/InputSource.kt b/src/main/kotlin/pw/binom/mirror/worker/convert/InputSource.kt new file mode 100644 index 0000000..9ab8302 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/convert/InputSource.kt @@ -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 +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/convert/JellyfinInput.kt b/src/main/kotlin/pw/binom/mirror/worker/convert/JellyfinInput.kt new file mode 100644 index 0000000..ae04841 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/convert/JellyfinInput.kt @@ -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() + } + } +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/convert/JobExecutor.kt b/src/main/kotlin/pw/binom/mirror/worker/convert/JobExecutor.kt new file mode 100644 index 0000000..19684fd --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/convert/JobExecutor.kt @@ -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, + val duration: Duration, + ) + + data class Progress( + val audioProgress: List, + 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() + 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() + val audioFiles = ArrayList() + 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" + } +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/convert/LocalStorageService.kt b/src/main/kotlin/pw/binom/mirror/worker/convert/LocalStorageService.kt new file mode 100644 index 0000000..54fa0cc --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/convert/LocalStorageService.kt @@ -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) + } +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/convert/OutputService.kt b/src/main/kotlin/pw/binom/mirror/worker/convert/OutputService.kt new file mode 100644 index 0000000..08d2584 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/convert/OutputService.kt @@ -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() + } + } + } +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/convert/S3Clients.kt b/src/main/kotlin/pw/binom/mirror/worker/convert/S3Clients.kt new file mode 100644 index 0000000..f45b56b --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/convert/S3Clients.kt @@ -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() +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/convert/S3Destination.kt b/src/main/kotlin/pw/binom/mirror/worker/convert/S3Destination.kt new file mode 100644 index 0000000..d8fa3ef --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/convert/S3Destination.kt @@ -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" +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/convert/S3Input.kt b/src/main/kotlin/pw/binom/mirror/worker/convert/S3Input.kt new file mode 100644 index 0000000..22565e8 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/convert/S3Input.kt @@ -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() + } + } +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/convert/ViewSizeCalculate.kt b/src/main/kotlin/pw/binom/mirror/worker/convert/ViewSizeCalculate.kt new file mode 100644 index 0000000..63dfa43 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/convert/ViewSizeCalculate.kt @@ -0,0 +1,20 @@ +package pw.binom.mirror.worker.convert + +object ViewSizeCalculate { + fun calc(width: Int, height: Int, maxWidth: Int?, maxHeight: Int?): Pair { + 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 + } +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/db/Job.kt b/src/main/kotlin/pw/binom/mirror/worker/db/Job.kt new file mode 100644 index 0000000..58cec44 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/db/Job.kt @@ -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 = 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" + } +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/db/JobRepository.kt b/src/main/kotlin/pw/binom/mirror/worker/db/JobRepository.kt new file mode 100644 index 0000000..11926c3 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/db/JobRepository.kt @@ -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) { + 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 { + 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(), + ) + } +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/worker/JobPoller.kt b/src/main/kotlin/pw/binom/mirror/worker/worker/JobPoller.kt new file mode 100644 index 0000000..7fdbb6a --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/worker/JobPoller.kt @@ -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) + } + } + } +} diff --git a/src/main/kotlin/pw/binom/mirror/worker/worker/JobProcessor.kt b/src/main/kotlin/pw/binom/mirror/worker/worker/JobProcessor.kt new file mode 100644 index 0000000..80726b4 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/worker/worker/JobProcessor.kt @@ -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() + } +} diff --git a/src/main/resources/application.yaml b/src/main/resources/application.yaml new file mode 100644 index 0000000..6fc66f7 --- /dev/null +++ b/src/main/resources/application.yaml @@ -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 diff --git a/src/main/resources/db/migration/V1__init.sql b/src/main/resources/db/migration/V1__init.sql new file mode 100644 index 0000000..22d0faf --- /dev/null +++ b/src/main/resources/db/migration/V1__init.sql @@ -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); diff --git a/src/test/kotlin/pw/binom/mirror/worker/AbstractMinioTest.kt b/src/test/kotlin/pw/binom/mirror/worker/AbstractMinioTest.kt new file mode 100644 index 0000000..f61dbdd --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/worker/AbstractMinioTest.kt @@ -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()) + } + } + } +} diff --git a/src/test/kotlin/pw/binom/mirror/worker/AbstractPostgresTest.kt b/src/test/kotlin/pw/binom/mirror/worker/AbstractPostgresTest.kt new file mode 100644 index 0000000..df13772 --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/worker/AbstractPostgresTest.kt @@ -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") + } +} diff --git a/src/test/kotlin/pw/binom/mirror/worker/JobPollerTest.kt b/src/test/kotlin/pw/binom/mirror/worker/JobPollerTest.kt new file mode 100644 index 0000000..d0dd645 --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/worker/JobPollerTest.kt @@ -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 await(timeoutMs: Long, condition: () -> T?): T? { + val deadline = System.currentTimeMillis() + timeoutMs + while (System.currentTimeMillis() < deadline) { + condition()?.let { return it } + Thread.sleep(100) + } + return null + } +} diff --git a/src/test/kotlin/pw/binom/mirror/worker/TestMedia.kt b/src/test/kotlin/pw/binom/mirror/worker/TestMedia.kt new file mode 100644 index 0000000..ef0012a --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/worker/TestMedia.kt @@ -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) { + 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" + } + } +} diff --git a/src/test/kotlin/pw/binom/mirror/worker/convert/DurationAsSecondsStringTest.kt b/src/test/kotlin/pw/binom/mirror/worker/convert/DurationAsSecondsStringTest.kt new file mode 100644 index 0000000..1435cb8 --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/worker/convert/DurationAsSecondsStringTest.kt @@ -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\"")) + } +} diff --git a/src/test/kotlin/pw/binom/mirror/worker/convert/FfmpegServiceTest.kt b/src/test/kotlin/pw/binom/mirror/worker/convert/FfmpegServiceTest.kt new file mode 100644 index 0000000..36b2ca1 --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/worker/convert/FfmpegServiceTest.kt @@ -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() + 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) + } +} diff --git a/src/test/kotlin/pw/binom/mirror/worker/convert/FfprobeServiceTest.kt b/src/test/kotlin/pw/binom/mirror/worker/convert/FfprobeServiceTest.kt new file mode 100644 index 0000000..a328ed9 --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/worker/convert/FfprobeServiceTest.kt @@ -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) + } +} diff --git a/src/test/kotlin/pw/binom/mirror/worker/convert/FileInfoTest.kt b/src/test/kotlin/pw/binom/mirror/worker/convert/FileInfoTest.kt new file mode 100644 index 0000000..4efbc70 --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/worker/convert/FileInfoTest.kt @@ -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) + } +} diff --git a/src/test/kotlin/pw/binom/mirror/worker/convert/InputServiceTest.kt b/src/test/kotlin/pw/binom/mirror/worker/convert/InputServiceTest.kt new file mode 100644 index 0000000..395c27c --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/worker/convert/InputServiceTest.kt @@ -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) + } +} diff --git a/src/test/kotlin/pw/binom/mirror/worker/convert/JobExecutorTest.kt b/src/test/kotlin/pw/binom/mirror/worker/convert/JobExecutorTest.kt new file mode 100644 index 0000000..423febd --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/worker/convert/JobExecutorTest.kt @@ -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(), + ) +} diff --git a/src/test/kotlin/pw/binom/mirror/worker/convert/LocalStorageServiceTest.kt b/src/test/kotlin/pw/binom/mirror/worker/convert/LocalStorageServiceTest.kt new file mode 100644 index 0000000..d0c2a21 --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/worker/convert/LocalStorageServiceTest.kt @@ -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) + } +} diff --git a/src/test/kotlin/pw/binom/mirror/worker/convert/OutputServiceTest.kt b/src/test/kotlin/pw/binom/mirror/worker/convert/OutputServiceTest.kt new file mode 100644 index 0000000..a3da797 --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/worker/convert/OutputServiceTest.kt @@ -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) + } +} diff --git a/src/test/kotlin/pw/binom/mirror/worker/convert/ViewSizeCalculateTest.kt b/src/test/kotlin/pw/binom/mirror/worker/convert/ViewSizeCalculateTest.kt new file mode 100644 index 0000000..11769e0 --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/worker/convert/ViewSizeCalculateTest.kt @@ -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) + } +} diff --git a/src/test/kotlin/pw/binom/mirror/worker/db/JobRepositoryTest.kt b/src/test/kotlin/pw/binom/mirror/worker/db/JobRepositoryTest.kt new file mode 100644 index 0000000..6340ba6 --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/worker/db/JobRepositoryTest.kt @@ -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() + 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 + } +}