use std::fmt;
use std::time::SystemTime;
use oxide_batch_core::{
BatchStatus, DefinitionIdentity, ExecutionVersion, FailureSummary, JobExecutionId,
JobInstanceId, JobInstanceKey, LifecycleError, OperatorRequestId,
};
use crate::{
ActorRef, AuthorizationClass, OperationId, OperatorAction, ReasonCode, RecoveryDisposition,
RecoveryProposal, RecoveryRequest, RecoveryRequestError, RepositoryError, RequestArguments,
RequestDigest, hex_digest, request_digest,
};
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct OperatorRequest {
action: OperatorAction,
operation_id: OperationId,
actor: ActorRef,
reason: Option<ReasonCode>,
target: OperatorTarget,
expected_version: Option<ExecutionVersion>,
arguments: RequestArguments,
digest: RequestDigest,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum RecoveryDirective {
MarkFailed(FailureSummary),
Abandon,
}
impl RecoveryDirective {
#[must_use]
pub const fn disposition(self) -> RecoveryDisposition {
match self {
Self::MarkFailed(_) => RecoveryDisposition::MarkFailed,
Self::Abandon => RecoveryDisposition::Abandon,
}
}
#[must_use]
pub const fn failure(self) -> Option<FailureSummary> {
match self {
Self::MarkFailed(failure) => Some(failure),
Self::Abandon => None,
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
enum OperatorTarget {
InstanceKey(Box<JobInstanceKey>),
Instance(JobInstanceId),
Execution(JobExecutionId),
}
impl OperatorTarget {
fn identity(&self) -> String {
match self {
Self::InstanceKey(key) => {
format!("instance-key:{}", hex_digest(&key.digest()))
}
Self::Instance(id) => format!("instance:{id}"),
Self::Execution(id) => format!("execution:{id}"),
}
}
}
impl OperatorRequest {
#[must_use]
pub fn launch(
operation_id: OperationId,
actor: ActorRef,
key: JobInstanceKey,
definition: DefinitionIdentity,
) -> Self {
Self::build(
OperatorAction::Launch,
operation_id,
actor,
None,
OperatorTarget::InstanceKey(Box::new(key)),
None,
RequestArguments::Definition(Box::new(definition)),
)
}
#[must_use]
pub fn restart(
operation_id: OperationId,
actor: ActorRef,
job_instance_id: JobInstanceId,
definition: DefinitionIdentity,
) -> Self {
Self::build(
OperatorAction::Restart,
operation_id,
actor,
None,
OperatorTarget::Instance(job_instance_id),
None,
RequestArguments::Definition(Box::new(definition)),
)
}
#[must_use]
pub fn stop(
operation_id: OperationId,
actor: ActorRef,
job_execution_id: JobExecutionId,
expected_version: ExecutionVersion,
) -> Self {
Self::build(
OperatorAction::Stop,
operation_id,
actor,
None,
OperatorTarget::Execution(job_execution_id),
Some(expected_version),
RequestArguments::None,
)
}
#[must_use]
pub fn abandon(
operation_id: OperationId,
actor: ActorRef,
reason: ReasonCode,
job_execution_id: JobExecutionId,
expected_version: ExecutionVersion,
) -> Self {
Self::build(
OperatorAction::Abandon,
operation_id,
actor,
Some(reason),
OperatorTarget::Execution(job_execution_id),
Some(expected_version),
RequestArguments::None,
)
}
#[must_use]
pub fn recover(
operation_id: OperationId,
actor: ActorRef,
reason: ReasonCode,
directive: RecoveryDirective,
proposal: &RecoveryProposal,
) -> Self {
let job_execution_id = proposal.evidence().execution_id();
let expected_version = proposal.observed_version();
let evidence_digest = *proposal.digest();
Self::build(
OperatorAction::Recover,
operation_id,
actor,
Some(reason),
OperatorTarget::Execution(job_execution_id),
Some(expected_version),
RequestArguments::Recovery {
directive,
evidence_digest,
unknown_commit: proposal.evidence().unknown_commit(),
},
)
}
fn build(
action: OperatorAction,
operation_id: OperationId,
actor: ActorRef,
reason: Option<ReasonCode>,
target: OperatorTarget,
expected_version: Option<ExecutionVersion>,
arguments: RequestArguments,
) -> Self {
let digest = request_digest(
action,
&target.identity(),
expected_version,
reason.as_ref(),
&arguments,
);
Self {
action,
operation_id,
actor,
reason,
target,
expected_version,
arguments,
digest,
}
}
#[must_use]
pub const fn action(&self) -> OperatorAction {
self.action
}
#[must_use]
pub const fn authorization_class(&self) -> AuthorizationClass {
self.action.authorization_class()
}
#[must_use]
pub const fn operation_id(&self) -> &OperationId {
&self.operation_id
}
#[must_use]
pub const fn actor(&self) -> &ActorRef {
&self.actor
}
#[must_use]
pub const fn reason(&self) -> Option<&ReasonCode> {
self.reason.as_ref()
}
#[must_use]
pub const fn expected_version(&self) -> Option<ExecutionVersion> {
self.expected_version
}
#[must_use]
pub const fn digest(&self) -> &RequestDigest {
&self.digest
}
#[must_use]
pub const fn job_execution_id(&self) -> Option<JobExecutionId> {
match self.target {
OperatorTarget::Execution(id) => Some(id),
_ => None,
}
}
#[must_use]
pub const fn job_instance_id(&self) -> Option<JobInstanceId> {
match self.target {
OperatorTarget::Instance(id) => Some(id),
_ => None,
}
}
#[must_use]
pub fn job_instance_key(&self) -> Option<&JobInstanceKey> {
match &self.target {
OperatorTarget::InstanceKey(key) => Some(key),
_ => None,
}
}
#[doc(hidden)]
#[must_use]
pub fn definition(&self) -> Option<&DefinitionIdentity> {
match &self.arguments {
RequestArguments::Definition(definition) => Some(definition),
_ => None,
}
}
#[doc(hidden)]
#[must_use]
pub fn recovery_request(&self) -> Option<Result<RecoveryRequest, RecoveryRequestError>> {
let RequestArguments::Recovery {
directive,
evidence_digest,
..
} = &self.arguments
else {
return None;
};
let expected_version = self.expected_version?;
let reason = self.reason.as_ref()?;
Some(match directive {
RecoveryDirective::Abandon => RecoveryRequest::abandon(
expected_version,
reason.as_str(),
self.actor.as_str(),
*evidence_digest,
),
RecoveryDirective::MarkFailed(failure) => RecoveryRequest::mark_failed(
expected_version,
reason.as_str(),
self.actor.as_str(),
*evidence_digest,
failure.category(),
failure.failure_id(),
),
})
}
#[doc(hidden)]
#[must_use]
pub fn recovery_guard(&self) -> Option<(RecoveryDirective, bool)> {
match &self.arguments {
RequestArguments::Recovery {
directive,
unknown_commit,
..
} => Some((*directive, *unknown_commit)),
_ => None,
}
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
#[non_exhaustive]
pub enum OperatorOutcomeClass {
Applied,
Replayed,
Rejected,
}
impl OperatorOutcomeClass {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Applied => "APPLIED",
Self::Replayed => "REPLAYED",
Self::Rejected => "REJECTED",
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum OperatorRejection {
OptimisticConflict {
current: ExecutionVersion,
},
InvalidState {
status: BatchStatus,
},
InstanceCompleted,
InstanceAbandoned,
ExecutionAlreadyActive {
execution_id: JobExecutionId,
status: BatchStatus,
},
IncompatibleDefinition,
RestartWithoutPriorAttempt,
StartLimitExceeded,
UnresolvedRecoveryRequired,
ExecutionNotFound,
InstanceNotFound,
UnsupportedAction,
}
impl OperatorRejection {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::OptimisticConflict { .. } => "OPTIMISTIC_CONFLICT",
Self::InvalidState { .. } => "INVALID_STATE",
Self::InstanceCompleted => "INSTANCE_COMPLETED",
Self::InstanceAbandoned => "INSTANCE_ABANDONED",
Self::ExecutionAlreadyActive { .. } => "EXECUTION_ALREADY_ACTIVE",
Self::IncompatibleDefinition => "INCOMPATIBLE_DEFINITION",
Self::RestartWithoutPriorAttempt => "RESTART_WITHOUT_PRIOR_ATTEMPT",
Self::StartLimitExceeded => "START_LIMIT_EXCEEDED",
Self::UnresolvedRecoveryRequired => "UNRESOLVED_RECOVERY_REQUIRED",
Self::ExecutionNotFound => "EXECUTION_NOT_FOUND",
Self::InstanceNotFound => "INSTANCE_NOT_FOUND",
Self::UnsupportedAction => "UNSUPPORTED_ACTION",
}
}
#[doc(hidden)]
#[must_use]
pub fn from_repository(error: &RepositoryError) -> Option<Self> {
match error {
RepositoryError::CompletedInstance { .. } => Some(Self::InstanceCompleted),
RepositoryError::AbandonedInstance { .. } => Some(Self::InstanceAbandoned),
RepositoryError::ExecutionAlreadyActive {
execution_id,
status,
..
} => Some(Self::ExecutionAlreadyActive {
execution_id: *execution_id,
status: *status,
}),
RepositoryError::IncompatibleDefinition { .. }
| RepositoryError::RestartStateNotFound { .. }
| RepositoryError::InvalidDefinitionUpgrade { .. } => {
Some(Self::IncompatibleDefinition)
}
RepositoryError::StartLimitExceeded { .. } => Some(Self::StartLimitExceeded),
RepositoryError::JobExecutionNotFound { .. } => Some(Self::ExecutionNotFound),
RepositoryError::JobInstanceNotFound { .. } => Some(Self::InstanceNotFound),
RepositoryError::Lifecycle(LifecycleError::StaleVersion { actual, .. }) => {
Some(Self::OptimisticConflict { current: *actual })
}
RepositoryError::Lifecycle(
LifecycleError::IllegalTransition { from, .. }
| LifecycleError::RestartRequiresNewAttempt { from },
) => Some(Self::InvalidState { status: *from }),
RepositoryError::RecoveryNotAllowed { status, .. }
| RepositoryError::Lifecycle(LifecycleError::NotRestartable { status }) => {
Some(Self::InvalidState { status: *status })
}
_ => None,
}
}
}
impl fmt::Display for OperatorRejection {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(self.as_str())
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct OperatorRecord {
id: OperatorRequestId,
action: OperatorAction,
operation_id: OperationId,
actor: ActorRef,
reason: Option<ReasonCode>,
digest: RequestDigest,
job_instance_id: Option<JobInstanceId>,
job_execution_id: Option<JobExecutionId>,
observed_version: Option<ExecutionVersion>,
prior_status: Option<BatchStatus>,
result_status: Option<BatchStatus>,
outcome: OperatorOutcomeClass,
rejection: Option<OperatorRejection>,
requested_at: SystemTime,
}
impl OperatorRecord {
#[must_use]
pub fn from_parts(id: OperatorRequestId, draft: OperatorRecordDraft) -> Self {
Self {
id,
action: draft.action,
operation_id: draft.operation_id,
actor: draft.actor,
reason: draft.reason,
digest: draft.digest,
job_instance_id: draft.job_instance_id,
job_execution_id: draft.job_execution_id,
observed_version: draft.observed_version,
prior_status: draft.prior_status,
result_status: draft.result_status,
outcome: draft.outcome,
rejection: draft.rejection,
requested_at: draft.requested_at,
}
}
#[must_use]
pub const fn id(&self) -> OperatorRequestId {
self.id
}
#[must_use]
pub const fn action(&self) -> OperatorAction {
self.action
}
#[must_use]
pub const fn operation_id(&self) -> &OperationId {
&self.operation_id
}
#[must_use]
pub const fn actor(&self) -> &ActorRef {
&self.actor
}
#[must_use]
pub const fn reason(&self) -> Option<&ReasonCode> {
self.reason.as_ref()
}
#[must_use]
pub const fn digest(&self) -> &RequestDigest {
&self.digest
}
#[must_use]
pub const fn job_instance_id(&self) -> Option<JobInstanceId> {
self.job_instance_id
}
#[must_use]
pub const fn job_execution_id(&self) -> Option<JobExecutionId> {
self.job_execution_id
}
#[must_use]
pub const fn observed_version(&self) -> Option<ExecutionVersion> {
self.observed_version
}
#[must_use]
pub const fn prior_status(&self) -> Option<BatchStatus> {
self.prior_status
}
#[must_use]
pub const fn result_status(&self) -> Option<BatchStatus> {
self.result_status
}
#[must_use]
pub const fn outcome(&self) -> OperatorOutcomeClass {
self.outcome
}
#[must_use]
pub const fn rejection(&self) -> Option<OperatorRejection> {
self.rejection
}
#[must_use]
pub const fn requested_at(&self) -> SystemTime {
self.requested_at
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct OperatorRecordDraft {
action: OperatorAction,
operation_id: OperationId,
actor: ActorRef,
reason: Option<ReasonCode>,
digest: RequestDigest,
job_instance_id: Option<JobInstanceId>,
job_execution_id: Option<JobExecutionId>,
observed_version: Option<ExecutionVersion>,
prior_status: Option<BatchStatus>,
result_status: Option<BatchStatus>,
outcome: OperatorOutcomeClass,
rejection: Option<OperatorRejection>,
requested_at: SystemTime,
}
impl OperatorRecordDraft {
#[must_use]
pub fn applied(
request: &OperatorRequest,
job_instance_id: Option<JobInstanceId>,
job_execution_id: Option<JobExecutionId>,
prior_status: Option<BatchStatus>,
result_status: Option<BatchStatus>,
requested_at: SystemTime,
) -> Self {
Self {
action: request.action(),
operation_id: request.operation_id().clone(),
actor: request.actor().clone(),
reason: request.reason().cloned(),
digest: *request.digest(),
job_instance_id,
job_execution_id,
observed_version: request.expected_version(),
prior_status,
result_status,
outcome: OperatorOutcomeClass::Applied,
rejection: None,
requested_at,
}
}
#[must_use]
pub fn rejected(
request: &OperatorRequest,
rejection: OperatorRejection,
requested_at: SystemTime,
) -> Self {
Self {
action: request.action(),
operation_id: request.operation_id().clone(),
actor: request.actor().clone(),
reason: request.reason().cloned(),
digest: *request.digest(),
job_instance_id: request.job_instance_id(),
job_execution_id: request.job_execution_id(),
observed_version: request.expected_version(),
prior_status: None,
result_status: None,
outcome: OperatorOutcomeClass::Rejected,
rejection: Some(rejection),
requested_at,
}
}
#[must_use]
#[allow(clippy::too_many_arguments)]
pub const fn from_durable(
action: OperatorAction,
operation_id: OperationId,
actor: ActorRef,
reason: Option<ReasonCode>,
digest: RequestDigest,
job_instance_id: Option<JobInstanceId>,
job_execution_id: Option<JobExecutionId>,
observed_version: Option<ExecutionVersion>,
prior_status: Option<BatchStatus>,
result_status: Option<BatchStatus>,
outcome: OperatorOutcomeClass,
rejection: Option<OperatorRejection>,
requested_at: SystemTime,
) -> Self {
Self {
action,
operation_id,
actor,
reason,
digest,
job_instance_id,
job_execution_id,
observed_version,
prior_status,
result_status,
outcome,
rejection,
requested_at,
}
}
#[must_use]
pub const fn action(&self) -> OperatorAction {
self.action
}
#[must_use]
pub const fn operation_id(&self) -> &OperationId {
&self.operation_id
}
#[must_use]
pub const fn actor(&self) -> &ActorRef {
&self.actor
}
#[must_use]
pub const fn reason(&self) -> Option<&ReasonCode> {
self.reason.as_ref()
}
#[must_use]
pub const fn digest(&self) -> &RequestDigest {
&self.digest
}
#[must_use]
pub const fn job_instance_id(&self) -> Option<JobInstanceId> {
self.job_instance_id
}
#[must_use]
pub const fn job_execution_id(&self) -> Option<JobExecutionId> {
self.job_execution_id
}
#[must_use]
pub const fn observed_version(&self) -> Option<ExecutionVersion> {
self.observed_version
}
#[must_use]
pub const fn prior_status(&self) -> Option<BatchStatus> {
self.prior_status
}
#[must_use]
pub const fn result_status(&self) -> Option<BatchStatus> {
self.result_status
}
#[must_use]
pub const fn outcome(&self) -> OperatorOutcomeClass {
self.outcome
}
#[must_use]
pub const fn rejection(&self) -> Option<OperatorRejection> {
self.rejection
}
#[must_use]
pub const fn requested_at(&self) -> SystemTime {
self.requested_at
}
}