use std::fmt;
use std::time::SystemTime;
use super::lifecycle::{validate_expected_version, validate_restart, validate_transition};
use super::{
DomainError, ExecutionVersion, ExitCode, FailureId, JobExecutionId, JobInstanceId,
JobInstanceKey, LifecycleError, LifecycleTransition, StepExecutionId, StepName,
};
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
#[non_exhaustive]
pub enum BatchStatus {
Starting,
Started,
Stopping,
Stopped,
Failed,
Completed,
Abandoned,
Unknown,
}
impl BatchStatus {
#[must_use]
pub const fn is_active(self) -> bool {
matches!(self, Self::Starting | Self::Started | Self::Stopping)
}
#[must_use]
pub const fn is_finished(self) -> bool {
matches!(
self,
Self::Stopped | Self::Failed | Self::Completed | Self::Abandoned
)
}
#[must_use]
pub const fn is_terminal(self) -> bool {
matches!(self, Self::Completed | Self::Abandoned)
}
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Starting => "STARTING",
Self::Started => "STARTED",
Self::Stopping => "STOPPING",
Self::Stopped => "STOPPED",
Self::Failed => "FAILED",
Self::Completed => "COMPLETED",
Self::Abandoned => "ABANDONED",
Self::Unknown => "UNKNOWN",
}
}
}
impl fmt::Display for BatchStatus {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(self.as_str())
}
}
#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct ExitStatus {
code: ExitCode,
}
impl ExitStatus {
#[must_use]
pub const fn new(code: ExitCode) -> Self {
Self { code }
}
#[must_use]
pub fn unknown() -> Self {
Self {
code: ExitCode::framework_owned("UNKNOWN"),
}
}
#[must_use]
pub fn completed() -> Self {
Self {
code: ExitCode::framework_owned("COMPLETED"),
}
}
#[must_use]
pub fn failed() -> Self {
Self {
code: ExitCode::framework_owned("FAILED"),
}
}
#[must_use]
pub fn stopped() -> Self {
Self {
code: ExitCode::framework_owned("STOPPED"),
}
}
#[must_use]
pub const fn code(&self) -> &ExitCode {
&self.code
}
}
impl fmt::Display for ExitStatus {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
self.code.fmt(formatter)
}
}
#[derive(Clone, Copy, Debug, Default, Eq, Hash, Ord, PartialEq, PartialOrd)]
#[non_exhaustive]
pub struct ExecutionCounts {
read: u64,
processed: u64,
written: u64,
filtered: u64,
committed: u64,
rolled_back: u64,
}
impl ExecutionCounts {
#[must_use]
pub const fn new(
read: u64,
processed: u64,
written: u64,
filtered: u64,
committed: u64,
rolled_back: u64,
) -> Self {
Self {
read,
processed,
written,
filtered,
committed,
rolled_back,
}
}
#[must_use]
pub const fn read(self) -> u64 {
self.read
}
#[must_use]
pub const fn processed(self) -> u64 {
self.processed
}
#[must_use]
pub const fn written(self) -> u64 {
self.written
}
#[must_use]
pub const fn filtered(self) -> u64 {
self.filtered
}
#[must_use]
pub const fn committed(self) -> u64 {
self.committed
}
#[must_use]
pub const fn rolled_back(self) -> u64 {
self.rolled_back
}
fn with_terminal_rollback(self) -> Result<Self, LifecycleError> {
let rolled_back = self
.rolled_back
.checked_add(1)
.ok_or(LifecycleError::CountExhausted)?;
Ok(Self {
rolled_back,
..self
})
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
#[allow(clippy::struct_field_names)]
pub struct ExecutionTimestamps {
created_at: SystemTime,
started_at: Option<SystemTime>,
ended_at: Option<SystemTime>,
}
impl ExecutionTimestamps {
pub fn new(
created_at: SystemTime,
started_at: Option<SystemTime>,
ended_at: Option<SystemTime>,
) -> Result<Self, DomainError> {
if started_at.is_some_and(|started| started < created_at)
|| ended_at.is_some_and(|ended| ended < started_at.unwrap_or(created_at))
{
return Err(DomainError::InvalidTimestampOrder);
}
Ok(Self {
created_at,
started_at,
ended_at,
})
}
#[must_use]
pub const fn created_at(self) -> SystemTime {
self.created_at
}
#[must_use]
pub const fn started_at(self) -> Option<SystemTime> {
self.started_at
}
#[must_use]
pub const fn ended_at(self) -> Option<SystemTime> {
self.ended_at
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
#[non_exhaustive]
pub enum FailureCategory {
InvalidDefinition,
DuplicateExecution,
IllegalTransition,
TransientInfrastructure,
PermanentInfrastructure,
UserComponent,
Cancelled,
Serialization,
Invariant,
OptimisticConflict,
Timeout,
UnsupportedCapability,
UnknownCommit,
ShutdownIncomplete,
StaleRecovered,
}
impl FailureCategory {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::InvalidDefinition => "invalid_definition",
Self::DuplicateExecution => "duplicate_execution",
Self::IllegalTransition => "illegal_transition",
Self::TransientInfrastructure => "transient_infrastructure",
Self::PermanentInfrastructure => "permanent_infrastructure",
Self::UserComponent => "user_component",
Self::Cancelled => "cancelled",
Self::Serialization => "serialization",
Self::Invariant => "invariant",
Self::OptimisticConflict => "optimistic_conflict",
Self::Timeout => "timeout",
Self::UnsupportedCapability => "unsupported_capability",
Self::UnknownCommit => "unknown_commit",
Self::ShutdownIncomplete => "shutdown_incomplete",
Self::StaleRecovered => "stale_recovered",
}
}
#[must_use]
pub const fn is_policy_eligible(self) -> bool {
matches!(
self,
Self::TransientInfrastructure
| Self::PermanentInfrastructure
| Self::UserComponent
| Self::OptimisticConflict
| Self::Timeout
)
}
#[must_use]
pub const fn durable_code(self) -> &'static str {
match self {
Self::InvalidDefinition => "INVALID_DEFINITION",
Self::DuplicateExecution => "DUPLICATE_EXECUTION",
Self::IllegalTransition => "ILLEGAL_TRANSITION",
Self::TransientInfrastructure => "TRANSIENT_INFRASTRUCTURE",
Self::PermanentInfrastructure => "PERMANENT_INFRASTRUCTURE",
Self::UserComponent => "USER_COMPONENT",
Self::Cancelled => "CANCELLED",
Self::Serialization => "SERIALIZATION",
Self::Invariant => "INVARIANT",
Self::OptimisticConflict => "OPTIMISTIC_CONFLICT",
Self::Timeout => "TIMEOUT",
Self::UnsupportedCapability => "UNSUPPORTED_CAPABILITY",
Self::UnknownCommit => "UNKNOWN_COMMIT",
Self::ShutdownIncomplete => "SHUTDOWN_INCOMPLETE",
Self::StaleRecovered => "STALE_RECOVERED",
}
}
#[must_use]
pub fn from_durable_code(value: &str) -> Option<Self> {
Some(match value {
"INVALID_DEFINITION" => Self::InvalidDefinition,
"DUPLICATE_EXECUTION" => Self::DuplicateExecution,
"ILLEGAL_TRANSITION" => Self::IllegalTransition,
"TRANSIENT_INFRASTRUCTURE" => Self::TransientInfrastructure,
"PERMANENT_INFRASTRUCTURE" => Self::PermanentInfrastructure,
"USER_COMPONENT" => Self::UserComponent,
"CANCELLED" => Self::Cancelled,
"SERIALIZATION" => Self::Serialization,
"INVARIANT" => Self::Invariant,
"OPTIMISTIC_CONFLICT" => Self::OptimisticConflict,
"TIMEOUT" => Self::Timeout,
"UNSUPPORTED_CAPABILITY" => Self::UnsupportedCapability,
"UNKNOWN_COMMIT" => Self::UnknownCommit,
"SHUTDOWN_INCOMPLETE" => Self::ShutdownIncomplete,
"STALE_RECOVERED" => Self::StaleRecovered,
_ => return None,
})
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct FailureSummary {
category: FailureCategory,
failure_id: FailureId,
}
impl FailureSummary {
#[must_use]
pub const fn new(category: FailureCategory, failure_id: FailureId) -> Self {
Self {
category,
failure_id,
}
}
#[must_use]
pub const fn category(self) -> FailureCategory {
self.category
}
#[must_use]
pub const fn failure_id(self) -> FailureId {
self.failure_id
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct ExecutionMetadata {
status: BatchStatus,
exit_status: ExitStatus,
timestamps: ExecutionTimestamps,
counts: ExecutionCounts,
failure: Option<FailureSummary>,
}
impl ExecutionMetadata {
pub fn new(
status: BatchStatus,
exit_status: ExitStatus,
timestamps: ExecutionTimestamps,
counts: ExecutionCounts,
failure: Option<FailureSummary>,
) -> Result<Self, DomainError> {
if status.is_active() && timestamps.ended_at().is_some() {
return Err(DomainError::ActiveExecutionHasEndTime);
}
if status.is_finished() && timestamps.ended_at().is_none() {
return Err(DomainError::FinishedExecutionMissingEndTime);
}
if matches!(status, BatchStatus::Failed) && failure.is_none() {
return Err(DomainError::FailedExecutionMissingFailure);
}
Ok(Self {
status,
exit_status,
timestamps,
counts,
failure,
})
}
#[must_use]
pub const fn status(&self) -> BatchStatus {
self.status
}
#[must_use]
pub const fn exit_status(&self) -> &ExitStatus {
&self.exit_status
}
#[must_use]
pub const fn timestamps(&self) -> ExecutionTimestamps {
self.timestamps
}
#[must_use]
pub const fn counts(&self) -> ExecutionCounts {
self.counts
}
#[must_use]
pub const fn failure(&self) -> Option<FailureSummary> {
self.failure
}
fn transition(&self, transition: LifecycleTransition) -> Result<Self, LifecycleError> {
validate_transition(self.status, transition)?;
let target = transition.target();
let transitioned_at = transition.transitioned_at();
let current_timestamps = self.timestamps;
if transitioned_at
< current_timestamps
.started_at()
.unwrap_or(current_timestamps.created_at())
|| current_timestamps
.ended_at()
.is_some_and(|ended_at| transitioned_at < ended_at)
{
return Err(LifecycleError::InvalidTransitionTime {
source: DomainError::InvalidTimestampOrder,
});
}
let started_at = if matches!(target, BatchStatus::Started) {
Some(transitioned_at)
} else {
current_timestamps.started_at()
};
let ended_at = if target.is_finished() {
current_timestamps.ended_at().or(Some(transitioned_at))
} else {
None
};
let timestamps =
ExecutionTimestamps::new(current_timestamps.created_at(), started_at, ended_at)
.map_err(|source| LifecycleError::InvalidTransitionTime { source })?;
let failure = transition.failure().or(self.failure);
let counts = if transition.terminal_rollback() {
self.counts.with_terminal_rollback()?
} else {
self.counts
};
Self::new(
target,
self.exit_status.clone(),
timestamps,
counts,
failure,
)
.map_err(|source| match source {
DomainError::FailedExecutionMissingFailure => {
LifecycleError::FailedTransitionMissingFailure
}
source => LifecycleError::InvalidTransitionTime { source },
})
}
fn with_exit_status(&self, exit_status: ExitStatus) -> Self {
Self {
status: self.status,
exit_status,
timestamps: self.timestamps,
counts: self.counts,
failure: self.failure,
}
}
fn starting(created_at: SystemTime) -> Self {
Self {
status: BatchStatus::Starting,
exit_status: ExitStatus::unknown(),
timestamps: ExecutionTimestamps {
created_at,
started_at: None,
ended_at: None,
},
counts: ExecutionCounts::default(),
failure: None,
}
}
}
#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct JobInstance {
id: JobInstanceId,
key: JobInstanceKey,
}
impl JobInstance {
#[must_use]
pub const fn new(id: JobInstanceId, key: JobInstanceKey) -> Self {
Self { id, key }
}
#[must_use]
pub const fn id(&self) -> JobInstanceId {
self.id
}
#[must_use]
pub const fn key(&self) -> &JobInstanceKey {
&self.key
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct JobExecution {
id: JobExecutionId,
job_instance_id: JobInstanceId,
metadata: ExecutionMetadata,
version: ExecutionVersion,
}
impl JobExecution {
#[must_use]
pub const fn new(
id: JobExecutionId,
job_instance_id: JobInstanceId,
metadata: ExecutionMetadata,
) -> Self {
Self {
id,
job_instance_id,
metadata,
version: ExecutionVersion::INITIAL,
}
}
#[must_use]
pub const fn from_snapshot(
id: JobExecutionId,
job_instance_id: JobInstanceId,
metadata: ExecutionMetadata,
version: ExecutionVersion,
) -> Self {
Self {
id,
job_instance_id,
metadata,
version,
}
}
#[must_use]
pub const fn id(&self) -> JobExecutionId {
self.id
}
#[must_use]
pub const fn job_instance_id(&self) -> JobInstanceId {
self.job_instance_id
}
#[must_use]
pub const fn metadata(&self) -> &ExecutionMetadata {
&self.metadata
}
#[must_use]
pub const fn version(&self) -> ExecutionVersion {
self.version
}
pub fn transition(
&mut self,
expected_version: ExecutionVersion,
transition: LifecycleTransition,
) -> Result<ExecutionVersion, LifecycleError> {
transition_execution(
&mut self.metadata,
&mut self.version,
expected_version,
transition,
)
}
pub fn enrich_exit_status(
&mut self,
expected_version: ExecutionVersion,
exit_status: ExitStatus,
) -> Result<ExecutionVersion, LifecycleError> {
enrich_execution_exit_status(
&mut self.metadata,
&mut self.version,
expected_version,
exit_status,
)
}
pub fn new_restart_attempt(
&self,
expected_version: ExecutionVersion,
new_execution_id: JobExecutionId,
created_at: SystemTime,
) -> Result<Self, LifecycleError> {
validate_expected_version(expected_version, self.version)?;
validate_restart(self.metadata.status())?;
if new_execution_id == self.id {
return Err(LifecycleError::AttemptIdentifierReused);
}
validate_restart_time(&self.metadata, created_at)?;
Ok(Self::new(
new_execution_id,
self.job_instance_id,
ExecutionMetadata::starting(created_at),
))
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct StepExecution {
id: StepExecutionId,
job_execution_id: JobExecutionId,
step_name: StepName,
metadata: ExecutionMetadata,
version: ExecutionVersion,
}
impl StepExecution {
#[must_use]
pub const fn new(
id: StepExecutionId,
job_execution_id: JobExecutionId,
step_name: StepName,
metadata: ExecutionMetadata,
) -> Self {
Self {
id,
job_execution_id,
step_name,
metadata,
version: ExecutionVersion::INITIAL,
}
}
#[must_use]
pub const fn from_snapshot(
id: StepExecutionId,
job_execution_id: JobExecutionId,
step_name: StepName,
metadata: ExecutionMetadata,
version: ExecutionVersion,
) -> Self {
Self {
id,
job_execution_id,
step_name,
metadata,
version,
}
}
#[must_use]
pub const fn id(&self) -> StepExecutionId {
self.id
}
#[must_use]
pub const fn job_execution_id(&self) -> JobExecutionId {
self.job_execution_id
}
#[must_use]
pub const fn step_name(&self) -> &StepName {
&self.step_name
}
#[must_use]
pub const fn metadata(&self) -> &ExecutionMetadata {
&self.metadata
}
#[must_use]
pub const fn version(&self) -> ExecutionVersion {
self.version
}
pub fn transition(
&mut self,
expected_version: ExecutionVersion,
transition: LifecycleTransition,
) -> Result<ExecutionVersion, LifecycleError> {
transition_execution(
&mut self.metadata,
&mut self.version,
expected_version,
transition,
)
}
pub fn enrich_exit_status(
&mut self,
expected_version: ExecutionVersion,
exit_status: ExitStatus,
) -> Result<ExecutionVersion, LifecycleError> {
enrich_execution_exit_status(
&mut self.metadata,
&mut self.version,
expected_version,
exit_status,
)
}
pub fn new_restart_attempt(
&self,
expected_version: ExecutionVersion,
new_execution_id: StepExecutionId,
new_job_execution_id: JobExecutionId,
created_at: SystemTime,
) -> Result<Self, LifecycleError> {
validate_expected_version(expected_version, self.version)?;
validate_restart(self.metadata.status())?;
if new_execution_id == self.id || new_job_execution_id == self.job_execution_id {
return Err(LifecycleError::AttemptIdentifierReused);
}
validate_restart_time(&self.metadata, created_at)?;
Ok(Self::new(
new_execution_id,
new_job_execution_id,
self.step_name.clone(),
ExecutionMetadata::starting(created_at),
))
}
}
fn transition_execution(
metadata: &mut ExecutionMetadata,
version: &mut ExecutionVersion,
expected_version: ExecutionVersion,
transition: LifecycleTransition,
) -> Result<ExecutionVersion, LifecycleError> {
validate_expected_version(expected_version, *version)?;
let updated_metadata = metadata.transition(transition)?;
let updated_version = version.next()?;
*metadata = updated_metadata;
*version = updated_version;
Ok(updated_version)
}
fn enrich_execution_exit_status(
metadata: &mut ExecutionMetadata,
version: &mut ExecutionVersion,
expected_version: ExecutionVersion,
exit_status: ExitStatus,
) -> Result<ExecutionVersion, LifecycleError> {
validate_expected_version(expected_version, *version)?;
let updated_version = version.next()?;
*metadata = metadata.with_exit_status(exit_status);
*version = updated_version;
Ok(updated_version)
}
fn validate_restart_time(
metadata: &ExecutionMetadata,
created_at: SystemTime,
) -> Result<(), LifecycleError> {
if metadata
.timestamps()
.ended_at()
.is_some_and(|ended_at| created_at < ended_at)
{
return Err(LifecycleError::InvalidTransitionTime {
source: DomainError::InvalidTimestampOrder,
});
}
Ok(())
}