use std::collections::HashMap;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::Mutex;
use sandbox_agent_error::SandboxError;
use crate::desktop_types::{DesktopProcessInfo, DesktopResolution, DesktopStreamStatusResponse};
use crate::process_runtime::{ProcessOwner, ProcessRuntime, ProcessStartSpec};
const NEKO_INTERNAL_PORT: u16 = 18100;
const NEKO_EPR: &str = "59050-59070";
const NEKO_READY_TIMEOUT: Duration = Duration::from_secs(15);
const NEKO_READY_POLL: Duration = Duration::from_millis(300);
#[derive(Debug, Clone)]
pub struct StreamingConfig {
pub video_codec: String,
pub audio_codec: String,
pub frame_rate: u32,
pub webrtc_port_range: String,
}
impl Default for StreamingConfig {
fn default() -> Self {
Self {
video_codec: "vp8".to_string(),
audio_codec: "opus".to_string(),
frame_rate: 30,
webrtc_port_range: NEKO_EPR.to_string(),
}
}
}
#[derive(Debug, Clone)]
pub struct DesktopStreamingManager {
inner: Arc<Mutex<DesktopStreamingState>>,
process_runtime: Arc<ProcessRuntime>,
}
#[derive(Debug)]
struct DesktopStreamingState {
active: bool,
process_id: Option<String>,
neko_base_url: Option<String>,
neko_session_cookie: Option<String>,
display: Option<String>,
resolution: Option<DesktopResolution>,
streaming_config: StreamingConfig,
window_id: Option<String>,
}
impl Default for DesktopStreamingState {
fn default() -> Self {
Self {
active: false,
process_id: None,
neko_base_url: None,
neko_session_cookie: None,
display: None,
resolution: None,
streaming_config: StreamingConfig::default(),
window_id: None,
}
}
}
impl DesktopStreamingManager {
pub fn new(process_runtime: Arc<ProcessRuntime>) -> Self {
Self {
inner: Arc::new(Mutex::new(DesktopStreamingState::default())),
process_runtime,
}
}
pub async fn start(
&self,
display: &str,
resolution: DesktopResolution,
environment: &HashMap<String, String>,
config: Option<StreamingConfig>,
window_id: Option<String>,
) -> Result<DesktopStreamStatusResponse, SandboxError> {
let config = config.unwrap_or_default();
let mut state = self.inner.lock().await;
if state.active {
return Ok(DesktopStreamStatusResponse {
active: true,
window_id: state.window_id.clone(),
process_id: state.process_id.clone(),
});
}
if let Some(ref old_id) = state.process_id {
let _ = self.process_runtime.stop_process(old_id, Some(2000)).await;
state.process_id = None;
state.neko_base_url = None;
state.neko_session_cookie = None;
}
let mut env = environment.clone();
env.insert("DISPLAY".to_string(), display.to_string());
let bind_addr = format!("0.0.0.0:{}", NEKO_INTERNAL_PORT);
let screen = format!(
"{}x{}@{}",
resolution.width, resolution.height, config.frame_rate
);
let snapshot = self
.process_runtime
.start_process(ProcessStartSpec {
command: "neko".to_string(),
args: vec![
"serve".to_string(),
"--server.bind".to_string(),
bind_addr,
"--desktop.screen".to_string(),
screen,
"--desktop.display".to_string(),
display.to_string(),
"--capture.video.display".to_string(),
display.to_string(),
"--capture.video.codec".to_string(),
config.video_codec.clone(),
"--capture.audio.codec".to_string(),
config.audio_codec.clone(),
"--webrtc.epr".to_string(),
config.webrtc_port_range.clone(),
"--webrtc.icelite".to_string(),
"--webrtc.nat1to1".to_string(),
"127.0.0.1".to_string(),
"--member.provider".to_string(),
"noauth".to_string(),
"--desktop.input.enabled=false".to_string(),
],
cwd: None,
env,
tty: false,
interactive: false,
owner: ProcessOwner::Desktop,
restart_policy: None,
})
.await
.map_err(|e| SandboxError::Conflict {
message: format!("failed to start neko streaming process: {e}"),
})?;
let neko_base = format!("http://127.0.0.1:{}", NEKO_INTERNAL_PORT);
let process_id_clone = snapshot.id.clone();
state.process_id = Some(snapshot.id.clone());
state.neko_base_url = Some(neko_base.clone());
state.display = Some(display.to_string());
state.resolution = Some(resolution);
state.streaming_config = config;
state.window_id = window_id;
state.active = true;
drop(state);
let deadline = tokio::time::Instant::now() + NEKO_READY_TIMEOUT;
let login_url = format!("{}/api/login", neko_base);
let client = reqwest::Client::builder()
.redirect(reqwest::redirect::Policy::none())
.build()
.unwrap_or_else(|_| reqwest::Client::new());
let mut session_cookie = None;
loop {
match client
.post(&login_url)
.json(&serde_json::json!({"username": "admin", "password": "admin"}))
.send()
.await
{
Ok(resp) if resp.status().is_success() => {
if let Some(set_cookie) = resp.headers().get("set-cookie") {
if let Ok(cookie_str) = set_cookie.to_str() {
if let Some(cookie_part) = cookie_str.split(';').next() {
session_cookie = Some(cookie_part.to_string());
}
}
}
tracing::info!("neko streaming process ready, session obtained");
let control_url = format!("{}/api/room/control/take", neko_base);
if let Some(ref cookie) = session_cookie {
let _ = client
.post(&control_url)
.header("Cookie", cookie.as_str())
.send()
.await;
tracing::info!("neko control taken");
}
break;
}
_ => {}
}
if tokio::time::Instant::now() >= deadline {
tracing::warn!("neko did not become ready within timeout, proceeding anyway");
break;
}
tokio::time::sleep(NEKO_READY_POLL).await;
}
if let Some(ref cookie) = session_cookie {
let mut state = self.inner.lock().await;
state.neko_session_cookie = Some(cookie.clone());
}
let state = self.inner.lock().await;
let state_window_id = state.window_id.clone();
drop(state);
Ok(DesktopStreamStatusResponse {
active: true,
window_id: state_window_id,
process_id: Some(process_id_clone),
})
}
pub async fn stop(&self) -> DesktopStreamStatusResponse {
let mut state = self.inner.lock().await;
if let Some(ref process_id) = state.process_id.take() {
let _ = self
.process_runtime
.stop_process(process_id, Some(3000))
.await;
}
state.active = false;
state.neko_base_url = None;
state.neko_session_cookie = None;
state.display = None;
state.resolution = None;
state.window_id = None;
DesktopStreamStatusResponse {
active: false,
window_id: None,
process_id: None,
}
}
pub async fn status(&self) -> DesktopStreamStatusResponse {
let state = self.inner.lock().await;
DesktopStreamStatusResponse {
active: state.active,
window_id: state.window_id.clone(),
process_id: state.process_id.clone(),
}
}
pub async fn ensure_active(&self) -> Result<(), SandboxError> {
if self.inner.lock().await.active {
Ok(())
} else {
Err(SandboxError::Conflict {
message: "desktop streaming is not active".to_string(),
})
}
}
pub async fn neko_ws_url(&self) -> Option<String> {
self.inner
.lock()
.await
.neko_base_url
.as_ref()
.map(|base| base.replace("http://", "ws://") + "/api/ws")
}
pub async fn neko_base_url(&self) -> Option<String> {
self.inner.lock().await.neko_base_url.clone()
}
pub async fn create_neko_session(&self) -> Option<String> {
let base_url = self.neko_base_url().await?;
let client = reqwest::Client::new();
let login_url = format!("{}/api/login", base_url);
let username = format!(
"user-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_nanos()
);
tracing::debug!(
"creating neko session: username={}, url={}",
username,
login_url
);
let resp = match client
.post(&login_url)
.json(&serde_json::json!({"username": username, "password": "admin"}))
.send()
.await
{
Ok(r) => r,
Err(e) => {
tracing::warn!("neko login request failed: {e}");
return None;
}
};
if !resp.status().is_success() {
tracing::warn!("neko login returned status {}", resp.status());
return None;
}
let cookie = resp
.headers()
.get("set-cookie")
.and_then(|v| v.to_str().ok())
.map(|v| v.split(';').next().unwrap_or(v).to_string());
let cookie = match cookie {
Some(c) => c,
None => {
tracing::warn!("neko login response missing set-cookie header");
return None;
}
};
tracing::debug!("neko session created: {}", username);
let control_url = format!("{}/api/room/control/take", base_url);
let _ = client
.post(&control_url)
.header("Cookie", &cookie)
.send()
.await;
Some(cookie)
}
pub async fn neko_session_cookie(&self) -> Option<String> {
self.inner.lock().await.neko_session_cookie.clone()
}
pub async fn resolution(&self) -> Option<DesktopResolution> {
self.inner.lock().await.resolution.clone()
}
pub async fn is_active(&self) -> bool {
self.inner.lock().await.active
}
pub async fn process_info(&self) -> Option<DesktopProcessInfo> {
let state = self.inner.lock().await;
let process_id = state.process_id.as_ref()?;
let snapshot = self.process_runtime.snapshot(process_id).await.ok()?;
Some(DesktopProcessInfo {
name: "neko".to_string(),
pid: snapshot.pid,
running: snapshot.status == crate::process_runtime::ProcessStatus::Running,
log_path: None,
})
}
}