feat: media-mirror-api — координатор зеркал (REST + очередь в Postgres)
- POST /api/mirror (дубль-контроль по itemId), GET статус/by-item/список, DELETE (cancelled/S3-очистка) - kotlinx-сериализация, Jackson выпилен из HTTP-стека (транзитив Flyway/AWS остался) - Spring Boot 4.1.0, Kotlin 2.4.10, Flyway, Testcontainers postgres+minio - тесты: 22, LINE 92.6%, smoke против живой БД прошёл
This commit is contained in:
@@ -0,0 +1,13 @@
|
||||
package pw.binom.mirror.api
|
||||
|
||||
import org.springframework.boot.autoconfigure.SpringBootApplication
|
||||
import org.springframework.boot.context.properties.ConfigurationPropertiesScan
|
||||
import org.springframework.boot.runApplication
|
||||
|
||||
@SpringBootApplication
|
||||
@ConfigurationPropertiesScan
|
||||
class MirrorApiApplication
|
||||
|
||||
fun main(args: Array<String>) {
|
||||
runApplication<MirrorApiApplication>(*args)
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
package pw.binom.mirror.api.config
|
||||
|
||||
import org.springframework.boot.context.properties.ConfigurationProperties
|
||||
import org.springframework.boot.context.properties.bind.DefaultValue
|
||||
|
||||
@ConfigurationProperties(prefix = "app")
|
||||
data class AppProperties(
|
||||
@param:DefaultValue
|
||||
val jellyfin: Jellyfin,
|
||||
@param:DefaultValue
|
||||
val s3: S3,
|
||||
) {
|
||||
data class Jellyfin(
|
||||
@param:DefaultValue("https://jellyfin.binom.pw/")
|
||||
val url: String,
|
||||
)
|
||||
|
||||
data class S3(
|
||||
@param:DefaultValue("https://s3.binom.pw")
|
||||
val url: String,
|
||||
@param:DefaultValue("")
|
||||
val accessKey: String,
|
||||
@param:DefaultValue("")
|
||||
val secretKey: String,
|
||||
@param:DefaultValue("media")
|
||||
val bucket: String,
|
||||
@param:DefaultValue("us-east-1")
|
||||
val region: String,
|
||||
@param:DefaultValue("mirror")
|
||||
val prefix: String,
|
||||
)
|
||||
}
|
||||
@@ -0,0 +1,21 @@
|
||||
package pw.binom.mirror.api.config
|
||||
|
||||
import kotlinx.serialization.json.Json
|
||||
import kotlinx.serialization.modules.SerializersModule
|
||||
import org.springframework.context.annotation.Bean
|
||||
import org.springframework.context.annotation.Configuration
|
||||
import pw.binom.mirror.api.serialization.UUIDSerializer
|
||||
import java.util.UUID
|
||||
|
||||
@Configuration
|
||||
class JsonConfig {
|
||||
|
||||
@Bean
|
||||
fun json(): Json = Json {
|
||||
ignoreUnknownKeys = true
|
||||
encodeDefaults = false
|
||||
serializersModule = SerializersModule {
|
||||
contextual(UUID::class, UUIDSerializer)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,91 @@
|
||||
package pw.binom.mirror.api.controller
|
||||
|
||||
import jakarta.validation.Validator
|
||||
import kotlinx.serialization.SerializationException
|
||||
import kotlinx.serialization.builtins.ListSerializer
|
||||
import kotlinx.serialization.json.Json
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.springframework.http.HttpStatus
|
||||
import org.springframework.http.MediaType
|
||||
import org.springframework.http.ResponseEntity
|
||||
import org.springframework.web.bind.annotation.DeleteMapping
|
||||
import org.springframework.web.bind.annotation.GetMapping
|
||||
import org.springframework.web.bind.annotation.PathVariable
|
||||
import org.springframework.web.bind.annotation.PostMapping
|
||||
import org.springframework.web.bind.annotation.RequestBody
|
||||
import org.springframework.web.bind.annotation.RequestMapping
|
||||
import org.springframework.web.bind.annotation.RequestParam
|
||||
import org.springframework.web.bind.annotation.RestController
|
||||
import pw.binom.mirror.api.dto.ByItemResponse
|
||||
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.service.MirrorService
|
||||
import java.util.UUID
|
||||
|
||||
@RestController
|
||||
@RequestMapping("/api/mirror")
|
||||
class MirrorController(
|
||||
private val service: MirrorService,
|
||||
private val json: Json,
|
||||
private val validator: Validator,
|
||||
) {
|
||||
private val logger = LoggerFactory.getLogger(MirrorController::class.java)
|
||||
|
||||
@PostMapping
|
||||
fun create(@RequestBody(required = false) body: String?): ResponseEntity<String> {
|
||||
val request = try {
|
||||
json.decodeFromString<CreateMirrorRequest>(body ?: "")
|
||||
} catch (e: SerializationException) {
|
||||
return badRequest("Invalid JSON body: ${e.message}")
|
||||
}
|
||||
val violation = validator.validate(request).firstOrNull()
|
||||
if (violation != null) {
|
||||
return badRequest(violation.message ?: "Validation failed")
|
||||
}
|
||||
val result = service.create(request)
|
||||
val status = if (result.created) HttpStatus.CREATED else HttpStatus.OK
|
||||
return jsonBody(json.encodeToString(result.response), status)
|
||||
}
|
||||
|
||||
@GetMapping("/{id}")
|
||||
fun status(@PathVariable id: UUID): ResponseEntity<String> {
|
||||
val status = service.status(id) ?: return ResponseEntity.notFound().build()
|
||||
return jsonBody(json.encodeToString(status), HttpStatus.OK)
|
||||
}
|
||||
|
||||
@GetMapping("/by-item/{itemId}")
|
||||
fun byItem(@PathVariable itemId: String): ResponseEntity<String> {
|
||||
val response = service.byItem(itemId) ?: return ResponseEntity.notFound().build()
|
||||
return jsonBody(json.encodeToString(response), HttpStatus.OK)
|
||||
}
|
||||
|
||||
@GetMapping
|
||||
fun list(
|
||||
@RequestParam(required = false) status: String?,
|
||||
@RequestParam(required = false, defaultValue = "50") limit: Int,
|
||||
): ResponseEntity<String> {
|
||||
val jobs = service.list(status, limit)
|
||||
val body = json.encodeToString(ListSerializer(MirrorStatus.serializer()), jobs)
|
||||
return jsonBody(body, HttpStatus.OK)
|
||||
}
|
||||
|
||||
@DeleteMapping("/{id}")
|
||||
fun delete(@PathVariable id: UUID): ResponseEntity<Unit> =
|
||||
if (service.delete(id)) {
|
||||
ResponseEntity.noContent().build()
|
||||
} else {
|
||||
ResponseEntity.notFound().build()
|
||||
}
|
||||
|
||||
private fun badRequest(message: String): ResponseEntity<String> {
|
||||
logger.warn("Bad request: {}", message)
|
||||
return jsonBody(json.encodeToString(ErrorResponse(message)), HttpStatus.BAD_REQUEST)
|
||||
}
|
||||
|
||||
private fun jsonBody(body: String, status: HttpStatus): ResponseEntity<String> =
|
||||
ResponseEntity.status(status)
|
||||
.contentType(MediaType.APPLICATION_JSON)
|
||||
.body(body)
|
||||
}
|
||||
@@ -0,0 +1,120 @@
|
||||
package pw.binom.mirror.api.db
|
||||
|
||||
import org.springframework.dao.DuplicateKeyException
|
||||
import org.springframework.jdbc.core.JdbcTemplate
|
||||
import org.springframework.jdbc.core.RowMapper
|
||||
import org.springframework.stereotype.Repository
|
||||
import java.sql.ResultSet
|
||||
import java.time.Instant
|
||||
import java.util.UUID
|
||||
|
||||
data class Job(
|
||||
val id: UUID,
|
||||
val itemId: String,
|
||||
val sourceUrl: String,
|
||||
val sourceType: String,
|
||||
val status: String,
|
||||
val progress: Int,
|
||||
val videoKey: String?,
|
||||
val audioKeys: String?,
|
||||
val error: String?,
|
||||
val createdAt: Instant,
|
||||
val updatedAt: Instant,
|
||||
val takenAt: Instant?,
|
||||
)
|
||||
|
||||
@Repository
|
||||
class JobRepository(private val jdbc: JdbcTemplate) {
|
||||
|
||||
fun insert(itemId: String, sourceUrl: String, sourceType: String): Job {
|
||||
val sql = """
|
||||
INSERT INTO media_mirror.jobs (item_id, source_url, source_type)
|
||||
VALUES (?, ?, ?)
|
||||
RETURNING id, item_id, source_url, source_type, status, progress,
|
||||
video_key, audio_keys, error, created_at, updated_at, taken_at
|
||||
""".trimIndent()
|
||||
return jdbc.queryForObject(sql, JOB_MAPPER, itemId, sourceUrl, sourceType)
|
||||
}
|
||||
|
||||
fun findById(id: UUID): Job? =
|
||||
findOne(
|
||||
"""
|
||||
SELECT id, item_id, source_url, source_type, status, progress,
|
||||
video_key, audio_keys, error, created_at, updated_at, taken_at
|
||||
FROM media_mirror.jobs WHERE id = ?
|
||||
""".trimIndent(),
|
||||
id,
|
||||
)
|
||||
|
||||
fun findByItemId(itemId: String): Job? =
|
||||
findOne(
|
||||
"""
|
||||
SELECT id, item_id, source_url, source_type, status, progress,
|
||||
video_key, audio_keys, error, created_at, updated_at, taken_at
|
||||
FROM media_mirror.jobs WHERE item_id = ?
|
||||
ORDER BY created_at DESC, id DESC LIMIT 1
|
||||
""".trimIndent(),
|
||||
itemId,
|
||||
)
|
||||
|
||||
fun findActiveByItemId(itemId: String): Job? =
|
||||
findOne(
|
||||
"""
|
||||
SELECT id, item_id, source_url, source_type, status, progress,
|
||||
video_key, audio_keys, error, created_at, updated_at, taken_at
|
||||
FROM media_mirror.jobs
|
||||
WHERE item_id = ? AND status NOT IN ('failed', 'cancelled')
|
||||
""".trimIndent(),
|
||||
itemId,
|
||||
)
|
||||
|
||||
fun list(status: String?, limit: Int): List<Job> {
|
||||
val base = """
|
||||
SELECT id, item_id, source_url, source_type, status, progress,
|
||||
video_key, audio_keys, error, created_at, updated_at, taken_at
|
||||
FROM media_mirror.jobs
|
||||
""".trimIndent()
|
||||
val sql = if (status == null) {
|
||||
"$base ORDER BY created_at LIMIT ?"
|
||||
} else {
|
||||
"$base WHERE status = ? ORDER BY created_at LIMIT ?"
|
||||
}
|
||||
return if (status == null) {
|
||||
jdbc.query(sql, JOB_MAPPER, limit)
|
||||
} else {
|
||||
jdbc.query(sql, JOB_MAPPER, status, limit)
|
||||
}
|
||||
}
|
||||
|
||||
fun updateStatus(id: UUID, status: String): Int =
|
||||
jdbc.update(
|
||||
"UPDATE media_mirror.jobs SET status = ?, updated_at = now() WHERE id = ?",
|
||||
status,
|
||||
id,
|
||||
)
|
||||
|
||||
fun delete(id: UUID): Int =
|
||||
jdbc.update("DELETE FROM media_mirror.jobs WHERE id = ?", id)
|
||||
|
||||
private fun findOne(sql: String, vararg args: Any): Job? =
|
||||
jdbc.query(sql, JOB_MAPPER, *args).firstOrNull()
|
||||
|
||||
companion object {
|
||||
private val JOB_MAPPER = RowMapper { rs: ResultSet, _: Int -> rs.toJob() }
|
||||
}
|
||||
}
|
||||
|
||||
private fun ResultSet.toJob(): Job = Job(
|
||||
id = getObject("id", UUID::class.java),
|
||||
itemId = getString("item_id"),
|
||||
sourceUrl = getString("source_url"),
|
||||
sourceType = getString("source_type"),
|
||||
status = getString("status"),
|
||||
progress = getInt("progress"),
|
||||
videoKey = getString("video_key"),
|
||||
audioKeys = getString("audio_keys"),
|
||||
error = getString("error"),
|
||||
createdAt = getTimestamp("created_at").toInstant(),
|
||||
updatedAt = getTimestamp("updated_at").toInstant(),
|
||||
takenAt = getTimestamp("taken_at")?.toInstant(),
|
||||
)
|
||||
@@ -0,0 +1,13 @@
|
||||
package pw.binom.mirror.api.dto
|
||||
|
||||
import jakarta.validation.constraints.NotBlank
|
||||
import kotlinx.serialization.Serializable
|
||||
|
||||
@Serializable
|
||||
data class CreateMirrorRequest(
|
||||
@field:NotBlank(message = "itemId must not be blank")
|
||||
val itemId: String,
|
||||
@field:NotBlank(message = "sourceUrl must not be blank")
|
||||
val sourceUrl: String,
|
||||
val sourceType: String = "jellyfin",
|
||||
)
|
||||
@@ -0,0 +1,8 @@
|
||||
package pw.binom.mirror.api.dto
|
||||
|
||||
import kotlinx.serialization.Serializable
|
||||
|
||||
@Serializable
|
||||
data class ErrorResponse(
|
||||
val error: String,
|
||||
)
|
||||
@@ -0,0 +1,39 @@
|
||||
package pw.binom.mirror.api.dto
|
||||
|
||||
import kotlinx.serialization.Serializable
|
||||
|
||||
@Serializable
|
||||
data class ByItemResponse(
|
||||
val itemId: String,
|
||||
val status: String,
|
||||
val files: MirrorFiles? = null,
|
||||
)
|
||||
|
||||
@Serializable
|
||||
data class MirrorFiles(
|
||||
val video: VideoFile,
|
||||
val audios: List<AudioFile> = emptyList(),
|
||||
)
|
||||
|
||||
@Serializable
|
||||
data class VideoFile(
|
||||
val key: String,
|
||||
val url: String,
|
||||
)
|
||||
|
||||
@Serializable
|
||||
data class AudioFile(
|
||||
val index: Int,
|
||||
val key: String,
|
||||
val title: String,
|
||||
val language: String,
|
||||
val url: String,
|
||||
)
|
||||
|
||||
@Serializable
|
||||
data class AudioKey(
|
||||
val key: String,
|
||||
val title: String,
|
||||
val language: String,
|
||||
val index: Int,
|
||||
)
|
||||
@@ -0,0 +1,21 @@
|
||||
package pw.binom.mirror.api.dto
|
||||
|
||||
import kotlinx.serialization.Contextual
|
||||
import kotlinx.serialization.Serializable
|
||||
import java.util.UUID
|
||||
|
||||
@Serializable
|
||||
data class CreateMirrorResponse(
|
||||
@Contextual
|
||||
val id: UUID,
|
||||
val status: String,
|
||||
)
|
||||
|
||||
@Serializable
|
||||
data class MirrorStatus(
|
||||
@Contextual
|
||||
val id: UUID,
|
||||
val itemId: String,
|
||||
val status: String,
|
||||
val progress: Int,
|
||||
)
|
||||
@@ -0,0 +1,64 @@
|
||||
package pw.binom.mirror.api.s3
|
||||
|
||||
import org.springframework.stereotype.Component
|
||||
import pw.binom.mirror.api.config.AppProperties
|
||||
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.model.DeleteObjectRequest
|
||||
import software.amazon.awssdk.services.s3.model.HeadObjectRequest
|
||||
import software.amazon.awssdk.services.s3.model.NoSuchKeyException
|
||||
import software.amazon.awssdk.services.s3.model.PutObjectRequest
|
||||
import java.net.URI
|
||||
|
||||
@Component
|
||||
class S3Storage(appProperties: AppProperties) {
|
||||
|
||||
private val s3 = appProperties.s3
|
||||
|
||||
private val client: S3Client = S3Client.builder()
|
||||
.endpointOverride(URI.create(s3.url))
|
||||
.region(Region.of(s3.region))
|
||||
.credentialsProvider(
|
||||
StaticCredentialsProvider.create(
|
||||
AwsBasicCredentials.create(s3.accessKey, s3.secretKey),
|
||||
),
|
||||
)
|
||||
.forcePathStyle(true)
|
||||
.build()
|
||||
|
||||
fun upload(key: String, content: ByteArray): String {
|
||||
client.putObject(
|
||||
PutObjectRequest.builder()
|
||||
.bucket(s3.bucket)
|
||||
.key(key)
|
||||
.build(),
|
||||
software.amazon.awssdk.core.sync.RequestBody.fromBytes(content),
|
||||
)
|
||||
return key
|
||||
}
|
||||
|
||||
fun deleteObject(key: String) {
|
||||
client.deleteObject(
|
||||
DeleteObjectRequest.builder()
|
||||
.bucket(s3.bucket)
|
||||
.key(key)
|
||||
.build(),
|
||||
)
|
||||
}
|
||||
|
||||
fun checkExists(key: String): Boolean = try {
|
||||
client.headObject(
|
||||
HeadObjectRequest.builder()
|
||||
.bucket(s3.bucket)
|
||||
.key(key)
|
||||
.build(),
|
||||
)
|
||||
true
|
||||
} catch (e: NoSuchKeyException) {
|
||||
false
|
||||
}
|
||||
|
||||
fun publicUrl(key: String): String = "${s3.url}/${s3.bucket}/$key"
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
package pw.binom.mirror.api.serialization
|
||||
|
||||
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 java.util.UUID
|
||||
|
||||
object UUIDSerializer : KSerializer<UUID> {
|
||||
override val descriptor: SerialDescriptor = PrimitiveSerialDescriptor("UUID", PrimitiveKind.STRING)
|
||||
|
||||
override fun deserialize(decoder: Decoder): UUID = UUID.fromString(decoder.decodeString())
|
||||
|
||||
override fun serialize(encoder: Encoder, value: UUID) {
|
||||
encoder.encodeString(value.toString())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,120 @@
|
||||
package pw.binom.mirror.api.service
|
||||
|
||||
import kotlinx.serialization.json.Json
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.springframework.dao.DuplicateKeyException
|
||||
import org.springframework.stereotype.Service
|
||||
import pw.binom.mirror.api.db.Job
|
||||
import pw.binom.mirror.api.db.JobRepository
|
||||
import pw.binom.mirror.api.dto.AudioFile
|
||||
import pw.binom.mirror.api.dto.AudioKey
|
||||
import pw.binom.mirror.api.dto.ByItemResponse
|
||||
import pw.binom.mirror.api.dto.CreateMirrorRequest
|
||||
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.VideoFile
|
||||
import pw.binom.mirror.api.s3.S3Storage
|
||||
import java.util.UUID
|
||||
|
||||
data class CreateResult(
|
||||
val response: CreateMirrorResponse,
|
||||
val created: Boolean,
|
||||
)
|
||||
|
||||
@Service
|
||||
class MirrorService(
|
||||
private val repository: JobRepository,
|
||||
private val s3: S3Storage,
|
||||
private val json: Json,
|
||||
) {
|
||||
private val logger = LoggerFactory.getLogger(MirrorService::class.java)
|
||||
|
||||
fun create(request: CreateMirrorRequest): CreateResult {
|
||||
repository.findActiveByItemId(request.itemId)?.let { existing ->
|
||||
return CreateResult(
|
||||
response = CreateMirrorResponse(id = existing.id, status = existing.status),
|
||||
created = false,
|
||||
)
|
||||
}
|
||||
val job = try {
|
||||
repository.insert(request.itemId, request.sourceUrl, request.sourceType)
|
||||
} catch (e: DuplicateKeyException) {
|
||||
logger.warn("Concurrent create for itemId={}, returning existing job", request.itemId)
|
||||
repository.findActiveByItemId(request.itemId)
|
||||
} ?: error("Failed to create job for itemId=${request.itemId}")
|
||||
return CreateResult(
|
||||
response = CreateMirrorResponse(id = job.id, status = job.status),
|
||||
created = true,
|
||||
)
|
||||
}
|
||||
|
||||
fun status(id: UUID): MirrorStatus? = repository.findById(id)?.toStatus()
|
||||
|
||||
fun byItem(itemId: String): ByItemResponse? {
|
||||
val job = repository.findByItemId(itemId) ?: return null
|
||||
val files = if (job.status == "done" && job.videoKey != null) {
|
||||
MirrorFiles(
|
||||
video = VideoFile(key = job.videoKey, url = s3.publicUrl(job.videoKey)),
|
||||
audios = parseAudios(job.audioKeys),
|
||||
)
|
||||
} else {
|
||||
null
|
||||
}
|
||||
return ByItemResponse(itemId = job.itemId, status = job.status, files = files)
|
||||
}
|
||||
|
||||
fun list(status: String?, limit: Int): List<MirrorStatus> =
|
||||
repository.list(status, limit.coerceIn(1, 500)).map { it.toStatus() }
|
||||
|
||||
fun delete(id: UUID): Boolean {
|
||||
val job = repository.findById(id) ?: return false
|
||||
return when (job.status) {
|
||||
"new", "processing" -> repository.updateStatus(id, "cancelled") > 0
|
||||
"done" -> {
|
||||
deleteFilesFromS3(job)
|
||||
repository.delete(id) > 0
|
||||
}
|
||||
else -> false
|
||||
}
|
||||
}
|
||||
|
||||
private fun deleteFilesFromS3(job: Job) {
|
||||
val keys = buildList {
|
||||
job.videoKey?.let { add(it) }
|
||||
parseAudioKeys(job.audioKeys).forEach { add(it.key) }
|
||||
}
|
||||
keys.filter { s3.checkExists(it) }.forEach { key ->
|
||||
runCatching { s3.deleteObject(key) }
|
||||
.onFailure { logger.warn("Failed to delete S3 object $key: ${it.message}") }
|
||||
}
|
||||
}
|
||||
|
||||
private fun parseAudios(raw: String?): List<AudioFile> =
|
||||
parseAudioKeys(raw).map {
|
||||
AudioFile(
|
||||
index = it.index,
|
||||
key = it.key,
|
||||
title = it.title,
|
||||
language = it.language,
|
||||
url = s3.publicUrl(it.key),
|
||||
)
|
||||
}
|
||||
|
||||
private fun parseAudioKeys(raw: String?): List<AudioKey> {
|
||||
if (raw.isNullOrBlank()) return emptyList()
|
||||
return try {
|
||||
json.decodeFromString<List<AudioKey>>(raw)
|
||||
} catch (e: Exception) {
|
||||
logger.warn("Failed to parse audio_keys for job: {}", e.message)
|
||||
emptyList()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun Job.toStatus(): MirrorStatus = MirrorStatus(
|
||||
id = id,
|
||||
itemId = itemId,
|
||||
status = status,
|
||||
progress = progress,
|
||||
)
|
||||
@@ -0,0 +1,21 @@
|
||||
app:
|
||||
jellyfin:
|
||||
url: https://jellyfin.binom.pw/
|
||||
s3:
|
||||
url: https://s3.binom.pw
|
||||
accessKey: ${S3_ACCESS_KEY}
|
||||
secretKey: ${S3_SECRET_KEY}
|
||||
bucket: media
|
||||
region: us-east-1
|
||||
prefix: mirror
|
||||
|
||||
spring:
|
||||
datasource:
|
||||
url: jdbc:postgresql://192.168.76.106:5432/glasses
|
||||
username: postgres
|
||||
password: postgres
|
||||
flyway:
|
||||
schemas: media_mirror
|
||||
|
||||
server:
|
||||
port: 8080
|
||||
@@ -0,0 +1,19 @@
|
||||
CREATE SCHEMA IF NOT EXISTS media_mirror;
|
||||
CREATE TABLE media_mirror.jobs (
|
||||
id uuid PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
item_id text NOT NULL, -- itemId из Jellyfin
|
||||
source_url text NOT NULL, -- URL исходника (jellyfin stream)
|
||||
source_type text NOT NULL DEFAULT 'jellyfin',
|
||||
status text NOT NULL DEFAULT 'new', -- new | processing | done | failed | cancelled
|
||||
progress int NOT NULL DEFAULT 0, -- 0..100
|
||||
video_key text, -- S3-ключ видео (после done)
|
||||
audio_keys jsonb, -- [{key, title, language, index}]
|
||||
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);
|
||||
CREATE UNIQUE INDEX IF NOT EXISTS uq_jobs_active_item
|
||||
ON media_mirror.jobs (item_id)
|
||||
WHERE status NOT IN ('failed', 'cancelled');
|
||||
@@ -0,0 +1,76 @@
|
||||
package pw.binom.mirror.api
|
||||
|
||||
import org.junit.jupiter.api.BeforeEach
|
||||
import org.springframework.beans.factory.annotation.Autowired
|
||||
import org.springframework.boot.test.context.SpringBootTest
|
||||
import org.springframework.boot.webmvc.test.autoconfigure.AutoConfigureMockMvc
|
||||
import org.springframework.jdbc.core.JdbcTemplate
|
||||
import org.springframework.test.context.DynamicPropertyRegistry
|
||||
import org.springframework.test.context.DynamicPropertySource
|
||||
import org.testcontainers.containers.MinIOContainer
|
||||
import org.testcontainers.containers.PostgreSQLContainer
|
||||
import pw.binom.mirror.api.s3.S3Storage
|
||||
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.model.S3Exception
|
||||
import java.net.URI
|
||||
|
||||
@SpringBootTest
|
||||
@AutoConfigureMockMvc
|
||||
abstract class AbstractIntegrationTest {
|
||||
|
||||
@Autowired
|
||||
lateinit var s3Storage: S3Storage
|
||||
|
||||
@Autowired
|
||||
lateinit var jdbcTemplate: JdbcTemplate
|
||||
|
||||
private var bucketReady = false
|
||||
|
||||
@BeforeEach
|
||||
fun cleanState() {
|
||||
ensureBucket()
|
||||
jdbcTemplate.update("DELETE FROM media_mirror.jobs")
|
||||
}
|
||||
|
||||
private fun ensureBucket() {
|
||||
if (bucketReady) return
|
||||
S3Client.builder()
|
||||
.endpointOverride(URI.create(minio.s3URL))
|
||||
.region(Region.US_EAST_1)
|
||||
.credentialsProvider(
|
||||
StaticCredentialsProvider.create(
|
||||
AwsBasicCredentials.create(minio.userName, minio.password),
|
||||
),
|
||||
)
|
||||
.forcePathStyle(true)
|
||||
.build()
|
||||
.use { client ->
|
||||
try {
|
||||
client.createBucket { it.bucket("media") }
|
||||
} catch (e: S3Exception) {
|
||||
if (e.statusCode() != 409) throw e
|
||||
}
|
||||
}
|
||||
bucketReady = true
|
||||
}
|
||||
|
||||
companion object {
|
||||
private val postgres: PostgreSQLContainer<*> = PostgreSQLContainer("postgres:16-alpine").apply { start() }
|
||||
private val minio: MinIOContainer = MinIOContainer("minio/minio:latest").apply { start() }
|
||||
|
||||
@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.s3.url") { minio.s3URL }
|
||||
registry.add("app.s3.accessKey") { minio.userName }
|
||||
registry.add("app.s3.secretKey") { minio.password }
|
||||
registry.add("app.s3.bucket") { "media" }
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,257 @@
|
||||
package pw.binom.mirror.api.controller
|
||||
|
||||
import kotlinx.serialization.json.Json
|
||||
import kotlinx.serialization.modules.SerializersModule
|
||||
import org.junit.jupiter.api.Assertions.assertEquals
|
||||
import org.junit.jupiter.api.Assertions.assertFalse
|
||||
import org.junit.jupiter.api.Assertions.assertNull
|
||||
import org.junit.jupiter.api.Assertions.assertTrue
|
||||
import org.junit.jupiter.api.Test
|
||||
import org.springframework.beans.factory.annotation.Autowired
|
||||
import org.springframework.http.MediaType
|
||||
import org.springframework.test.web.servlet.MockMvc
|
||||
import org.springframework.test.web.servlet.MvcResult
|
||||
import org.springframework.test.web.servlet.request.MockMvcRequestBuilders.delete
|
||||
import org.springframework.test.web.servlet.request.MockMvcRequestBuilders.get
|
||||
import org.springframework.test.web.servlet.request.MockMvcRequestBuilders.post
|
||||
import org.springframework.test.web.servlet.result.MockMvcResultMatchers.status
|
||||
import pw.binom.mirror.api.AbstractIntegrationTest
|
||||
import pw.binom.mirror.api.dto.AudioFile
|
||||
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.VideoFile
|
||||
import pw.binom.mirror.api.db.JobRepository
|
||||
import pw.binom.mirror.api.s3.S3Storage
|
||||
import pw.binom.mirror.api.serialization.UUIDSerializer
|
||||
import java.nio.charset.StandardCharsets
|
||||
import java.util.UUID
|
||||
|
||||
class MirrorControllerTest : AbstractIntegrationTest() {
|
||||
|
||||
@Autowired
|
||||
lateinit var mockMvc: MockMvc
|
||||
|
||||
@Autowired
|
||||
lateinit var repository: JobRepository
|
||||
|
||||
@Autowired
|
||||
lateinit var s3: S3Storage
|
||||
|
||||
private val json = Json {
|
||||
ignoreUnknownKeys = true
|
||||
serializersModule = SerializersModule {
|
||||
contextual(UUID::class, UUIDSerializer)
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `POST mirror creates a job and returns 201`() {
|
||||
val result = postCreate("create-id", "https://jellyfin.binom.pw/Videos/1/stream")
|
||||
.andExpect(status().isCreated)
|
||||
.andReturn()
|
||||
|
||||
val response = json.decodeFromString<CreateMirrorResponse>(bodyOf(result))
|
||||
assertEquals("new", response.status)
|
||||
assertTrue(response.id != UUID(0L, 0L))
|
||||
|
||||
val status = json.decodeFromString<MirrorStatus>(
|
||||
bodyOf(mockMvc.perform(get("/api/mirror/${response.id}")).andExpect(status().isOk).andReturn()),
|
||||
)
|
||||
assertEquals("create-id", status.itemId)
|
||||
assertEquals("new", status.status)
|
||||
assertEquals(0, status.progress)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `POST mirror returns 200 with same task on duplicate itemId`() {
|
||||
val first = postCreate("dup-id").andExpect(status().isCreated).andReturn()
|
||||
val firstResponse = json.decodeFromString<CreateMirrorResponse>(bodyOf(first))
|
||||
|
||||
val second = postCreate("dup-id").andExpect(status().isOk).andReturn()
|
||||
val secondResponse = json.decodeFromString<CreateMirrorResponse>(bodyOf(second))
|
||||
|
||||
assertEquals(firstResponse.id, secondResponse.id)
|
||||
assertEquals("new", secondResponse.status)
|
||||
|
||||
val count = jdbcTemplate.queryForObject(
|
||||
"SELECT count(*) FROM media_mirror.jobs WHERE item_id = 'dup-id'",
|
||||
Int::class.java,
|
||||
)
|
||||
assertEquals(1, count)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `POST mirror returns 400 on blank itemId`() {
|
||||
postCreate("", "https://source.example/stream")
|
||||
.andExpect(status().isBadRequest)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `POST mirror returns 400 on invalid json body`() {
|
||||
mockMvc.perform(
|
||||
post("/api/mirror").contentType(MediaType.APPLICATION_JSON).content("""{"itemId": """),
|
||||
).andExpect(status().isBadRequest)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `POST mirror returns 400 on missing body`() {
|
||||
mockMvc.perform(
|
||||
post("/api/mirror").contentType(MediaType.APPLICATION_JSON),
|
||||
).andExpect(status().isBadRequest)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `POST mirror does not duplicate a processing job`() {
|
||||
val created = json.decodeFromString<CreateMirrorResponse>(
|
||||
bodyOf(postCreate("busy-id").andExpect(status().isCreated).andReturn()),
|
||||
)
|
||||
jdbcTemplate.update(
|
||||
"UPDATE media_mirror.jobs SET status = 'processing', progress = 42 WHERE id = ?",
|
||||
created.id,
|
||||
)
|
||||
val duplicate = json.decodeFromString<CreateMirrorResponse>(
|
||||
bodyOf(postCreate("busy-id").andExpect(status().isOk).andReturn()),
|
||||
)
|
||||
assertEquals(created.id, duplicate.id)
|
||||
assertEquals("processing", duplicate.status)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `GET by-item returns files when done`() {
|
||||
val itemId = "done-item"
|
||||
val job = repository.insert(itemId, "https://source.example/stream", "jellyfin")
|
||||
markDone(job.id, itemId)
|
||||
|
||||
val response = json.decodeFromString<ByItemResponse>(
|
||||
bodyOf(mockMvc.perform(get("/api/mirror/by-item/$itemId")).andExpect(status().isOk).andReturn()),
|
||||
)
|
||||
|
||||
assertEquals(itemId, response.itemId)
|
||||
assertEquals("done", response.status)
|
||||
val files = requireNotNull(response.files)
|
||||
assertEquals(
|
||||
VideoFile("mirror/$itemId/video.mkv", s3.publicUrl("mirror/$itemId/video.mkv")),
|
||||
files.video,
|
||||
)
|
||||
assertEquals(
|
||||
listOf(
|
||||
AudioFile(0, "mirror/$itemId/audio-0.ogg", "original", "jpn", s3.publicUrl("mirror/$itemId/audio-0.ogg")),
|
||||
AudioFile(1, "mirror/$itemId/audio-1.ogg", "dub", "rus", s3.publicUrl("mirror/$itemId/audio-1.ogg")),
|
||||
),
|
||||
files.audios,
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `GET by-item returns files null when not done`() {
|
||||
val itemId = "new-item"
|
||||
repository.insert(itemId, "https://source.example/stream", "jellyfin")
|
||||
|
||||
val response = json.decodeFromString<ByItemResponse>(
|
||||
bodyOf(mockMvc.perform(get("/api/mirror/by-item/$itemId")).andExpect(status().isOk).andReturn()),
|
||||
)
|
||||
|
||||
assertEquals(itemId, response.itemId)
|
||||
assertEquals("new", response.status)
|
||||
assertNull(response.files)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `GET by-item returns 404 when missing`() {
|
||||
mockMvc.perform(get("/api/mirror/by-item/unknown-item")).andExpect(status().isNotFound)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `GET status returns 404 when missing`() {
|
||||
mockMvc.perform(get("/api/mirror/${UUID.randomUUID()}")).andExpect(status().isNotFound)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `GET list filters by status and applies limit`() {
|
||||
val newItem = repository.insert("list-new-1", "https://s/1", "jellyfin")
|
||||
val doneItem = repository.insert("list-done-1", "https://s/2", "jellyfin")
|
||||
markDone(doneItem.id, "list-done-1")
|
||||
val doneItem2 = repository.insert("list-done-2", "https://s/3", "jellyfin")
|
||||
markDone(doneItem2.id, "list-done-2")
|
||||
repository.insert("list-new-2", "https://s/4", "jellyfin")
|
||||
assertEquals("new", newItem.status)
|
||||
|
||||
val doneList = json.decodeFromString<List<MirrorStatus>>(
|
||||
bodyOf(mockMvc.perform(get("/api/mirror").param("status", "done")).andExpect(status().isOk).andReturn()),
|
||||
)
|
||||
assertEquals(listOf("list-done-1", "list-done-2"), doneList.map { it.itemId })
|
||||
assertTrue(doneList.all { it.status == "done" })
|
||||
|
||||
val limited = json.decodeFromString<List<MirrorStatus>>(
|
||||
bodyOf(mockMvc.perform(get("/api/mirror").param("limit", "1")).andExpect(status().isOk).andReturn()),
|
||||
)
|
||||
assertEquals(1, limited.size)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `DELETE cancels a new job and is idempotent`() {
|
||||
val job = repository.insert("delete-new", "https://s/1", "jellyfin")
|
||||
|
||||
mockMvc.perform(delete("/api/mirror/${job.id}")).andExpect(status().isNoContent)
|
||||
mockMvc.perform(delete("/api/mirror/${job.id}")).andExpect(status().isNotFound)
|
||||
|
||||
val status = json.decodeFromString<MirrorStatus>(
|
||||
bodyOf(mockMvc.perform(get("/api/mirror/${job.id}")).andExpect(status().isOk).andReturn()),
|
||||
)
|
||||
assertEquals("cancelled", status.status)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `DELETE removes done job and its s3 files`() {
|
||||
val itemId = "delete-done"
|
||||
val job = repository.insert(itemId, "https://s/1", "jellyfin")
|
||||
markDone(job.id, itemId)
|
||||
assertTrue(s3.checkExists("mirror/$itemId/video.mkv"))
|
||||
|
||||
mockMvc.perform(delete("/api/mirror/${job.id}")).andExpect(status().isNoContent)
|
||||
|
||||
assertFalse(s3.checkExists("mirror/$itemId/video.mkv"))
|
||||
assertFalse(s3.checkExists("mirror/$itemId/audio-0.ogg"))
|
||||
mockMvc.perform(get("/api/mirror/${job.id}")).andExpect(status().isNotFound)
|
||||
mockMvc.perform(delete("/api/mirror/${job.id}")).andExpect(status().isNotFound)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `DELETE returns 404 for missing job`() {
|
||||
mockMvc.perform(delete("/api/mirror/${UUID.randomUUID()}")).andExpect(status().isNotFound)
|
||||
}
|
||||
|
||||
private fun markDone(id: UUID, itemId: String) {
|
||||
jdbcTemplate.update(
|
||||
"""
|
||||
UPDATE media_mirror.jobs
|
||||
SET status = 'done', progress = 100, video_key = ?, audio_keys = ?::jsonb
|
||||
WHERE id = ?
|
||||
""".trimIndent(),
|
||||
"mirror/$itemId/video.mkv",
|
||||
"""[
|
||||
{"key": "mirror/$itemId/audio-0.ogg", "title": "original", "language": "jpn", "index": 0},
|
||||
{"key": "mirror/$itemId/audio-1.ogg", "title": "dub", "language": "rus", "index": 1}
|
||||
]""",
|
||||
id,
|
||||
)
|
||||
s3.upload("mirror/$itemId/video.mkv", byteArrayOf(1, 2, 3))
|
||||
s3.upload("mirror/$itemId/audio-0.ogg", byteArrayOf(4, 5))
|
||||
s3.upload("mirror/$itemId/audio-1.ogg", byteArrayOf(6, 7))
|
||||
}
|
||||
|
||||
private fun postCreate(
|
||||
itemId: String,
|
||||
sourceUrl: String = "https://jellyfin.binom.pw/Videos/1/stream",
|
||||
): org.springframework.test.web.servlet.ResultActions =
|
||||
mockMvc.perform(
|
||||
post("/api/mirror")
|
||||
.contentType(MediaType.APPLICATION_JSON)
|
||||
.content("""{"itemId":"$itemId","sourceUrl":"$sourceUrl","sourceType":"jellyfin"}"""),
|
||||
)
|
||||
|
||||
private fun bodyOf(result: MvcResult): String =
|
||||
result.response.getContentAsString(StandardCharsets.UTF_8)
|
||||
}
|
||||
@@ -0,0 +1,107 @@
|
||||
package pw.binom.mirror.api.db
|
||||
|
||||
import org.junit.jupiter.api.Assertions.assertEquals
|
||||
import org.junit.jupiter.api.Assertions.assertNull
|
||||
import org.junit.jupiter.api.Assertions.assertThrows
|
||||
import org.junit.jupiter.api.Assertions.assertTrue
|
||||
import org.junit.jupiter.api.Test
|
||||
import org.springframework.beans.factory.annotation.Autowired
|
||||
import org.springframework.dao.DuplicateKeyException
|
||||
import pw.binom.mirror.api.AbstractIntegrationTest
|
||||
|
||||
class JobRepositoryTest : AbstractIntegrationTest() {
|
||||
|
||||
@Autowired
|
||||
lateinit var repository: JobRepository
|
||||
|
||||
@Test
|
||||
fun `insert then findById returns the same job`() {
|
||||
val inserted = repository.insert("crud-id", "https://s/1", "jellyfin")
|
||||
|
||||
val found = repository.findById(inserted.id)
|
||||
assertEquals(inserted, found)
|
||||
assertEquals("crud-id", found?.itemId)
|
||||
assertEquals("new", found?.status)
|
||||
assertEquals(0, found?.progress)
|
||||
assertNull(found?.videoKey)
|
||||
assertNull(found?.audioKeys)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `updateStatus changes status and updated_at`() {
|
||||
val inserted = repository.insert("upd-id", "https://s/1", "jellyfin")
|
||||
|
||||
assertEquals(1, repository.updateStatus(inserted.id, "processing"))
|
||||
assertEquals("processing", repository.findById(inserted.id)?.status)
|
||||
|
||||
assertEquals(1, repository.updateStatus(inserted.id, "cancelled"))
|
||||
assertEquals("cancelled", repository.findById(inserted.id)?.status)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `delete removes the row`() {
|
||||
val inserted = repository.insert("del-id", "https://s/1", "jellyfin")
|
||||
|
||||
assertEquals(1, repository.delete(inserted.id))
|
||||
assertNull(repository.findById(inserted.id))
|
||||
assertEquals(0, repository.delete(inserted.id))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `list filters by status and respects limit`() {
|
||||
repository.insert("l-new-1", "https://s/1", "jellyfin")
|
||||
repository.insert("l-new-2", "https://s/2", "jellyfin")
|
||||
val done = repository.insert("l-done-1", "https://s/3", "jellyfin")
|
||||
jdbcTemplate.update("UPDATE media_mirror.jobs SET status = 'done' WHERE id = ?", done.id)
|
||||
|
||||
val all = repository.list(null, 100)
|
||||
assertEquals(3, all.size)
|
||||
|
||||
val doneOnly = repository.list("done", 100)
|
||||
assertEquals(listOf("l-done-1"), doneOnly.map { it.itemId })
|
||||
|
||||
val limited = repository.list(null, 2)
|
||||
assertEquals(2, limited.size)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `duplicate active itemId violates partial unique index`() {
|
||||
repository.insert("uniq-id", "https://s/1", "jellyfin")
|
||||
assertThrows(DuplicateKeyException::class.java) {
|
||||
repository.insert("uniq-id", "https://s/2", "jellyfin")
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `failed job allows creating a new one for the same itemId`() {
|
||||
val first = repository.insert("reuse-id", "https://s/1", "jellyfin")
|
||||
repository.updateStatus(first.id, "failed")
|
||||
|
||||
val second = repository.insert("reuse-id", "https://s/2", "jellyfin")
|
||||
assertTrue(second.id != first.id)
|
||||
assertEquals("new", second.status)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `findActiveByItemId only matches non-failed non-cancelled jobs`() {
|
||||
val active = repository.insert("active-id", "https://s/1", "jellyfin")
|
||||
assertEquals(active.id, repository.findActiveByItemId("active-id")?.id)
|
||||
|
||||
repository.updateStatus(active.id, "cancelled")
|
||||
assertNull(repository.findActiveByItemId("active-id"))
|
||||
|
||||
repository.insert("active2-id", "https://s/2", "jellyfin")
|
||||
val done = repository.insert("active3-id", "https://s/3", "jellyfin")
|
||||
jdbcTemplate.update("UPDATE media_mirror.jobs SET status = 'done' WHERE id = ?", done.id)
|
||||
assertEquals(done.id, repository.findActiveByItemId("active3-id")?.id)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `findByItemId returns the latest job for an item`() {
|
||||
val first = repository.insert("byitem-id", "https://s/1", "jellyfin")
|
||||
repository.updateStatus(first.id, "failed")
|
||||
val second = repository.insert("byitem-id", "https://s/2", "jellyfin")
|
||||
|
||||
assertEquals(second.id, repository.findByItemId("byitem-id")?.id)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user