diff --git a/Cargo.lock b/Cargo.lock index a8461bb..bac5b07 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -695,6 +695,10 @@ dependencies = [ "tempfile", ] +[[package]] +name = "invariant-coordinator" +version = "0.1.0" + [[package]] name = "invariant-engine" version = "0.1.0" @@ -720,6 +724,10 @@ dependencies = [ "tracing-subscriber", ] +[[package]] +name = "invariant-scheduler" +version = "0.1.0" + [[package]] name = "invariant-types" version = "0.1.0" diff --git a/Cargo.toml b/Cargo.toml index 35a3fe1..1a30a31 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,8 +1,10 @@ [workspace] resolver = "3" members = [ + "crates/invariant-coordinator", "crates/invariant-engine", "crates/invariant-journal", + "crates/invariant-scheduler", "crates/invariant-types", ] diff --git a/crates/invariant-coordinator/Cargo.toml b/crates/invariant-coordinator/Cargo.toml new file mode 100644 index 0000000..ce4d9fb --- /dev/null +++ b/crates/invariant-coordinator/Cargo.toml @@ -0,0 +1,6 @@ +[package] +name = "invariant-coordinator" +version = "0.1.0" +edition = "2024" + +[dependencies] diff --git a/crates/invariant-coordinator/src/domain/error.rs b/crates/invariant-coordinator/src/domain/error.rs new file mode 100644 index 0000000..cfe7019 --- /dev/null +++ b/crates/invariant-coordinator/src/domain/error.rs @@ -0,0 +1,81 @@ +use std::fmt; + +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum DomainError { + EmptyJobId, + EmptyTargetRef, + EmptyWorkerId, + EmptyLeaseToken, + ZeroLeaseDuration, + TimeOverflow, + LeaseEpochOverflow, + AttemptOverflow, + JobAlreadyLeased, + JobNotLeased, + JobTerminal, + LeaseHeldByDifferentWorker, + StaleLease, + LeaseExpired, +} + +impl fmt::Display for DomainError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + let message = match self { + Self::EmptyJobId => "job id cannot be empty", + Self::EmptyTargetRef => "target ref cannot be empty", + Self::EmptyWorkerId => "worker id cannot be empty", + Self::EmptyLeaseToken => "lease token cannot be empty", + Self::ZeroLeaseDuration => "lease duration must be greater than zero", + Self::TimeOverflow => "scheduler time overflow", + Self::LeaseEpochOverflow => "lease epoch overflow", + Self::AttemptOverflow => "attempt number overflow", + Self::JobAlreadyLeased => "job already has an active lease", + Self::JobNotLeased => "job is not leased", + Self::JobTerminal => "job is terminal", + Self::LeaseHeldByDifferentWorker => "lease is held by a different worker", + Self::StaleLease => "lease credential is stale", + Self::LeaseExpired => "lease has expired", + }; + f.write_str(message) + } +} + +impl std::error::Error for DomainError {} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn display_should_render_stable_messages_for_all_domain_errors() { + let cases = [ + (DomainError::EmptyJobId, "job id cannot be empty"), + (DomainError::EmptyTargetRef, "target ref cannot be empty"), + (DomainError::EmptyWorkerId, "worker id cannot be empty"), + (DomainError::EmptyLeaseToken, "lease token cannot be empty"), + ( + DomainError::ZeroLeaseDuration, + "lease duration must be greater than zero", + ), + (DomainError::TimeOverflow, "scheduler time overflow"), + (DomainError::LeaseEpochOverflow, "lease epoch overflow"), + (DomainError::AttemptOverflow, "attempt number overflow"), + ( + DomainError::JobAlreadyLeased, + "job already has an active lease", + ), + (DomainError::JobNotLeased, "job is not leased"), + (DomainError::JobTerminal, "job is terminal"), + ( + DomainError::LeaseHeldByDifferentWorker, + "lease is held by a different worker", + ), + (DomainError::StaleLease, "lease credential is stale"), + (DomainError::LeaseExpired, "lease has expired"), + ]; + + for (error, expected) in cases { + assert_eq!(error.to_string(), expected); + } + } +} diff --git a/crates/invariant-coordinator/src/domain/job.rs b/crates/invariant-coordinator/src/domain/job.rs new file mode 100644 index 0000000..66a218f --- /dev/null +++ b/crates/invariant-coordinator/src/domain/job.rs @@ -0,0 +1,519 @@ +use crate::domain::{ + error::DomainError, + lease::{JobLease, LeaseCredential, LeaseDuration, LeaseEpoch, LeaseToken}, + time::SchedulerTime, + worker::WorkerId, +}; + +#[derive(Clone, Debug, PartialEq, Eq, Hash)] +pub struct JobId(String); + +#[derive(Clone, Debug, PartialEq, Eq, Hash)] +pub struct TargetRef(String); + +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub struct AttemptNumber(u32); + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum JobStatus { + Queued, + Leased, + Completed, + Failed, + Cancelled, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct Job { + id: JobId, + target: TargetRef, + status: JobStatus, + current_lease: Option, + attempt: AttemptNumber, + last_lease_epoch: LeaseEpoch, +} + +impl JobId { + pub fn new(value: impl Into) -> Result { + let value = value.into(); + if value.is_empty() { + return Err(DomainError::EmptyJobId); + } + Ok(Self(value)) + } + + pub fn as_str(&self) -> &str { + &self.0 + } +} + +impl TargetRef { + pub fn new(value: impl Into) -> Result { + let value = value.into(); + if value.is_empty() { + return Err(DomainError::EmptyTargetRef); + } + Ok(Self(value)) + } + + pub fn as_str(&self) -> &str { + &self.0 + } +} + +impl AttemptNumber { + pub const fn zero() -> Self { + Self(0) + } + + pub const fn value(self) -> u32 { + self.0 + } + + pub fn checked_next(self) -> Result { + self.0 + .checked_add(1) + .map(Self) + .ok_or(DomainError::AttemptOverflow) + } +} + +impl JobStatus { + pub fn is_terminal(self) -> bool { + matches!(self, Self::Completed | Self::Failed | Self::Cancelled) + } +} + +impl Job { + pub fn new(id: JobId, target: TargetRef) -> Self { + Self { + id, + target, + status: JobStatus::Queued, + current_lease: None, + attempt: AttemptNumber::zero(), + last_lease_epoch: LeaseEpoch::zero(), + } + } + + pub fn id(&self) -> &JobId { + &self.id + } + + pub fn target(&self) -> &TargetRef { + &self.target + } + + pub fn status(&self) -> JobStatus { + self.status + } + + pub fn current_lease(&self) -> Option<&JobLease> { + self.current_lease.as_ref() + } + + pub fn attempt(&self) -> AttemptNumber { + self.attempt + } + + pub fn last_lease_epoch(&self) -> LeaseEpoch { + self.last_lease_epoch + } + + pub fn claim( + &mut self, + holder: WorkerId, + token: LeaseToken, + now: SchedulerTime, + duration: LeaseDuration, + ) -> Result { + if self.status.is_terminal() { + return Err(DomainError::JobTerminal); + } + + if let Some(current) = &self.current_lease + && !current.is_expired_at(now) + { + return Err(DomainError::JobAlreadyLeased); + } + + let next_epoch = self.last_lease_epoch.checked_next()?; + let next_attempt = self.attempt.checked_next()?; + let lease = JobLease::new(holder, token, next_epoch, now, duration)?; + + self.current_lease = Some(lease.clone()); + self.last_lease_epoch = next_epoch; + self.attempt = next_attempt; + self.status = JobStatus::Leased; + + Ok(lease) + } + + pub fn renew_lease( + &mut self, + credential: &LeaseCredential, + now: SchedulerTime, + duration: LeaseDuration, + ) -> Result<(), DomainError> { + if self.status.is_terminal() { + return Err(DomainError::JobTerminal); + } + + let lease = self + .current_lease + .as_mut() + .ok_or(DomainError::JobNotLeased)?; + + lease.renew(credential, now, duration) + } + + pub fn complete( + &mut self, + credential: &LeaseCredential, + now: SchedulerTime, + ) -> Result<(), DomainError> { + if self.status.is_terminal() { + return Err(DomainError::JobTerminal); + } + + let lease = self + .current_lease + .as_ref() + .ok_or(DomainError::JobNotLeased)?; + + lease.require_current(credential, now)?; + + self.status = JobStatus::Completed; + self.current_lease = None; + + Ok(()) + } + + pub fn release( + &mut self, + credential: &LeaseCredential, + now: SchedulerTime, + ) -> Result<(), DomainError> { + if self.status.is_terminal() { + return Err(DomainError::JobTerminal); + } + + let lease = self + .current_lease + .as_ref() + .ok_or(DomainError::JobNotLeased)?; + + lease.require_current(credential, now)?; + + self.current_lease = None; + self.status = JobStatus::Queued; + + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn job_id(value: &str) -> JobId { + JobId::new(value).unwrap() + } + + fn target(value: &str) -> TargetRef { + TargetRef::new(value).unwrap() + } + + fn worker(value: &str) -> WorkerId { + WorkerId::new(value).unwrap() + } + + fn token(value: &str) -> LeaseToken { + LeaseToken::new(value).unwrap() + } + + fn duration(value: u64) -> LeaseDuration { + LeaseDuration::from_millis(value).unwrap() + } + + fn new_job() -> Job { + Job::new(job_id("job-1"), target("target-1")) + } + + #[test] + fn new_job_starts_queued_without_lease() { + let job = new_job(); + + assert_eq!(job.status(), JobStatus::Queued); + assert_eq!(job.current_lease(), None); + assert_eq!(job.attempt(), AttemptNumber::zero()); + assert_eq!(job.last_lease_epoch(), LeaseEpoch::zero()); + } + + #[test] + fn constructors_should_reject_empty_ids() { + assert_eq!(JobId::new(""), Err(DomainError::EmptyJobId)); + assert_eq!(TargetRef::new(""), Err(DomainError::EmptyTargetRef)); + } + + #[test] + fn claim_queued_job_creates_lease_epoch_one_attempt_one() { + let mut job = new_job(); + let now = SchedulerTime::from_millis_since_epoch(10); + + let lease = job + .claim(worker("worker-1"), token("token-1"), now, duration(5)) + .unwrap(); + + assert_eq!(job.status(), JobStatus::Leased); + assert_eq!(job.attempt(), AttemptNumber(1)); + assert_eq!(job.last_lease_epoch(), LeaseEpoch::from_value(1)); + assert_eq!(lease.epoch(), LeaseEpoch::from_value(1)); + assert_eq!(lease.holder().as_str(), "worker-1"); + assert_eq!( + lease.expires_at(), + SchedulerTime::from_millis_since_epoch(15) + ); + assert_eq!(job.current_lease(), Some(&lease)); + } + + #[test] + fn claim_before_expiry_is_rejected() { + let mut job = new_job(); + let now = SchedulerTime::from_millis_since_epoch(10); + job.claim(worker("worker-1"), token("token-1"), now, duration(5)) + .unwrap(); + + let result = job.claim( + worker("worker-2"), + token("token-2"), + SchedulerTime::from_millis_since_epoch(14), + duration(5), + ); + + assert_eq!(result, Err(DomainError::JobAlreadyLeased)); + } + + #[test] + fn renew_valid_lease_keeps_token_epoch_and_attempt() { + let mut job = new_job(); + let now = SchedulerTime::from_millis_since_epoch(10); + let lease = job + .claim(worker("worker-1"), token("token-1"), now, duration(5)) + .unwrap(); + let credential = lease.credential(); + + let result = job.renew_lease( + &credential, + SchedulerTime::from_millis_since_epoch(12), + duration(10), + ); + + assert_eq!(result, Ok(())); + assert_eq!(job.attempt(), AttemptNumber(1)); + let renewed = job.current_lease().unwrap(); + assert_eq!(renewed.token().as_str(), "token-1"); + assert_eq!(renewed.epoch(), LeaseEpoch::from_value(1)); + assert_eq!( + renewed.expires_at(), + SchedulerTime::from_millis_since_epoch(22) + ); + } + + #[test] + fn renew_with_wrong_holder_is_rejected() { + let mut job = new_job(); + let now = SchedulerTime::from_millis_since_epoch(10); + let lease = job + .claim(worker("worker-1"), token("token-1"), now, duration(5)) + .unwrap(); + let credential = + LeaseCredential::new(worker("worker-2"), lease.token().clone(), lease.epoch()); + + let result = job.renew_lease(&credential, now, duration(5)); + + assert_eq!(result, Err(DomainError::LeaseHeldByDifferentWorker)); + } + + #[test] + fn renew_with_wrong_token_is_rejected() { + let mut job = new_job(); + let now = SchedulerTime::from_millis_since_epoch(10); + let lease = job + .claim(worker("worker-1"), token("token-1"), now, duration(5)) + .unwrap(); + let credential = + LeaseCredential::new(lease.holder().clone(), token("token-2"), lease.epoch()); + + let result = job.renew_lease(&credential, now, duration(5)); + + assert_eq!(result, Err(DomainError::StaleLease)); + } + + #[test] + fn renew_with_wrong_epoch_is_rejected() { + let mut job = new_job(); + let now = SchedulerTime::from_millis_since_epoch(10); + let lease = job + .claim(worker("worker-1"), token("token-1"), now, duration(5)) + .unwrap(); + let credential = LeaseCredential::new( + lease.holder().clone(), + lease.token().clone(), + LeaseEpoch::from_value(2), + ); + + let result = job.renew_lease(&credential, now, duration(5)); + + assert_eq!(result, Err(DomainError::StaleLease)); + } + + #[test] + fn expired_lease_can_be_reclaimed_with_new_epoch() { + let mut job = new_job(); + let lease = job + .claim( + worker("worker-1"), + token("token-1"), + SchedulerTime::from_millis_since_epoch(10), + duration(5), + ) + .unwrap(); + + let new_lease = job + .claim( + worker("worker-2"), + token("token-2"), + SchedulerTime::from_millis_since_epoch(15), + duration(5), + ) + .unwrap(); + + assert_eq!(lease.epoch(), LeaseEpoch::from_value(1)); + assert_eq!(new_lease.epoch(), LeaseEpoch::from_value(2)); + assert_eq!(new_lease.holder().as_str(), "worker-2"); + assert_eq!(job.attempt(), AttemptNumber(2)); + } + + #[test] + fn stale_complete_after_reclaim_is_rejected() { + let mut job = new_job(); + let old_lease = job + .claim( + worker("worker-1"), + token("token-1"), + SchedulerTime::from_millis_since_epoch(10), + duration(5), + ) + .unwrap(); + job.claim( + worker("worker-2"), + token("token-2"), + SchedulerTime::from_millis_since_epoch(15), + duration(5), + ) + .unwrap(); + + let result = job.complete( + &old_lease.credential(), + SchedulerTime::from_millis_since_epoch(16), + ); + + assert_eq!(result, Err(DomainError::LeaseHeldByDifferentWorker)); + assert_eq!(job.status(), JobStatus::Leased); + assert_eq!(job.current_lease().unwrap().holder().as_str(), "worker-2"); + } + + #[test] + fn expired_lease_cannot_complete_even_without_reclaim() { + let mut job = new_job(); + let lease = job + .claim( + worker("worker-1"), + token("token-1"), + SchedulerTime::from_millis_since_epoch(10), + duration(5), + ) + .unwrap(); + + let result = job.complete( + &lease.credential(), + SchedulerTime::from_millis_since_epoch(15), + ); + + assert_eq!(result, Err(DomainError::LeaseExpired)); + assert_eq!(job.status(), JobStatus::Leased); + } + + #[test] + fn valid_complete_marks_completed_and_clears_lease() { + let mut job = new_job(); + let lease = job + .claim( + worker("worker-1"), + token("token-1"), + SchedulerTime::from_millis_since_epoch(10), + duration(5), + ) + .unwrap(); + + let result = job.complete( + &lease.credential(), + SchedulerTime::from_millis_since_epoch(14), + ); + + assert_eq!(result, Ok(())); + assert_eq!(job.status(), JobStatus::Completed); + assert_eq!(job.current_lease(), None); + } + + #[test] + fn terminal_job_cannot_be_claimed() { + let mut job = new_job(); + let lease = job + .claim( + worker("worker-1"), + token("token-1"), + SchedulerTime::from_millis_since_epoch(10), + duration(5), + ) + .unwrap(); + job.complete( + &lease.credential(), + SchedulerTime::from_millis_since_epoch(14), + ) + .unwrap(); + + let result = job.claim( + worker("worker-2"), + token("token-2"), + SchedulerTime::from_millis_since_epoch(15), + duration(5), + ); + + assert_eq!(result, Err(DomainError::JobTerminal)); + } + + #[test] + fn valid_release_returns_to_queued_and_preserves_epoch() { + let mut job = new_job(); + let lease = job + .claim( + worker("worker-1"), + token("token-1"), + SchedulerTime::from_millis_since_epoch(10), + duration(5), + ) + .unwrap(); + + let result = job.release( + &lease.credential(), + SchedulerTime::from_millis_since_epoch(14), + ); + + assert_eq!(result, Ok(())); + assert_eq!(job.status(), JobStatus::Queued); + assert_eq!(job.current_lease(), None); + assert_eq!(job.last_lease_epoch(), LeaseEpoch::from_value(1)); + } +} diff --git a/crates/invariant-coordinator/src/domain/lease.rs b/crates/invariant-coordinator/src/domain/lease.rs new file mode 100644 index 0000000..e2e3e6c --- /dev/null +++ b/crates/invariant-coordinator/src/domain/lease.rs @@ -0,0 +1,327 @@ +use std::time::Duration; + +use crate::domain::{error::DomainError, time::SchedulerTime, worker::WorkerId}; + +#[derive(Clone, Debug, PartialEq, Eq, Hash)] +pub struct LeaseToken(String); + +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub struct LeaseEpoch(u64); + +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub struct LeaseDuration(Duration); + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct LeaseCredential { + holder: WorkerId, + token: LeaseToken, + epoch: LeaseEpoch, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct JobLease { + holder: WorkerId, + token: LeaseToken, + epoch: LeaseEpoch, + acquired_at: SchedulerTime, + renewed_at: SchedulerTime, + expires_at: SchedulerTime, +} + +impl LeaseToken { + pub fn new(value: impl Into) -> Result { + let value = value.into(); + if value.is_empty() { + return Err(DomainError::EmptyLeaseToken); + } + Ok(Self(value)) + } + + pub fn as_str(&self) -> &str { + &self.0 + } +} + +impl LeaseEpoch { + pub const fn zero() -> Self { + Self(0) + } + + pub const fn from_value(value: u64) -> Self { + Self(value) + } + + pub const fn value(self) -> u64 { + self.0 + } + + pub fn checked_next(self) -> Result { + self.0 + .checked_add(1) + .map(Self) + .ok_or(DomainError::LeaseEpochOverflow) + } +} + +impl LeaseDuration { + pub fn from_duration(value: Duration) -> Result { + if value.is_zero() { + return Err(DomainError::ZeroLeaseDuration); + } + Ok(Self(value)) + } + + pub fn from_millis(value: u64) -> Result { + Self::from_duration(Duration::from_millis(value)) + } + + pub fn from_secs(value: u64) -> Result { + Self::from_duration(Duration::from_secs(value)) + } + + pub fn as_duration(self) -> Duration { + self.0 + } +} + +impl LeaseCredential { + pub fn new(holder: WorkerId, token: LeaseToken, epoch: LeaseEpoch) -> Self { + Self { + holder, + token, + epoch, + } + } + + pub fn holder(&self) -> &WorkerId { + &self.holder + } + + pub fn token(&self) -> &LeaseToken { + &self.token + } + + pub fn epoch(&self) -> LeaseEpoch { + self.epoch + } +} + +impl JobLease { + pub fn new( + holder: WorkerId, + token: LeaseToken, + epoch: LeaseEpoch, + now: SchedulerTime, + duration: LeaseDuration, + ) -> Result { + let expires_at = now.checked_add_duration(duration.as_duration())?; + Ok(Self { + holder, + token, + epoch, + acquired_at: now, + renewed_at: now, + expires_at, + }) + } + + pub fn holder(&self) -> &WorkerId { + &self.holder + } + + pub fn token(&self) -> &LeaseToken { + &self.token + } + + pub fn epoch(&self) -> LeaseEpoch { + self.epoch + } + + pub fn acquired_at(&self) -> SchedulerTime { + self.acquired_at + } + + pub fn renewed_at(&self) -> SchedulerTime { + self.renewed_at + } + + pub fn expires_at(&self) -> SchedulerTime { + self.expires_at + } + + pub fn credential(&self) -> LeaseCredential { + LeaseCredential::new(self.holder.clone(), self.token.clone(), self.epoch) + } + + pub fn is_expired_at(&self, now: SchedulerTime) -> bool { + now >= self.expires_at + } + + pub fn require_current( + &self, + credential: &LeaseCredential, + now: SchedulerTime, + ) -> Result<(), DomainError> { + if self.holder != *credential.holder() { + return Err(DomainError::LeaseHeldByDifferentWorker); + } + if self.token != *credential.token() || self.epoch != credential.epoch() { + return Err(DomainError::StaleLease); + } + if self.is_expired_at(now) { + return Err(DomainError::LeaseExpired); + } + Ok(()) + } + + pub fn renew( + &mut self, + credential: &LeaseCredential, + now: SchedulerTime, + duration: LeaseDuration, + ) -> Result<(), DomainError> { + self.require_current(credential, now)?; + self.renewed_at = now; + self.expires_at = now.checked_add_duration(duration.as_duration())?; + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn worker(value: &str) -> WorkerId { + WorkerId::new(value).unwrap() + } + + fn token(value: &str) -> LeaseToken { + LeaseToken::new(value).unwrap() + } + + fn duration(value: u64) -> LeaseDuration { + LeaseDuration::from_millis(value).unwrap() + } + + #[test] + fn lease_duration_should_reject_zero() { + let result = LeaseDuration::from_millis(0); + + assert_eq!(result, Err(DomainError::ZeroLeaseDuration)); + } + + #[test] + fn lease_epoch_should_increment_from_zero() { + let result = LeaseEpoch::zero().checked_next(); + + assert_eq!(result, Ok(LeaseEpoch(1))); + } + + #[test] + fn lease_epoch_should_reject_overflow() { + let result = LeaseEpoch(u64::MAX).checked_next(); + + assert_eq!(result, Err(DomainError::LeaseEpochOverflow)); + } + + #[test] + fn new_should_compute_expiry_and_initial_renewed_at() { + let now = SchedulerTime::from_millis_since_epoch(10); + + let lease = JobLease::new( + worker("worker-1"), + token("token-1"), + LeaseEpoch(1), + now, + duration(5), + ) + .unwrap(); + + assert_eq!(lease.acquired_at(), now); + assert_eq!(lease.renewed_at(), now); + assert_eq!( + lease.expires_at(), + SchedulerTime::from_millis_since_epoch(15) + ); + } + + #[test] + fn require_current_should_reject_wrong_holder() { + let now = SchedulerTime::from_millis_since_epoch(10); + let lease = JobLease::new( + worker("worker-1"), + token("token-1"), + LeaseEpoch(1), + now, + duration(5), + ) + .unwrap(); + let credential = LeaseCredential::new(worker("worker-2"), token("token-1"), LeaseEpoch(1)); + + let result = lease.require_current(&credential, now); + + assert_eq!(result, Err(DomainError::LeaseHeldByDifferentWorker)); + } + + #[test] + fn require_current_should_reject_wrong_token() { + let now = SchedulerTime::from_millis_since_epoch(10); + let lease = JobLease::new( + worker("worker-1"), + token("token-1"), + LeaseEpoch(1), + now, + duration(5), + ) + .unwrap(); + let credential = LeaseCredential::new(worker("worker-1"), token("token-2"), LeaseEpoch(1)); + + let result = lease.require_current(&credential, now); + + assert_eq!(result, Err(DomainError::StaleLease)); + } + + #[test] + fn require_current_should_reject_expired_lease_at_deadline() { + let now = SchedulerTime::from_millis_since_epoch(10); + let lease = JobLease::new( + worker("worker-1"), + token("token-1"), + LeaseEpoch(1), + now, + duration(5), + ) + .unwrap(); + let credential = lease.credential(); + + let result = lease.require_current(&credential, SchedulerTime::from_millis_since_epoch(15)); + + assert_eq!(result, Err(DomainError::LeaseExpired)); + } + + #[test] + fn renew_should_keep_token_and_epoch_and_extend_expiry() { + let now = SchedulerTime::from_millis_since_epoch(10); + let mut lease = JobLease::new( + worker("worker-1"), + token("token-1"), + LeaseEpoch(1), + now, + duration(5), + ) + .unwrap(); + let credential = lease.credential(); + let renewed_at = SchedulerTime::from_millis_since_epoch(12); + + let result = lease.renew(&credential, renewed_at, duration(10)); + + assert_eq!(result, Ok(())); + assert_eq!(lease.token().as_str(), "token-1"); + assert_eq!(lease.epoch(), LeaseEpoch(1)); + assert_eq!(lease.renewed_at(), renewed_at); + assert_eq!( + lease.expires_at(), + SchedulerTime::from_millis_since_epoch(22) + ); + } +} diff --git a/crates/invariant-coordinator/src/domain/mod.rs b/crates/invariant-coordinator/src/domain/mod.rs new file mode 100644 index 0000000..42eb2f6 --- /dev/null +++ b/crates/invariant-coordinator/src/domain/mod.rs @@ -0,0 +1,11 @@ +pub mod error; +pub mod job; +pub mod lease; +pub mod time; +pub mod worker; + +pub use error::DomainError; +pub use job::{AttemptNumber, Job, JobId, JobStatus, TargetRef}; +pub use lease::{JobLease, LeaseCredential, LeaseDuration, LeaseEpoch, LeaseToken}; +pub use time::SchedulerTime; +pub use worker::WorkerId; diff --git a/crates/invariant-coordinator/src/domain/time.rs b/crates/invariant-coordinator/src/domain/time.rs new file mode 100644 index 0000000..50a9985 --- /dev/null +++ b/crates/invariant-coordinator/src/domain/time.rs @@ -0,0 +1,59 @@ +use std::time::Duration; + +use crate::domain::error::DomainError; + +/// A point in time on the scheduler's timeline. +/// +/// `SchedulerTime` is measured as elapsed duration since the scheduler epoch. +/// Unlike `std::time::Instant`, it is a domain value rather than a process-local +/// clock reading. +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub struct SchedulerTime(Duration); + +impl SchedulerTime { + pub fn from_duration_since_epoch(value: Duration) -> Self { + Self(value) + } + + pub fn from_millis_since_epoch(value: u64) -> Self { + Self(Duration::from_millis(value)) + } + + pub fn from_secs_since_epoch(value: u64) -> Self { + Self(Duration::from_secs(value)) + } + + pub fn duration_since_epoch(self) -> Duration { + self.0 + } + + pub fn checked_add_duration(self, duration: Duration) -> Result { + self.0 + .checked_add(duration) + .map(Self) + .ok_or(DomainError::TimeOverflow) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn checked_add_duration_should_return_new_scheduler_time() { + let now = SchedulerTime::from_millis_since_epoch(10); + + let result = now.checked_add_duration(Duration::from_millis(5)); + + assert_eq!(result, Ok(SchedulerTime::from_millis_since_epoch(15))); + } + + #[test] + fn checked_add_duration_should_reject_overflow() { + let now = SchedulerTime::from_duration_since_epoch(Duration::MAX); + + let result = now.checked_add_duration(Duration::from_nanos(1)); + + assert_eq!(result, Err(DomainError::TimeOverflow)); + } +} diff --git a/crates/invariant-coordinator/src/domain/worker.rs b/crates/invariant-coordinator/src/domain/worker.rs new file mode 100644 index 0000000..ec9562c --- /dev/null +++ b/crates/invariant-coordinator/src/domain/worker.rs @@ -0,0 +1,37 @@ +use crate::domain::error::DomainError; + +#[derive(Clone, Debug, PartialEq, Eq, Hash)] +pub struct WorkerId(String); + +impl WorkerId { + pub fn new(value: impl Into) -> Result { + let value = value.into(); + if value.is_empty() { + return Err(DomainError::EmptyWorkerId); + } + Ok(Self(value)) + } + + pub fn as_str(&self) -> &str { + &self.0 + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn new_should_reject_empty_worker_id() { + let result = WorkerId::new(""); + + assert_eq!(result, Err(DomainError::EmptyWorkerId)); + } + + #[test] + fn as_str_should_return_worker_id() { + let worker_id = WorkerId::new("worker-1").unwrap(); + + assert_eq!(worker_id.as_str(), "worker-1"); + } +} diff --git a/crates/invariant-coordinator/src/lib.rs b/crates/invariant-coordinator/src/lib.rs new file mode 100644 index 0000000..2d87f10 --- /dev/null +++ b/crates/invariant-coordinator/src/lib.rs @@ -0,0 +1,6 @@ +//! invariant-coordinator — GLOBAL durable-ownership / coordination layer (WIP). +//! +//! Parked here during the local-scheduler milestone: leases, fencing epochs, and +//! the lease-laden `Job` aggregate. The global layer is not yet fully designed. +//! `invariant-scheduler` (local) does NOT depend on this crate — work flows down. +pub mod domain; diff --git a/crates/invariant-scheduler/Cargo.toml b/crates/invariant-scheduler/Cargo.toml new file mode 100644 index 0000000..8b95538 --- /dev/null +++ b/crates/invariant-scheduler/Cargo.toml @@ -0,0 +1,6 @@ +[package] +name = "invariant-scheduler" +version = "0.1.0" +edition = "2024" + +[dependencies] diff --git a/crates/invariant-scheduler/src/domain/error.rs b/crates/invariant-scheduler/src/domain/error.rs new file mode 100644 index 0000000..a0f20cc --- /dev/null +++ b/crates/invariant-scheduler/src/domain/error.rs @@ -0,0 +1,48 @@ +use std::fmt; + +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum DomainError { + EmptyJobId, + EmptyTargetRef, + TimeOverflow, + AttemptOverflow, + JobTerminal, + QueueFull, +} + +impl fmt::Display for DomainError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + let message = match self { + Self::EmptyJobId => "job id cannot be empty", + Self::EmptyTargetRef => "target ref cannot be empty", + Self::TimeOverflow => "scheduler time overflow", + Self::AttemptOverflow => "attempt number overflow", + Self::JobTerminal => "job is terminal", + Self::QueueFull => "ready queue is full", + }; + f.write_str(message) + } +} + +impl std::error::Error for DomainError {} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn display_should_render_stable_messages_for_all_domain_errors() { + let cases = [ + (DomainError::EmptyJobId, "job id cannot be empty"), + (DomainError::EmptyTargetRef, "target ref cannot be empty"), + (DomainError::TimeOverflow, "scheduler time overflow"), + (DomainError::AttemptOverflow, "attempt number overflow"), + (DomainError::JobTerminal, "job is terminal"), + (DomainError::QueueFull, "ready queue is full"), + ]; + + for (error, expected) in cases { + assert_eq!(error.to_string(), expected); + } + } +} diff --git a/crates/invariant-scheduler/src/domain/mod.rs b/crates/invariant-scheduler/src/domain/mod.rs new file mode 100644 index 0000000..c7cc940 --- /dev/null +++ b/crates/invariant-scheduler/src/domain/mod.rs @@ -0,0 +1,5 @@ +pub mod error; +pub mod time; + +pub use error::DomainError; +pub use time::SchedulerTime; diff --git a/crates/invariant-scheduler/src/domain/time.rs b/crates/invariant-scheduler/src/domain/time.rs new file mode 100644 index 0000000..50a9985 --- /dev/null +++ b/crates/invariant-scheduler/src/domain/time.rs @@ -0,0 +1,59 @@ +use std::time::Duration; + +use crate::domain::error::DomainError; + +/// A point in time on the scheduler's timeline. +/// +/// `SchedulerTime` is measured as elapsed duration since the scheduler epoch. +/// Unlike `std::time::Instant`, it is a domain value rather than a process-local +/// clock reading. +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub struct SchedulerTime(Duration); + +impl SchedulerTime { + pub fn from_duration_since_epoch(value: Duration) -> Self { + Self(value) + } + + pub fn from_millis_since_epoch(value: u64) -> Self { + Self(Duration::from_millis(value)) + } + + pub fn from_secs_since_epoch(value: u64) -> Self { + Self(Duration::from_secs(value)) + } + + pub fn duration_since_epoch(self) -> Duration { + self.0 + } + + pub fn checked_add_duration(self, duration: Duration) -> Result { + self.0 + .checked_add(duration) + .map(Self) + .ok_or(DomainError::TimeOverflow) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn checked_add_duration_should_return_new_scheduler_time() { + let now = SchedulerTime::from_millis_since_epoch(10); + + let result = now.checked_add_duration(Duration::from_millis(5)); + + assert_eq!(result, Ok(SchedulerTime::from_millis_since_epoch(15))); + } + + #[test] + fn checked_add_duration_should_reject_overflow() { + let now = SchedulerTime::from_duration_since_epoch(Duration::MAX); + + let result = now.checked_add_duration(Duration::from_nanos(1)); + + assert_eq!(result, Err(DomainError::TimeOverflow)); + } +} diff --git a/crates/invariant-scheduler/src/lib.rs b/crates/invariant-scheduler/src/lib.rs new file mode 100644 index 0000000..d7abca1 --- /dev/null +++ b/crates/invariant-scheduler/src/lib.rs @@ -0,0 +1 @@ +pub mod domain;