diff --git a/src/main/kotlin/pw/binom/mirror/api/config/AppProperties.kt b/src/main/kotlin/pw/binom/mirror/api/config/AppProperties.kt index c328752..dbe8469 100644 --- a/src/main/kotlin/pw/binom/mirror/api/config/AppProperties.kt +++ b/src/main/kotlin/pw/binom/mirror/api/config/AppProperties.kt @@ -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( diff --git a/src/main/kotlin/pw/binom/mirror/api/controller/MirrorController.kt b/src/main/kotlin/pw/binom/mirror/api/controller/MirrorController.kt index 19caa27..b7c68df 100644 --- a/src/main/kotlin/pw/binom/mirror/api/controller/MirrorController.kt +++ b/src/main/kotlin/pw/binom/mirror/api/controller/MirrorController.kt @@ -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 { + 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 { logger.warn("Bad request: {}", message) return jsonBody(json.encodeToString(ErrorResponse(message)), HttpStatus.BAD_REQUEST) diff --git a/src/main/kotlin/pw/binom/mirror/api/dto/ScanResult.kt b/src/main/kotlin/pw/binom/mirror/api/dto/ScanResult.kt new file mode 100644 index 0000000..3441368 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/api/dto/ScanResult.kt @@ -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, +) diff --git a/src/main/kotlin/pw/binom/mirror/api/service/JellyfinScanner.kt b/src/main/kotlin/pw/binom/mirror/api/service/JellyfinScanner.kt new file mode 100644 index 0000000..818ea21 --- /dev/null +++ b/src/main/kotlin/pw/binom/mirror/api/service/JellyfinScanner.kt @@ -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 { + val body = get(baseUrl, "/Users", apiKey) + return json.decodeFromString>(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(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, +) { + @Serializable + data class Item( + @SerialName("Id") val id: String, + ) +} diff --git a/src/main/resources/application.yaml b/src/main/resources/application.yaml index 9c328b5..b95db91 100644 --- a/src/main/resources/application.yaml +++ b/src/main/resources/application.yaml @@ -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} diff --git a/src/test/kotlin/pw/binom/mirror/api/AbstractIntegrationTest.kt b/src/test/kotlin/pw/binom/mirror/api/AbstractIntegrationTest.kt index 61f28a2..5f323ca 100644 --- a/src/test/kotlin/pw/binom/mirror/api/AbstractIntegrationTest.kt +++ b/src/test/kotlin/pw/binom/mirror/api/AbstractIntegrationTest.kt @@ -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" } } } } diff --git a/src/test/kotlin/pw/binom/mirror/api/FakeJellyfin.kt b/src/test/kotlin/pw/binom/mirror/api/FakeJellyfin.kt new file mode 100644 index 0000000..050e1fe --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/api/FakeJellyfin.kt @@ -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 = 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) } + } +} diff --git a/src/test/kotlin/pw/binom/mirror/api/controller/MirrorControllerTest.kt b/src/test/kotlin/pw/binom/mirror/api/controller/MirrorControllerTest.kt index 9356a0a..8fc5e5d 100644 --- a/src/test/kotlin/pw/binom/mirror/api/controller/MirrorControllerTest.kt +++ b/src/test/kotlin/pw/binom/mirror/api/controller/MirrorControllerTest.kt @@ -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( + 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( """ diff --git a/src/test/kotlin/pw/binom/mirror/api/service/JellyfinScannerTest.kt b/src/test/kotlin/pw/binom/mirror/api/service/JellyfinScannerTest.kt new file mode 100644 index 0000000..8c0a647 --- /dev/null +++ b/src/test/kotlin/pw/binom/mirror/api/service/JellyfinScannerTest.kt @@ -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) + } +}