refactor: 精简运行时重复初始化

统一服务与 CLI 的数据库初始化入口,由视频流管理器负责 WebRTC 采集源同步,并将 MSD 目录创建收敛到控制器内部。
This commit is contained in:
mofeng-git
2026-08-26 11:37:39 +08:00
parent 486f3887c2
commit 0054100414
6 changed files with 70 additions and 81 deletions

View File

@@ -1,3 +1,36 @@
mod pool; mod pool;
use std::path::Path;
use crate::error::Result;
pub use pool::DatabasePool; pub use pool::DatabasePool;
/// Open the application database stored in `data_dir` and ensure its schema exists.
pub async fn open_database_pool(data_dir: &Path) -> Result<DatabasePool> {
let db = DatabasePool::new(&data_dir.join("one-kvm.db")).await?;
db.init_schema().await?;
Ok(db)
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn open_database_pool_creates_data_dir_and_initializes_schema() {
let temp_dir = tempfile::tempdir().unwrap();
let data_dir = temp_dir.path().join("nested").join("data");
let db = open_database_pool(&data_dir).await.unwrap();
assert!(data_dir.join("one-kvm.db").is_file());
let users_table: Option<String> = sqlx::query_scalar(
"SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'users'",
)
.fetch_optional(db.pool())
.await
.unwrap();
assert_eq!(users_table.as_deref(), Some("users"));
}
}

View File

@@ -2,7 +2,7 @@ use std::collections::HashSet;
use std::future::Future; use std::future::Future;
use std::io::Write; use std::io::Write;
use std::net::{IpAddr, SocketAddr}; use std::net::{IpAddr, SocketAddr};
use std::path::{Path, PathBuf}; use std::path::PathBuf;
use axum_server::tls_rustls::RustlsConfig; use axum_server::tls_rustls::RustlsConfig;
use clap::{Args, Parser, Subcommand, ValueEnum}; use clap::{Args, Parser, Subcommand, ValueEnum};
@@ -12,7 +12,7 @@ use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};
use one_kvm::auth::{SessionStore, TwoFactorService, UserStore}; use one_kvm::auth::{SessionStore, TwoFactorService, UserStore};
use one_kvm::config; use one_kvm::config;
use one_kvm::db::DatabasePool; use one_kvm::db::open_database_pool;
use one_kvm::platform::PlatformCapabilities; use one_kvm::platform::PlatformCapabilities;
use one_kvm::runtime::{RuntimeBuilder, WebConfigOverrides}; use one_kvm::runtime::{RuntimeBuilder, WebConfigOverrides};
use one_kvm::state::ShutdownAction; use one_kvm::state::ShutdownAction;
@@ -298,13 +298,6 @@ async fn shutdown_signal() -> anyhow::Result<()> {
Ok(()) Ok(())
} }
async fn open_database_pool(data_dir: &Path) -> anyhow::Result<DatabasePool> {
let db_path = data_dir.join("one-kvm.db");
let db = DatabasePool::new(&db_path).await?;
db.init_schema().await?;
Ok(db)
}
async fn run_servers_until_shutdown<F, E>( async fn run_servers_until_shutdown<F, E>(
mut servers: FuturesUnordered<F>, mut servers: FuturesUnordered<F>,
shutdown_signal: impl Future<Output = ShutdownAction>, shutdown_signal: impl Future<Output = ShutdownAction>,
@@ -348,7 +341,6 @@ fn restart_current_process(exe_path: Option<PathBuf>) -> anyhow::Result<()> {
} }
async fn run_cli_command(command: CliCommand, data_dir: PathBuf) -> anyhow::Result<()> { async fn run_cli_command(command: CliCommand, data_dir: PathBuf) -> anyhow::Result<()> {
tokio::fs::create_dir_all(&data_dir).await?;
let db = open_database_pool(&data_dir).await?; let db = open_database_pool(&data_dir).await?;
let users = UserStore::new(db.clone_pool()); let users = UserStore::new(db.clone_pool());
let two_factor = TwoFactorService::new(db.clone_pool()); let two_factor = TwoFactorService::new(db.clone_pool());

View File

@@ -62,12 +62,8 @@ impl MsdController {
), ),
} }
if let Err(e) = std::fs::create_dir_all(&self.images_path) { tokio::fs::create_dir_all(&self.images_path).await?;
warn!("Failed to create images directory: {}", e); tokio::fs::create_dir_all(&self.ventoy_dir).await?;
}
if let Err(e) = std::fs::create_dir_all(&self.ventoy_dir) {
warn!("Failed to create ventoy directory: {}", e);
}
info!("Fetching MSD function from OtgService"); info!("Fetching MSD function from OtgService");
let msd_func = self let msd_func = self

