use std::error::Error;
use std::fmt;
use std::sync::Arc;
use std::time::SystemTime;
use crate::{
BatchStatus, Clock, ExecutionVersion, JobExecution, JobExecutionId, JobInstanceId,
JobRepository, LifecycleTransition, OperationId, OperatorAction, OperatorOutcomeClass,
OperatorRecord, OperatorRecordDraft, OperatorRejection, OperatorRequest, ReasonCode,
RecoveryDirective, RecoveryRequestError, RepositoryError, RepositoryUnitOfWork,
TelemetryEventKind, TelemetryEventSink, TelemetryRecord,
};
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct OperatorOutcome {
class: OperatorOutcomeClass,
record: OperatorRecord,
execution: Option<JobExecution>,
changed: bool,
}
impl OperatorOutcome {
const fn new(
class: OperatorOutcomeClass,
record: OperatorRecord,
execution: Option<JobExecution>,
changed: bool,
) -> Self {
Self {
class,
record,
execution,
changed,
}
}
#[must_use]
pub const fn class(&self) -> OperatorOutcomeClass {
self.class
}
#[must_use]
pub const fn record(&self) -> &OperatorRecord {
&self.record
}
#[must_use]
pub const fn execution(&self) -> Option<&JobExecution> {
self.execution.as_ref()
}
#[must_use]
pub const fn rejection(&self) -> Option<OperatorRejection> {
self.record.rejection()
}
#[must_use]
pub const fn changed_state(&self) -> bool {
self.changed
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum OperatorError {
OperationIdConflict {
action: OperatorAction,
operation_id: OperationId,
},
OperationOutcomeUnknown,
InvalidRecoveryRequest(RecoveryRequestError),
Repository(RepositoryError),
}
impl fmt::Display for OperatorError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::OperationIdConflict {
action,
operation_id,
} => write!(
formatter,
"operation identifier {operation_id} was already recorded for {action} with a different request"
),
Self::OperationOutcomeUnknown => {
formatter.write_str("the operator commit outcome is unknown")
}
Self::InvalidRecoveryRequest(error) => error.fmt(formatter),
Self::Repository(error) => error.fmt(formatter),
}
}
}
impl Error for OperatorError {
fn source(&self) -> Option<&(dyn Error + 'static)> {
match self {
Self::InvalidRecoveryRequest(error) => Some(error),
Self::Repository(error) => Some(error),
_ => None,
}
}
}
impl From<RepositoryError> for OperatorError {
fn from(value: RepositoryError) -> Self {
match value {
RepositoryError::CommitOutcomeUnknown => Self::OperationOutcomeUnknown,
other => Self::Repository(other),
}
}
}
#[derive(Clone)]
pub struct JobOperator<R> {
repository: R,
clock: Arc<dyn Clock>,
event_sinks: Vec<Arc<dyn TelemetryEventSink>>,
}
impl<R> fmt::Debug for JobOperator<R> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("JobOperator")
.finish_non_exhaustive()
}
}
impl<R: JobRepository> JobOperator<R> {
pub const fn new(repository: R, clock: Arc<dyn Clock>) -> Self {
Self {
repository,
clock,
event_sinks: Vec::new(),
}
}
#[must_use]
pub fn with_event_sink(mut self, sink: Arc<dyn TelemetryEventSink>) -> Self {
self.event_sinks.push(sink);
self
}
pub const fn repository(&self) -> &R {
&self.repository
}
pub async fn execute(
&self,
request: &OperatorRequest,
) -> Result<OperatorOutcome, OperatorError> {
if let Some(recorded) = self.replay(request).await? {
self.emit_outcome(request, &recorded);
return Ok(recorded);
}
let requested_at = self.clock.now();
let mut unit = self.repository.begin().await?;
let effect = match self.apply(unit.as_mut(), request).await {
Ok(effect) => effect,
Err(EffectFailure::Rejected(rejection)) => {
let _ = unit.rollback().await;
let outcome = self
.audit_rejection(request, rejection, requested_at)
.await?;
self.emit_outcome(request, &outcome);
return Ok(outcome);
}
Err(EffectFailure::Failed(error)) => {
let _ = unit.rollback().await;
return Err(error);
}
};
let draft = OperatorRecordDraft::applied(
request,
effect.job_instance_id,
effect.job_execution_id,
effect.prior_status,
effect.result_status,
requested_at,
);
let record = match unit.append_operator_request(&draft).await {
Ok(record) => record,
Err(RepositoryError::ConcurrentModification) => {
let _ = unit.rollback().await;
return self.replay(request).await?.ok_or(OperatorError::Repository(
RepositoryError::ConcurrentModification,
));
}
Err(error) => return Err(error.into()),
};
unit.commit().await?;
let outcome = OperatorOutcome::new(
OperatorOutcomeClass::Applied,
record,
effect.execution,
effect.changed,
);
self.emit_outcome(request, &outcome);
Ok(outcome)
}
fn emit_outcome(&self, request: &OperatorRequest, outcome: &OperatorOutcome) {
let primary = match outcome.class() {
OperatorOutcomeClass::Applied | OperatorOutcomeClass::Replayed => {
TelemetryEventKind::OperatorRequestAccepted
}
_ => TelemetryEventKind::OperatorRequestRejected,
};
self.emit_record(&TelemetryRecord::operator(
primary,
request,
Some(outcome.class()),
outcome.rejection(),
));
if request.action() == OperatorAction::Recover {
let recovery = if outcome.class() == OperatorOutcomeClass::Rejected {
TelemetryEventKind::RecoveryRejected
} else {
TelemetryEventKind::RecoveryApplied
};
self.emit_record(&TelemetryRecord::operator(
recovery,
request,
Some(outcome.class()),
outcome.rejection(),
));
}
self.emit_record(&TelemetryRecord::operator(
TelemetryEventKind::OperatorRequestCompleted,
request,
Some(outcome.class()),
outcome.rejection(),
));
}
fn emit_record(&self, record: &TelemetryRecord) {
for sink in &self.event_sinks {
crate::telemetry::emit_safely(Some(sink), record);
}
}
async fn replay(
&self,
request: &OperatorRequest,
) -> Result<Option<OperatorOutcome>, OperatorError> {
let mut unit = self.repository.begin().await?;
let recorded = unit
.find_operator_request(request.action(), request.operation_id())
.await?;
unit.rollback().await?;
let Some(record) = recorded else {
return Ok(None);
};
if record.digest() != request.digest() {
return Err(OperatorError::OperationIdConflict {
action: request.action(),
operation_id: request.operation_id().clone(),
});
}
Ok(Some(OperatorOutcome::new(
OperatorOutcomeClass::Replayed,
record,
None,
false,
)))
}
async fn audit_rejection(
&self,
request: &OperatorRequest,
rejection: OperatorRejection,
requested_at: SystemTime,
) -> Result<OperatorOutcome, OperatorError> {
let draft = OperatorRecordDraft::rejected(request, rejection, requested_at);
let mut unit = self.repository.begin().await?;
let record = match unit.append_operator_request(&draft).await {
Ok(record) => record,
Err(RepositoryError::ConcurrentModification) => {
let _ = unit.rollback().await;
return self.replay(request).await?.ok_or(OperatorError::Repository(
RepositoryError::ConcurrentModification,
));
}
Err(error) => return Err(error.into()),
};
unit.commit().await?;
Ok(OperatorOutcome::new(
OperatorOutcomeClass::Rejected,
record,
None,
false,
))
}
async fn apply(
&self,
unit: &mut dyn RepositoryUnitOfWork,
request: &OperatorRequest,
) -> Result<AppliedEffect, EffectFailure> {
match request.action() {
OperatorAction::Launch => self.launch(unit, request).await,
OperatorAction::Restart => self.restart(unit, request).await,
OperatorAction::Stop => self.stop(unit, request).await,
OperatorAction::Abandon => self.abandon(unit, request).await,
OperatorAction::Recover => self.recover(unit, request).await,
_ => Err(EffectFailure::Rejected(
OperatorRejection::UnsupportedAction,
)),
}
}
async fn launch(
&self,
unit: &mut dyn RepositoryUnitOfWork,
request: &OperatorRequest,
) -> Result<AppliedEffect, EffectFailure> {
let Some(key) = request.job_instance_key() else {
return Err(EffectFailure::Rejected(OperatorRejection::InstanceNotFound));
};
let definition = request
.definition()
.ok_or(EffectFailure::Rejected(
OperatorRejection::IncompatibleDefinition,
))?
.clone();
let selection = unit
.select_or_create_job_instance(key)
.await
.map_err(EffectFailure::classify)?;
let instance_id = selection.instance().id();
let execution = unit
.create_job_execution_with_definition(instance_id, &definition)
.await
.map_err(EffectFailure::classify)?;
Ok(AppliedEffect::created(instance_id, execution))
}
async fn restart(
&self,
unit: &mut dyn RepositoryUnitOfWork,
request: &OperatorRequest,
) -> Result<AppliedEffect, EffectFailure> {
let Some(instance_id) = request.job_instance_id() else {
return Err(EffectFailure::Rejected(OperatorRejection::InstanceNotFound));
};
let definition = request
.definition()
.ok_or(EffectFailure::Rejected(
OperatorRejection::IncompatibleDefinition,
))?
.clone();
let prior = unit
.job_executions(instance_id)
.await
.map_err(EffectFailure::classify)?;
let latest = prior.last().ok_or(EffectFailure::Rejected(
OperatorRejection::RestartWithoutPriorAttempt,
))?;
let prior_status = latest.metadata().status();
if matches!(prior_status, BatchStatus::Unknown) {
return Err(EffectFailure::Rejected(OperatorRejection::InvalidState {
status: prior_status,
}));
}
let execution = unit
.create_job_execution_with_definition(instance_id, &definition)
.await
.map_err(EffectFailure::classify)?;
Ok(AppliedEffect::created(instance_id, execution).with_prior(prior_status))
}
async fn stop(
&self,
unit: &mut dyn RepositoryUnitOfWork,
request: &OperatorRequest,
) -> Result<AppliedEffect, EffectFailure> {
let (id, expected_version) = execution_target(request)?;
let observed = unit
.get_job_execution(id)
.await
.map_err(EffectFailure::classify)?
.ok_or(EffectFailure::Rejected(
OperatorRejection::ExecutionNotFound,
))?;
let status = observed.metadata().status();
if !matches!(status, BatchStatus::Starting | BatchStatus::Started) {
if matches!(status, BatchStatus::Stopping) || status.is_finished() {
return Ok(AppliedEffect::unchanged(&observed));
}
return Err(EffectFailure::Rejected(OperatorRejection::InvalidState {
status,
}));
}
let execution = unit
.request_execution_stop(id, expected_version, request.actor(), self.clock.now())
.await
.map_err(EffectFailure::classify)?;
Ok(AppliedEffect::updated(&execution, status))
}
async fn abandon(
&self,
unit: &mut dyn RepositoryUnitOfWork,
request: &OperatorRequest,
) -> Result<AppliedEffect, EffectFailure> {
let (id, expected_version) = execution_target(request)?;
let observed = unit
.get_job_execution(id)
.await
.map_err(EffectFailure::classify)?
.ok_or(EffectFailure::Rejected(
OperatorRejection::ExecutionNotFound,
))?;
let status = observed.metadata().status();
match status {
BatchStatus::Abandoned => return Ok(AppliedEffect::unchanged(&observed)),
BatchStatus::Stopped | BatchStatus::Failed => {}
BatchStatus::Unknown => {
let decision = unit
.recovery_decision(id)
.await
.map_err(EffectFailure::classify)?;
if decision.is_none() {
return Err(EffectFailure::Rejected(
OperatorRejection::UnresolvedRecoveryRequired,
));
}
}
other => {
return Err(EffectFailure::Rejected(OperatorRejection::InvalidState {
status: other,
}));
}
}
let transition = LifecycleTransition::new(BatchStatus::Abandoned, self.clock.now());
let execution = unit
.transition_job_execution(id, expected_version, transition)
.await
.map_err(EffectFailure::classify)?;
Ok(AppliedEffect::updated(&execution, status))
}
async fn recover(
&self,
unit: &mut dyn RepositoryUnitOfWork,
request: &OperatorRequest,
) -> Result<AppliedEffect, EffectFailure> {
let (id, expected_version) = execution_target(request)?;
let execution = unit
.get_job_execution(id)
.await
.map_err(EffectFailure::classify)?
.ok_or(EffectFailure::Rejected(
OperatorRejection::ExecutionNotFound,
))?;
if execution.version() != expected_version {
return Err(EffectFailure::Rejected(
OperatorRejection::OptimisticConflict {
current: execution.version(),
},
));
}
let (directive, unknown_commit) = request.recovery_guard().ok_or(
EffectFailure::Rejected(OperatorRejection::InvalidState {
status: execution.metadata().status(),
}),
)?;
if unknown_commit
&& matches!(directive, RecoveryDirective::MarkFailed(_))
&& request.reason().map(ReasonCode::as_str) != Some("UNKNOWN_EFFECT")
{
return Err(EffectFailure::Rejected(
OperatorRejection::UnresolvedRecoveryRequired,
));
}
let recovery = request
.recovery_request()
.ok_or(EffectFailure::Rejected(OperatorRejection::InvalidState {
status: BatchStatus::Unknown,
}))?
.map_err(|error| EffectFailure::Failed(OperatorError::InvalidRecoveryRequest(error)))?;
let result = unit
.recover_job_execution(id, &recovery)
.await
.map_err(EffectFailure::classify)?;
Ok(AppliedEffect::updated(
result.execution(),
result.decision().prior_status(),
))
}
}
fn execution_target(
request: &OperatorRequest,
) -> Result<(JobExecutionId, ExecutionVersion), EffectFailure> {
let Some(id) = request.job_execution_id() else {
return Err(EffectFailure::Rejected(
OperatorRejection::ExecutionNotFound,
));
};
let expected_version = request.expected_version().ok_or(EffectFailure::Rejected(
OperatorRejection::InvalidState {
status: BatchStatus::Unknown,
},
))?;
Ok((id, expected_version))
}
struct AppliedEffect {
job_instance_id: Option<JobInstanceId>,
job_execution_id: Option<JobExecutionId>,
prior_status: Option<BatchStatus>,
result_status: Option<BatchStatus>,
execution: Option<JobExecution>,
changed: bool,
}
impl AppliedEffect {
fn created(instance_id: JobInstanceId, execution: JobExecution) -> Self {
Self {
job_instance_id: Some(instance_id),
job_execution_id: Some(execution.id()),
prior_status: None,
result_status: Some(execution.metadata().status()),
execution: Some(execution),
changed: true,
}
}
fn updated(execution: &JobExecution, prior_status: BatchStatus) -> Self {
Self {
job_instance_id: Some(execution.job_instance_id()),
job_execution_id: Some(execution.id()),
prior_status: Some(prior_status),
result_status: Some(execution.metadata().status()),
execution: Some(execution.clone()),
changed: true,
}
}
fn unchanged(execution: &JobExecution) -> Self {
let status = execution.metadata().status();
Self {
job_instance_id: Some(execution.job_instance_id()),
job_execution_id: Some(execution.id()),
prior_status: Some(status),
result_status: Some(status),
execution: Some(execution.clone()),
changed: false,
}
}
const fn with_prior(mut self, prior_status: BatchStatus) -> Self {
self.prior_status = Some(prior_status);
self
}
}
enum EffectFailure {
Rejected(OperatorRejection),
Failed(OperatorError),
}
impl EffectFailure {
fn classify(error: RepositoryError) -> Self {
OperatorRejection::from_repository(&error)
.map_or_else(|| Self::Failed(OperatorError::from(error)), Self::Rejected)
}
}