Archive v0.1.0 — 사진/영상 라이브러리 관리 프로그램
Tauri v2 + SolidJS + SQLite + ffmpeg 사이드카로 구현한 크로스플랫폼 사진·영상 관리 앱. 로컬/NAS(SFTP·WebDAV·FTP) 소스, 가상 그리드, 인앱 재생, 태그/이동/삭제/undo, 중복 탐지, 포터블 배포. - archive-db: SQLite 스키마·마이그레이션·단일 writer 스레드 + FTS5 trigram - archive-vfs: VFS 4백엔드(local/sftp/ftp/webdav) + 자격증명(키체인/볼트) - archive-indexer: 스캔·해시·썸네일·중복탐지·태그·파일작업·유지보수 - archive-media: localhost HTTP 미디어 서버(Range) + ffmpeg 스트림 잡 - 프론트: 3-pane UI, justified 가상 그리드, 라이트박스, 중복 검토 패널 Rust 테스트 49개 통과. CI: win x64/arm64 포터블 zip + macOS universal dmg. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,20 @@
|
||||
[package]
|
||||
name = "archive-media"
|
||||
version = "0.1.0"
|
||||
edition = "2021"
|
||||
|
||||
[dependencies]
|
||||
axum = "0.8"
|
||||
async-trait = { workspace = true }
|
||||
tokio = { workspace = true }
|
||||
tokio-util = { version = "0.7", features = ["io"] }
|
||||
bytes = { workspace = true }
|
||||
serde = { workspace = true }
|
||||
mime_guess = "2"
|
||||
rand = "0.9"
|
||||
thiserror = { workspace = true }
|
||||
tracing = { workspace = true }
|
||||
|
||||
[dev-dependencies]
|
||||
reqwest = { version = "0.12", default-features = false, features = ["rustls-tls"] }
|
||||
tempfile = "3"
|
||||
@@ -0,0 +1,385 @@
|
||||
//! localhost 미디어 스트리밍 서버.
|
||||
//!
|
||||
//! Tauri asset/커스텀 프로토콜은 대용량 영상 Range/스트리밍에 결함이 있어
|
||||
//! (tauri#6375, #7355, #11371) 모든 미디어 바이트는 이 서버를 통해 전달한다.
|
||||
//! 소비자: 웹뷰 `<video>/<img>`, (M2+) ffmpeg 썸네일러/트랜스코더.
|
||||
|
||||
pub mod range;
|
||||
pub mod stream;
|
||||
|
||||
use axum::body::Body;
|
||||
use axum::extract::{Path as AxumPath, State};
|
||||
use axum::http::{header, HeaderMap, StatusCode};
|
||||
use axum::response::{IntoResponse, Response};
|
||||
use axum::routing::get;
|
||||
use axum::Router;
|
||||
use range::{parse_range, RangeSpec};
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
use tokio::io::{AsyncReadExt, AsyncSeekExt};
|
||||
use tokio::sync::oneshot;
|
||||
use tokio_util::io::ReaderStream;
|
||||
|
||||
/// 원격 파일 스트리밍 핸들 (전체 길이 + AsyncRead).
|
||||
pub struct RemoteStream {
|
||||
pub total_len: u64,
|
||||
pub reader: Box<dyn tokio::io::AsyncRead + Send + Unpin>,
|
||||
pub mime: String,
|
||||
}
|
||||
|
||||
/// file_id → 실제 경로 해석. 앱 쪽에서 DB 조회로 구현한다.
|
||||
#[async_trait::async_trait]
|
||||
pub trait PathResolver: Send + Sync + 'static {
|
||||
fn resolve(&self, file_id: i64) -> Option<PathBuf>;
|
||||
|
||||
/// 썸네일 캐시 파일 경로 해석 (size_class: 0=grid, 1=preview).
|
||||
fn resolve_thumb(&self, file_id: i64, size_class: u8) -> Option<PathBuf> {
|
||||
let _ = (file_id, size_class);
|
||||
None
|
||||
}
|
||||
|
||||
/// 원격 파일 여부 (true면 resolve 대신 open_remote 사용).
|
||||
fn is_remote(&self, file_id: i64) -> bool {
|
||||
let _ = file_id;
|
||||
false
|
||||
}
|
||||
|
||||
/// 원격 파일의 지정 범위를 스트리밍한다.
|
||||
async fn open_remote(&self, file_id: i64, range: Option<std::ops::Range<u64>>) -> Option<RemoteStream> {
|
||||
let _ = (file_id, range);
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
pub struct MediaServer {
|
||||
pub port: u16,
|
||||
pub token: String,
|
||||
pub streams: Arc<stream::StreamManager>,
|
||||
shutdown_tx: Option<oneshot::Sender<()>>,
|
||||
}
|
||||
|
||||
impl MediaServer {
|
||||
pub fn media_url(&self, file_id: i64) -> String {
|
||||
format!("http://127.0.0.1:{}/media/{}/{}", self.port, self.token, file_id)
|
||||
}
|
||||
|
||||
pub fn stream_url(&self, job_id: u64) -> String {
|
||||
format!("http://127.0.0.1:{}/stream/{}/{}", self.port, self.token, job_id)
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for MediaServer {
|
||||
fn drop(&mut self) {
|
||||
if let Some(tx) = self.shutdown_tx.take() {
|
||||
let _ = tx.send(());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
struct ServerState {
|
||||
token: String,
|
||||
resolver: Arc<dyn PathResolver>,
|
||||
streams: Arc<stream::StreamManager>,
|
||||
}
|
||||
|
||||
/// 서버를 127.0.0.1의 임의 포트에 기동한다. `ffmpeg`가 None이면 스트림 라우트는 501.
|
||||
pub async fn start(
|
||||
resolver: Arc<dyn PathResolver>,
|
||||
ffmpeg: Option<PathBuf>,
|
||||
) -> std::io::Result<MediaServer> {
|
||||
let token = generate_token();
|
||||
let streams = Arc::new(stream::StreamManager::new(ffmpeg));
|
||||
let state = Arc::new(ServerState {
|
||||
token: token.clone(),
|
||||
resolver,
|
||||
streams: streams.clone(),
|
||||
});
|
||||
|
||||
let app = Router::new()
|
||||
.route("/health", get(|| async { "ok" }))
|
||||
.route("/media/{token}/{id}", get(serve_media))
|
||||
.route("/thumb/{token}/{id}", get(serve_thumb))
|
||||
.route("/stream/{token}/{id}", get(serve_stream))
|
||||
.with_state(state);
|
||||
|
||||
let listener = tokio::net::TcpListener::bind(("127.0.0.1", 0)).await?;
|
||||
let port = listener.local_addr()?.port();
|
||||
let (shutdown_tx, shutdown_rx) = oneshot::channel::<()>();
|
||||
|
||||
tokio::spawn(async move {
|
||||
let result = axum::serve(listener, app)
|
||||
.with_graceful_shutdown(async {
|
||||
let _ = shutdown_rx.await;
|
||||
})
|
||||
.await;
|
||||
if let Err(e) = result {
|
||||
tracing::error!("미디어 서버 종료 오류: {e}");
|
||||
}
|
||||
});
|
||||
|
||||
tracing::info!(port, "미디어 서버 시작");
|
||||
Ok(MediaServer {
|
||||
port,
|
||||
token,
|
||||
streams,
|
||||
shutdown_tx: Some(shutdown_tx),
|
||||
})
|
||||
}
|
||||
|
||||
async fn serve_stream(
|
||||
State(state): State<Arc<ServerState>>,
|
||||
AxumPath((token, job_id)): AxumPath<(String, u64)>,
|
||||
) -> Response {
|
||||
if token != state.token {
|
||||
return StatusCode::NOT_FOUND.into_response();
|
||||
}
|
||||
if !state.streams.available() {
|
||||
return StatusCode::NOT_IMPLEMENTED.into_response();
|
||||
}
|
||||
let Some(stdout) = state.streams.spawn(job_id) else {
|
||||
return StatusCode::NOT_FOUND.into_response();
|
||||
};
|
||||
let body = Body::from_stream(ReaderStream::with_capacity(stdout, 256 * 1024));
|
||||
Response::builder()
|
||||
.status(StatusCode::OK)
|
||||
.header(header::CONTENT_TYPE, "video/mp4")
|
||||
.header(header::CACHE_CONTROL, "no-store")
|
||||
.body(body)
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
fn generate_token() -> String {
|
||||
let v: u128 = rand::random();
|
||||
format!("{v:032x}")
|
||||
}
|
||||
|
||||
async fn serve_media(
|
||||
State(state): State<Arc<ServerState>>,
|
||||
AxumPath((token, file_id)): AxumPath<(String, i64)>,
|
||||
headers: HeaderMap,
|
||||
) -> Response {
|
||||
// 토큰 불일치는 404 — 존재 여부를 노출하지 않는다
|
||||
if token != state.token {
|
||||
return StatusCode::NOT_FOUND.into_response();
|
||||
}
|
||||
|
||||
// 원격 파일은 VFS 스트리밍 경로로
|
||||
if state.resolver.is_remote(file_id) {
|
||||
return serve_remote(&state.resolver, file_id, &headers).await;
|
||||
}
|
||||
|
||||
let resolver = state.resolver.clone();
|
||||
let path = tokio::task::spawn_blocking(move || resolver.resolve(file_id))
|
||||
.await
|
||||
.ok()
|
||||
.flatten();
|
||||
let Some(path) = path else {
|
||||
return StatusCode::NOT_FOUND.into_response();
|
||||
};
|
||||
|
||||
let Ok(mut file) = tokio::fs::File::open(&path).await else {
|
||||
return StatusCode::NOT_FOUND.into_response();
|
||||
};
|
||||
let Ok(meta) = file.metadata().await else {
|
||||
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
|
||||
};
|
||||
let len = meta.len();
|
||||
|
||||
let mime = mime_guess::from_path(&path).first_or_octet_stream();
|
||||
let range_header = headers
|
||||
.get(header::RANGE)
|
||||
.and_then(|v| v.to_str().ok());
|
||||
|
||||
match parse_range(range_header, len) {
|
||||
RangeSpec::Unsatisfiable => Response::builder()
|
||||
.status(StatusCode::RANGE_NOT_SATISFIABLE)
|
||||
.header(header::CONTENT_RANGE, format!("bytes */{len}"))
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
|
||||
RangeSpec::Full => {
|
||||
let stream = ReaderStream::new(file);
|
||||
Response::builder()
|
||||
.status(StatusCode::OK)
|
||||
.header(header::CONTENT_TYPE, mime.as_ref())
|
||||
.header(header::CONTENT_LENGTH, len)
|
||||
.header(header::ACCEPT_RANGES, "bytes")
|
||||
.body(Body::from_stream(stream))
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
RangeSpec::Partial(start, end) => {
|
||||
if file.seek(std::io::SeekFrom::Start(start)).await.is_err() {
|
||||
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
|
||||
}
|
||||
let window = end - start + 1;
|
||||
let stream = ReaderStream::new(file.take(window));
|
||||
Response::builder()
|
||||
.status(StatusCode::PARTIAL_CONTENT)
|
||||
.header(header::CONTENT_TYPE, mime.as_ref())
|
||||
.header(header::CONTENT_LENGTH, window)
|
||||
.header(header::CONTENT_RANGE, format!("bytes {start}-{end}/{len}"))
|
||||
.header(header::ACCEPT_RANGES, "bytes")
|
||||
.body(Body::from_stream(stream))
|
||||
.unwrap()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// 원격 파일 스트리밍 — VFS의 open_remote로 범위를 받아 서빙한다.
|
||||
async fn serve_remote(
|
||||
resolver: &Arc<dyn PathResolver>,
|
||||
file_id: i64,
|
||||
headers: &HeaderMap,
|
||||
) -> Response {
|
||||
// 먼저 전체 길이를 알아야 Range 계산이 가능 — HEAD 격으로 0바이트 요청
|
||||
let head = resolver.open_remote(file_id, Some(0..0)).await;
|
||||
let Some(head) = head else {
|
||||
return StatusCode::NOT_FOUND.into_response();
|
||||
};
|
||||
let len = head.total_len;
|
||||
let mime = head.mime.clone();
|
||||
drop(head);
|
||||
|
||||
let range_header = headers.get(header::RANGE).and_then(|v| v.to_str().ok());
|
||||
match parse_range(range_header, len) {
|
||||
RangeSpec::Unsatisfiable => Response::builder()
|
||||
.status(StatusCode::RANGE_NOT_SATISFIABLE)
|
||||
.header(header::CONTENT_RANGE, format!("bytes */{len}"))
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
RangeSpec::Full => {
|
||||
let Some(stream) = resolver.open_remote(file_id, None).await else {
|
||||
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
|
||||
};
|
||||
Response::builder()
|
||||
.status(StatusCode::OK)
|
||||
.header(header::CONTENT_TYPE, mime)
|
||||
.header(header::CONTENT_LENGTH, len)
|
||||
.header(header::ACCEPT_RANGES, "bytes")
|
||||
.body(Body::from_stream(ReaderStream::new(stream.reader)))
|
||||
.unwrap()
|
||||
}
|
||||
RangeSpec::Partial(start, end) => {
|
||||
let Some(stream) = resolver.open_remote(file_id, Some(start..end + 1)).await else {
|
||||
return StatusCode::INTERNAL_SERVER_ERROR.into_response();
|
||||
};
|
||||
let window = end - start + 1;
|
||||
Response::builder()
|
||||
.status(StatusCode::PARTIAL_CONTENT)
|
||||
.header(header::CONTENT_TYPE, mime)
|
||||
.header(header::CONTENT_LENGTH, window)
|
||||
.header(header::CONTENT_RANGE, format!("bytes {start}-{end}/{len}"))
|
||||
.header(header::ACCEPT_RANGES, "bytes")
|
||||
.body(Body::from_stream(ReaderStream::new(stream.reader)))
|
||||
.unwrap()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(serde::Deserialize)]
|
||||
struct ThumbQuery {
|
||||
#[serde(default)]
|
||||
s: u8,
|
||||
}
|
||||
|
||||
async fn serve_thumb(
|
||||
State(state): State<Arc<ServerState>>,
|
||||
AxumPath((token, file_id)): AxumPath<(String, i64)>,
|
||||
axum::extract::Query(q): axum::extract::Query<ThumbQuery>,
|
||||
) -> Response {
|
||||
if token != state.token {
|
||||
return StatusCode::NOT_FOUND.into_response();
|
||||
}
|
||||
let resolver = state.resolver.clone();
|
||||
let path = tokio::task::spawn_blocking(move || resolver.resolve_thumb(file_id, q.s))
|
||||
.await
|
||||
.ok()
|
||||
.flatten();
|
||||
let Some(path) = path else {
|
||||
return StatusCode::NOT_FOUND.into_response();
|
||||
};
|
||||
let Ok(bytes) = tokio::fs::read(&path).await else {
|
||||
return StatusCode::NOT_FOUND.into_response();
|
||||
};
|
||||
Response::builder()
|
||||
.status(StatusCode::OK)
|
||||
.header(header::CONTENT_TYPE, "image/jpeg")
|
||||
.header(header::CONTENT_LENGTH, bytes.len())
|
||||
.header(header::CACHE_CONTROL, "max-age=3600")
|
||||
.body(Body::from(bytes))
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use std::collections::HashMap;
|
||||
|
||||
struct StaticResolver(HashMap<i64, PathBuf>);
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl PathResolver for StaticResolver {
|
||||
fn resolve(&self, file_id: i64) -> Option<PathBuf> {
|
||||
self.0.get(&file_id).cloned()
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn serves_full_and_range_requests() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let p = dir.path().join("video.mp4");
|
||||
let data: Vec<u8> = (0..=255u8).cycle().take(100_000).collect();
|
||||
std::fs::write(&p, &data).unwrap();
|
||||
|
||||
let mut map = HashMap::new();
|
||||
map.insert(42i64, p);
|
||||
let server = start(Arc::new(StaticResolver(map)), None).await.unwrap();
|
||||
let base = server.media_url(42);
|
||||
let client = reqwest::Client::new();
|
||||
|
||||
// 전체 응답
|
||||
let res = client.get(&base).send().await.unwrap();
|
||||
assert_eq!(res.status(), 200);
|
||||
assert_eq!(res.headers()["accept-ranges"], "bytes");
|
||||
assert_eq!(res.headers()["content-type"], "video/mp4");
|
||||
assert_eq!(res.bytes().await.unwrap().len(), 100_000);
|
||||
|
||||
// Range 응답 (영상 시킹 시나리오)
|
||||
let res = client
|
||||
.get(&base)
|
||||
.header("Range", "bytes=1000-1999")
|
||||
.send()
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(res.status(), 206);
|
||||
assert_eq!(res.headers()["content-range"], "bytes 1000-1999/100000");
|
||||
let body = res.bytes().await.unwrap();
|
||||
assert_eq!(body.len(), 1000);
|
||||
assert_eq!(&body[..], &data[1000..2000]);
|
||||
|
||||
// suffix range (mp4 moov-at-end 케이스)
|
||||
let res = client
|
||||
.get(&base)
|
||||
.header("Range", "bytes=-500")
|
||||
.send()
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(res.status(), 206);
|
||||
assert_eq!(&res.bytes().await.unwrap()[..], &data[99_500..]);
|
||||
|
||||
// 416
|
||||
let res = client
|
||||
.get(&base)
|
||||
.header("Range", "bytes=200000-")
|
||||
.send()
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(res.status(), 416);
|
||||
assert_eq!(res.headers()["content-range"], "bytes */100000");
|
||||
|
||||
// 잘못된 토큰 → 404
|
||||
let bad = format!("http://127.0.0.1:{}/media/wrongtoken/42", server.port);
|
||||
assert_eq!(client.get(&bad).send().await.unwrap().status(), 404);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,105 @@
|
||||
//! HTTP Range 헤더 파싱 (단일 범위만 지원 — 웹뷰 `<video>`/`<img>`는 단일 범위만 보낸다).
|
||||
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
pub enum RangeSpec {
|
||||
/// Range 헤더 없음 또는 무시 가능한 형식 → 전체 응답 (200)
|
||||
Full,
|
||||
/// 206 — [start, end] (inclusive)
|
||||
Partial(u64, u64),
|
||||
/// 416 — 만족 불가
|
||||
Unsatisfiable,
|
||||
}
|
||||
|
||||
/// `Range: bytes=a-b` / `bytes=a-` / `bytes=-suffix` 파싱.
|
||||
pub fn parse_range(header: Option<&str>, len: u64) -> RangeSpec {
|
||||
let Some(header) = header else {
|
||||
return RangeSpec::Full;
|
||||
};
|
||||
let Some(spec) = header.strip_prefix("bytes=") else {
|
||||
return RangeSpec::Full; // 알 수 없는 단위는 무시하고 전체 응답
|
||||
};
|
||||
// 다중 범위는 첫 범위만 사용
|
||||
let first = spec.split(',').next().unwrap_or("").trim();
|
||||
let Some((start_s, end_s)) = first.split_once('-') else {
|
||||
return RangeSpec::Full;
|
||||
};
|
||||
|
||||
if len == 0 {
|
||||
return RangeSpec::Unsatisfiable;
|
||||
}
|
||||
|
||||
match (start_s.is_empty(), end_s.is_empty()) {
|
||||
// bytes=-suffix : 마지막 suffix 바이트
|
||||
(true, false) => match end_s.parse::<u64>() {
|
||||
Ok(0) => RangeSpec::Unsatisfiable,
|
||||
Ok(suffix) => {
|
||||
let start = len.saturating_sub(suffix);
|
||||
RangeSpec::Partial(start, len - 1)
|
||||
}
|
||||
Err(_) => RangeSpec::Full,
|
||||
},
|
||||
// bytes=start- : start부터 끝까지
|
||||
(false, true) => match start_s.parse::<u64>() {
|
||||
Ok(start) if start < len => RangeSpec::Partial(start, len - 1),
|
||||
Ok(_) => RangeSpec::Unsatisfiable,
|
||||
Err(_) => RangeSpec::Full,
|
||||
},
|
||||
// bytes=start-end
|
||||
(false, false) => match (start_s.parse::<u64>(), end_s.parse::<u64>()) {
|
||||
(Ok(start), Ok(end)) => {
|
||||
if start > end || start >= len {
|
||||
RangeSpec::Unsatisfiable
|
||||
} else {
|
||||
RangeSpec::Partial(start, end.min(len - 1))
|
||||
}
|
||||
}
|
||||
_ => RangeSpec::Full,
|
||||
},
|
||||
(true, true) => RangeSpec::Full,
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn no_header_is_full() {
|
||||
assert_eq!(parse_range(None, 100), RangeSpec::Full);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn simple_range() {
|
||||
assert_eq!(parse_range(Some("bytes=0-49"), 100), RangeSpec::Partial(0, 49));
|
||||
assert_eq!(parse_range(Some("bytes=50-"), 100), RangeSpec::Partial(50, 99));
|
||||
assert_eq!(parse_range(Some("bytes=-10"), 100), RangeSpec::Partial(90, 99));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn end_clamped_to_len() {
|
||||
assert_eq!(parse_range(Some("bytes=90-200"), 100), RangeSpec::Partial(90, 99));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unsatisfiable() {
|
||||
assert_eq!(parse_range(Some("bytes=100-"), 100), RangeSpec::Unsatisfiable);
|
||||
assert_eq!(parse_range(Some("bytes=200-300"), 100), RangeSpec::Unsatisfiable);
|
||||
assert_eq!(parse_range(Some("bytes=5-2"), 100), RangeSpec::Unsatisfiable);
|
||||
assert_eq!(parse_range(Some("bytes=0-"), 0), RangeSpec::Unsatisfiable);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn garbage_falls_back_to_full() {
|
||||
assert_eq!(parse_range(Some("items=0-5"), 100), RangeSpec::Full);
|
||||
assert_eq!(parse_range(Some("bytes=abc-def"), 100), RangeSpec::Full);
|
||||
assert_eq!(parse_range(Some("bytes=-"), 100), RangeSpec::Full);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn multi_range_uses_first() {
|
||||
assert_eq!(
|
||||
parse_range(Some("bytes=0-10, 20-30"), 100),
|
||||
RangeSpec::Partial(0, 10)
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,160 @@
|
||||
//! ffmpeg 리먹스/트랜스코딩 스트림 잡 관리.
|
||||
//!
|
||||
//! 웹뷰가 재생 못 하는 컨테이너/코덱은 ffmpeg로 fMP4를 만들어 chunked로
|
||||
//! 흘려보낸다. 시킹은 잡 재생성(-ss 입력 시킹) + 프론트 seekBase 오프셋 방식.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::path::PathBuf;
|
||||
use std::process::Stdio;
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::sync::Mutex;
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Deserialize)]
|
||||
#[serde(rename_all = "lowercase")]
|
||||
pub enum VideoMode {
|
||||
/// 스트림 카피 (컨테이너만 교체)
|
||||
Copy,
|
||||
/// 카피 + hvc1 태그 (macOS HEVC)
|
||||
CopyHvc1,
|
||||
/// H.264 하드웨어 인코더로 변환
|
||||
Transcode,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Deserialize)]
|
||||
#[serde(rename_all = "lowercase")]
|
||||
pub enum AudioMode {
|
||||
Copy,
|
||||
/// AAC 변환 (AC-3/DTS/Vorbis 등)
|
||||
Aac,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct StreamSpec {
|
||||
pub input: PathBuf,
|
||||
pub start_seconds: f64,
|
||||
pub video: VideoMode,
|
||||
pub audio: AudioMode,
|
||||
}
|
||||
|
||||
pub struct StreamManager {
|
||||
ffmpeg: Option<PathBuf>,
|
||||
specs: Mutex<HashMap<u64, StreamSpec>>,
|
||||
procs: Mutex<HashMap<u64, tokio::process::Child>>,
|
||||
next_id: AtomicU64,
|
||||
}
|
||||
|
||||
impl StreamManager {
|
||||
pub fn new(ffmpeg: Option<PathBuf>) -> Self {
|
||||
StreamManager {
|
||||
ffmpeg,
|
||||
specs: Mutex::new(HashMap::new()),
|
||||
procs: Mutex::new(HashMap::new()),
|
||||
next_id: AtomicU64::new(1),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn available(&self) -> bool {
|
||||
self.ffmpeg.is_some()
|
||||
}
|
||||
|
||||
/// 잡을 등록하고 id를 돌려준다 (프로세스는 GET 시점에 기동).
|
||||
pub fn create(&self, spec: StreamSpec) -> u64 {
|
||||
let id = self.next_id.fetch_add(1, Ordering::SeqCst);
|
||||
self.specs.lock().unwrap().insert(id, spec);
|
||||
id
|
||||
}
|
||||
|
||||
pub fn stop(&self, id: u64) {
|
||||
self.specs.lock().unwrap().remove(&id);
|
||||
if let Some(mut child) = self.procs.lock().unwrap().remove(&id) {
|
||||
let _ = child.start_kill();
|
||||
}
|
||||
}
|
||||
|
||||
pub fn stop_all(&self) {
|
||||
self.specs.lock().unwrap().clear();
|
||||
let mut procs = self.procs.lock().unwrap();
|
||||
for (_, mut child) in procs.drain() {
|
||||
let _ = child.start_kill();
|
||||
}
|
||||
}
|
||||
|
||||
/// GET 핸들러에서 호출 — ffmpeg를 기동하고 stdout을 돌려준다.
|
||||
pub(crate) fn spawn(&self, id: u64) -> Option<tokio::process::ChildStdout> {
|
||||
let spec = self.specs.lock().unwrap().get(&id).cloned()?;
|
||||
let ffmpeg = self.ffmpeg.clone()?;
|
||||
|
||||
// 같은 잡의 이전 프로세스는 종료 (시킹 재요청 등)
|
||||
if let Some(mut old) = self.procs.lock().unwrap().remove(&id) {
|
||||
let _ = old.start_kill();
|
||||
}
|
||||
|
||||
let mut cmd = tokio::process::Command::new(&ffmpeg);
|
||||
cmd.args(["-v", "error"]);
|
||||
if spec.start_seconds > 0.01 {
|
||||
cmd.args(["-ss", &format!("{:.3}", spec.start_seconds)]);
|
||||
}
|
||||
cmd.arg("-i").arg(&spec.input);
|
||||
cmd.args(["-map", "0:v:0", "-map", "0:a:0?"]);
|
||||
|
||||
match spec.video {
|
||||
VideoMode::Copy => {
|
||||
cmd.args(["-c:v", "copy"]);
|
||||
}
|
||||
VideoMode::CopyHvc1 => {
|
||||
cmd.args(["-c:v", "copy", "-tag:v", "hvc1"]);
|
||||
}
|
||||
VideoMode::Transcode => {
|
||||
#[cfg(windows)]
|
||||
cmd.args(["-c:v", "h264_mf", "-b:v", "6M"]);
|
||||
#[cfg(target_os = "macos")]
|
||||
cmd.args(["-c:v", "h264_videotoolbox", "-b:v", "6M"]);
|
||||
#[cfg(not(any(windows, target_os = "macos")))]
|
||||
cmd.args(["-c:v", "libx264", "-preset", "veryfast", "-crf", "23"]);
|
||||
// H.264는 짝수 해상도 필수
|
||||
cmd.args(["-vf", "scale=trunc(iw/2)*2:trunc(ih/2)*2"]);
|
||||
}
|
||||
}
|
||||
match spec.audio {
|
||||
AudioMode::Copy => {
|
||||
cmd.args(["-c:a", "copy"]);
|
||||
}
|
||||
AudioMode::Aac => {
|
||||
cmd.args(["-c:a", "aac", "-b:a", "192k", "-ac", "2"]);
|
||||
}
|
||||
}
|
||||
cmd.args([
|
||||
"-movflags",
|
||||
"frag_keyframe+empty_moov+default_base_moof",
|
||||
"-f",
|
||||
"mp4",
|
||||
"pipe:1",
|
||||
]);
|
||||
cmd.stdout(Stdio::piped()).stderr(Stdio::null()).stdin(Stdio::null());
|
||||
cmd.kill_on_drop(true);
|
||||
#[cfg(windows)]
|
||||
{
|
||||
cmd.creation_flags(0x0800_0000); // CREATE_NO_WINDOW
|
||||
}
|
||||
|
||||
let mut child = match cmd.spawn() {
|
||||
Ok(c) => c,
|
||||
Err(e) => {
|
||||
tracing::error!("ffmpeg 스트림 기동 실패: {e}");
|
||||
return None;
|
||||
}
|
||||
};
|
||||
let stdout = child.stdout.take()?;
|
||||
self.procs.lock().unwrap().insert(id, child);
|
||||
Some(stdout)
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for StreamManager {
|
||||
fn drop(&mut self) {
|
||||
let mut procs = self.procs.lock().unwrap();
|
||||
for (_, mut child) in procs.drain() {
|
||||
let _ = child.start_kill();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,98 @@
|
||||
//! ffmpeg 스트림(리먹스) 실기 테스트 — 사이드카가 있을 때만 실행.
|
||||
|
||||
use archive_media::stream::{AudioMode, StreamSpec, VideoMode};
|
||||
use archive_media::PathResolver;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
|
||||
fn ffmpeg_path() -> Option<PathBuf> {
|
||||
let bin = Path::new(env!("CARGO_MANIFEST_DIR")).join("../../binaries");
|
||||
let triple = if cfg!(all(windows, target_arch = "x86_64")) {
|
||||
"x86_64-pc-windows-msvc"
|
||||
} else if cfg!(all(windows, target_arch = "aarch64")) {
|
||||
"aarch64-pc-windows-msvc"
|
||||
} else if cfg!(all(target_os = "macos", target_arch = "aarch64")) {
|
||||
"aarch64-apple-darwin"
|
||||
} else {
|
||||
"x86_64-apple-darwin"
|
||||
};
|
||||
let ext = if cfg!(windows) { ".exe" } else { "" };
|
||||
let p = bin.join(format!("ffmpeg-{triple}{ext}"));
|
||||
p.is_file().then_some(p)
|
||||
}
|
||||
|
||||
struct NoResolver;
|
||||
#[async_trait::async_trait]
|
||||
impl PathResolver for NoResolver {
|
||||
fn resolve(&self, _: i64) -> Option<PathBuf> {
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
fn make_mkv(ffmpeg: &Path, dir: &Path) -> PathBuf {
|
||||
let out = dir.join("입력.mkv");
|
||||
let vcodec = if cfg!(windows) { "h264_mf" } else { "h264_videotoolbox" };
|
||||
let status = std::process::Command::new(ffmpeg)
|
||||
.args([
|
||||
"-y", "-v", "error",
|
||||
"-f", "lavfi", "-i", "testsrc2=duration=3:size=320x240:rate=25",
|
||||
"-f", "lavfi", "-i", "sine=frequency=440:duration=3",
|
||||
"-c:v", vcodec, "-b:v", "300k", "-c:a", "ac3", "-shortest",
|
||||
])
|
||||
.arg(&out)
|
||||
.status()
|
||||
.unwrap();
|
||||
assert!(status.success());
|
||||
out
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn remux_mkv_to_fmp4_stream() {
|
||||
let Some(ffmpeg) = ffmpeg_path() else {
|
||||
eprintln!("사이드카 없음 — 건너뜀");
|
||||
return;
|
||||
};
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let mkv = make_mkv(&ffmpeg, dir.path());
|
||||
|
||||
let server = archive_media::start(Arc::new(NoResolver), Some(ffmpeg)).await.unwrap();
|
||||
|
||||
// MKV(h264 + ac3) → 영상 카피 + 오디오 AAC 변환 (모드 3)
|
||||
let job = server.streams.create(StreamSpec {
|
||||
input: mkv,
|
||||
start_seconds: 0.0,
|
||||
video: VideoMode::Copy,
|
||||
audio: AudioMode::Aac,
|
||||
});
|
||||
let url = server.stream_url(job);
|
||||
|
||||
let res = reqwest::get(&url).await.unwrap();
|
||||
assert_eq!(res.status(), 200);
|
||||
assert_eq!(res.headers()["content-type"], "video/mp4");
|
||||
|
||||
let body = res.bytes().await.unwrap();
|
||||
assert!(body.len() > 10_000, "스트림이 너무 짧음: {} bytes", body.len());
|
||||
// fMP4 시그니처: 4번째 바이트부터 'ftyp'
|
||||
assert_eq!(&body[4..8], b"ftyp", "fMP4 헤더가 아님");
|
||||
// fragmented mp4는 moof 박스를 포함
|
||||
assert!(
|
||||
body.windows(4).any(|w| w == b"moof"),
|
||||
"fragmented MP4가 아님 (moof 없음)"
|
||||
);
|
||||
|
||||
server.streams.stop(job);
|
||||
|
||||
// 시킹 시나리오: 1.5초 지점에서 재기동
|
||||
let job2 = server.streams.create(StreamSpec {
|
||||
input: dir.path().join("입력.mkv"),
|
||||
start_seconds: 1.5,
|
||||
video: VideoMode::Copy,
|
||||
audio: AudioMode::Aac,
|
||||
});
|
||||
let res2 = reqwest::get(server.stream_url(job2)).await.unwrap();
|
||||
assert_eq!(res2.status(), 200);
|
||||
let body2 = res2.bytes().await.unwrap();
|
||||
assert!(body2.len() > 5_000);
|
||||
assert!(body2.len() < body.len(), "1.5초 이후 스트림이 전체보다 짧아야 함");
|
||||
server.streams.stop(job2);
|
||||
}
|
||||
Reference in New Issue
Block a user