1pub mod audit_log;
11pub mod backpressure;
12pub mod checkpointing;
13pub mod circuit_breaker;
14pub mod cluster;
15pub mod compaction;
16pub mod connection_pool;
17pub mod consensus;
18pub mod coordinator;
19pub mod discovery;
20pub mod distributed_enhancements;
21pub mod fault_tolerance;
22pub mod heartbeat;
23pub mod job_dag;
24pub mod job_preemption;
25pub mod job_tracker;
26pub mod leader_election;
27pub mod lease;
28pub mod load_balancer;
29pub mod membership;
30pub mod message_bus;
31pub mod message_queue;
32pub mod metrics_aggregator;
33pub mod node_health;
34pub mod node_registry;
35pub mod node_topology;
36pub mod partition;
37pub mod pb;
38pub mod raft_primitives;
39pub mod replication;
40pub mod resource_quota;
41pub mod scheduler;
42pub mod segment;
43pub mod segment_merge;
44pub mod shard;
45pub mod shard_map;
46pub mod snapshot_store;
47pub mod task_distribution;
48pub mod task_priority_queue;
49pub mod task_queue;
50pub mod task_retry;
51pub mod twopc;
52pub mod weighted_round_robin;
53pub mod work_stealing;
54pub mod worker;
55pub mod worker_draining;
56
57use std::collections::HashMap;
58use std::sync::Arc;
59use std::time::{Duration, Instant};
60use thiserror::Error;
61use tokio::sync::RwLock;
62use uuid::Uuid;
63
64pub type Result<T> = std::result::Result<T, DistributedError>;
66
67#[derive(Debug, Error)]
69pub enum DistributedError {
70 #[error("Worker error: {0}")]
71 Worker(String),
72
73 #[error("Coordinator error: {0}")]
74 Coordinator(String),
75
76 #[error("Job error: {0}")]
77 Job(String),
78
79 #[error("Network error: {0}")]
80 Network(#[from] tonic::transport::Error),
81
82 #[error("gRPC status error: {0}")]
83 Status(#[from] tonic::Status),
84
85 #[error("Serialization error: {0}")]
86 Serialization(#[from] serde_json::Error),
87
88 #[error("IO error: {0}")]
89 Io(#[from] std::io::Error),
90
91 #[error("Discovery error: {0}")]
92 Discovery(String),
93
94 #[error("Scheduling error: {0}")]
95 Scheduling(String),
96
97 #[error("Segmentation error: {0}")]
98 Segmentation(String),
99
100 #[error("Timeout error")]
101 Timeout,
102
103 #[error("Invalid configuration: {0}")]
104 InvalidConfig(String),
105
106 #[error("Resource exhausted: {0}")]
107 ResourceExhausted(String),
108
109 #[error("Error: {0}")]
110 Other(Box<dyn std::error::Error + Send + Sync>),
111}
112
113impl From<Box<dyn std::error::Error + Send + Sync>> for DistributedError {
114 fn from(err: Box<dyn std::error::Error + Send + Sync>) -> Self {
115 DistributedError::Other(err)
116 }
117}
118
119#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
121pub struct DistributedConfig {
122 pub coordinator_addr: String,
124
125 pub max_retries: u32,
127
128 pub heartbeat_interval: Duration,
130
131 pub job_timeout: Duration,
133
134 pub max_concurrent_jobs: u32,
136
137 pub fault_tolerance: bool,
139
140 pub discovery_method: DiscoveryMethod,
142}
143
144impl Default for DistributedConfig {
145 fn default() -> Self {
146 Self {
147 coordinator_addr: "127.0.0.1:50051".to_string(),
148 max_retries: 3,
149 heartbeat_interval: Duration::from_secs(30),
150 job_timeout: Duration::from_secs(3600),
151 max_concurrent_jobs: 4,
152 fault_tolerance: true,
153 discovery_method: DiscoveryMethod::Static,
154 }
155 }
156}
157
158#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
160#[allow(dead_code)]
161pub enum DiscoveryMethod {
162 Static,
164 #[allow(clippy::upper_case_acronyms)]
166 MDNS,
167 Etcd,
169 Consul,
171}
172
173#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
175pub enum SplitStrategy {
176 SegmentBased,
178 TileBased,
180 GopBased,
182}
183
184impl From<SplitStrategy> for i32 {
185 fn from(strategy: SplitStrategy) -> Self {
186 match strategy {
187 SplitStrategy::SegmentBased => 0,
188 SplitStrategy::TileBased => 1,
189 SplitStrategy::GopBased => 2,
190 }
191 }
192}
193
194#[derive(
196 Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, serde::Serialize, serde::Deserialize,
197)]
198pub enum JobPriority {
199 Low = 0,
200 Normal = 1,
201 High = 2,
202 Critical = 3,
203}
204
205impl From<JobPriority> for u32 {
206 fn from(priority: JobPriority) -> Self {
207 priority as u32
208 }
209}
210
211#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
213pub struct DistributedJob {
214 pub id: Uuid,
216
217 pub task_id: Uuid,
219
220 pub source_url: String,
222
223 pub codec: String,
225
226 pub strategy: SplitStrategy,
228
229 pub priority: JobPriority,
231
232 pub params: EncodingParams,
234
235 pub output_url: String,
237
238 pub deadline: Option<i64>,
240}
241
242#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
244pub struct EncodingParams {
245 pub bitrate: Option<u32>,
246 pub width: Option<u32>,
247 pub height: Option<u32>,
248 pub preset: Option<String>,
249 pub profile: Option<String>,
250 pub crf: Option<u32>,
251 pub extra_params: std::collections::HashMap<String, String>,
252}
253
254impl Default for EncodingParams {
255 fn default() -> Self {
256 Self {
257 bitrate: None,
258 width: None,
259 height: None,
260 preset: Some("medium".to_string()),
261 profile: None,
262 crf: Some(23),
263 extra_params: std::collections::HashMap::new(),
264 }
265 }
266}
267
268#[derive(Debug, Clone)]
270struct JobRecord {
271 #[allow(dead_code)]
273 job: DistributedJob,
274 status: JobStatus,
276 submitted_at: Instant,
278 retries: u32,
280}
281
282pub struct DistributedEncoder {
289 config: DistributedConfig,
290 jobs: Arc<RwLock<HashMap<Uuid, JobRecord>>>,
292}
293
294impl DistributedEncoder {
295 #[must_use]
297 pub fn new(config: DistributedConfig) -> Self {
298 Self {
299 config,
300 jobs: Arc::new(RwLock::new(HashMap::new())),
301 }
302 }
303
304 #[must_use]
306 pub fn with_defaults() -> Self {
307 Self::new(DistributedConfig::default())
308 }
309
310 #[must_use]
312 pub fn config(&self) -> &DistributedConfig {
313 &self.config
314 }
315
316 pub async fn job_count(&self) -> usize {
318 self.jobs.read().await.len()
319 }
320
321 pub async fn active_job_count(&self) -> usize {
323 self.jobs
324 .read()
325 .await
326 .values()
327 .filter(|r| {
328 matches!(
329 r.status,
330 JobStatus::Pending | JobStatus::Assigned | JobStatus::InProgress
331 )
332 })
333 .count()
334 }
335
336 pub async fn submit_job(&self, job: DistributedJob) -> Result<Uuid> {
355 if job.source_url.is_empty() {
357 return Err(DistributedError::InvalidConfig(
358 "source_url must not be empty".to_string(),
359 ));
360 }
361 if job.output_url.is_empty() {
362 return Err(DistributedError::InvalidConfig(
363 "output_url must not be empty".to_string(),
364 ));
365 }
366 if job.codec.is_empty() {
367 return Err(DistributedError::InvalidConfig(
368 "codec must not be empty".to_string(),
369 ));
370 }
371
372 if let Some(deadline) = job.deadline {
374 let now = std::time::SystemTime::now()
375 .duration_since(std::time::UNIX_EPOCH)
376 .map_err(|e| DistributedError::Job(format!("System time error: {e}")))?;
377 if deadline < now.as_secs() as i64 {
378 return Err(DistributedError::Job(
379 "Job deadline is already in the past".to_string(),
380 ));
381 }
382 }
383
384 let mut jobs = self.jobs.write().await;
385
386 if jobs.contains_key(&job.id) {
388 return Err(DistributedError::Job(format!(
389 "Job with ID {} already exists",
390 job.id
391 )));
392 }
393
394 let active_count = jobs
396 .values()
397 .filter(|r| {
398 matches!(
399 r.status,
400 JobStatus::Pending | JobStatus::Assigned | JobStatus::InProgress
401 )
402 })
403 .count();
404
405 if active_count >= self.config.max_concurrent_jobs as usize {
406 return Err(DistributedError::ResourceExhausted(format!(
407 "Maximum concurrent jobs ({}) reached",
408 self.config.max_concurrent_jobs
409 )));
410 }
411
412 let job_id = job.id;
413
414 tracing::info!(
415 "Submitting job {} (codec={}, strategy={:?}, priority={:?}) to coordinator at {}",
416 job_id,
417 job.codec,
418 job.strategy,
419 job.priority,
420 self.config.coordinator_addr
421 );
422
423 jobs.insert(
424 job_id,
425 JobRecord {
426 job,
427 status: JobStatus::Pending,
428 submitted_at: Instant::now(),
429 retries: 0,
430 },
431 );
432
433 Ok(job_id)
434 }
435
436 pub async fn job_status(&self, job_id: Uuid) -> Result<JobStatus> {
446 tracing::debug!("Querying status for job {}", job_id);
447
448 let mut jobs = self.jobs.write().await;
449 let record = jobs
450 .get_mut(&job_id)
451 .ok_or_else(|| DistributedError::Job(format!("Job {job_id} not found")))?;
452
453 if matches!(
455 record.status,
456 JobStatus::Pending | JobStatus::Assigned | JobStatus::InProgress
457 ) && record.submitted_at.elapsed() > self.config.job_timeout
458 {
459 tracing::warn!(
460 "Job {} has timed out after {:?}",
461 job_id,
462 self.config.job_timeout
463 );
464 record.status = JobStatus::Failed;
465 }
466
467 Ok(record.status)
468 }
469
470 pub async fn cancel_job(&self, job_id: Uuid) -> Result<()> {
480 tracing::info!("Cancelling job {}", job_id);
481
482 let mut jobs = self.jobs.write().await;
483 let record = jobs
484 .get_mut(&job_id)
485 .ok_or_else(|| DistributedError::Job(format!("Job {job_id} not found")))?;
486
487 match record.status {
488 JobStatus::Completed => {
489 return Err(DistributedError::Job(format!(
490 "Job {job_id} is already completed and cannot be cancelled"
491 )));
492 }
493 JobStatus::Failed => {
494 return Err(DistributedError::Job(format!(
495 "Job {job_id} has already failed and cannot be cancelled"
496 )));
497 }
498 JobStatus::Cancelled => {
499 return Err(DistributedError::Job(format!(
500 "Job {job_id} is already cancelled"
501 )));
502 }
503 _ => {}
504 }
505
506 record.status = JobStatus::Cancelled;
507 Ok(())
508 }
509
510 pub async fn advance_job(&self, job_id: Uuid) -> Result<JobStatus> {
518 let mut jobs = self.jobs.write().await;
519 let record = jobs
520 .get_mut(&job_id)
521 .ok_or_else(|| DistributedError::Job(format!("Job {job_id} not found")))?;
522
523 record.status = match record.status {
524 JobStatus::Pending => JobStatus::Assigned,
525 JobStatus::Assigned => JobStatus::InProgress,
526 JobStatus::InProgress => JobStatus::Completed,
527 other => {
528 return Err(DistributedError::Job(format!(
529 "Cannot advance job in terminal state: {other:?}"
530 )));
531 }
532 };
533
534 Ok(record.status)
535 }
536
537 pub async fn fail_job(&self, job_id: Uuid) -> Result<()> {
543 let mut jobs = self.jobs.write().await;
544 let record = jobs
545 .get_mut(&job_id)
546 .ok_or_else(|| DistributedError::Job(format!("Job {job_id} not found")))?;
547
548 if matches!(record.status, JobStatus::Completed | JobStatus::Cancelled) {
549 return Err(DistributedError::Job(format!(
550 "Cannot fail job {job_id} in terminal state: {:?}",
551 record.status
552 )));
553 }
554
555 if self.config.fault_tolerance && record.retries < self.config.max_retries {
557 record.retries += 1;
558 record.status = JobStatus::Pending;
559 tracing::info!(
560 "Retrying job {} (attempt {}/{})",
561 job_id,
562 record.retries,
563 self.config.max_retries
564 );
565 } else {
566 record.status = JobStatus::Failed;
567 }
568
569 Ok(())
570 }
571
572 pub async fn job_retries(&self, job_id: Uuid) -> Result<u32> {
578 let jobs = self.jobs.read().await;
579 let record = jobs
580 .get(&job_id)
581 .ok_or_else(|| DistributedError::Job(format!("Job {job_id} not found")))?;
582 Ok(record.retries)
583 }
584
585 pub async fn list_jobs(&self) -> Vec<(Uuid, JobStatus)> {
587 self.jobs
588 .read()
589 .await
590 .iter()
591 .map(|(id, record)| (*id, record.status))
592 .collect()
593 }
594}
595
596#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
598pub enum JobStatus {
599 Pending,
600 Assigned,
601 InProgress,
602 Completed,
603 Failed,
604 Cancelled,
605}
606
607#[cfg(test)]
608mod tests {
609 use super::*;
610
611 fn make_job() -> DistributedJob {
612 DistributedJob {
613 id: Uuid::new_v4(),
614 task_id: Uuid::new_v4(),
615 source_url: "s3://bucket/input.mp4".to_string(),
616 codec: "av1".to_string(),
617 strategy: SplitStrategy::SegmentBased,
618 priority: JobPriority::Normal,
619 params: EncodingParams::default(),
620 output_url: "s3://bucket/output.mp4".to_string(),
621 deadline: None,
622 }
623 }
624
625 #[test]
626 fn test_default_config() {
627 let config = DistributedConfig::default();
628 assert_eq!(config.coordinator_addr, "127.0.0.1:50051");
629 assert_eq!(config.max_retries, 3);
630 assert_eq!(config.max_concurrent_jobs, 4);
631 }
632
633 #[test]
634 fn test_encoder_creation() {
635 let encoder = DistributedEncoder::with_defaults();
636 assert_eq!(encoder.config().coordinator_addr, "127.0.0.1:50051");
637 }
638
639 #[test]
640 fn test_job_priority_ordering() {
641 assert!(JobPriority::Critical > JobPriority::High);
642 assert!(JobPriority::High > JobPriority::Normal);
643 assert!(JobPriority::Normal > JobPriority::Low);
644 }
645
646 #[tokio::test]
647 async fn test_submit_and_query_job() {
648 let encoder = DistributedEncoder::with_defaults();
649 let job = make_job();
650 let job_id = job.id;
651
652 let returned_id = encoder
653 .submit_job(job)
654 .await
655 .expect("submit should succeed");
656 assert_eq!(returned_id, job_id);
657
658 let status = encoder
659 .job_status(job_id)
660 .await
661 .expect("status should succeed");
662 assert_eq!(status, JobStatus::Pending);
663 }
664
665 #[tokio::test]
666 async fn test_submit_rejects_empty_source_url() {
667 let encoder = DistributedEncoder::with_defaults();
668 let mut job = make_job();
669 job.source_url = String::new();
670
671 let result = encoder.submit_job(job).await;
672 assert!(result.is_err());
673 }
674
675 #[tokio::test]
676 async fn test_submit_rejects_empty_codec() {
677 let encoder = DistributedEncoder::with_defaults();
678 let mut job = make_job();
679 job.codec = String::new();
680
681 let result = encoder.submit_job(job).await;
682 assert!(result.is_err());
683 }
684
685 #[tokio::test]
686 async fn test_submit_rejects_empty_output_url() {
687 let encoder = DistributedEncoder::with_defaults();
688 let mut job = make_job();
689 job.output_url = String::new();
690
691 let result = encoder.submit_job(job).await;
692 assert!(result.is_err());
693 }
694
695 #[tokio::test]
696 async fn test_submit_rejects_duplicate_id() {
697 let encoder = DistributedEncoder::with_defaults();
698 let job = make_job();
699 let dup = job.clone();
700
701 encoder
702 .submit_job(job)
703 .await
704 .expect("first submit should succeed");
705 let result = encoder.submit_job(dup).await;
706 assert!(result.is_err());
707 }
708
709 #[tokio::test]
710 async fn test_cancel_job() {
711 let encoder = DistributedEncoder::with_defaults();
712 let job = make_job();
713 let job_id = job.id;
714
715 encoder
716 .submit_job(job)
717 .await
718 .expect("submit should succeed");
719 encoder
720 .cancel_job(job_id)
721 .await
722 .expect("cancel should succeed");
723
724 let status = encoder
725 .job_status(job_id)
726 .await
727 .expect("status should succeed");
728 assert_eq!(status, JobStatus::Cancelled);
729 }
730
731 #[tokio::test]
732 async fn test_cancel_nonexistent_job_fails() {
733 let encoder = DistributedEncoder::with_defaults();
734 let result = encoder.cancel_job(Uuid::new_v4()).await;
735 assert!(result.is_err());
736 }
737
738 #[tokio::test]
739 async fn test_cancel_completed_job_fails() {
740 let encoder = DistributedEncoder::with_defaults();
741 let job = make_job();
742 let job_id = job.id;
743
744 encoder
745 .submit_job(job)
746 .await
747 .expect("submit should succeed");
748 encoder
750 .advance_job(job_id)
751 .await
752 .expect("advance should succeed"); encoder
754 .advance_job(job_id)
755 .await
756 .expect("advance should succeed"); encoder
758 .advance_job(job_id)
759 .await
760 .expect("advance should succeed"); let result = encoder.cancel_job(job_id).await;
763 assert!(result.is_err());
764 }
765
766 #[tokio::test]
767 async fn test_advance_job_lifecycle() {
768 let encoder = DistributedEncoder::with_defaults();
769 let job = make_job();
770 let job_id = job.id;
771
772 encoder
773 .submit_job(job)
774 .await
775 .expect("submit should succeed");
776
777 let s1 = encoder
778 .advance_job(job_id)
779 .await
780 .expect("advance should succeed");
781 assert_eq!(s1, JobStatus::Assigned);
782
783 let s2 = encoder
784 .advance_job(job_id)
785 .await
786 .expect("advance should succeed");
787 assert_eq!(s2, JobStatus::InProgress);
788
789 let s3 = encoder
790 .advance_job(job_id)
791 .await
792 .expect("advance should succeed");
793 assert_eq!(s3, JobStatus::Completed);
794
795 let result = encoder.advance_job(job_id).await;
797 assert!(result.is_err());
798 }
799
800 #[tokio::test]
801 async fn test_fail_job_with_retry() {
802 let config = DistributedConfig {
803 max_retries: 2,
804 fault_tolerance: true,
805 ..DistributedConfig::default()
806 };
807 let encoder = DistributedEncoder::new(config);
808 let job = make_job();
809 let job_id = job.id;
810
811 encoder
812 .submit_job(job)
813 .await
814 .expect("submit should succeed");
815
816 encoder.fail_job(job_id).await.expect("fail should succeed");
818 let status = encoder
819 .job_status(job_id)
820 .await
821 .expect("status should succeed");
822 assert_eq!(status, JobStatus::Pending);
823 let retries = encoder
824 .job_retries(job_id)
825 .await
826 .expect("retries should succeed");
827 assert_eq!(retries, 1);
828
829 encoder.fail_job(job_id).await.expect("fail should succeed");
831 let retries = encoder
832 .job_retries(job_id)
833 .await
834 .expect("retries should succeed");
835 assert_eq!(retries, 2);
836
837 encoder.fail_job(job_id).await.expect("fail should succeed");
839 let status = encoder
840 .job_status(job_id)
841 .await
842 .expect("status should succeed");
843 assert_eq!(status, JobStatus::Failed);
844 }
845
846 #[tokio::test]
847 async fn test_fail_without_fault_tolerance() {
848 let config = DistributedConfig {
849 fault_tolerance: false,
850 ..DistributedConfig::default()
851 };
852 let encoder = DistributedEncoder::new(config);
853 let job = make_job();
854 let job_id = job.id;
855
856 encoder
857 .submit_job(job)
858 .await
859 .expect("submit should succeed");
860 encoder.fail_job(job_id).await.expect("fail should succeed");
861
862 let status = encoder
863 .job_status(job_id)
864 .await
865 .expect("status should succeed");
866 assert_eq!(status, JobStatus::Failed);
867 }
868
869 #[tokio::test]
870 async fn test_concurrency_limit() {
871 let config = DistributedConfig {
872 max_concurrent_jobs: 2,
873 ..DistributedConfig::default()
874 };
875 let encoder = DistributedEncoder::new(config);
876
877 encoder
878 .submit_job(make_job())
879 .await
880 .expect("first should succeed");
881 encoder
882 .submit_job(make_job())
883 .await
884 .expect("second should succeed");
885
886 let result = encoder.submit_job(make_job()).await;
888 assert!(result.is_err());
889 }
890
891 #[tokio::test]
892 async fn test_concurrency_freed_after_cancel() {
893 let config = DistributedConfig {
894 max_concurrent_jobs: 1,
895 ..DistributedConfig::default()
896 };
897 let encoder = DistributedEncoder::new(config);
898
899 let job = make_job();
900 let job_id = job.id;
901 encoder.submit_job(job).await.expect("first should succeed");
902
903 assert!(encoder.submit_job(make_job()).await.is_err());
905
906 encoder
908 .cancel_job(job_id)
909 .await
910 .expect("cancel should succeed");
911
912 encoder
914 .submit_job(make_job())
915 .await
916 .expect("after cancel should succeed");
917 }
918
919 #[tokio::test]
920 async fn test_list_jobs() {
921 let encoder = DistributedEncoder::with_defaults();
922 let j1 = make_job();
923 let j2 = make_job();
924 let id1 = j1.id;
925 let id2 = j2.id;
926
927 encoder.submit_job(j1).await.expect("submit should succeed");
928 encoder.submit_job(j2).await.expect("submit should succeed");
929
930 let jobs = encoder.list_jobs().await;
931 assert_eq!(jobs.len(), 2);
932
933 let ids: Vec<Uuid> = jobs.iter().map(|(id, _)| *id).collect();
934 assert!(ids.contains(&id1));
935 assert!(ids.contains(&id2));
936 }
937
938 #[tokio::test]
939 async fn test_job_count() {
940 let encoder = DistributedEncoder::with_defaults();
941 assert_eq!(encoder.job_count().await, 0);
942
943 encoder
944 .submit_job(make_job())
945 .await
946 .expect("submit should succeed");
947 assert_eq!(encoder.job_count().await, 1);
948 assert_eq!(encoder.active_job_count().await, 1);
949 }
950
951 #[tokio::test]
952 async fn test_status_nonexistent_job_fails() {
953 let encoder = DistributedEncoder::with_defaults();
954 let result = encoder.job_status(Uuid::new_v4()).await;
955 assert!(result.is_err());
956 }
957
958 #[tokio::test]
959 async fn test_job_timeout_detection() {
960 let config = DistributedConfig {
961 job_timeout: Duration::from_millis(1),
962 ..DistributedConfig::default()
963 };
964 let encoder = DistributedEncoder::new(config);
965 let job = make_job();
966 let job_id = job.id;
967
968 encoder
969 .submit_job(job)
970 .await
971 .expect("submit should succeed");
972
973 tokio::time::sleep(Duration::from_millis(10)).await;
975
976 let status = encoder
977 .job_status(job_id)
978 .await
979 .expect("status should succeed");
980 assert_eq!(status, JobStatus::Failed);
981 }
982
983 #[tokio::test]
984 async fn test_submit_past_deadline_rejected() {
985 let encoder = DistributedEncoder::with_defaults();
986 let mut job = make_job();
987 job.deadline = Some(0); let result = encoder.submit_job(job).await;
990 assert!(result.is_err());
991 }
992}