1pub 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
28pub use array::{DistributedArray, DistributedArrayManager};
30
31pub 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
40pub use communication::{
42 CommunicationEndpoint, CommunicationManager, DistributedMessage, HeartbeatHandler,
43 MessageHandler,
44};
45
46pub use fault_tolerance::{
48 initialize_fault_tolerance, ClusterHealthSummary, FaultDetectionStrategy,
49 FaultToleranceManager, NodeHealth as FaultNodeHealth, NodeInfo as FaultNodeInfo,
50 RecoveryStrategy,
51};
52
53pub use load_balancing::{
55 LoadBalancer as DistributedLoadBalancer, LoadBalancingStats, LoadBalancingStrategy,
56 NodeLoad as LoadBalancerNodeLoad, TaskAssignment as LoadBalancerTaskAssignment,
57};
58
59pub use orchestration::{
61 OrchestrationEngine, OrchestrationStats, OrchestratorNode, Task as OrchestrationTask,
62 TaskPriority as OrchestrationTaskPriority, TaskStatus as OrchestrationTaskStatus, Workflow,
63 WorkflowStatus,
64};
65
66pub 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
74pub 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
81pub use lock_free::{LockFreeCounter, LockFreeQueue, LockFreeStack};
83
84pub use task_graph::TaskGraph;
86
87pub 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
94pub 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#[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#[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}