use std::collections::BTreeSet;
use std::error::Error;
use std::fmt;
use std::future::Future;
use std::num::NonZeroU64;
use std::pin::Pin;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::SystemTime;
use oxide_batch_core::{
BatchStatus, DefinitionIdentity, DefinitionRevision, DefinitionUpgrade, DomainError,
ExecutionMetadata, ExecutionTimestamps, ExecutionVersion, ExitStatus, FailureCategory,
FailureId, FailureSummary, IdentifierKind, JobExecution, JobExecutionId, JobInstance,
JobInstanceId, JobInstanceKey, JobName, LifecycleError, LifecycleTransition, NodeId,
RecoveryDecisionId, StartLimit, StepExecution, StepExecutionId, StepName, StepPartitionId,
};
use crate::{
ActorRef, FlowDecision, FlowDecisionRequest, FlowStepState, FlowTransitionKind, OperationId,
OperatorAction, OperatorRecord, OperatorRecordDraft, OwnerToken, PartitionAggregate,
PartitionAggregationError, PartitionPlanEntry, PurgeCounts, PurgePlan, PurgePlanRequest,
PurgeSurvey, ReasonCode, RetentionAction, RetentionHold, RetentionRecord, RetentionRecordDraft,
StepPartition,
};
const MAX_RECOVERY_REASON_BYTES: usize = 64;
const MAX_OPERATOR_REFERENCE_BYTES: usize = 128;
pub type BoxFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
pub trait Clock: Send + Sync {
fn now(&self) -> SystemTime;
}
#[derive(Clone, Copy, Debug, Default)]
pub struct SystemClock;
impl Clock for SystemClock {
fn now(&self) -> SystemTime {
SystemTime::now()
}
}
pub trait IdGenerator: Send + Sync {
fn next_job_instance_id(&self) -> Result<JobInstanceId, IdGenerationError>;
fn next_job_execution_id(&self) -> Result<JobExecutionId, IdGenerationError>;
fn next_step_execution_id(&self) -> Result<StepExecutionId, IdGenerationError>;
fn next_failure_id(&self) -> Result<FailureId, IdGenerationError>;
}
#[derive(Debug)]
pub struct SequentialIdGenerator {
next: AtomicU64,
}
impl SequentialIdGenerator {
#[must_use]
pub const fn new(first: NonZeroU64) -> Self {
Self {
next: AtomicU64::new(first.get()),
}
}
fn next_raw(&self, kind: IdentifierKind) -> Result<u64, IdGenerationError> {
self.next
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
if current == 0 {
None
} else {
Some(current.checked_add(1).unwrap_or(0))
}
})
.map_err(|_| IdGenerationError::Exhausted { kind })
}
}
impl IdGenerator for SequentialIdGenerator {
fn next_job_instance_id(&self) -> Result<JobInstanceId, IdGenerationError> {
JobInstanceId::new(self.next_raw(IdentifierKind::JobInstance)?)
.map_err(IdGenerationError::Invalid)
}
fn next_job_execution_id(&self) -> Result<JobExecutionId, IdGenerationError> {
JobExecutionId::new(self.next_raw(IdentifierKind::JobExecution)?)
.map_err(IdGenerationError::Invalid)
}
fn next_step_execution_id(&self) -> Result<StepExecutionId, IdGenerationError> {
StepExecutionId::new(self.next_raw(IdentifierKind::StepExecution)?)
.map_err(IdGenerationError::Invalid)
}
fn next_failure_id(&self) -> Result<FailureId, IdGenerationError> {
FailureId::new(self.next_raw(IdentifierKind::Failure)?).map_err(IdGenerationError::Invalid)
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum IdGenerationError {
Exhausted {
kind: IdentifierKind,
},
Invalid(DomainError),
}
impl fmt::Display for IdGenerationError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Exhausted { kind } => write!(formatter, "{kind} identifier source is exhausted"),
Self::Invalid(error) => write!(formatter, "generated identifier was invalid: {error}"),
}
}
}
impl Error for IdGenerationError {
fn source(&self) -> Option<&(dyn Error + 'static)> {
match self {
Self::Invalid(error) => Some(error),
Self::Exhausted { .. } => None,
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum JobInstanceSelection {
Created(JobInstance),
Existing(JobInstance),
}
impl JobInstanceSelection {
#[must_use]
pub const fn instance(&self) -> &JobInstance {
match self {
Self::Created(instance) | Self::Existing(instance) => instance,
}
}
#[must_use]
pub const fn was_created(&self) -> bool {
matches!(self, Self::Created(_))
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum RecoveryDisposition {
MarkFailed,
Abandon,
}
impl RecoveryDisposition {
#[must_use]
pub const fn resulting_status(self) -> BatchStatus {
match self {
Self::MarkFailed => BatchStatus::Failed,
Self::Abandon => BatchStatus::Abandoned,
}
}
}
#[derive(Clone, Eq, PartialEq)]
pub struct RecoveryRequest {
expected_version: ExecutionVersion,
disposition: RecoveryDisposition,
reason_code: String,
operator_reference: String,
evidence_digest: [u8; 32],
failure: Option<FailureSummary>,
}
impl RecoveryRequest {
pub fn mark_failed(
expected_version: ExecutionVersion,
reason_code: impl Into<String>,
operator_reference: impl Into<String>,
evidence_digest: [u8; 32],
failure_category: FailureCategory,
failure_id: FailureId,
) -> Result<Self, RecoveryRequestError> {
Self::new(
expected_version,
RecoveryDisposition::MarkFailed,
reason_code,
operator_reference,
evidence_digest,
Some(FailureSummary::new(failure_category, failure_id)),
)
}
pub fn abandon(
expected_version: ExecutionVersion,
reason_code: impl Into<String>,
operator_reference: impl Into<String>,
evidence_digest: [u8; 32],
) -> Result<Self, RecoveryRequestError> {
Self::new(
expected_version,
RecoveryDisposition::Abandon,
reason_code,
operator_reference,
evidence_digest,
None,
)
}
fn new(
expected_version: ExecutionVersion,
disposition: RecoveryDisposition,
reason_code: impl Into<String>,
operator_reference: impl Into<String>,
evidence_digest: [u8; 32],
failure: Option<FailureSummary>,
) -> Result<Self, RecoveryRequestError> {
let reason_code = reason_code.into();
validate_recovery_text(
&reason_code,
RecoveryField::ReasonCode,
MAX_RECOVERY_REASON_BYTES,
)?;
let operator_reference = operator_reference.into();
validate_recovery_text(
&operator_reference,
RecoveryField::OperatorReference,
MAX_OPERATOR_REFERENCE_BYTES,
)?;
Ok(Self {
expected_version,
disposition,
reason_code,
operator_reference,
evidence_digest,
failure,
})
}
#[must_use]
pub const fn expected_version(&self) -> ExecutionVersion {
self.expected_version
}
#[must_use]
pub const fn disposition(&self) -> RecoveryDisposition {
self.disposition
}
#[must_use]
pub fn reason_code(&self) -> &str {
&self.reason_code
}
#[must_use]
pub fn operator_reference(&self) -> &str {
&self.operator_reference
}
#[must_use]
pub const fn evidence_digest(&self) -> &[u8; 32] {
&self.evidence_digest
}
#[must_use]
pub const fn failure(&self) -> Option<FailureSummary> {
self.failure
}
}
impl fmt::Debug for RecoveryRequest {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("RecoveryRequest")
.field("expected_version", &self.expected_version)
.field("disposition", &self.disposition)
.field("reason_code", &self.reason_code)
.field("operator_reference", &self.operator_reference)
.field("evidence_digest", &"<redacted>")
.field("failure", &self.failure)
.finish()
}
}
#[derive(Clone, Eq, PartialEq)]
pub struct RecoveryDecision {
id: RecoveryDecisionId,
job_execution_id: JobExecutionId,
execution_version: ExecutionVersion,
prior_status: BatchStatus,
resulting_status: BatchStatus,
reason_code: String,
operator_reference: String,
evidence_digest: [u8; 32],
decided_at: SystemTime,
}
impl RecoveryDecision {
#[allow(clippy::too_many_arguments)]
#[doc(hidden)]
#[must_use]
pub fn new(
id: RecoveryDecisionId,
job_execution_id: JobExecutionId,
execution_version: ExecutionVersion,
prior_status: BatchStatus,
resulting_status: BatchStatus,
reason_code: String,
operator_reference: String,
evidence_digest: [u8; 32],
decided_at: SystemTime,
) -> Self {
Self {
id,
job_execution_id,
execution_version,
prior_status,
resulting_status,
reason_code,
operator_reference,
evidence_digest,
decided_at,
}
}
#[must_use]
pub const fn id(&self) -> RecoveryDecisionId {
self.id
}
#[must_use]
pub const fn job_execution_id(&self) -> JobExecutionId {
self.job_execution_id
}
#[must_use]
pub const fn execution_version(&self) -> ExecutionVersion {
self.execution_version
}
#[must_use]
pub const fn prior_status(&self) -> BatchStatus {
self.prior_status
}
#[must_use]
pub const fn resulting_status(&self) -> BatchStatus {
self.resulting_status
}
#[must_use]
pub fn reason_code(&self) -> &str {
&self.reason_code
}
#[must_use]
pub fn operator_reference(&self) -> &str {
&self.operator_reference
}
#[must_use]
pub const fn evidence_digest(&self) -> &[u8; 32] {
&self.evidence_digest
}
#[must_use]
pub const fn decided_at(&self) -> SystemTime {
self.decided_at
}
}
impl fmt::Debug for RecoveryDecision {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("RecoveryDecision")
.field("id", &self.id)
.field("job_execution_id", &self.job_execution_id)
.field("execution_version", &self.execution_version)
.field("prior_status", &self.prior_status)
.field("resulting_status", &self.resulting_status)
.field("reason_code", &self.reason_code)
.field("operator_reference", &self.operator_reference)
.field("evidence_digest", &"<redacted>")
.field("decided_at", &self.decided_at)
.finish()
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct RecoveryResult {
execution: JobExecution,
decision: RecoveryDecision,
}
impl RecoveryResult {
#[doc(hidden)]
#[must_use]
pub const fn new(execution: JobExecution, decision: RecoveryDecision) -> Self {
Self {
execution,
decision,
}
}
#[must_use]
pub const fn execution(&self) -> &JobExecution {
&self.execution
}
#[must_use]
pub const fn decision(&self) -> &RecoveryDecision {
&self.decision
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum RecoveryField {
ReasonCode,
OperatorReference,
}
#[derive(Clone, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum RecoveryRequestError {
Empty {
field: RecoveryField,
},
TooLong {
field: RecoveryField,
max_bytes: usize,
},
SurroundingWhitespace {
field: RecoveryField,
},
ControlCharacter {
field: RecoveryField,
},
}
impl fmt::Display for RecoveryRequestError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Empty { field } => write!(formatter, "{field:?} must not be empty"),
Self::TooLong { field, max_bytes } => {
write!(formatter, "{field:?} exceeds {max_bytes} bytes")
}
Self::SurroundingWhitespace { field } => {
write!(formatter, "{field:?} has surrounding whitespace")
}
Self::ControlCharacter { field } => {
write!(formatter, "{field:?} contains a control character")
}
}
}
}
impl Error for RecoveryRequestError {}
fn validate_recovery_text(
value: &str,
field: RecoveryField,
max_bytes: usize,
) -> Result<(), RecoveryRequestError> {
if value.is_empty() {
return Err(RecoveryRequestError::Empty { field });
}
if value.len() > max_bytes {
return Err(RecoveryRequestError::TooLong { field, max_bytes });
}
if value.trim() != value {
return Err(RecoveryRequestError::SurroundingWhitespace { field });
}
if value.chars().any(char::is_control) {
return Err(RecoveryRequestError::ControlCharacter { field });
}
Ok(())
}
#[doc(hidden)]
pub fn recovered_execution(
prior: &JobExecution,
request: &RecoveryRequest,
decided_at: SystemTime,
) -> Result<JobExecution, RepositoryError> {
if prior.version() != request.expected_version() {
return Err(RepositoryError::Lifecycle(LifecycleError::StaleVersion {
expected: request.expected_version(),
actual: prior.version(),
}));
}
let prior_status = prior.metadata().status();
if !matches!(
prior_status,
BatchStatus::Starting | BatchStatus::Started | BatchStatus::Stopping | BatchStatus::Unknown
) {
return Err(RepositoryError::RecoveryNotAllowed {
id: prior.id(),
status: prior_status,
});
}
let current_time = prior.metadata().timestamps();
let timestamps = ExecutionTimestamps::new(
current_time.created_at(),
current_time.started_at(),
Some(decided_at),
)?;
let resulting_status = request.disposition().resulting_status();
let metadata = ExecutionMetadata::new(
resulting_status,
prior.metadata().exit_status().clone(),
timestamps,
prior.metadata().counts(),
request.failure(),
)?;
Ok(JobExecution::from_snapshot(
prior.id(),
prior.job_instance_id(),
metadata,
prior.version().next()?,
))
}
pub trait JobRepository: Send + Sync {
fn connection_capacity(&self) -> u32 {
1
}
fn descriptor(&self) -> RepositoryDescriptor {
RepositoryDescriptor::new(0, [])
}
fn begin<'a>(
&'a self,
) -> BoxFuture<'a, Result<Box<dyn RepositoryUnitOfWork + 'a>, RepositoryError>>;
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct ExecutionControl {
execution: JobExecution,
owner_matches: bool,
stop_requested: bool,
}
impl ExecutionControl {
#[doc(hidden)]
#[must_use]
pub const fn new(execution: JobExecution, owner_matches: bool, stop_requested: bool) -> Self {
Self {
execution,
owner_matches,
stop_requested,
}
}
#[must_use]
pub const fn execution(&self) -> &JobExecution {
&self.execution
}
#[must_use]
pub const fn owner_matches(&self) -> bool {
self.owner_matches
}
#[must_use]
pub const fn stop_requested(&self) -> bool {
self.stop_requested
}
}
pub trait RepositoryUnitOfWork: Send {
fn register_definition_upgrade<'a>(
&'a mut self,
job_name: &'a JobName,
upgrade: &'a DefinitionUpgrade,
) -> BoxFuture<'a, Result<(), RepositoryError>>;
fn select_or_create_job_instance<'a>(
&'a mut self,
key: &'a JobInstanceKey,
) -> BoxFuture<'a, Result<JobInstanceSelection, RepositoryError>>;
fn create_job_execution(
&mut self,
job_instance_id: JobInstanceId,
) -> BoxFuture<'_, Result<JobExecution, RepositoryError>>;
fn create_job_execution_with_definition<'a>(
&'a mut self,
job_instance_id: JobInstanceId,
definition: &'a DefinitionIdentity,
) -> BoxFuture<'a, Result<JobExecution, RepositoryError>>;
fn create_step_execution<'a>(
&'a mut self,
job_execution_id: JobExecutionId,
step_name: &'a StepName,
) -> BoxFuture<'a, Result<StepExecution, RepositoryError>>;
fn create_flow_step_execution<'a>(
&'a mut self,
_job_execution_id: JobExecutionId,
_step_name: &'a StepName,
_node_id: &'a NodeId,
_start_limit: StartLimit,
) -> BoxFuture<'a, Result<StepExecution, RepositoryError>> {
Box::pin(async { Err(RepositoryError::FlowStateCorrupt) })
}
fn transition_job_execution(
&mut self,
id: JobExecutionId,
expected_version: ExecutionVersion,
transition: LifecycleTransition,
) -> BoxFuture<'_, Result<JobExecution, RepositoryError>>;
fn enrich_job_exit_status<'a>(
&'a mut self,
id: JobExecutionId,
expected_version: ExecutionVersion,
exit_status: &'a ExitStatus,
) -> BoxFuture<'a, Result<JobExecution, RepositoryError>>;
fn transition_step_execution(
&mut self,
id: StepExecutionId,
expected_version: ExecutionVersion,
transition: LifecycleTransition,
) -> BoxFuture<'_, Result<StepExecution, RepositoryError>>;
fn enrich_step_exit_status<'a>(
&'a mut self,
id: StepExecutionId,
expected_version: ExecutionVersion,
exit_status: &'a ExitStatus,
) -> BoxFuture<'a, Result<StepExecution, RepositoryError>>;
fn find_job_instance<'a>(
&'a mut self,
key: &'a JobInstanceKey,
) -> BoxFuture<'a, Result<Option<JobInstance>, RepositoryError>>;
fn get_job_instance(
&mut self,
id: JobInstanceId,
) -> BoxFuture<'_, Result<Option<JobInstance>, RepositoryError>>;
fn get_job_execution(
&mut self,
id: JobExecutionId,
) -> BoxFuture<'_, Result<Option<JobExecution>, RepositoryError>>;
fn job_executions(
&mut self,
job_instance_id: JobInstanceId,
) -> BoxFuture<'_, Result<Vec<JobExecution>, RepositoryError>>;
fn get_step_execution(
&mut self,
id: StepExecutionId,
) -> BoxFuture<'_, Result<Option<StepExecution>, RepositoryError>>;
fn step_executions(
&mut self,
job_execution_id: JobExecutionId,
) -> BoxFuture<'_, Result<Vec<StepExecution>, RepositoryError>>;
fn latest_flow_step<'a>(
&'a mut self,
_job_instance_id: JobInstanceId,
_node_id: &'a NodeId,
) -> BoxFuture<'a, Result<Option<FlowStepState>, RepositoryError>> {
Box::pin(async { Err(RepositoryError::FlowStateCorrupt) })
}
fn append_flow_decision<'a>(
&'a mut self,
_request: &'a FlowDecisionRequest,
) -> BoxFuture<'a, Result<FlowDecision, RepositoryError>> {
Box::pin(async { Err(RepositoryError::FlowStateCorrupt) })
}
fn find_reusable_flow_decision<'a>(
&'a mut self,
_job_instance_id: JobInstanceId,
_node_id: &'a NodeId,
_plan_fingerprint: &'a [u8; 32],
_input_digest: &'a [u8; 32],
_kind: FlowTransitionKind,
) -> BoxFuture<'a, Result<Option<FlowDecision>, RepositoryError>> {
Box::pin(async { Err(RepositoryError::FlowStateCorrupt) })
}
fn flow_decisions(
&mut self,
_job_execution_id: JobExecutionId,
) -> BoxFuture<'_, Result<Vec<FlowDecision>, RepositoryError>> {
Box::pin(async { Err(RepositoryError::FlowStateCorrupt) })
}
fn create_step_partition_plan<'a>(
&'a mut self,
_step_execution_id: StepExecutionId,
_entries: &'a [PartitionPlanEntry],
) -> BoxFuture<'a, Result<Vec<StepPartition>, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::StepPartitions,
})
})
}
fn step_partition_plan(
&mut self,
_step_execution_id: StepExecutionId,
) -> BoxFuture<'_, Result<Vec<StepPartition>, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::StepPartitions,
})
})
}
fn restart_step_partition_plan(
&mut self,
_source_step_execution_id: StepExecutionId,
_target_step_execution_id: StepExecutionId,
) -> BoxFuture<'_, Result<Vec<StepPartition>, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::StepPartitions,
})
})
}
fn assign_step_partition(
&mut self,
_id: StepPartitionId,
_expected_version: ExecutionVersion,
_worker_step_execution_id: StepExecutionId,
) -> BoxFuture<'_, Result<StepPartition, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::StepPartitions,
})
})
}
fn complete_step_partition(
&mut self,
_id: StepPartitionId,
_expected_version: ExecutionVersion,
_worker_step_execution_id: StepExecutionId,
) -> BoxFuture<'_, Result<StepPartition, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::StepPartitions,
})
})
}
fn aggregate_step_partitions(
&mut self,
_step_execution_id: StepExecutionId,
_expected_version: ExecutionVersion,
_transitioned_at: SystemTime,
) -> BoxFuture<'_, Result<StepExecution, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::StepPartitions,
})
})
}
fn recover_job_execution<'a>(
&'a mut self,
id: JobExecutionId,
request: &'a RecoveryRequest,
) -> BoxFuture<'a, Result<RecoveryResult, RepositoryError>>;
fn recovery_decision(
&mut self,
id: JobExecutionId,
) -> BoxFuture<'_, Result<Option<RecoveryDecision>, RepositoryError>>;
fn find_operator_request<'a>(
&'a mut self,
_action: OperatorAction,
_operation_id: &'a OperationId,
) -> BoxFuture<'a, Result<Option<OperatorRecord>, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::OperatorRequests,
})
})
}
fn append_operator_request<'a>(
&'a mut self,
_draft: &'a OperatorRecordDraft,
) -> BoxFuture<'a, Result<OperatorRecord, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::OperatorRequests,
})
})
}
fn request_execution_stop<'a>(
&'a mut self,
_id: JobExecutionId,
_expected_version: ExecutionVersion,
_actor: &'a ActorRef,
_requested_at: SystemTime,
) -> BoxFuture<'a, Result<JobExecution, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::StopRequests,
})
})
}
fn claim_execution_owner<'a>(
&'a mut self,
_id: JobExecutionId,
_expected_version: ExecutionVersion,
_owner: &'a OwnerToken,
_claimed_at: SystemTime,
) -> BoxFuture<'a, Result<JobExecution, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::ExecutionOwnership,
})
})
}
fn observe_execution_control<'a>(
&'a mut self,
_id: JobExecutionId,
_owner: &'a OwnerToken,
_observed_at: SystemTime,
) -> BoxFuture<'a, Result<ExecutionControl, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::ExecutionOwnership,
})
})
}
fn job_instance_hold(
&mut self,
_id: JobInstanceId,
) -> BoxFuture<'_, Result<Option<RetentionHold>, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::InstanceHolds,
})
})
}
fn place_instance_hold<'a>(
&'a mut self,
_id: JobInstanceId,
_actor: &'a ActorRef,
_reason: &'a ReasonCode,
_placed_at: SystemTime,
) -> BoxFuture<'a, Result<RetentionHold, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::InstanceHolds,
})
})
}
fn release_instance_hold(
&mut self,
_id: JobInstanceId,
) -> BoxFuture<'_, Result<Option<RetentionHold>, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::InstanceHolds,
})
})
}
fn find_retention_action<'a>(
&'a mut self,
_action: RetentionAction,
_operation_id: &'a OperationId,
) -> BoxFuture<'a, Result<Option<RetentionRecord>, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::RetentionPurge,
})
})
}
fn append_retention_action<'a>(
&'a mut self,
_draft: &'a RetentionRecordDraft,
) -> BoxFuture<'a, Result<RetentionRecord, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::RetentionPurge,
})
})
}
fn purge_survey<'a>(
&'a mut self,
_request: &'a PurgePlanRequest,
) -> BoxFuture<'a, Result<PurgeSurvey, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::RetentionPurge,
})
})
}
fn apply_purge<'a>(
&'a mut self,
_plan: &'a PurgePlan,
) -> BoxFuture<'a, Result<PurgeCounts, RepositoryError>> {
Box::pin(async {
Err(RepositoryError::UnsupportedCapability {
capability: RepositoryCapability::RetentionPurge,
})
})
}
fn commit<'a>(self: Box<Self>) -> BoxFuture<'a, Result<(), RepositoryError>>
where
Self: 'a;
fn rollback<'a>(self: Box<Self>) -> BoxFuture<'a, Result<(), RepositoryError>>
where
Self: 'a;
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
#[non_exhaustive]
pub enum RepositoryCapability {
OperatorRequests,
StopRequests,
ExecutionOwnership,
InstanceHolds,
RetentionPurge,
StepPartitions,
}
impl RepositoryCapability {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::OperatorRequests => "operator requests",
Self::StopRequests => "durable stop requests",
Self::ExecutionOwnership => "execution ownership evidence",
Self::InstanceHolds => "instance holds",
Self::RetentionPurge => "retention purge",
Self::StepPartitions => "durable step partitions",
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct RepositoryDescriptor {
descriptor_version: u32,
schema_version: u32,
capabilities: BTreeSet<RepositoryCapability>,
}
impl RepositoryDescriptor {
pub const CURRENT_VERSION: u32 = 1;
#[must_use]
pub fn new(
schema_version: u32,
capabilities: impl IntoIterator<Item = RepositoryCapability>,
) -> Self {
Self {
descriptor_version: Self::CURRENT_VERSION,
schema_version,
capabilities: capabilities.into_iter().collect(),
}
}
#[must_use]
pub const fn descriptor_version(&self) -> u32 {
self.descriptor_version
}
#[must_use]
pub const fn schema_version(&self) -> u32 {
self.schema_version
}
#[must_use]
pub fn declares(&self, capability: RepositoryCapability) -> bool {
self.capabilities.contains(&capability)
}
#[must_use]
pub fn capabilities(&self) -> impl ExactSizeIterator<Item = RepositoryCapability> + '_ {
self.capabilities.iter().copied()
}
pub fn require(&self, capability: RepositoryCapability) -> Result<(), RepositoryError> {
if self.declares(capability) {
Ok(())
} else {
Err(RepositoryError::UnsupportedCapability { capability })
}
}
}
impl fmt::Display for RepositoryCapability {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(self.as_str())
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum RepositoryError {
SchemaUninitialized,
MigrationRequired {
current: u32,
supported: u32,
},
NewerSchema {
current: u32,
supported: u32,
},
IdentifierOutOfRange {
kind: IdentifierKind,
value: u64,
},
JobInstanceNotFound {
id: JobInstanceId,
},
JobExecutionNotFound {
id: JobExecutionId,
},
StepExecutionNotFound {
id: StepExecutionId,
},
StepPartitionNotFound {
id: StepPartitionId,
},
EmptyPartitionPlan,
PartitionPlanTooLarge {
max: usize,
},
DuplicatePartitionKey,
PartitionPlanExists {
step_execution_id: StepExecutionId,
},
PartitionPlanNotCommitted {
step_execution_id: StepExecutionId,
},
PartitionUpdateNotAllowed {
id: StepPartitionId,
status: BatchStatus,
},
PartitionWorkerMismatch {
partition_id: StepPartitionId,
worker_step_execution_id: StepExecutionId,
},
PartitionWorkerAlreadyAssigned {
worker_step_execution_id: StepExecutionId,
},
PartitionWorkerStale {
partition_id: StepPartitionId,
worker_step_execution_id: StepExecutionId,
},
PartitionParentNotActive {
step_execution_id: StepExecutionId,
status: BatchStatus,
},
PartitionAggregationIncomplete {
step_execution_id: StepExecutionId,
status: BatchStatus,
},
PartitionStateCorrupt,
DuplicateIdentifier {
kind: IdentifierKind,
value: u64,
},
CompletedInstance {
id: JobInstanceId,
},
AbandonedInstance {
id: JobInstanceId,
},
ExecutionAlreadyActive {
instance_id: JobInstanceId,
execution_id: JobExecutionId,
status: BatchStatus,
},
DefinitionDrift {
job_name: JobName,
revision: DefinitionRevision,
},
DefinitionJobMismatch {
expected: JobName,
actual: JobName,
},
IncompatibleDefinition {
instance_id: JobInstanceId,
},
UnsupportedManifestVersion {
format: u16,
},
InvalidDefinitionUpgrade {
execution_id: JobExecutionId,
},
DefinitionUpgradeConflict {
job_name: JobName,
},
RestartStateNotFound {
execution_id: JobExecutionId,
step_name: StepName,
},
FaultStateCorrupt,
StartLimitExceeded {
instance_id: JobInstanceId,
node_id: NodeId,
limit: StartLimit,
},
FlowStateCorrupt,
RecoveryNotAllowed {
id: JobExecutionId,
status: BatchStatus,
},
ExecutionOwned {
id: JobExecutionId,
},
ExecutionOwnershipNotAllowed {
id: JobExecutionId,
status: BatchStatus,
},
Domain(DomainError),
Identifier(IdGenerationError),
Lifecycle(LifecycleError),
RetentionPlanStale,
UnsupportedCapability {
capability: RepositoryCapability,
},
ConcurrentModification,
CommitOutcomeUnknown,
Unavailable,
}
impl fmt::Display for RepositoryError {
#[allow(clippy::too_many_lines)]
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::SchemaUninitialized => {
formatter.write_str("PostgreSQL metadata schema is not initialized")
}
Self::MigrationRequired { current, supported } => write!(
formatter,
"PostgreSQL metadata schema version {current} requires migration to {supported}"
),
Self::NewerSchema { current, supported } => write!(
formatter,
"PostgreSQL metadata schema version {current} is newer than supported version {supported}"
),
Self::IdentifierOutOfRange { kind, value } => {
write!(
formatter,
"{kind} identifier {value} exceeds PostgreSQL bigint"
)
}
Self::JobInstanceNotFound { id } => {
write!(formatter, "job instance {id} was not found")
}
Self::JobExecutionNotFound { id } => {
write!(formatter, "job execution {id} was not found")
}
Self::StepExecutionNotFound { id } => {
write!(formatter, "step execution {id} was not found")
}
Self::StepPartitionNotFound { id } => {
write!(formatter, "step partition {id} was not found")
}
Self::EmptyPartitionPlan => {
formatter.write_str("partition plan must contain at least one entry")
}
Self::PartitionPlanTooLarge { max } => {
write!(formatter, "partition plan exceeds {max} entries")
}
Self::DuplicatePartitionKey => {
formatter.write_str("partition plan contains a duplicate key")
}
Self::PartitionPlanExists { step_execution_id } => write!(
formatter,
"step execution {step_execution_id} already has a partition plan"
),
Self::PartitionPlanNotCommitted { step_execution_id } => write!(
formatter,
"step execution {step_execution_id} partition plan must commit before assignment"
),
Self::PartitionUpdateNotAllowed { id, status } => write!(
formatter,
"step partition {id} cannot be updated from {status}"
),
Self::PartitionWorkerMismatch {
partition_id,
worker_step_execution_id,
} => write!(
formatter,
"worker step execution {worker_step_execution_id} does not belong to partition {partition_id}"
),
Self::PartitionWorkerAlreadyAssigned {
worker_step_execution_id,
} => write!(
formatter,
"worker step execution {worker_step_execution_id} is already assigned to a partition"
),
Self::PartitionWorkerStale {
partition_id,
worker_step_execution_id,
} => write!(
formatter,
"worker step execution {worker_step_execution_id} is not the current worker for partition {partition_id}"
),
Self::PartitionParentNotActive {
step_execution_id,
status,
} => write!(
formatter,
"partition parent step execution {step_execution_id} cannot mutate children from {status}"
),
Self::PartitionAggregationIncomplete {
step_execution_id,
status,
} => write!(
formatter,
"step execution {step_execution_id} cannot aggregate a child in {status}"
),
Self::PartitionStateCorrupt => {
formatter.write_str("durable partition state is unusable and no work may begin")
}
Self::DuplicateIdentifier { kind, value } => {
write!(formatter, "{kind} identifier {value} already exists")
}
Self::CompletedInstance { id } => {
write!(formatter, "job instance {id} is already completed")
}
Self::AbandonedInstance { id } => {
write!(formatter, "job instance {id} is abandoned")
}
Self::ExecutionAlreadyActive {
instance_id,
execution_id,
status,
} => write!(
formatter,
"job instance {instance_id} already has execution {execution_id} in {status}"
),
Self::DefinitionDrift { job_name, revision } => write!(
formatter,
"job {job_name} definition revision {} has drifted",
revision.as_str()
),
Self::DefinitionJobMismatch { expected, actual } => write!(
formatter,
"definition for job {actual} cannot be used for job {expected}"
),
Self::IncompatibleDefinition { instance_id } => write!(
formatter,
"job instance {instance_id} has no direct compatible definition"
),
Self::UnsupportedManifestVersion { format } => {
write!(
formatter,
"definition manifest format {format} is unsupported"
)
}
Self::InvalidDefinitionUpgrade { execution_id } => write!(
formatter,
"definition upgrade for execution {execution_id} is incomplete"
),
Self::DefinitionUpgradeConflict { job_name } => {
write!(formatter, "job {job_name} definition upgrade conflicts")
}
Self::RestartStateNotFound {
execution_id,
step_name,
} => write!(
formatter,
"restart execution {execution_id} has no durable source for step {step_name}"
),
Self::FaultStateCorrupt => {
formatter.write_str("durable fault state is unusable and no work may begin")
}
Self::StartLimitExceeded {
instance_id,
node_id,
limit,
} => write!(
formatter,
"job instance {instance_id} exhausted start limit {} for node {}",
limit.get(),
node_id.as_str()
),
Self::FlowStateCorrupt => {
formatter.write_str("durable flow history is unusable and no work may begin")
}
Self::RecoveryNotAllowed { id, status } => {
write!(
formatter,
"job execution {id} in {status} cannot be recovered"
)
}
Self::ExecutionOwned { id } => {
write!(formatter, "job execution {id} is owned by another process")
}
Self::ExecutionOwnershipNotAllowed { id, status } => write!(
formatter,
"job execution {id} in {status} cannot acquire process ownership"
),
Self::Domain(error) => write!(formatter, "invalid repository domain value: {error}"),
Self::Identifier(error) => write!(formatter, "identifier generation failed: {error}"),
Self::Lifecycle(error) => error.fmt(formatter),
Self::RetentionPlanStale => {
formatter.write_str("the purge plan is stale and nothing was deleted")
}
Self::UnsupportedCapability { capability } => {
write!(formatter, "the adapter does not support {capability}")
}
Self::ConcurrentModification => {
formatter.write_str("repository unit of work is based on a stale snapshot")
}
Self::CommitOutcomeUnknown => formatter.write_str(
"PostgreSQL commit outcome is unknown; inspect durable metadata before recovery",
),
Self::Unavailable => formatter.write_str("repository is unavailable"),
}
}
}
#[doc(hidden)]
pub fn aggregate_partition_parent(
parent: &StepExecution,
expected_version: ExecutionVersion,
aggregate: &PartitionAggregate,
transitioned_at: SystemTime,
failure: Option<FailureSummary>,
) -> Result<StepExecution, RepositoryError> {
let transition = if aggregate.status() == BatchStatus::Failed {
LifecycleTransition::failed(
transitioned_at,
failure.ok_or(LifecycleError::FailedTransitionMissingFailure)?,
)
} else {
LifecycleTransition::new(aggregate.status(), transitioned_at)
};
let mut transitioned = parent.clone();
transitioned.transition(expected_version, transition)?;
let metadata = ExecutionMetadata::new(
aggregate.status(),
aggregate.exit_status().clone(),
transitioned.metadata().timestamps(),
aggregate.counts(),
transitioned.metadata().failure(),
)?;
Ok(StepExecution::from_snapshot(
transitioned.id(),
transitioned.job_execution_id(),
transitioned.step_name().clone(),
metadata,
transitioned.version(),
))
}
#[doc(hidden)]
#[must_use]
pub fn map_partition_aggregation(
step_execution_id: StepExecutionId,
error: PartitionAggregationError,
) -> RepositoryError {
match error {
PartitionAggregationError::Incomplete { status } => {
RepositoryError::PartitionAggregationIncomplete {
step_execution_id,
status,
}
}
PartitionAggregationError::CountExhausted => {
RepositoryError::Lifecycle(LifecycleError::CountExhausted)
}
PartitionAggregationError::EmptyPlan
| PartitionAggregationError::PlanTooLarge { .. }
| PartitionAggregationError::DuplicateKey => RepositoryError::PartitionStateCorrupt,
}
}
impl Error for RepositoryError {
fn source(&self) -> Option<&(dyn Error + 'static)> {
match self {
Self::Domain(error) => Some(error),
Self::Identifier(error) => Some(error),
Self::Lifecycle(error) => Some(error),
_ => None,
}
}
}
impl From<DomainError> for RepositoryError {
fn from(error: DomainError) -> Self {
Self::Domain(error)
}
}
impl From<IdGenerationError> for RepositoryError {
fn from(error: IdGenerationError) -> Self {
Self::Identifier(error)
}
}
impl From<LifecycleError> for RepositoryError {
fn from(error: LifecycleError) -> Self {
Self::Lifecycle(error)
}
}