Skip to main content

oximedia_distributed/
lib.rs

1//! Distributed encoding coordinator for `OxiMedia`.
2//!
3//! This crate provides a distributed video encoding system with:
4//! - Central coordinator for job management
5//! - Worker nodes for distributed encoding
6//! - Multiple splitting strategies (segment, tile, GOP-based)
7//! - Load balancing and fault tolerance
8//! - gRPC-based communication
9
10pub 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
64/// Result type for distributed operations
65pub type Result<T> = std::result::Result<T, DistributedError>;
66
67/// Errors that can occur in distributed encoding
68#[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/// Configuration for the distributed encoder
120#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
121pub struct DistributedConfig {
122    /// Coordinator address
123    pub coordinator_addr: String,
124
125    /// Maximum number of retry attempts
126    pub max_retries: u32,
127
128    /// Heartbeat interval
129    pub heartbeat_interval: Duration,
130
131    /// Job timeout
132    pub job_timeout: Duration,
133
134    /// Maximum concurrent jobs per worker
135    pub max_concurrent_jobs: u32,
136
137    /// Enable fault tolerance
138    pub fault_tolerance: bool,
139
140    /// Worker discovery method
141    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/// Worker discovery methods
159#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
160#[allow(dead_code)]
161pub enum DiscoveryMethod {
162    /// Static configuration
163    Static,
164    /// Multicast DNS
165    #[allow(clippy::upper_case_acronyms)]
166    MDNS,
167    /// etcd-based discovery
168    Etcd,
169    /// Consul-based discovery
170    Consul,
171}
172
173/// Job splitting strategy
174#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
175pub enum SplitStrategy {
176    /// Split by time segments
177    SegmentBased,
178    /// Split by spatial tiles
179    TileBased,
180    /// Split by GOP (Group of Pictures)
181    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/// Job priority levels
195#[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/// Represents a distributed encoding job
212#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
213pub struct DistributedJob {
214    /// Unique job identifier
215    pub id: Uuid,
216
217    /// Task identifier (multiple jobs can belong to same task)
218    pub task_id: Uuid,
219
220    /// Source video URL
221    pub source_url: String,
222
223    /// Target codec
224    pub codec: String,
225
226    /// Splitting strategy
227    pub strategy: SplitStrategy,
228
229    /// Job priority
230    pub priority: JobPriority,
231
232    /// Encoding parameters
233    pub params: EncodingParams,
234
235    /// Output destination
236    pub output_url: String,
237
238    /// Deadline timestamp (Unix epoch)
239    pub deadline: Option<i64>,
240}
241
242/// Encoding parameters
243#[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/// Internal record for a submitted job, tracking its lifecycle.
269#[derive(Debug, Clone)]
270struct JobRecord {
271    /// The original job definition.
272    #[allow(dead_code)]
273    job: DistributedJob,
274    /// Current status.
275    status: JobStatus,
276    /// When the job was submitted.
277    submitted_at: Instant,
278    /// Number of retry attempts so far.
279    retries: u32,
280}
281
282/// Main distributed encoder interface.
283///
284/// Maintains an in-process job store so that `submit_job`, `job_status`, and
285/// `cancel_job` operate on real state. In a production deployment the store
286/// would be backed by the gRPC coordinator; this implementation provides a
287/// fully functional local fallback that exercises the complete lifecycle.
288pub struct DistributedEncoder {
289    config: DistributedConfig,
290    /// Job store keyed by job UUID.
291    jobs: Arc<RwLock<HashMap<Uuid, JobRecord>>>,
292}
293
294impl DistributedEncoder {
295    /// Create a new distributed encoder with the given configuration
296    #[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    /// Create a new distributed encoder with default configuration
305    #[must_use]
306    pub fn with_defaults() -> Self {
307        Self::new(DistributedConfig::default())
308    }
309
310    /// Get the current configuration
311    #[must_use]
312    pub fn config(&self) -> &DistributedConfig {
313        &self.config
314    }
315
316    /// Return the number of currently tracked jobs.
317    pub async fn job_count(&self) -> usize {
318        self.jobs.read().await.len()
319    }
320
321    /// Return the number of active (non-terminal) jobs.
322    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    /// Submit a job for distributed encoding.
337    ///
338    /// Validates the job, checks concurrency limits, registers it in the
339    /// internal store, and returns the job ID on success.
340    ///
341    /// # Arguments
342    ///
343    /// * `job` - The encoding job to submit
344    ///
345    /// # Returns
346    ///
347    /// Returns the job ID on success
348    ///
349    /// # Errors
350    ///
351    /// Returns `DistributedError::InvalidConfig` if the job definition is
352    /// invalid, or `DistributedError::ResourceExhausted` if the maximum
353    /// concurrent job limit has been reached.
354    pub async fn submit_job(&self, job: DistributedJob) -> Result<Uuid> {
355        // --- validation ---
356        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        // Check deadline is not already in the past
373        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        // Check for duplicate job ID
387        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        // Enforce concurrency limit
395        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    /// Query the status of a previously submitted job.
437    ///
438    /// In addition to returning the stored status, this method performs
439    /// timeout checking: if a job has been active longer than the configured
440    /// `job_timeout` it is automatically marked as `Failed`.
441    ///
442    /// # Errors
443    ///
444    /// Returns `DistributedError::Job` if the job ID is not found.
445    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        // Check for timeout on active jobs
454        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    /// Cancel a previously submitted job.
471    ///
472    /// Only jobs that are not yet in a terminal state (`Completed`, `Failed`,
473    /// `Cancelled`) can be cancelled.
474    ///
475    /// # Errors
476    ///
477    /// Returns `DistributedError::Job` if the job ID is not found or the job
478    /// is already in a terminal state.
479    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    /// Advance a job to the next logical status (for internal/testing use).
511    ///
512    /// Transitions: Pending -> Assigned -> InProgress -> Completed
513    ///
514    /// # Errors
515    ///
516    /// Returns error if the job is not found or is in a terminal state.
517    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    /// Mark a job as failed (for internal/testing use).
538    ///
539    /// # Errors
540    ///
541    /// Returns error if the job is not found or already in a terminal state.
542    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        // Check if we should retry
556        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    /// Get the retry count for a job.
573    ///
574    /// # Errors
575    ///
576    /// Returns error if the job is not found.
577    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    /// List all job IDs with their current statuses.
586    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/// Job execution status
597#[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        // Advance to Completed
749        encoder
750            .advance_job(job_id)
751            .await
752            .expect("advance should succeed"); // Assigned
753        encoder
754            .advance_job(job_id)
755            .await
756            .expect("advance should succeed"); // InProgress
757        encoder
758            .advance_job(job_id)
759            .await
760            .expect("advance should succeed"); // Completed
761
762        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        // Cannot advance past Completed
796        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        // First failure: should retry (back to Pending)
817        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        // Second failure: should retry again
830        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        // Third failure: max retries exhausted, should be Failed
838        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        // Third should be rejected
887        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        // Cannot submit another
904        assert!(encoder.submit_job(make_job()).await.is_err());
905
906        // Cancel the first
907        encoder
908            .cancel_job(job_id)
909            .await
910            .expect("cancel should succeed");
911
912        // Now we can submit
913        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        // Wait briefly for timeout
974        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); // epoch = far in the past
988
989        let result = encoder.submit_job(job).await;
990        assert!(result.is_err());
991    }
992}