Skip to main content

scirs2_core/distributed/
mod.rs

1//! Production-grade distributed computing infrastructure
2//!
3//! This module provides comprehensive distributed computing capabilities
4//! for SciRS2 Core 1.0, including distributed arrays, cluster management,
5//! fault tolerance, and scalable computation orchestration.
6//!
7//! ## Distributed Primitives
8//!
9//! Low-level local primitives (work queues, worker pools, map-reduce) are
10//! provided in [`primitives`].  These complement the cluster-management
11//! machinery with simple, channel-based concurrency helpers.
12
13pub mod array;
14pub mod cluster;
15pub mod communication;
16pub mod compression;
17pub mod fault_tolerance;
18pub mod load_balancing;
19pub mod lock_free;
20pub mod orchestration;
21pub mod par_iter;
22pub mod parallel_scan;
23pub mod param_server;
24pub mod primitives;
25pub mod scheduler;
26pub mod task_graph;
27
28// Array operations
29pub use array::{DistributedArray, DistributedArrayManager};
30
31// Cluster management
32pub use cluster::{
33    initialize_cluster_manager, BackoffStrategy, ClusterConfiguration, ClusterEventLog,
34    ClusterHealth, ClusterManager, ClusterState, ComputeCapacity, DistributedTask,
35    NodeCapabilities, NodeInfo as ClusterNodeInfo, NodeMetadata, NodeStatus, NodeType,
36    ResourceRequirements, RetryPolicy, TaskId, TaskParameters, TaskPriority as ClusterTaskPriority,
37    TaskType,
38};
39
40// Communication
41pub use communication::{
42    CommunicationEndpoint, CommunicationManager, DistributedMessage, HeartbeatHandler,
43    MessageHandler,
44};
45
46// Fault tolerance
47pub use fault_tolerance::{
48    initialize_fault_tolerance, ClusterHealthSummary, FaultDetectionStrategy,
49    FaultToleranceManager, NodeHealth as FaultNodeHealth, NodeInfo as FaultNodeInfo,
50    RecoveryStrategy,
51};
52
53// Load balancing
54pub use load_balancing::{
55    LoadBalancer as DistributedLoadBalancer, LoadBalancingStats, LoadBalancingStrategy,
56    NodeLoad as LoadBalancerNodeLoad, TaskAssignment as LoadBalancerTaskAssignment,
57};
58
59// Orchestration
60pub use orchestration::{
61    OrchestrationEngine, OrchestrationStats, OrchestratorNode, Task as OrchestrationTask,
62    TaskPriority as OrchestrationTaskPriority, TaskStatus as OrchestrationTaskStatus, Workflow,
63    WorkflowStatus,
64};
65
66// Scheduler
67pub use scheduler::{
68    initialize_distributed_scheduler, CompletedTask, DistributedScheduler, ExecutionTracker,
69    FailedTask, LoadBalancer as SchedulerLoadBalancer,
70    LoadBalancingStrategy as SchedulerLoadBalancingStrategy, NodeLoad as SchedulerNodeLoad,
71    SchedulingAlgorithm, SchedulingPolicies, TaskAssignment as SchedulerTaskAssignment, TaskQueue,
72};
73
74// Distributed primitives (work queue, worker pool, map-reduce, resource monitor)
75pub use primitives::{
76    chunked_parallel_process, distributed_map, distributed_map_reduce, try_distributed_map,
77    try_distributed_map_reduce, DistributedError, DistributedSliceExt, ResourceMonitor, WorkQueue,
78    WorkReceiver, WorkerPool,
79};
80
81// Lock-free data structures
82pub use lock_free::{LockFreeCounter, LockFreeQueue, LockFreeStack};
83
84// Task graph executor
85pub use task_graph::TaskGraph;
86
87// Parallel scan / prefix sum
88pub use parallel_scan::{
89    parallel_prefix_max, parallel_prefix_min, parallel_prefix_sum, parallel_prefix_sum_exclusive,
90    parallel_prefix_sum_f64, parallel_prefix_sum_i64, parallel_scan, parallel_scan_exclusive,
91    segmented_prefix_sum, try_parallel_prefix_sum, try_parallel_scan,
92};
93
94// Parallel iterator combinators
95pub use par_iter::{
96    par_all, par_any, par_filter, par_filter_map, par_fold, par_for_each, par_map, par_sort,
97    par_sort_by, try_par_fold, try_par_map,
98};
99
100/// Initialize distributed computing infrastructure
101#[allow(dead_code)]
102pub fn initialize_distributed_computing() -> crate::error::CoreResult<()> {
103    cluster::initialize_cluster_manager()?;
104    scheduler::initialize_distributed_scheduler()?;
105    fault_tolerance::initialize_fault_tolerance()?;
106    Ok(())
107}
108
109/// Get distributed system status
110#[allow(dead_code)]
111pub fn get_distributed_status() -> crate::error::CoreResult<DistributedSystemStatus> {
112    let cluster_manager = cluster::ClusterManager::global()?;
113    let scheduler = scheduler::DistributedScheduler::global()?;
114
115    Ok(DistributedSystemStatus {
116        cluster_health: cluster_manager.get_health()?,
117        active_nodes: cluster_manager.get_active_nodes()?.len(),
118        pending_tasks: scheduler.get_pending_task_count()?,
119        total_capacity: cluster_manager.get_total_capacity()?,
120    })
121}
122
123#[derive(Debug, Clone)]
124pub struct DistributedSystemStatus {
125    pub cluster_health: ClusterHealth,
126    pub active_nodes: usize,
127    pub pending_tasks: usize,
128    pub total_capacity: ComputeCapacity,
129}