Files
spacedrive/core/src/lib.rs
T
Jamie Pine 469fc7cbe2 refactor: complete migration to modular action system and remove ActionOutput enum
- Eliminated the centralized `ActionOutput` enum, achieving true modularity in action handling.
- Introduced the `ActionType` trait, allowing each action to define its own output type directly.
- Updated multiple actions to utilize native output types, enhancing type safety and reducing complexity.
- Enhanced the `ActionManager` with a generic `dispatch_action` method for improved action execution.
- Preserved existing infrastructure for validation, audit logging, and error handling throughout the migration.

These changes significantly improve the modularity and maintainability of the action system while ensuring a consistent API surface.
2025-09-09 02:03:17 -04:00

528 lines
15 KiB
Rust

#![allow(warnings)]
//! Spacedrive Core v2
//!
//! A unified, simplified architecture for cross-platform file management.
pub mod common;
pub mod config;
pub mod context;
pub mod cqrs;
pub mod crypto;
pub mod device;
pub mod domain;
pub mod filetype;
pub mod infra;
pub mod library;
pub mod location;
pub mod ops;
pub mod service;
pub mod testing;
pub mod volume;
use service::network::protocol::pairing::PairingProtocolHandler;
use service::network::utils::logging::NetworkLogger;
// Compatibility module for legacy networking references
pub mod networking {
pub use crate::service::network::*;
}
use crate::config::AppConfig;
use crate::context::CoreContext;
use crate::cqrs::{Query, QueryManager};
use crate::device::DeviceManager;
use crate::infra::action::manager::ActionManager;
use crate::infra::event::{Event, EventBus};
use crate::library::LibraryManager;
use crate::service::Services;
use crate::volume::{VolumeDetectionConfig, VolumeManager};
use std::path::PathBuf;
use std::sync::Arc;
use tokio::sync::{mpsc, RwLock};
use tracing::{error, info};
/// Pending pairing request information
#[derive(Debug, Clone)]
pub struct PendingPairingRequest {
pub request_id: uuid::Uuid,
pub device_id: uuid::Uuid,
pub device_name: String,
pub received_at: chrono::DateTime<chrono::Utc>,
}
/// Spacedrop request message
#[derive(serde::Serialize, serde::Deserialize)]
struct SpacedropRequest {
transfer_id: uuid::Uuid,
file_path: String,
sender_name: String,
message: Option<String>,
file_size: u64,
}
// NOTE: SimplePairingUI has been moved to CLI infrastructure
// See: src/infrastructure/cli/pairing_ui.rs for CLI-specific implementations
/// Bridge between networking events and core events
pub struct NetworkEventBridge {
network_events: mpsc::UnboundedReceiver<service::network::NetworkEvent>,
core_events: Arc<EventBus>,
}
impl NetworkEventBridge {
pub fn new(
network_events: mpsc::UnboundedReceiver<service::network::NetworkEvent>,
core_events: Arc<EventBus>,
) -> Self {
Self {
network_events,
core_events,
}
}
pub async fn run(mut self) {
while let Some(event) = self.network_events.recv().await {
if let Some(core_event) = self.translate_event(event) {
self.core_events.emit(core_event);
}
}
}
fn translate_event(&self, event: service::network::NetworkEvent) -> Option<Event> {
match event {
service::network::NetworkEvent::ConnectionEstablished { device_id, .. } => {
Some(Event::DeviceConnected {
device_id,
device_name: "Connected Device".to_string(),
})
}
service::network::NetworkEvent::ConnectionLost { device_id, .. } => {
Some(Event::DeviceDisconnected { device_id })
}
service::network::NetworkEvent::PairingCompleted {
device_id,
device_info,
} => Some(Event::DeviceConnected {
device_id,
device_name: device_info.device_name,
}),
_ => None, // Some events don't map to core events
}
}
}
/// The main context for all core operations
pub struct Core {
/// Application configuration
pub config: Arc<RwLock<AppConfig>>,
/// Device manager
pub device: Arc<DeviceManager>,
/// Library manager
pub libraries: Arc<LibraryManager>,
/// Volume manager
pub volumes: Arc<VolumeManager>,
/// Event bus for state changes
pub events: Arc<EventBus>,
/// Container for high-level services
pub services: Services,
/// Shared context for core components
pub context: Arc<CoreContext>,
}
impl Core {
/// Initialize a new Core instance with default data directory
pub async fn new() -> Result<Self, Box<dyn std::error::Error>> {
let data_dir = crate::config::default_data_dir()?;
Self::new_with_config(data_dir).await
}
/// Initialize a new Core instance with custom data directory
pub async fn new_with_config(data_dir: PathBuf) -> Result<Self, Box<dyn std::error::Error>> {
info!("Initializing Spacedrive Core at {:?}", data_dir);
// 1. Load or create app config
let config = AppConfig::load_or_create(&data_dir)?;
config.ensure_directories()?;
let config = Arc::new(RwLock::new(config));
// 2. Initialize device manager
let device = Arc::new(DeviceManager::init_with_path(&data_dir)?);
// Set the global device ID for legacy compatibility
common::utils::set_current_device_id(device.device_id()?);
// 3. Create event bus
let events = Arc::new(EventBus::default());
// 4. Initialize volume manager
let volume_config = VolumeDetectionConfig::default();
let device_id = device.device_id()?;
let volumes = Arc::new(VolumeManager::new(device_id, volume_config, events.clone()));
// 5. Initialize volume detection
// info!("Initializing volume detection...");
// match volumes.initialize().await {
// Ok(()) => info!("Volume manager initialized"),
// Err(e) => error!("Failed to initialize volume manager: {}", e),
// }
// 6. Initialize library manager with libraries directory
let libraries_dir = config.read().await.libraries_dir();
let libraries = Arc::new(LibraryManager::new_with_dir(libraries_dir, events.clone()));
// 7. Initialize library key manager
let library_key_manager =
Arc::new(crate::crypto::library_key_manager::LibraryKeyManager::new()?);
// 8. Register all job types
info!("Registering job types...");
crate::ops::register_all_jobs();
info!("Job types registered");
// 9. Create the context that will be shared with services
let mut context_inner = CoreContext::new(
events.clone(),
device.clone(),
libraries.clone(),
volumes.clone(),
library_key_manager.clone(),
);
// Set job logging configuration if enabled
let app_config = config.read().await;
if app_config.job_logging.enabled {
context_inner
.set_job_logging(app_config.job_logging.clone(), app_config.job_logs_dir());
}
drop(app_config);
let context = Arc::new(context_inner);
// 10. Initialize services first, passing them the context
let services = Services::new(context.clone());
// 11. Auto-load all libraries with context for job manager initialization
info!("Loading existing libraries...");
let loaded_libraries: Vec<Arc<crate::library::Library>> =
match libraries.load_all_with_context(context.clone()).await {
Ok(count) => {
info!("Loaded {} libraries", count);
libraries.list().await
}
Err(e) => {
error!("Failed to load libraries: {}", e);
vec![]
}
};
// Initialize sidecar manager for each loaded library
for library in &loaded_libraries {
info!("Initializing sidecar manager for library {}", library.id());
if let Err(e) = services.sidecar_manager.init_library(&library).await {
error!(
"Failed to initialize sidecar manager for library {}: {}",
library.id(),
e
);
} else {
// Run bootstrap scan
if let Err(e) = services.sidecar_manager.bootstrap_scan(&library).await {
error!(
"Failed to run sidecar bootstrap scan for library {}: {}",
library.id(),
e
);
}
}
}
info!("Starting background services...");
match services.start_all().await {
Ok(()) => info!("Background services started"),
Err(e) => error!("Failed to start services: {}", e),
}
// 12. Initialize ActionManager and set it in context
let action_manager = Arc::new(crate::infra::action::manager::ActionManager::new(
context.clone(),
));
context.set_action_manager(action_manager).await;
// 13. Emit startup event
events.emit(Event::CoreStarted);
Ok(Self {
config,
device,
libraries,
volumes,
events,
services,
context,
})
}
/// Get the application configuration
pub fn config(&self) -> Arc<RwLock<AppConfig>> {
self.config.clone()
}
/// Initialize networking using master key
pub async fn init_networking(&mut self) -> Result<(), Box<dyn std::error::Error>> {
self.init_networking_with_logger(Arc::new(service::network::SilentLogger))
.await
}
/// Initialize networking with custom logger
pub async fn init_networking_with_logger(
&mut self,
logger: Arc<dyn service::network::NetworkLogger>,
) -> Result<(), Box<dyn std::error::Error>> {
logger.info("Initializing networking...").await;
// Initialize networking service through the services container
let data_dir = self.config.read().await.data_dir.clone();
self.services
.init_networking(
self.device.clone(),
self.services.library_key_manager.clone(),
data_dir,
)
.await?;
// Start the networking service
self.services.start_networking().await?;
// Get the networking service for protocol registration
if let Some(networking_service) = self.services.networking() {
// Register default protocol handlers
self.register_default_protocols(&networking_service).await?;
// Set up event bridge to integrate with core event system
let event_bridge = NetworkEventBridge::new(
networking_service
.subscribe_events()
.await
.unwrap_or_else(|| {
let (_, rx) = tokio::sync::mpsc::unbounded_channel();
rx
}),
self.events.clone(),
);
tokio::spawn(event_bridge.run());
// Make networking service available to the context for other services
self.context.set_networking(networking_service).await;
}
logger.info("Networking initialized successfully").await;
Ok(())
}
/// Register default protocol handlers
async fn register_default_protocols(
&self,
networking: &service::network::NetworkingService,
) -> Result<(), Box<dyn std::error::Error>> {
let logger = std::sync::Arc::new(service::network::utils::logging::ConsoleLogger);
// Get command sender for the pairing handler's state machine
let command_sender = networking
.command_sender()
.ok_or("NetworkingEventLoop command sender not available")?
.clone();
// Get data directory from config
let data_dir = {
let config = self.config.read().await;
config.data_dir.clone()
};
let pairing_handler = Arc::new(
service::network::protocol::PairingProtocolHandler::new_with_persistence(
networking.identity().clone(),
networking.device_registry(),
logger.clone(),
command_sender,
data_dir,
),
);
// Try to load persisted sessions, but don't fail if there's an error
if let Err(e) = pairing_handler.load_persisted_sessions().await {
logger
.warn(&format!(
"Failed to load persisted pairing sessions: {}. Starting with empty sessions.",
e
))
.await;
}
// Start the state machine task for pairing
service::network::protocol::PairingProtocolHandler::start_state_machine_task(
pairing_handler.clone(),
);
// Start cleanup task for expired sessions
service::network::protocol::PairingProtocolHandler::start_cleanup_task(
pairing_handler.clone(),
);
let messaging_handler = service::network::protocol::MessagingProtocolHandler::new();
let mut file_transfer_handler =
service::network::protocol::FileTransferProtocolHandler::new_default(logger.clone());
// Inject device registry into file transfer handler for encryption
file_transfer_handler.set_device_registry(networking.device_registry());
let protocol_registry = networking.protocol_registry();
{
let mut registry = protocol_registry.write().await;
registry.register_handler(pairing_handler)?;
registry.register_handler(Arc::new(messaging_handler))?;
registry.register_handler(Arc::new(file_transfer_handler))?;
}
Ok(())
}
/// Initialize networking from Arc<Core> - for daemon use
pub async fn init_networking_shared(
core: Arc<Core>,
) -> Result<Arc<Core>, Box<dyn std::error::Error>> {
info!("Initializing networking for shared core...");
// Create a new Core with networking enabled
let mut new_core =
Core::new_with_config(core.config().read().await.data_dir.clone()).await?;
// Initialize networking on the new core
new_core.init_networking().await?;
info!("Networking initialized successfully for shared core");
Ok(Arc::new(new_core))
}
/// Get the networking service (if initialized)
pub fn networking(&self) -> Option<Arc<service::network::NetworkingService>> {
self.services.networking()
}
/// Get list of connected devices
pub async fn get_connected_devices(
&self,
) -> Result<Vec<uuid::Uuid>, Box<dyn std::error::Error>> {
Ok(self.services.device.get_connected_devices().await?)
}
/// Get detailed information about connected devices
pub async fn get_connected_devices_info(
&self,
) -> Result<Vec<service::network::DeviceInfo>, Box<dyn std::error::Error>> {
Ok(self.services.device.get_connected_devices_info().await?)
}
/// Add a location to the file system watcher
pub async fn add_watched_location(
&self,
location_id: uuid::Uuid,
library_id: uuid::Uuid,
path: std::path::PathBuf,
enabled: bool,
) -> Result<(), Box<dyn std::error::Error>> {
use crate::service::watcher::WatchedLocation;
let watched_location = WatchedLocation {
id: location_id,
library_id,
path,
enabled,
};
Ok(self
.services
.location_watcher
.add_location(watched_location)
.await?)
}
/// Remove a location from the file system watcher
pub async fn remove_watched_location(
&self,
location_id: uuid::Uuid,
) -> Result<(), Box<dyn std::error::Error>> {
Ok(self
.services
.location_watcher
.remove_location(location_id)
.await?)
}
/// Update file watching settings for a location
pub async fn update_watched_location(
&self,
location_id: uuid::Uuid,
enabled: bool,
) -> Result<(), Box<dyn std::error::Error>> {
Ok(self
.services
.location_watcher
.update_location(location_id, enabled)
.await?)
}
/// Get all currently watched locations
pub async fn get_watched_locations(&self) -> Vec<crate::service::watcher::WatchedLocation> {
self.services.location_watcher.get_watched_locations().await
}
/// Execute an action using the unified CQRS API.
///
/// This method provides a unified, type-safe entry point for all operations.
/// Actions return their natural output types - domain objects, job handles, etc.
pub async fn execute_action<A: crate::infra::action::ActionTrait>(&self, action: A) -> anyhow::Result<A::Output> {
let action_manager = ActionManager::new(self.context.clone());
action_manager.dispatch(action).await
.map_err(|e| anyhow::anyhow!("Action execution failed: {}", e))
}
/// Execute a query using the CQRS API.
///
/// This method provides a unified, type-safe entry point for all read operations.
/// It uses the QueryManager for consistent infrastructure (validation, logging, etc.).
pub async fn execute_query<Q: Query>(&self, query: Q) -> anyhow::Result<Q::Output> {
let query_manager = QueryManager::new(self.context.clone());
query_manager.dispatch(query).await
}
/// Shutdown the core gracefully
pub async fn shutdown(&self) -> Result<(), Box<dyn std::error::Error>> {
info!("Shutting down Spacedrive Core...");
// Networking service is stopped by services.stop_all()
// Stop all services
self.services.stop_all().await?;
// Stop volume monitoring
self.volumes.stop_monitoring().await;
// Close all libraries
self.libraries.close_all().await?;
// Save configuration
self.config.write().await.save()?;
// Emit shutdown event
self.events.emit(Event::CoreShutdown);
info!("Spacedrive Core shutdown complete");
Ok(())
}
}