use anyhow::{Context, Result};
use axum::{
Router,
body::Body,
extract::{Path, State},
http::{HeaderMap, StatusCode, header},
response::{IntoResponse, Response},
routing::get,
};
use crate::config::Quality;
use rand::Rng;
use std::path::PathBuf;
use std::process::Stdio;
use std::sync::Arc;
use tokio::fs::File;
use tokio::process::{Child, Command};
use tokio::sync::Mutex;
use tokio_util::io::ReaderStream;
const MASTER_FILENAME: &str = "master.m3u8";
const PLAYLIST_FILENAME: &str = "playlist.m3u8";
const SEGMENT_DURATION_SECS: u32 = 4;
#[allow(dead_code)]
const PLAYLIST_SIZE: u32 = 12;
#[derive(Debug, Clone)]
pub struct Session {
pub token: String,
pub video_url: String,
pub audio_url: String,
pub quality_label: String,
pub work_dir: PathBuf,
pub duration_secs: u64,
}
#[derive(Clone)]
pub struct StreamServer {
bind_addr: String,
port: u16,
inner: Arc<Inner>,
}
struct Inner {
session: Mutex<Option<Session>>,
active: Mutex<Option<Child>>,
}
impl StreamServer {
pub fn new(bind_addr: impl Into<String>, port: u16) -> Self {
Self {
bind_addr: bind_addr.into(),
port,
inner: Arc::new(Inner {
session: Mutex::new(None),
active: Mutex::new(None),
}),
}
}
pub async fn set_session(
&self,
video_url: String,
audio_url: String,
quality_label: String,
quality: Quality,
duration_secs: u64,
public_host: &str,
) -> Result<String> {
if let Some(mut child) = self.inner.active.lock().await.take() {
let _ = child.start_kill();
let _ = child.wait().await;
}
if let Some(prev) = self.inner.session.lock().await.take() {
let _ = std::fs::remove_dir_all(&prev.work_dir);
}
let token = random_token();
let work_dir = std::env::temp_dir().join(format!("grod-stream-{token}"));
std::fs::create_dir_all(&work_dir)
.with_context(|| format!("creating stream work dir {}", work_dir.display()))?;
let playlist_path = work_dir.join(PLAYLIST_FILENAME);
let master_path = work_dir.join(MASTER_FILENAME);
let (codecs, resolution, bandwidth) = master_hints(quality);
let master_body = format!(
"#EXTM3U\n\
#EXT-X-VERSION:6\n\
#EXT-X-INDEPENDENT-SEGMENTS\n\
#EXT-X-STREAM-INF:BANDWIDTH={bandwidth},CODECS=\"{codecs}\",RESOLUTION={resolution}\n\
{PLAYLIST_FILENAME}\n"
);
std::fs::write(&master_path, master_body)
.with_context(|| format!("writing master playlist {}", master_path.display()))?;
let mut cmd = Command::new("ffmpeg");
let keyint = SEGMENT_DURATION_SECS * 24; let x264_params = format!(
"keyint={keyint}:min-keyint={keyint}:scenecut=0:repeat-headers=1"
);
cmd.arg("-hide_banner")
.arg("-loglevel")
.arg("warning")
.arg("-fflags")
.arg("+genpts")
.arg("-reconnect").arg("1")
.arg("-reconnect_at_eof").arg("1")
.arg("-reconnect_streamed").arg("1")
.arg("-reconnect_on_network_error").arg("1")
.arg("-reconnect_on_http_error").arg("4xx,5xx")
.arg("-reconnect_delay_max").arg("5")
.arg("-i")
.arg(&video_url)
.arg("-reconnect").arg("1")
.arg("-reconnect_at_eof").arg("1")
.arg("-reconnect_streamed").arg("1")
.arg("-reconnect_on_network_error").arg("1")
.arg("-reconnect_on_http_error").arg("4xx,5xx")
.arg("-reconnect_delay_max").arg("5")
.arg("-i")
.arg(&audio_url)
.arg("-map")
.arg("0:v:0")
.arg("-map")
.arg("1:a:0")
.arg("-c:v")
.arg("libx264")
.arg("-preset")
.arg("veryfast")
.arg("-crf")
.arg("20")
.arg("-x264-params")
.arg(&x264_params)
.arg("-c:a")
.arg("copy")
.arg("-copyts")
.arg("-start_at_zero")
.arg("-muxdelay")
.arg("0")
.arg("-muxpreload")
.arg("0")
.arg("-hls_time")
.arg(SEGMENT_DURATION_SECS.to_string())
.arg("-hls_list_size")
.arg("0") .arg("-hls_flags")
.arg("independent_segments")
.arg("-hls_segment_type")
.arg("mpegts")
.arg("-hls_playlist_type")
.arg("event")
.arg("-hls_segment_filename")
.arg(work_dir.join("seg_%05d.ts").to_string_lossy().to_string())
.arg("-f")
.arg("hls")
.arg(playlist_path.to_string_lossy().to_string())
.stdout(Stdio::null())
.stderr(Stdio::inherit())
.kill_on_drop(true);
let child = cmd
.spawn()
.context("failed to spawn ffmpeg — is ffmpeg installed?")?;
let session = Session {
token: token.clone(),
video_url,
audio_url,
quality_label,
work_dir,
duration_secs,
};
*self.inner.session.lock().await = Some(session);
*self.inner.active.lock().await = Some(child);
let inner_clone = self.inner.clone();
tokio::spawn(async move {
loop {
tokio::time::sleep(std::time::Duration::from_secs(2)).await;
let mut guard = inner_clone.active.lock().await;
let exited = match guard.as_mut() {
Some(c) => matches!(c.try_wait(), Ok(Some(_))),
None => true,
};
if exited {
if let Some(mut c) = guard.take() {
let _ = c.wait().await;
}
break;
}
}
});
Ok(format!(
"http://{}:{}/stream/{}/{}",
public_host, self.port, token, MASTER_FILENAME
))
}
pub async fn current(&self) -> Option<Session> {
self.inner.session.lock().await.clone()
}
pub async fn clear(&self) {
if let Some(mut child) = self.inner.active.lock().await.take() {
let _ = child.start_kill();
let _ = child.wait().await;
}
if let Some(prev) = self.inner.session.lock().await.take() {
let _ = std::fs::remove_dir_all(&prev.work_dir);
}
}
pub async fn run(self) -> Result<()> {
let app = Router::new()
.route("/stream/{token}/{file}", get(file_handler))
.with_state(self.inner.clone());
let bind = format!("{}:{}", self.bind_addr, self.port);
let listener = tokio::net::TcpListener::bind(&bind)
.await
.with_context(|| format!("binding stream server on {bind}"))?;
eprintln!("Stream server listening on {bind}");
axum::serve(listener, app)
.await
.context("stream server failed")?;
Ok(())
}
}
async fn file_handler(
Path((token, file)): Path<(String, String)>,
State(inner): State<Arc<Inner>>,
) -> Response {
if file.contains('/') || file.contains("..") {
return (StatusCode::BAD_REQUEST, "invalid filename").into_response();
}
let session = {
let guard = inner.session.lock().await;
match guard.as_ref() {
Some(s) if s.token == token => s.clone(),
_ => return (StatusCode::NOT_FOUND, "no session for token").into_response(),
}
};
let path = session.work_dir.join(&file);
let is_media_playlist = file == PLAYLIST_FILENAME;
let mut tries = 0;
let file_ready = loop {
let exists = path.exists();
if exists && !is_media_playlist {
break true;
}
if exists && is_media_playlist {
let seg_count = std::fs::read_dir(&session.work_dir)
.map(|d| d.filter_map(|e| e.ok()).filter(|e| {
e.file_name().to_string_lossy().ends_with(".ts")
}).count())
.unwrap_or(0);
if seg_count >= 3 {
break true;
}
}
if tries >= 50 {
break false;
}
tries += 1;
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
};
if !file_ready {
return (StatusCode::NOT_FOUND, "file not ready").into_response();
}
let f = match File::open(&path).await {
Ok(f) => f,
Err(e) => {
return (
StatusCode::INTERNAL_SERVER_ERROR,
format!("open {}: {e}", path.display()),
)
.into_response();
}
};
let content_type = if file.ends_with(".m3u8") {
"application/vnd.apple.mpegurl"
} else if file.ends_with(".ts") {
"video/mp2t"
} else {
"application/octet-stream"
};
let stream = ReaderStream::new(f);
let body = Body::from_stream(stream);
let mut headers = HeaderMap::new();
headers.insert(header::CONTENT_TYPE, content_type.parse().unwrap());
headers.insert(header::CACHE_CONTROL, "no-store".parse().unwrap());
headers.insert(header::ACCESS_CONTROL_ALLOW_ORIGIN, "*".parse().unwrap());
(StatusCode::OK, headers, body).into_response()
}
fn master_hints(q: Quality) -> (&'static str, &'static str, u32) {
match q {
Quality::Best | Quality::P1080 => ("avc1.640028,mp4a.40.2", "1920x1080", 6_000_000),
Quality::P720 => ("avc1.4d401f,mp4a.40.2", "1280x720", 3_000_000),
Quality::P480 => ("avc1.4d401e,mp4a.40.2", "854x480", 1_500_000),
Quality::P360 => ("avc1.42c01e,mp4a.40.2", "640x360", 800_000),
}
}
fn random_token() -> String {
let mut rng = rand::thread_rng();
let mut bytes = [0u8; 12];
rng.fill(&mut bytes);
bytes.iter().map(|b| format!("{b:02x}")).collect()
}