Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ff72389ea6 |
@@ -13,6 +13,8 @@ data class AppProperties(
|
||||
data class Jellyfin(
|
||||
@param:DefaultValue("https://jellyfin.binom.pw/")
|
||||
val url: String,
|
||||
@param:DefaultValue("")
|
||||
val apiKey: String,
|
||||
)
|
||||
|
||||
data class S3(
|
||||
|
||||
@@ -21,6 +21,9 @@ import pw.binom.mirror.api.dto.CreateMirrorRequest
|
||||
import pw.binom.mirror.api.dto.CreateMirrorResponse
|
||||
import pw.binom.mirror.api.dto.ErrorResponse
|
||||
import pw.binom.mirror.api.dto.MirrorStatus
|
||||
import pw.binom.mirror.api.dto.ScanResult
|
||||
import pw.binom.mirror.api.service.JellyfinScanner
|
||||
import pw.binom.mirror.api.service.JellyfinUnavailableException
|
||||
import pw.binom.mirror.api.service.MirrorService
|
||||
import java.util.UUID
|
||||
|
||||
@@ -28,6 +31,7 @@ import java.util.UUID
|
||||
@RequestMapping("/api/mirror")
|
||||
class MirrorController(
|
||||
private val service: MirrorService,
|
||||
private val scanner: JellyfinScanner,
|
||||
private val json: Json,
|
||||
private val validator: Validator,
|
||||
) {
|
||||
@@ -79,6 +83,20 @@ class MirrorController(
|
||||
ResponseEntity.notFound().build()
|
||||
}
|
||||
|
||||
@PostMapping("/scan")
|
||||
fun scan(): ResponseEntity<String> {
|
||||
val result = try {
|
||||
scanner.scan()
|
||||
} catch (e: JellyfinUnavailableException) {
|
||||
logger.warn("Jellyfin scan failed: {}", e.message)
|
||||
return jsonBody(
|
||||
json.encodeToString(ErrorResponse(e.message ?: "Jellyfin unavailable")),
|
||||
HttpStatus.BAD_GATEWAY,
|
||||
)
|
||||
}
|
||||
return jsonBody(json.encodeToString(ScanResult.serializer(), result), HttpStatus.OK)
|
||||
}
|
||||
|
||||
private fun badRequest(message: String): ResponseEntity<String> {
|
||||
logger.warn("Bad request: {}", message)
|
||||
return jsonBody(json.encodeToString(ErrorResponse(message)), HttpStatus.BAD_REQUEST)
|
||||
|
||||
@@ -0,0 +1,10 @@
|
||||
package pw.binom.mirror.api.dto
|
||||
|
||||
import kotlinx.serialization.Serializable
|
||||
|
||||
@Serializable
|
||||
data class ScanResult(
|
||||
val found: Int,
|
||||
val created: Int,
|
||||
val skipped: Int,
|
||||
)
|
||||
@@ -0,0 +1,114 @@
|
||||
package pw.binom.mirror.api.service
|
||||
|
||||
import kotlinx.serialization.SerialName
|
||||
import kotlinx.serialization.Serializable
|
||||
import kotlinx.serialization.json.Json
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.springframework.dao.DuplicateKeyException
|
||||
import org.springframework.stereotype.Service
|
||||
import pw.binom.mirror.api.config.AppProperties
|
||||
import pw.binom.mirror.api.db.JobRepository
|
||||
import pw.binom.mirror.api.dto.ScanResult
|
||||
import java.net.HttpURLConnection
|
||||
import java.net.URL
|
||||
import java.nio.charset.StandardCharsets
|
||||
|
||||
@Service
|
||||
class JellyfinScanner(
|
||||
private val repository: JobRepository,
|
||||
private val properties: AppProperties,
|
||||
private val json: Json,
|
||||
) {
|
||||
private val logger = LoggerFactory.getLogger(JellyfinScanner::class.java)
|
||||
|
||||
fun scan(): ScanResult {
|
||||
val baseUrl = properties.jellyfin.url.trimEnd('/')
|
||||
val apiKey = properties.jellyfin.apiKey
|
||||
val userId = fetchUsers(baseUrl, apiKey).first().id
|
||||
var found = 0
|
||||
var created = 0
|
||||
var skipped = 0
|
||||
var startIndex = 0
|
||||
while (true) {
|
||||
val page = fetchItems(baseUrl, userId, apiKey, startIndex, PAGE_SIZE)
|
||||
found += page.items.size
|
||||
for (item in page.items) {
|
||||
if (repository.findActiveByItemId(item.id) != null) {
|
||||
skipped++
|
||||
} else {
|
||||
try {
|
||||
repository.insert(item.id, "", "jellyfin")
|
||||
created++
|
||||
} catch (e: DuplicateKeyException) {
|
||||
logger.warn("Concurrent insert for itemId={}, skipping", item.id)
|
||||
skipped++
|
||||
}
|
||||
}
|
||||
}
|
||||
if (page.items.isEmpty()) break
|
||||
startIndex += page.items.size
|
||||
if (startIndex >= page.totalRecordCount) break
|
||||
}
|
||||
logger.info("Scan finished: found={}, created={}, skipped={}", found, created, skipped)
|
||||
return ScanResult(found = found, created = created, skipped = skipped)
|
||||
}
|
||||
|
||||
private fun fetchUsers(baseUrl: String, apiKey: String): List<JellyfinUser> {
|
||||
val body = get(baseUrl, "/Users", apiKey)
|
||||
return json.decodeFromString<List<JellyfinUser>>(body)
|
||||
}
|
||||
|
||||
private fun fetchItems(
|
||||
baseUrl: String,
|
||||
userId: String,
|
||||
apiKey: String,
|
||||
startIndex: Int,
|
||||
limit: Int,
|
||||
): ItemsPage {
|
||||
val path = "/Users/$userId/Items?Recursive=true&IncludeItemTypes=Movie,Episode&StartIndex=$startIndex&Limit=$limit"
|
||||
val body = get(baseUrl, path, apiKey)
|
||||
return json.decodeFromString<ItemsPage>(body)
|
||||
}
|
||||
|
||||
private fun get(baseUrl: String, path: String, apiKey: String): String {
|
||||
val connection = URL("$baseUrl$path").openConnection() as HttpURLConnection
|
||||
connection.requestMethod = "GET"
|
||||
connection.setRequestProperty("X-Emby-Token", apiKey)
|
||||
connection.connectTimeout = 10_000
|
||||
connection.readTimeout = 30_000
|
||||
try {
|
||||
val code = connection.responseCode
|
||||
if (code != HttpURLConnection.HTTP_OK) {
|
||||
val detail = connection.errorStream
|
||||
?.use { it.bufferedReader(StandardCharsets.UTF_8).readText() }
|
||||
.orEmpty()
|
||||
throw JellyfinUnavailableException("Jellyfin returned HTTP $code: $detail")
|
||||
}
|
||||
return connection.inputStream.use { it.bufferedReader(StandardCharsets.UTF_8).readText() }
|
||||
} finally {
|
||||
connection.disconnect()
|
||||
}
|
||||
}
|
||||
|
||||
private companion object {
|
||||
const val PAGE_SIZE = 200
|
||||
}
|
||||
}
|
||||
|
||||
class JellyfinUnavailableException(message: String) : RuntimeException(message)
|
||||
|
||||
@Serializable
|
||||
private data class JellyfinUser(
|
||||
@SerialName("Id") val id: String,
|
||||
)
|
||||
|
||||
@Serializable
|
||||
private data class ItemsPage(
|
||||
@SerialName("TotalRecordCount") val totalRecordCount: Int,
|
||||
@SerialName("Items") val items: List<Item>,
|
||||
) {
|
||||
@Serializable
|
||||
data class Item(
|
||||
@SerialName("Id") val id: String,
|
||||
)
|
||||
}
|
||||
@@ -1,6 +1,7 @@
|
||||
app:
|
||||
jellyfin:
|
||||
url: https://jellyfin.binom.pw/
|
||||
api-key: ${JELLYFIN_API_KEY}
|
||||
s3:
|
||||
url: https://s3.binom.pw
|
||||
accessKey: ${S3_ACCESS_KEY}
|
||||
|
||||
@@ -33,6 +33,8 @@ abstract class AbstractIntegrationTest {
|
||||
fun cleanState() {
|
||||
ensureBucket()
|
||||
jdbcTemplate.update("DELETE FROM media_mirror.jobs")
|
||||
fakeJellyfin.items = emptyList()
|
||||
fakeJellyfin.fail = false
|
||||
}
|
||||
|
||||
private fun ensureBucket() {
|
||||
@@ -61,6 +63,8 @@ abstract class AbstractIntegrationTest {
|
||||
private val postgres: PostgreSQLContainer<*> = PostgreSQLContainer("postgres:16-alpine").apply { start() }
|
||||
private val minio: MinIOContainer = MinIOContainer("minio/minio:latest").apply { start() }
|
||||
|
||||
val fakeJellyfin = FakeJellyfin()
|
||||
|
||||
@JvmStatic
|
||||
@DynamicPropertySource
|
||||
fun properties(registry: DynamicPropertyRegistry) {
|
||||
@@ -71,6 +75,8 @@ abstract class AbstractIntegrationTest {
|
||||
registry.add("app.s3.accessKey") { minio.userName }
|
||||
registry.add("app.s3.secretKey") { minio.password }
|
||||
registry.add("app.s3.bucket") { "media" }
|
||||
registry.add("app.jellyfin.url") { fakeJellyfin.url }
|
||||
registry.add("app.jellyfin.apiKey") { "test-api-key" }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
package pw.binom.mirror.api
|
||||
|
||||
import com.sun.net.httpserver.HttpExchange
|
||||
import com.sun.net.httpserver.HttpServer
|
||||
import java.net.InetSocketAddress
|
||||
import java.nio.charset.StandardCharsets
|
||||
|
||||
class FakeJellyfin {
|
||||
|
||||
@Volatile
|
||||
var items: List<String> = emptyList()
|
||||
|
||||
@Volatile
|
||||
var fail: Boolean = false
|
||||
|
||||
private val server: HttpServer = HttpServer.create(InetSocketAddress(0), 0).also { s ->
|
||||
s.createContext("/") { exchange -> handle(exchange) }
|
||||
s.start()
|
||||
}
|
||||
|
||||
val url: String
|
||||
get() = "http://localhost:${server.address.port}"
|
||||
|
||||
fun close() = server.stop(0)
|
||||
|
||||
private fun handle(exchange: HttpExchange) {
|
||||
try {
|
||||
if (fail) {
|
||||
exchange.sendResponseHeaders(500, -1)
|
||||
return
|
||||
}
|
||||
val path = exchange.requestURI.path
|
||||
when {
|
||||
path == "/Users" -> respond(exchange, 200, """[{"Id": "user-1", "Name": "Test User"}]""")
|
||||
path.startsWith("/Users/") && path.endsWith("/Items") -> respondItems(exchange)
|
||||
else -> respond(exchange, 404, """{"error": "not found"}""")
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
runCatching { exchange.sendResponseHeaders(500, -1) }
|
||||
} finally {
|
||||
exchange.close()
|
||||
}
|
||||
}
|
||||
|
||||
private fun respondItems(exchange: HttpExchange) {
|
||||
val query = exchange.requestURI.rawQuery.orEmpty()
|
||||
val params = query.split("&").mapNotNull { part ->
|
||||
val eq = part.indexOf('=')
|
||||
if (eq < 0) null else part.substring(0, eq) to part.substring(eq + 1)
|
||||
}.toMap()
|
||||
val startIndex = params["StartIndex"]?.toIntOrNull() ?: 0
|
||||
val limit = params["Limit"]?.toIntOrNull() ?: 200
|
||||
val page = items.drop(startIndex).take(limit)
|
||||
val body = buildString {
|
||||
append("{\"TotalRecordCount\": ${items.size}, \"Items\": [")
|
||||
append(page.joinToString(",") { id -> """{"Id": "$id", "Name": "Item $id", "Type": "Movie"}""" })
|
||||
append("]}")
|
||||
}
|
||||
respond(exchange, 200, body)
|
||||
}
|
||||
|
||||
private fun respond(exchange: HttpExchange, code: Int, body: String) {
|
||||
val bytes = body.toByteArray(StandardCharsets.UTF_8)
|
||||
exchange.sendResponseHeaders(code, bytes.size.toLong())
|
||||
exchange.responseBody.use { it.write(bytes) }
|
||||
}
|
||||
}
|
||||
@@ -21,6 +21,7 @@ import pw.binom.mirror.api.dto.ByItemResponse
|
||||
import pw.binom.mirror.api.dto.CreateMirrorResponse
|
||||
import pw.binom.mirror.api.dto.MirrorFiles
|
||||
import pw.binom.mirror.api.dto.MirrorStatus
|
||||
import pw.binom.mirror.api.dto.ScanResult
|
||||
import pw.binom.mirror.api.dto.VideoFile
|
||||
import pw.binom.mirror.api.db.JobRepository
|
||||
import pw.binom.mirror.api.s3.S3Storage
|
||||
@@ -223,6 +224,27 @@ class MirrorControllerTest : AbstractIntegrationTest() {
|
||||
mockMvc.perform(delete("/api/mirror/${UUID.randomUUID()}")).andExpect(status().isNotFound)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `POST scan returns counts of found created skipped`() {
|
||||
fakeJellyfin.items = listOf("scan-1", "scan-2", "scan-3")
|
||||
repository.insert("scan-2", "", "jellyfin")
|
||||
|
||||
val result = json.decodeFromString<ScanResult>(
|
||||
bodyOf(mockMvc.perform(post("/api/mirror/scan")).andExpect(status().isOk).andReturn()),
|
||||
)
|
||||
|
||||
assertEquals(3, result.found)
|
||||
assertEquals(2, result.created)
|
||||
assertEquals(1, result.skipped)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `POST scan returns 502 when jellyfin is unavailable`() {
|
||||
fakeJellyfin.fail = true
|
||||
|
||||
mockMvc.perform(post("/api/mirror/scan")).andExpect(status().isBadGateway)
|
||||
}
|
||||
|
||||
private fun markDone(id: UUID, itemId: String) {
|
||||
jdbcTemplate.update(
|
||||
"""
|
||||
|
||||
@@ -0,0 +1,97 @@
|
||||
package pw.binom.mirror.api.service
|
||||
|
||||
import org.junit.jupiter.api.Assertions.assertEquals
|
||||
import org.junit.jupiter.api.Assertions.assertThrows
|
||||
import org.junit.jupiter.api.Test
|
||||
import org.springframework.beans.factory.annotation.Autowired
|
||||
import pw.binom.mirror.api.AbstractIntegrationTest
|
||||
import pw.binom.mirror.api.db.JobRepository
|
||||
import java.util.UUID
|
||||
|
||||
class JellyfinScannerTest : AbstractIntegrationTest() {
|
||||
|
||||
@Autowired
|
||||
lateinit var repository: JobRepository
|
||||
|
||||
@Autowired
|
||||
lateinit var scanner: JellyfinScanner
|
||||
|
||||
@Test
|
||||
fun `scan finds all items across multiple pages`() {
|
||||
fakeJellyfin.items = (0 until 450).map { "item-$it" }
|
||||
|
||||
val result = scanner.scan()
|
||||
|
||||
assertEquals(450, result.found)
|
||||
assertEquals(450, result.created)
|
||||
assertEquals(0, result.skipped)
|
||||
val count = jdbcTemplate.queryForObject("SELECT count(*) FROM media_mirror.jobs", Int::class.java)
|
||||
assertEquals(450, count)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `scan skips active jobs and recreates failed ones`() {
|
||||
val done = repository.insert("done-item", "", "jellyfin")
|
||||
markStatus(done.id, "done")
|
||||
val failed = repository.insert("failed-item", "", "jellyfin")
|
||||
markStatus(failed.id, "failed")
|
||||
val cancelled = repository.insert("cancelled-item", "", "jellyfin")
|
||||
markStatus(cancelled.id, "cancelled")
|
||||
fakeJellyfin.items = listOf("done-item", "failed-item", "cancelled-item", "fresh-item")
|
||||
|
||||
val result = scanner.scan()
|
||||
|
||||
assertEquals(4, result.found)
|
||||
assertEquals(3, result.created)
|
||||
assertEquals(1, result.skipped)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `scan is idempotent`() {
|
||||
fakeJellyfin.items = (0 until 10).map { "idem-$it" }
|
||||
|
||||
assertEquals(10, scanner.scan().created)
|
||||
val second = scanner.scan()
|
||||
assertEquals(10, second.found)
|
||||
assertEquals(0, second.created)
|
||||
assertEquals(10, second.skipped)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `scan recreates job after done job is deleted`() {
|
||||
val itemId = "recreate-item"
|
||||
val job = repository.insert(itemId, "", "jellyfin")
|
||||
markStatus(job.id, "done")
|
||||
fakeJellyfin.items = listOf(itemId)
|
||||
|
||||
assertEquals(0, scanner.scan().created)
|
||||
|
||||
repository.delete(job.id)
|
||||
val result = scanner.scan()
|
||||
assertEquals(1, result.found)
|
||||
assertEquals(1, result.created)
|
||||
assertEquals(0, result.skipped)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `scan throws when jellyfin returns an error`() {
|
||||
fakeJellyfin.fail = true
|
||||
|
||||
assertThrows(JellyfinUnavailableException::class.java) {
|
||||
scanner.scan()
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `scan with empty library creates nothing`() {
|
||||
val result = scanner.scan()
|
||||
|
||||
assertEquals(0, result.found)
|
||||
assertEquals(0, result.created)
|
||||
assertEquals(0, result.skipped)
|
||||
}
|
||||
|
||||
private fun markStatus(id: UUID, status: String) {
|
||||
jdbcTemplate.update("UPDATE media_mirror.jobs SET status = ? WHERE id = ?", status, id)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user