pub mod backpressure;
pub mod checkpointing;
pub mod circuit_breaker;
pub mod cluster;
pub mod consensus;
pub mod coordinator;
pub mod discovery;
pub mod fault_tolerance;
pub mod heartbeat;
pub mod job_tracker;
pub mod leader_election;
pub mod load_balancer;
pub mod message_bus;
pub mod message_queue;
pub mod metrics_aggregator;
pub mod node_health;
pub mod node_registry;
pub mod node_topology;
pub mod partition;
pub mod pb;
pub mod raft_primitives;
pub mod replication;
pub mod resource_quota;
pub mod scheduler;
pub mod segment;
pub mod shard;
pub mod shard_map;
pub mod snapshot_store;
pub mod task_distribution;
pub mod task_priority_queue;
pub mod task_queue;
pub mod task_retry;
pub mod work_stealing;
pub mod worker;
use std::time::Duration;
use thiserror::Error;
use uuid::Uuid;
pub type Result<T> = std::result::Result<T, DistributedError>;
#[derive(Debug, Error)]
pub enum DistributedError {
#[error("Worker error: {0}")]
Worker(String),
#[error("Coordinator error: {0}")]
Coordinator(String),
#[error("Job error: {0}")]
Job(String),
#[error("Network error: {0}")]
Network(#[from] tonic::transport::Error),
#[error("gRPC status error: {0}")]
Status(#[from] tonic::Status),
#[error("Serialization error: {0}")]
Serialization(#[from] serde_json::Error),
#[error("IO error: {0}")]
Io(#[from] std::io::Error),
#[error("Discovery error: {0}")]
Discovery(String),
#[error("Scheduling error: {0}")]
Scheduling(String),
#[error("Segmentation error: {0}")]
Segmentation(String),
#[error("Timeout error")]
Timeout,
#[error("Invalid configuration: {0}")]
InvalidConfig(String),
#[error("Resource exhausted: {0}")]
ResourceExhausted(String),
#[error("Error: {0}")]
Other(Box<dyn std::error::Error + Send + Sync>),
}
impl From<Box<dyn std::error::Error + Send + Sync>> for DistributedError {
fn from(err: Box<dyn std::error::Error + Send + Sync>) -> Self {
DistributedError::Other(err)
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct DistributedConfig {
pub coordinator_addr: String,
pub max_retries: u32,
pub heartbeat_interval: Duration,
pub job_timeout: Duration,
pub max_concurrent_jobs: u32,
pub fault_tolerance: bool,
pub discovery_method: DiscoveryMethod,
}
impl Default for DistributedConfig {
fn default() -> Self {
Self {
coordinator_addr: "127.0.0.1:50051".to_string(),
max_retries: 3,
heartbeat_interval: Duration::from_secs(30),
job_timeout: Duration::from_secs(3600),
max_concurrent_jobs: 4,
fault_tolerance: true,
discovery_method: DiscoveryMethod::Static,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[allow(dead_code)]
pub enum DiscoveryMethod {
Static,
#[allow(clippy::upper_case_acronyms)]
MDNS,
Etcd,
Consul,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub enum SplitStrategy {
SegmentBased,
TileBased,
GopBased,
}
impl From<SplitStrategy> for i32 {
fn from(strategy: SplitStrategy) -> Self {
match strategy {
SplitStrategy::SegmentBased => 0,
SplitStrategy::TileBased => 1,
SplitStrategy::GopBased => 2,
}
}
}
#[derive(
Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, serde::Serialize, serde::Deserialize,
)]
pub enum JobPriority {
Low = 0,
Normal = 1,
High = 2,
Critical = 3,
}
impl From<JobPriority> for u32 {
fn from(priority: JobPriority) -> Self {
priority as u32
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct DistributedJob {
pub id: Uuid,
pub task_id: Uuid,
pub source_url: String,
pub codec: String,
pub strategy: SplitStrategy,
pub priority: JobPriority,
pub params: EncodingParams,
pub output_url: String,
pub deadline: Option<i64>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct EncodingParams {
pub bitrate: Option<u32>,
pub width: Option<u32>,
pub height: Option<u32>,
pub preset: Option<String>,
pub profile: Option<String>,
pub crf: Option<u32>,
pub extra_params: std::collections::HashMap<String, String>,
}
impl Default for EncodingParams {
fn default() -> Self {
Self {
bitrate: None,
width: None,
height: None,
preset: Some("medium".to_string()),
profile: None,
crf: Some(23),
extra_params: std::collections::HashMap::new(),
}
}
}
pub struct DistributedEncoder {
config: DistributedConfig,
}
impl DistributedEncoder {
#[must_use]
pub fn new(config: DistributedConfig) -> Self {
Self { config }
}
#[must_use]
pub fn with_defaults() -> Self {
Self::new(DistributedConfig::default())
}
#[must_use]
pub fn config(&self) -> &DistributedConfig {
&self.config
}
pub async fn submit_job(&self, job: DistributedJob) -> Result<Uuid> {
tracing::info!(
"Submitting job {} to coordinator at {}",
job.id,
self.config.coordinator_addr
);
Ok(job.id)
}
pub async fn job_status(&self, job_id: Uuid) -> Result<JobStatus> {
tracing::debug!("Querying status for job {}", job_id);
Ok(JobStatus::Pending)
}
pub async fn cancel_job(&self, job_id: Uuid) -> Result<()> {
tracing::info!("Cancelling job {}", job_id);
Ok(())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub enum JobStatus {
Pending,
Assigned,
InProgress,
Completed,
Failed,
Cancelled,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_default_config() {
let config = DistributedConfig::default();
assert_eq!(config.coordinator_addr, "127.0.0.1:50051");
assert_eq!(config.max_retries, 3);
assert_eq!(config.max_concurrent_jobs, 4);
}
#[test]
fn test_encoder_creation() {
let encoder = DistributedEncoder::with_defaults();
assert_eq!(encoder.config().coordinator_addr, "127.0.0.1:50051");
}
#[test]
fn test_job_priority_ordering() {
assert!(JobPriority::Critical > JobPriority::High);
assert!(JobPriority::High > JobPriority::Normal);
assert!(JobPriority::Normal > JobPriority::Low);
}
}