Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 96a9a0536d | |||
| 68eefff884 |
@@ -50,6 +50,7 @@ dependencies {
|
||||
implementation("org.jetbrains.kotlinx:kotlinx-serialization-json:1.11.0")
|
||||
implementation(platform("software.amazon.awssdk:bom:2.46.21"))
|
||||
implementation("software.amazon.awssdk:s3")
|
||||
implementation("com.squareup.okhttp3:okhttp:4.12.0")
|
||||
|
||||
testImplementation("org.springframework.boot:spring-boot-starter-test")
|
||||
testImplementation(kotlin("test-junit5"))
|
||||
|
||||
@@ -8,6 +8,7 @@ import java.time.Duration
|
||||
data class AppProperties(
|
||||
val pollInterval: Duration = Duration.ofSeconds(2),
|
||||
val staleTimeout: Duration = Duration.ofMinutes(30),
|
||||
val heartbeatInterval: Duration = Duration.ofSeconds(10),
|
||||
val jellyfin: Jellyfin = Jellyfin(),
|
||||
val s3: S3 = S3(),
|
||||
val ffmpegPath: String = "ffmpeg",
|
||||
|
||||
@@ -28,6 +28,7 @@ class InputService(
|
||||
logger.info("Download cancelled: {}", source)
|
||||
throw e
|
||||
} catch (e: Throwable) {
|
||||
logger.warn("Download failed: {}", e.message, e)
|
||||
throw IllegalStateException("Can't download video from $source", e)
|
||||
}
|
||||
logger.info("Downloaded success in {} ms", (System.nanoTime() - startTime) / 1_000_000)
|
||||
|
||||
@@ -1,15 +1,18 @@
|
||||
package pw.binom.mirror.worker.convert
|
||||
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.withContext
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.Request
|
||||
import org.slf4j.Logger
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.springframework.stereotype.Component
|
||||
import java.net.HttpURLConnection
|
||||
import java.net.URL
|
||||
import java.io.IOException
|
||||
import java.nio.file.Files
|
||||
import java.nio.file.Path
|
||||
import java.nio.file.StandardCopyOption
|
||||
import java.util.concurrent.TimeUnit
|
||||
|
||||
@Component
|
||||
class JellyfinInput(
|
||||
@@ -17,24 +20,46 @@ class JellyfinInput(
|
||||
) : InputService.InputImplementation {
|
||||
private val logger: Logger = LoggerFactory.getLogger(JellyfinInput::class.java)
|
||||
|
||||
private val client: OkHttpClient by lazy {
|
||||
OkHttpClient.Builder()
|
||||
.connectTimeout(30, TimeUnit.SECONDS)
|
||||
.readTimeout(0, TimeUnit.SECONDS)
|
||||
.followRedirects(true)
|
||||
.followSslRedirects(true)
|
||||
.build()
|
||||
}
|
||||
|
||||
override fun isSupport(source: InputSource): Boolean = source is InputSource.Jellyfin
|
||||
|
||||
override 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}")
|
||||
var lastError: IOException? = null
|
||||
for (attempt in 1..3) {
|
||||
try {
|
||||
return@withContext downloadOnce(source)
|
||||
} catch (e: IOException) {
|
||||
lastError = e
|
||||
logger.warn("Download attempt {}/3 failed: {}", attempt, e.message)
|
||||
if (attempt < 3) {
|
||||
delay(5_000)
|
||||
}
|
||||
}
|
||||
}
|
||||
throw lastError!!
|
||||
}
|
||||
|
||||
private fun downloadOnce(source: InputSource.Jellyfin): Path {
|
||||
val request = Request.Builder().url(source.streamUrl).get().build()
|
||||
val response = client.newCall(request).execute()
|
||||
response.use {
|
||||
if (!response.isSuccessful) {
|
||||
throw IllegalStateException("Invalid response code ${response.code}")
|
||||
}
|
||||
val resultFile = localStorageService.genFile()
|
||||
logger.info("Downloading file in to ${resultFile}...")
|
||||
try {
|
||||
connection.inputStream.use { input ->
|
||||
logger.info("Downloading file in to {}...", resultFile)
|
||||
val path = try {
|
||||
val body = response.body ?: throw IllegalStateException("Empty response body")
|
||||
body.byteStream().use { input ->
|
||||
Files.copy(input, resultFile, StandardCopyOption.REPLACE_EXISTING)
|
||||
}
|
||||
logger.info("File success downloaded!")
|
||||
@@ -43,8 +68,7 @@ class JellyfinInput(
|
||||
runCatching { Files.deleteIfExists(resultFile) }
|
||||
throw e
|
||||
}
|
||||
} finally {
|
||||
connection.disconnect()
|
||||
return path
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -42,6 +42,17 @@ open class JobRepository(
|
||||
)
|
||||
}
|
||||
|
||||
open fun heartbeat(id: UUID) {
|
||||
jdbcTemplate.update(
|
||||
"""
|
||||
UPDATE media_mirror.jobs
|
||||
SET taken_at = now(), updated_at = now()
|
||||
WHERE id = ? AND status = 'processing'
|
||||
""".trimIndent(),
|
||||
id,
|
||||
)
|
||||
}
|
||||
|
||||
open fun markDone(id: UUID, videoKey: String, audioKeys: List<AudioKey>) {
|
||||
val audioJson = json.encodeToString(audioKeysSerializer, audioKeys)
|
||||
jdbcTemplate.update(
|
||||
|
||||
@@ -1,6 +1,13 @@
|
||||
package pw.binom.mirror.worker.worker
|
||||
|
||||
import kotlinx.coroutines.CancellationException
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.isActive
|
||||
import kotlinx.coroutines.launch
|
||||
import org.slf4j.Logger
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.springframework.stereotype.Service
|
||||
@@ -29,6 +36,13 @@ class JobProcessor(
|
||||
region = props.s3.region,
|
||||
prefix = props.s3.prefix,
|
||||
)
|
||||
val heartbeatScope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
|
||||
val heartbeatJob = heartbeatScope.launch {
|
||||
while (isActive) {
|
||||
delay(props.heartbeatInterval.toMillis())
|
||||
runCatching { jobRepository.heartbeat(job.id) }
|
||||
}
|
||||
}
|
||||
try {
|
||||
val result = jobExecutor.execute(
|
||||
job = job,
|
||||
@@ -45,6 +59,9 @@ class JobProcessor(
|
||||
} catch (e: Throwable) {
|
||||
logger.warn("Job {} finished with error: {}", job.id, e.message)
|
||||
jobRepository.markFailed(job.id, e.message ?: e.toString())
|
||||
} finally {
|
||||
heartbeatJob.cancel()
|
||||
heartbeatScope.cancel()
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -11,6 +11,7 @@ spring:
|
||||
app:
|
||||
poll-interval: 2s
|
||||
stale-timeout: 30m
|
||||
heartbeat-interval: 10s
|
||||
jellyfin:
|
||||
url: https://jellyfin.binom.pw/
|
||||
apiKey: ${JELLYFIN_API_KEY:}
|
||||
|
||||
@@ -146,6 +146,32 @@ class JobRepositoryTest : AbstractPostgresTest() {
|
||||
assertEquals(1, repository.findByStatus("processing").size)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `heartbeat refreshes takenAt of a processing job`() {
|
||||
val id = insertJob(
|
||||
itemId = "item-1",
|
||||
status = "processing",
|
||||
takenAt = Instant.now().minus(1, ChronoUnit.HOURS),
|
||||
)
|
||||
|
||||
repository.heartbeat(id)
|
||||
|
||||
val job = repository.findById(id)!!
|
||||
assertEquals("processing", job.status)
|
||||
assertTrue(job.takenAt!!.isAfter(Instant.now().minus(5, ChronoUnit.SECONDS)))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `heartbeat does not touch a non-processing job`() {
|
||||
val id = insertJob(itemId = "item-1", status = "new")
|
||||
|
||||
repository.heartbeat(id)
|
||||
|
||||
val job = repository.findById(id)!!
|
||||
assertEquals("new", job.status)
|
||||
assertNull(job.takenAt)
|
||||
}
|
||||
|
||||
private fun insertJob(
|
||||
itemId: String,
|
||||
status: String = "new",
|
||||
|
||||
Reference in New Issue
Block a user