Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| b9c248e359 | |||
| 12233e9bfa | |||
| d10759b26b | |||
| 96a9a0536d |
@@ -50,6 +50,7 @@ dependencies {
|
|||||||
implementation("org.jetbrains.kotlinx:kotlinx-serialization-json:1.11.0")
|
implementation("org.jetbrains.kotlinx:kotlinx-serialization-json:1.11.0")
|
||||||
implementation(platform("software.amazon.awssdk:bom:2.46.21"))
|
implementation(platform("software.amazon.awssdk:bom:2.46.21"))
|
||||||
implementation("software.amazon.awssdk:s3")
|
implementation("software.amazon.awssdk:s3")
|
||||||
|
implementation("com.squareup.okhttp3:okhttp:4.12.0")
|
||||||
|
|
||||||
testImplementation("org.springframework.boot:spring-boot-starter-test")
|
testImplementation("org.springframework.boot:spring-boot-starter-test")
|
||||||
testImplementation(kotlin("test-junit5"))
|
testImplementation(kotlin("test-junit5"))
|
||||||
|
|||||||
@@ -82,7 +82,9 @@ class FfmpegService(
|
|||||||
}
|
}
|
||||||
if (width != null || height != null) {
|
if (width != null || height != null) {
|
||||||
args += "-vf"
|
args += "-vf"
|
||||||
args += "scale=${width ?: "-1"}:${height ?: "-1"}"
|
// Высота прибита (480), ширина — по пропорциям; pad доводит кадр до кратности 16
|
||||||
|
// (H.264 макроблоки 16x16): чёрные поля по краям, картинка не искажается.
|
||||||
|
args += "scale=${width ?: "-1"}:${height ?: "-1"},pad=ceil(iw/16)*16:ceil(ih/16)*16:(ow-iw)/2:(oh-ih)/2"
|
||||||
}
|
}
|
||||||
if (!removeAudio) {
|
if (!removeAudio) {
|
||||||
if (audioCodec != null) {
|
if (audioCodec != null) {
|
||||||
|
|||||||
@@ -28,6 +28,7 @@ class InputService(
|
|||||||
logger.info("Download cancelled: {}", source)
|
logger.info("Download cancelled: {}", source)
|
||||||
throw e
|
throw e
|
||||||
} catch (e: Throwable) {
|
} catch (e: Throwable) {
|
||||||
|
logger.warn("Download failed: {}", e.message, e)
|
||||||
throw IllegalStateException("Can't download video from $source", e)
|
throw IllegalStateException("Can't download video from $source", e)
|
||||||
}
|
}
|
||||||
logger.info("Downloaded success in {} ms", (System.nanoTime() - startTime) / 1_000_000)
|
logger.info("Downloaded success in {} ms", (System.nanoTime() - startTime) / 1_000_000)
|
||||||
|
|||||||
@@ -1,15 +1,18 @@
|
|||||||
package pw.binom.mirror.worker.convert
|
package pw.binom.mirror.worker.convert
|
||||||
|
|
||||||
import kotlinx.coroutines.Dispatchers
|
import kotlinx.coroutines.Dispatchers
|
||||||
|
import kotlinx.coroutines.delay
|
||||||
import kotlinx.coroutines.withContext
|
import kotlinx.coroutines.withContext
|
||||||
|
import okhttp3.OkHttpClient
|
||||||
|
import okhttp3.Request
|
||||||
import org.slf4j.Logger
|
import org.slf4j.Logger
|
||||||
import org.slf4j.LoggerFactory
|
import org.slf4j.LoggerFactory
|
||||||
import org.springframework.stereotype.Component
|
import org.springframework.stereotype.Component
|
||||||
import java.net.HttpURLConnection
|
import java.io.IOException
|
||||||
import java.net.URL
|
|
||||||
import java.nio.file.Files
|
import java.nio.file.Files
|
||||||
import java.nio.file.Path
|
import java.nio.file.Path
|
||||||
import java.nio.file.StandardCopyOption
|
import java.nio.file.StandardCopyOption
|
||||||
|
import java.util.concurrent.TimeUnit
|
||||||
|
|
||||||
@Component
|
@Component
|
||||||
class JellyfinInput(
|
class JellyfinInput(
|
||||||
@@ -17,24 +20,46 @@ class JellyfinInput(
|
|||||||
) : InputService.InputImplementation {
|
) : InputService.InputImplementation {
|
||||||
private val logger: Logger = LoggerFactory.getLogger(JellyfinInput::class.java)
|
private val logger: Logger = LoggerFactory.getLogger(JellyfinInput::class.java)
|
||||||
|
|
||||||
|
private val client: OkHttpClient by lazy {
|
||||||
|
OkHttpClient.Builder()
|
||||||
|
.connectTimeout(30, TimeUnit.SECONDS)
|
||||||
|
.readTimeout(0, TimeUnit.SECONDS)
|
||||||
|
.followRedirects(true)
|
||||||
|
.followSslRedirects(true)
|
||||||
|
.build()
|
||||||
|
}
|
||||||
|
|
||||||
override fun isSupport(source: InputSource): Boolean = source is InputSource.Jellyfin
|
override fun isSupport(source: InputSource): Boolean = source is InputSource.Jellyfin
|
||||||
|
|
||||||
override suspend fun download(source: InputSource): Path = withContext(Dispatchers.IO) {
|
override suspend fun download(source: InputSource): Path = withContext(Dispatchers.IO) {
|
||||||
source as InputSource.Jellyfin
|
source as InputSource.Jellyfin
|
||||||
val url = URL(source.streamUrl)
|
var lastError: IOException? = null
|
||||||
val connection = url.openConnection() as HttpURLConnection
|
for (attempt in 1..3) {
|
||||||
connection.requestMethod = "GET"
|
try {
|
||||||
connection.instanceFollowRedirects = true
|
return@withContext downloadOnce(source)
|
||||||
connection.connectTimeout = 30_000
|
} catch (e: IOException) {
|
||||||
connection.readTimeout = 120_000
|
lastError = e
|
||||||
try {
|
logger.warn("Download attempt {}/3 failed: {}", attempt, e.message)
|
||||||
if (connection.responseCode != HttpURLConnection.HTTP_OK) {
|
if (attempt < 3) {
|
||||||
throw IllegalStateException("Invalid response code ${connection.responseCode}")
|
delay(5_000)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
throw lastError!!
|
||||||
|
}
|
||||||
|
|
||||||
|
private fun downloadOnce(source: InputSource.Jellyfin): Path {
|
||||||
|
val request = Request.Builder().url(source.streamUrl).get().build()
|
||||||
|
val response = client.newCall(request).execute()
|
||||||
|
response.use {
|
||||||
|
if (!response.isSuccessful) {
|
||||||
|
throw IllegalStateException("Invalid response code ${response.code}")
|
||||||
}
|
}
|
||||||
val resultFile = localStorageService.genFile()
|
val resultFile = localStorageService.genFile()
|
||||||
logger.info("Downloading file in to ${resultFile}...")
|
logger.info("Downloading file in to {}...", resultFile)
|
||||||
try {
|
val path = try {
|
||||||
connection.inputStream.use { input ->
|
val body = response.body ?: throw IllegalStateException("Empty response body")
|
||||||
|
body.byteStream().use { input ->
|
||||||
Files.copy(input, resultFile, StandardCopyOption.REPLACE_EXISTING)
|
Files.copy(input, resultFile, StandardCopyOption.REPLACE_EXISTING)
|
||||||
}
|
}
|
||||||
logger.info("File success downloaded!")
|
logger.info("File success downloaded!")
|
||||||
@@ -43,8 +68,7 @@ class JellyfinInput(
|
|||||||
runCatching { Files.deleteIfExists(resultFile) }
|
runCatching { Files.deleteIfExists(resultFile) }
|
||||||
throw e
|
throw e
|
||||||
}
|
}
|
||||||
} finally {
|
return path
|
||||||
connection.disconnect()
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -26,6 +26,7 @@ data class Job(
|
|||||||
val createdAt: Instant,
|
val createdAt: Instant,
|
||||||
val updatedAt: Instant,
|
val updatedAt: Instant,
|
||||||
val takenAt: Instant? = null,
|
val takenAt: Instant? = null,
|
||||||
|
val retryCount: Int = 0,
|
||||||
) {
|
) {
|
||||||
companion object {
|
companion object {
|
||||||
const val STATUS_NEW = "new"
|
const val STATUS_NEW = "new"
|
||||||
@@ -33,5 +34,6 @@ data class Job(
|
|||||||
const val STATUS_DONE = "done"
|
const val STATUS_DONE = "done"
|
||||||
const val STATUS_FAILED = "failed"
|
const val STATUS_FAILED = "failed"
|
||||||
const val STATUS_CANCELLED = "cancelled"
|
const val STATUS_CANCELLED = "cancelled"
|
||||||
|
const val MAX_RETRIES = 3
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -71,7 +71,7 @@ open class JobRepository(
|
|||||||
jdbcTemplate.update(
|
jdbcTemplate.update(
|
||||||
"""
|
"""
|
||||||
UPDATE media_mirror.jobs
|
UPDATE media_mirror.jobs
|
||||||
SET status = 'failed', error = ?, updated_at = now()
|
SET status = 'failed', error = ?, retry_count = retry_count + 1, updated_at = now()
|
||||||
WHERE id = ?
|
WHERE id = ?
|
||||||
""".trimIndent(),
|
""".trimIndent(),
|
||||||
error.take(4000),
|
error.take(4000),
|
||||||
@@ -110,8 +110,16 @@ open class JobRepository(
|
|||||||
val jobs = jdbcTemplate.query(
|
val jobs = jdbcTemplate.query(
|
||||||
"""
|
"""
|
||||||
SELECT * FROM media_mirror.jobs
|
SELECT * FROM media_mirror.jobs
|
||||||
WHERE status = 'new'
|
WHERE (
|
||||||
ORDER BY created_at
|
status = 'new'
|
||||||
|
OR (status = 'failed' AND retry_count < 3)
|
||||||
|
OR (status = 'processing' AND taken_at < now() - interval '30 seconds')
|
||||||
|
)
|
||||||
|
ORDER BY
|
||||||
|
CASE WHEN status = 'new' THEN 0
|
||||||
|
WHEN status = 'failed' THEN 1
|
||||||
|
ELSE 2 END,
|
||||||
|
created_at
|
||||||
LIMIT 1
|
LIMIT 1
|
||||||
FOR UPDATE SKIP LOCKED
|
FOR UPDATE SKIP LOCKED
|
||||||
""".trimIndent(),
|
""".trimIndent(),
|
||||||
@@ -152,6 +160,7 @@ open class JobRepository(
|
|||||||
createdAt = rs.getTimestamp("created_at").toInstant(),
|
createdAt = rs.getTimestamp("created_at").toInstant(),
|
||||||
updatedAt = rs.getTimestamp("updated_at").toInstant(),
|
updatedAt = rs.getTimestamp("updated_at").toInstant(),
|
||||||
takenAt = rs.getTimestamp("taken_at")?.toInstant(),
|
takenAt = rs.getTimestamp("taken_at")?.toInstant(),
|
||||||
|
retryCount = rs.getInt("retry_count"),
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -58,6 +58,12 @@ class JobProcessor(
|
|||||||
throw e
|
throw e
|
||||||
} catch (e: Throwable) {
|
} catch (e: Throwable) {
|
||||||
logger.warn("Job {} finished with error: {}", job.id, e.message)
|
logger.warn("Job {} finished with error: {}", job.id, e.message)
|
||||||
|
val currentRetries = jobRepository.findById(job.id)?.retryCount ?: 0
|
||||||
|
if (currentRetries >= Job.MAX_RETRIES) {
|
||||||
|
logger.warn("Job {} exceeded max retries ({}), permanently failed", job.id, currentRetries)
|
||||||
|
} else {
|
||||||
|
logger.info("Job {} failed (retry {}/{}), will be retried", job.id, currentRetries, Job.MAX_RETRIES)
|
||||||
|
}
|
||||||
jobRepository.markFailed(job.id, e.message ?: e.toString())
|
jobRepository.markFailed(job.id, e.message ?: e.toString())
|
||||||
} finally {
|
} finally {
|
||||||
heartbeatJob.cancel()
|
heartbeatJob.cancel()
|
||||||
|
|||||||
@@ -0,0 +1 @@
|
|||||||
|
ALTER TABLE media_mirror.jobs ADD COLUMN IF NOT EXISTS retry_count int NOT NULL DEFAULT 0;
|
||||||
@@ -122,6 +122,81 @@ class JobRepositoryTest : AbstractPostgresTest() {
|
|||||||
assertEquals("boom", job.error)
|
assertEquals("boom", job.error)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `markFailed increments retryCount`() {
|
||||||
|
val id = insertJob("item-1")
|
||||||
|
repository.takeNextJob()
|
||||||
|
|
||||||
|
repository.markFailed(id, "boom")
|
||||||
|
|
||||||
|
val job = repository.findByStatus("failed").single()
|
||||||
|
assertEquals(1, job.retryCount)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `takeNextJob picks failed job with retryCount below max`() {
|
||||||
|
val id = insertJob(itemId = "item-1", status = "failed", retryCount = 0)
|
||||||
|
|
||||||
|
val taken = repository.takeNextJob()
|
||||||
|
|
||||||
|
assertNotNull(taken)
|
||||||
|
assertEquals(id, taken.id)
|
||||||
|
assertEquals("processing", taken.status)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `takeNextJob does not pick failed job with retryCount at max`() {
|
||||||
|
insertJob(itemId = "item-1", status = "failed", retryCount = 3)
|
||||||
|
|
||||||
|
val taken = repository.takeNextJob()
|
||||||
|
|
||||||
|
assertNull(taken)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `takeNextJob picks stale processing job`() {
|
||||||
|
val id = insertJob(
|
||||||
|
itemId = "item-1",
|
||||||
|
status = "processing",
|
||||||
|
takenAt = Instant.now().minus(31, ChronoUnit.SECONDS),
|
||||||
|
)
|
||||||
|
|
||||||
|
val taken = repository.takeNextJob()
|
||||||
|
|
||||||
|
assertNotNull(taken)
|
||||||
|
assertEquals(id, taken.id)
|
||||||
|
assertEquals("processing", taken.status)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `takeNextJob does not pick fresh processing job`() {
|
||||||
|
insertJob(itemId = "item-1", status = "processing", takenAt = Instant.now())
|
||||||
|
|
||||||
|
val taken = repository.takeNextJob()
|
||||||
|
|
||||||
|
assertNull(taken)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `takeNextJob priority new over failed over stale`() {
|
||||||
|
val newId = insertJob(itemId = "item-new")
|
||||||
|
val failedId = insertJob(itemId = "item-failed", status = "failed", retryCount = 0)
|
||||||
|
insertJob(
|
||||||
|
itemId = "item-stale",
|
||||||
|
status = "processing",
|
||||||
|
takenAt = Instant.now().minus(31, ChronoUnit.SECONDS),
|
||||||
|
)
|
||||||
|
|
||||||
|
val first = repository.takeNextJob()
|
||||||
|
assertEquals(newId, first?.id)
|
||||||
|
|
||||||
|
val second = repository.takeNextJob()
|
||||||
|
assertEquals(failedId, second?.id)
|
||||||
|
|
||||||
|
val third = repository.takeNextJob()
|
||||||
|
assertEquals("item-stale", third?.itemId)
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
fun `stale processing job returns to new`() {
|
fun `stale processing job returns to new`() {
|
||||||
val id = insertJob(
|
val id = insertJob(
|
||||||
@@ -177,18 +252,20 @@ class JobRepositoryTest : AbstractPostgresTest() {
|
|||||||
status: String = "new",
|
status: String = "new",
|
||||||
sourceUrl: String = "http://localhost/video.mp4",
|
sourceUrl: String = "http://localhost/video.mp4",
|
||||||
takenAt: Instant? = null,
|
takenAt: Instant? = null,
|
||||||
|
retryCount: Int = 0,
|
||||||
): UUID {
|
): UUID {
|
||||||
val id = UUID.randomUUID()
|
val id = UUID.randomUUID()
|
||||||
JdbcTemplate(dataSource).update(
|
JdbcTemplate(dataSource).update(
|
||||||
"""
|
"""
|
||||||
INSERT INTO media_mirror.jobs (id, item_id, source_url, status, taken_at)
|
INSERT INTO media_mirror.jobs (id, item_id, source_url, status, taken_at, retry_count)
|
||||||
VALUES (?, ?, ?, ?, ?)
|
VALUES (?, ?, ?, ?, ?, ?)
|
||||||
""".trimIndent(),
|
""".trimIndent(),
|
||||||
id,
|
id,
|
||||||
itemId,
|
itemId,
|
||||||
sourceUrl,
|
sourceUrl,
|
||||||
status,
|
status,
|
||||||
takenAt?.let { java.sql.Timestamp.from(it) },
|
takenAt?.let { java.sql.Timestamp.from(it) },
|
||||||
|
retryCount,
|
||||||
)
|
)
|
||||||
return id
|
return id
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user