View File

@@ -8,7 +8,7 @@ use crate::audio::{AudioController, AudioControllerConfig, AudioQuality};
use crate::auth::{SessionStore, TwoFactorService, UserStore}; use crate::auth::{SessionStore, TwoFactorService, UserStore};
use crate::computer_use::ComputerUseManager; use crate::computer_use::ComputerUseManager;
use crate::config::{self, AppConfig, ConfigStore}; use crate::config::{self, AppConfig, ConfigStore};
use crate::db::DatabasePool; use crate::db::{open_database_pool, DatabasePool};
use crate::events::EventBus; use crate::events::EventBus;
use crate::extensions::ExtensionManager; use crate::extensions::ExtensionManager;
use crate::hid::{HidBackendType, HidController}; use crate::hid::{HidBackendType, HidController};
@@ -126,8 +126,6 @@ impl RuntimeBuilder {
} }
} }
connect_capture_to_webrtc(&streamer, &webrtc).await;
let stream_manager = VideoStreamManager::with_webrtc_streamer( let stream_manager = VideoStreamManager::with_webrtc_streamer(
streamer.clone(), streamer.clone(),
webrtc.clone() as Arc<dyn crate::video::traits::VideoOutput>, webrtc.clone() as Arc<dyn crate::video::traits::VideoOutput>,
@@ -226,23 +224,19 @@ impl ApplicationRuntime {
async fn load_runtime_config( async fn load_runtime_config(
data_dir: &Path, data_dir: &Path,
) -> anyhow::Result<(DatabasePool, ConfigStore, AppConfig)> { ) -> anyhow::Result<(DatabasePool, ConfigStore, AppConfig)> {
tokio::fs::create_dir_all(data_dir).await?; let db = open_database_pool(data_dir).await?;
let db_path = data_dir.join("one-kvm.db");
let db = DatabasePool::new(&db_path).await?;
db.init_schema().await?;
let config_store = ConfigStore::new(db.clone_pool())?; let config_store = ConfigStore::new(db.clone_pool())?;
config_store.load().await?; config_store.load().await?;
let mut config = (*config_store.get()).clone(); let mut config = (*config_store.get()).clone();
config.apply_platform_defaults(); config.apply_platform_defaults();
prepare_linux_runtime_dirs(data_dir, &config_store, &mut config).await?; normalize_msd_config(data_dir, &config_store, &mut config).await?;
Ok((db, config_store, config)) Ok((db, config_store, config))
} }
#[cfg(unix)] #[cfg(unix)]
async fn prepare_linux_runtime_dirs( async fn normalize_msd_config(
data_dir: &Path, data_dir: &Path,
config_store: &ConfigStore, config_store: &ConfigStore,
config: &mut AppConfig, config: &mut AppConfig,
@@ -263,19 +257,11 @@ async fn prepare_linux_runtime_dirs(
if msd_dir_updated { if msd_dir_updated {
config_store.set(config.clone()).await?; config_store.set(config.clone()).await?;
} }
let msd_dir = PathBuf::from(&config.msd.msd_dir);
if let Err(error) = tokio::fs::create_dir_all(msd_dir.join("images")).await {
tracing::warn!("Failed to create MSD images directory: {}", error);
}
if let Err(error) = tokio::fs::create_dir_all(msd_dir.join("ventoy")).await {
tracing::warn!("Failed to create MSD ventoy directory: {}", error);
}
Ok(()) Ok(())
} }
#[cfg(not(unix))] #[cfg(not(unix))]
async fn prepare_linux_runtime_dirs( async fn normalize_msd_config(
_data_dir: &Path, _data_dir: &Path,
_config_store: &ConfigStore, _config_store: &ConfigStore,
_config: &mut AppConfig, _config: &mut AppConfig,
@@ -499,27 +485,6 @@ async fn build_audio(config: &AppConfig, events: &Arc<EventBus>) -> Arc<AudioCon
controller controller
} }
async fn connect_capture_to_webrtc(streamer: &Arc<Streamer>, webrtc: &Arc<WebRtcStreamer>) {
let (device_path, resolution, format, fps, jpeg_quality) =
streamer.current_capture_config().await;
tracing::debug!(
"Initial video config: {}x{} {:?} @ {}fps",
resolution.width,
resolution.height,
format,
fps
);
webrtc.update_video_config(resolution, format, fps).await;
if let Some(device_path) = device_path {
webrtc
.set_capture_device(device_path, jpeg_quality, streamer.current_device().await)
.await;
tracing::debug!("WebRTC streamer configured for direct capture");
} else {
tracing::warn!("No capture device configured for WebRTC");
}
}
async fn connect_audio_recovery( async fn connect_audio_recovery(
audio: &Arc<AudioController>, audio: &Arc<AudioController>,
stream_manager: &Arc<VideoStreamManager>, stream_manager: &Arc<VideoStreamManager>,
@@ -595,4 +560,25 @@ mod tests {
assert_eq!(config.web.https_port, original_https_port); assert_eq!(config.web.https_port, original_https_port);
assert!(config.web.https_enabled); assert!(config.web.https_enabled);
} }
#[cfg(unix)]
#[tokio::test]
async fn normalizing_disabled_msd_does_not_create_module_directories() {
let temp_dir = tempfile::tempdir().unwrap();
let data_dir = temp_dir.path().join("data");
let msd_dir = temp_dir.path().join("disabled-msd");
let db = open_database_pool(&data_dir).await.unwrap();
let config_store = ConfigStore::new(db.clone_pool()).unwrap();
config_store.load().await.unwrap();
let mut config = (*config_store.get()).clone();
config.msd.enabled = false;
config.msd.msd_dir = msd_dir.to_string_lossy().into_owned();
normalize_msd_config(&data_dir, &config_store, &mut config)
.await
.unwrap();
assert!(!msd_dir.join("images").exists());
assert!(!msd_dir.join("ventoy").exists());
}
} }

View File

@@ -258,14 +258,6 @@ impl UsbCoordinator {
let inquiry_changed = old_config.flash_inquiry_string != new_config.flash_inquiry_string let inquiry_changed = old_config.flash_inquiry_string != new_config.flash_inquiry_string
|| old_config.cdrom_inquiry_string != new_config.cdrom_inquiry_string; || old_config.cdrom_inquiry_string != new_config.cdrom_inquiry_string;
let msd_dir = new_config.msd_dir_path();
if let Err(error) = std::fs::create_dir_all(msd_dir.join("images")) {
tracing::warn!("Failed to create MSD images directory: {}", error);
}
if let Err(error) = std::fs::create_dir_all(msd_dir.join("ventoy")) {
tracing::warn!("Failed to create MSD ventoy directory: {}", error);
}
if !options.force && old_enabled == new_enabled && !directory_changed && !inquiry_changed { if !options.force && old_enabled == new_enabled && !directory_changed && !inquiry_changed {
tracing::info!("MSD configuration unchanged, no reload needed"); tracing::info!("MSD configuration unchanged, no reload needed");
return Ok(()); return Ok(());

View File

@@ -195,25 +195,15 @@ impl VideoStreamManager {
info!("Initializing video stream manager with mode: {:?}", mode); info!("Initializing video stream manager with mode: {:?}", mode);
*self.mode.write().await = mode.clone(); *self.mode.write().await = mode.clone();
// Check if streamer is already initialized (capturer exists) // A failed fixed-device configuration can leave the streamer in a transient
let needs_init = self.streamer.state().await == StreamerState::Uninitialized; // state without a capture device. Treat that the same as an uninitialized
// streamer so the advertised auto-detection fallback actually runs.
let state = self.streamer.state().await;
let (device_path, _, _, _, _) = self.streamer.current_capture_config().await;
let needs_init = state == StreamerState::Uninitialized || device_path.is_none();
if needs_init { if needs_init {
match mode { self.streamer.init_auto().await?;
StreamMode::Mjpeg => {
// Initialize MJPEG streamer
if let Err(e) = self.streamer.init_auto().await {
warn!("Failed to auto-initialize MJPEG streamer: {}", e);
}
}
StreamMode::WebRTC => {
// WebRTC is initialized on-demand when clients connect
// But we still need to initialize the video capture
if let Err(e) = self.streamer.init_auto().await {
warn!("Failed to auto-initialize video capture for WebRTC: {}", e);
}
}
}
} }
self.sync_webrtc_capture_source("after init").await; self.sync_webrtc_capture_source("after init").await;