use std::collections::BTreeSet;
use std::error::Error;
use std::fmt;
use std::time::{Duration, SystemTime};
use oxide_batch_core::{
BatchStatus, ExecutionVersion, JobExecutionId, JobInstanceId, JobName, RetentionActionId,
};
use crate::{ActorRef, CanonicalWriter, OperationId, ReasonCode, RepositoryError, hex_digest};
pub const MAX_PURGE_BATCH: u32 = 1000;
pub const MIN_PURGE_AGE: Duration = Duration::from_hours(1);
pub const DEFAULT_PURGE_AGE: Duration = Duration::from_hours(30 * 24);
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct PurgeBatchBound(u32);
impl PurgeBatchBound {
pub const fn new(value: u32) -> Result<Self, RetentionError> {
if value == 0 || value > MAX_PURGE_BATCH {
return Err(RetentionError::BatchBoundOutOfRange { requested: value });
}
Ok(Self(value))
}
#[must_use]
pub const fn get(self) -> u32 {
self.0
}
}
impl Default for PurgeBatchBound {
fn default() -> Self {
Self(MAX_PURGE_BATCH)
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct TerminalStatusSet(BTreeSet<BatchStatus>);
impl TerminalStatusSet {
pub fn new(statuses: impl IntoIterator<Item = BatchStatus>) -> Result<Self, RetentionError> {
let mut set = BTreeSet::new();
for status in statuses {
if !status.is_finished() {
return Err(RetentionError::NonTerminalStatus { status });
}
set.insert(status);
}
if set.is_empty() {
return Err(RetentionError::EmptyStatusSet);
}
Ok(Self(set))
}
#[must_use]
pub fn all() -> Self {
Self(
[
BatchStatus::Completed,
BatchStatus::Failed,
BatchStatus::Stopped,
BatchStatus::Abandoned,
]
.into_iter()
.collect(),
)
}
#[must_use]
pub fn contains(&self, status: BatchStatus) -> bool {
self.0.contains(&status)
}
#[must_use]
pub fn iter(&self) -> impl ExactSizeIterator<Item = BatchStatus> + '_ {
self.0.iter().copied()
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PurgePlanRequest {
job_name: JobName,
statuses: TerminalStatusSet,
minimum_age: Duration,
batch: PurgeBatchBound,
}
impl PurgePlanRequest {
pub fn new(
job_name: JobName,
statuses: TerminalStatusSet,
minimum_age: Duration,
batch: PurgeBatchBound,
) -> Result<Self, RetentionError> {
if minimum_age.as_secs() < MIN_PURGE_AGE.as_secs() {
return Err(RetentionError::AgeBoundTooSmall {
minimum: MIN_PURGE_AGE,
});
}
Ok(Self {
job_name,
statuses,
minimum_age,
batch,
})
}
#[must_use]
pub const fn job_name(&self) -> &JobName {
&self.job_name
}
#[must_use]
pub const fn statuses(&self) -> &TerminalStatusSet {
&self.statuses
}
#[must_use]
pub const fn minimum_age(&self) -> Duration {
self.minimum_age
}
#[must_use]
pub const fn batch(&self) -> PurgeBatchBound {
self.batch
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct PurgeCandidate {
job_instance_id: JobInstanceId,
job_execution_id: JobExecutionId,
version: ExecutionVersion,
}
impl PurgeCandidate {
#[must_use]
pub const fn new(
job_instance_id: JobInstanceId,
job_execution_id: JobExecutionId,
version: ExecutionVersion,
) -> Self {
Self {
job_instance_id,
job_execution_id,
version,
}
}
#[must_use]
pub const fn job_instance_id(&self) -> JobInstanceId {
self.job_instance_id
}
#[must_use]
pub const fn job_execution_id(&self) -> JobExecutionId {
self.job_execution_id
}
#[must_use]
pub const fn version(&self) -> ExecutionVersion {
self.version
}
}
#[derive(Clone, Copy, Debug, Default, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct PurgeCounts {
flow_decisions: u64,
recovery_decisions: u64,
operator_requests: u64,
step_partitions: u64,
step_executions: u64,
job_executions: u64,
job_instances: u64,
}
impl PurgeCounts {
#[must_use]
pub const fn new(
flow_decisions: u64,
recovery_decisions: u64,
operator_requests: u64,
step_partitions: u64,
step_executions: u64,
job_executions: u64,
job_instances: u64,
) -> Self {
Self {
flow_decisions,
recovery_decisions,
operator_requests,
step_partitions,
step_executions,
job_executions,
job_instances,
}
}
#[must_use]
pub const fn flow_decisions(self) -> u64 {
self.flow_decisions
}
#[must_use]
pub const fn recovery_decisions(self) -> u64 {
self.recovery_decisions
}
#[must_use]
pub const fn operator_requests(self) -> u64 {
self.operator_requests
}
#[must_use]
pub const fn step_partitions(self) -> u64 {
self.step_partitions
}
#[must_use]
pub const fn step_executions(self) -> u64 {
self.step_executions
}
#[must_use]
pub const fn job_executions(self) -> u64 {
self.job_executions
}
#[must_use]
pub const fn job_instances(self) -> u64 {
self.job_instances
}
}
#[derive(Clone, Debug, Default, Eq, PartialEq)]
pub struct PurgeSurvey {
candidates: Vec<PurgeCandidate>,
counts: PurgeCounts,
}
impl PurgeSurvey {
#[must_use]
pub const fn new(candidates: Vec<PurgeCandidate>, counts: PurgeCounts) -> Self {
Self { candidates, counts }
}
#[must_use]
pub fn candidates(&self) -> &[PurgeCandidate] {
&self.candidates
}
#[must_use]
pub const fn counts(&self) -> PurgeCounts {
self.counts
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PurgePlan {
request: PurgePlanRequest,
candidates: Vec<PurgeCandidate>,
counts: PurgeCounts,
digest: [u8; 32],
}
impl PurgePlan {
#[doc(hidden)]
#[must_use]
pub fn new(request: PurgePlanRequest, survey: PurgeSurvey) -> Self {
let digest = plan_digest(&request, survey.candidates());
Self {
request,
candidates: survey.candidates,
counts: survey.counts,
digest,
}
}
#[must_use]
pub const fn request(&self) -> &PurgePlanRequest {
&self.request
}
#[must_use]
pub fn candidates(&self) -> &[PurgeCandidate] {
&self.candidates
}
#[must_use]
pub const fn counts(&self) -> PurgeCounts {
self.counts
}
#[must_use]
pub const fn digest(&self) -> &[u8; 32] {
&self.digest
}
#[must_use]
pub fn digest_hex(&self) -> String {
hex_digest(&self.digest)
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.candidates.is_empty()
}
}
fn plan_digest(request: &PurgePlanRequest, candidates: &[PurgeCandidate]) -> [u8; 32] {
let mut writer = CanonicalWriter::new("oxide-batch.retention-plan.v1");
writer.push_str(request.job_name().as_str());
for status in request.statuses().iter() {
writer.push_str(status.as_str());
}
writer.push_u64(request.minimum_age().as_secs());
writer.push_u64(u64::from(request.batch().get()));
for candidate in candidates {
writer.push_u64(candidate.job_instance_id().get());
writer.push_u64(candidate.job_execution_id().get());
writer.push_u64(candidate.version().get());
}
writer.digest()
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
#[non_exhaustive]
pub enum RetentionAction {
Hold,
ReleaseHold,
ApplyPurge,
}
impl RetentionAction {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Hold => "HOLD",
Self::ReleaseHold => "RELEASE_HOLD",
Self::ApplyPurge => "APPLY_PURGE",
}
}
}
impl fmt::Display for RetentionAction {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(self.as_str())
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct RetentionHold {
job_instance_id: JobInstanceId,
actor: ActorRef,
reason: ReasonCode,
placed_at: SystemTime,
}
impl RetentionHold {
#[must_use]
pub const fn new(
job_instance_id: JobInstanceId,
actor: ActorRef,
reason: ReasonCode,
placed_at: SystemTime,
) -> Self {
Self {
job_instance_id,
actor,
reason,
placed_at,
}
}
#[must_use]
pub const fn job_instance_id(&self) -> JobInstanceId {
self.job_instance_id
}
#[must_use]
pub const fn actor(&self) -> &ActorRef {
&self.actor
}
#[must_use]
pub const fn reason(&self) -> &ReasonCode {
&self.reason
}
#[must_use]
pub const fn placed_at(&self) -> SystemTime {
self.placed_at
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
#[non_exhaustive]
pub enum RetentionOutcome {
Applied,
Replayed,
Rejected,
}
impl RetentionOutcome {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Applied => "APPLIED",
Self::Replayed => "REPLAYED",
Self::Rejected => "REJECTED",
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct RetentionRecord {
id: RetentionActionId,
action: RetentionAction,
operation_id: OperationId,
actor: ActorRef,
reason: ReasonCode,
job_instance_id: Option<JobInstanceId>,
plan_digest: Option<[u8; 32]>,
counts: PurgeCounts,
batch_bound: Option<PurgeBatchBound>,
outcome: RetentionOutcome,
applied_at: SystemTime,
}
impl RetentionRecord {
#[must_use]
pub fn from_parts(id: RetentionActionId, draft: RetentionRecordDraft) -> Self {
Self {
id,
action: draft.action,
operation_id: draft.operation_id,
actor: draft.actor,
reason: draft.reason,
job_instance_id: draft.job_instance_id,
plan_digest: draft.plan_digest,
counts: draft.counts,
batch_bound: draft.batch_bound,
outcome: draft.outcome,
applied_at: draft.applied_at,
}
}
#[must_use]
pub const fn id(&self) -> RetentionActionId {
self.id
}
#[must_use]
pub const fn action(&self) -> RetentionAction {
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) -> &ReasonCode {
&self.reason
}
#[must_use]
pub const fn job_instance_id(&self) -> Option<JobInstanceId> {
self.job_instance_id
}
#[must_use]
pub const fn plan_digest(&self) -> Option<&[u8; 32]> {
self.plan_digest.as_ref()
}
#[must_use]
pub const fn counts(&self) -> PurgeCounts {
self.counts
}
#[must_use]
pub const fn batch_bound(&self) -> Option<PurgeBatchBound> {
self.batch_bound
}
#[must_use]
pub const fn outcome(&self) -> RetentionOutcome {
self.outcome
}
#[must_use]
pub const fn applied_at(&self) -> SystemTime {
self.applied_at
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct RetentionRecordDraft {
action: RetentionAction,
operation_id: OperationId,
actor: ActorRef,
reason: ReasonCode,
job_instance_id: Option<JobInstanceId>,
plan_digest: Option<[u8; 32]>,
counts: PurgeCounts,
batch_bound: Option<PurgeBatchBound>,
outcome: RetentionOutcome,
applied_at: SystemTime,
}
impl RetentionRecordDraft {
#[must_use]
pub fn instance_action(
action: RetentionAction,
operation_id: OperationId,
actor: ActorRef,
reason: ReasonCode,
job_instance_id: JobInstanceId,
applied_at: SystemTime,
) -> Self {
Self {
action,
operation_id,
actor,
reason,
job_instance_id: Some(job_instance_id),
plan_digest: None,
counts: PurgeCounts::default(),
batch_bound: None,
outcome: RetentionOutcome::Applied,
applied_at,
}
}
#[must_use]
pub const fn purge(
operation_id: OperationId,
actor: ActorRef,
reason: ReasonCode,
plan_digest: [u8; 32],
counts: PurgeCounts,
batch_bound: PurgeBatchBound,
applied_at: SystemTime,
) -> Self {
Self {
action: RetentionAction::ApplyPurge,
operation_id,
actor,
reason,
job_instance_id: None,
plan_digest: Some(plan_digest),
counts,
batch_bound: Some(batch_bound),
outcome: RetentionOutcome::Applied,
applied_at,
}
}
#[must_use]
#[allow(clippy::too_many_arguments)]
pub const fn from_durable(
action: RetentionAction,
operation_id: OperationId,
actor: ActorRef,
reason: ReasonCode,
job_instance_id: Option<JobInstanceId>,
plan_digest: Option<[u8; 32]>,
counts: PurgeCounts,
batch_bound: Option<PurgeBatchBound>,
outcome: RetentionOutcome,
applied_at: SystemTime,
) -> Self {
Self {
action,
operation_id,
actor,
reason,
job_instance_id,
plan_digest,
counts,
batch_bound,
outcome,
applied_at,
}
}
#[must_use]
pub const fn action(&self) -> RetentionAction {
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) -> &ReasonCode {
&self.reason
}
#[must_use]
pub const fn job_instance_id(&self) -> Option<JobInstanceId> {
self.job_instance_id
}
#[must_use]
pub const fn plan_digest(&self) -> Option<&[u8; 32]> {
self.plan_digest.as_ref()
}
#[must_use]
pub const fn counts(&self) -> PurgeCounts {
self.counts
}
#[must_use]
pub const fn batch_bound(&self) -> Option<PurgeBatchBound> {
self.batch_bound
}
#[must_use]
pub const fn outcome(&self) -> RetentionOutcome {
self.outcome
}
#[must_use]
pub const fn applied_at(&self) -> SystemTime {
self.applied_at
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum RetentionError {
BatchBoundOutOfRange {
requested: u32,
},
AgeBoundTooSmall {
minimum: Duration,
},
NonTerminalStatus {
status: BatchStatus,
},
EmptyStatusSet,
RetentionPlanStale,
InstanceHeld {
job_instance_id: JobInstanceId,
},
OperationIdConflict {
action: RetentionAction,
operation_id: OperationId,
},
OperationOutcomeUnknown,
Repository(RepositoryError),
}
impl fmt::Display for RetentionError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::BatchBoundOutOfRange { requested } => write!(
formatter,
"purge batch bound {requested} is outside 1..={MAX_PURGE_BATCH}"
),
Self::AgeBoundTooSmall { minimum } => write!(
formatter,
"the minimum age must be at least {} seconds",
minimum.as_secs()
),
Self::NonTerminalStatus { status } => {
write!(formatter, "{status} is not a finished status")
}
Self::EmptyStatusSet => formatter.write_str("a purge must target at least one status"),
Self::RetentionPlanStale => {
formatter.write_str("the purge plan is stale and nothing was deleted")
}
Self::InstanceHeld { job_instance_id } => {
write!(formatter, "job instance {job_instance_id} is held")
}
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 retention commit outcome is unknown")
}
Self::Repository(error) => error.fmt(formatter),
}
}
}
impl Error for RetentionError {
fn source(&self) -> Option<&(dyn Error + 'static)> {
match self {
Self::Repository(error) => Some(error),
_ => None,
}
}
}
impl From<RepositoryError> for RetentionError {
fn from(value: RepositoryError) -> Self {
match value {
RepositoryError::CommitOutcomeUnknown => Self::OperationOutcomeUnknown,
RepositoryError::RetentionPlanStale => Self::RetentionPlanStale,
other => Self::Repository(other),
}
}
}