use std::{num::NonZeroU32, num::NonZeroU64, time::Duration};
use runifold_core::CheckpointId;
use crate::{
WorkerId, WorkflowStoreError, WorkflowStoreErrorKind, WorkflowStoreFuture,
WorkflowTaskCleanupLease, WorkflowTaskRetentionStore, WorkflowTaskTombstoneCursor,
WorkflowTenantId,
};
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct WorkflowTaskTombstoneRetention(NonZeroU64);
impl WorkflowTaskTombstoneRetention {
pub fn new(duration: Duration) -> Result<Self, WorkflowStoreError> {
let millis = u64::try_from(duration.as_millis())
.ok()
.and_then(NonZeroU64::new)
.ok_or_else(|| {
invalid_input("Task tombstone retention must fit in positive whole milliseconds")
})?;
Ok(Self(millis))
}
pub const fn as_millis(self) -> u64 {
self.0.get()
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct WorkflowTaskTombstonePurgeLimit(NonZeroU32);
impl WorkflowTaskTombstonePurgeLimit {
pub fn new(value: u32) -> Result<Self, WorkflowStoreError> {
let value =
NonZeroU32::new(value).ok_or_else(|| invalid_input("purge limit must be positive"))?;
if value.get() > 1_000 {
return Err(invalid_input("purge limit cannot exceed 1,000"));
}
Ok(Self(value))
}
pub const fn get(self) -> u32 {
self.0.get()
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct WorkflowTaskTombstoneApprovalWindow(NonZeroU64);
impl WorkflowTaskTombstoneApprovalWindow {
pub fn new(duration: Duration) -> Result<Self, WorkflowStoreError> {
let millis = u64::try_from(duration.as_millis())
.ok()
.and_then(NonZeroU64::new)
.ok_or_else(|| {
invalid_input("purge approval window must fit in positive whole milliseconds")
})?;
Ok(Self(millis))
}
pub const fn as_millis(self) -> u64 {
self.0.get()
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct WorkflowTaskTombstoneApprovalInboxLimit(NonZeroU32);
impl WorkflowTaskTombstoneApprovalInboxLimit {
pub fn new(value: u32) -> Result<Self, WorkflowStoreError> {
let value = NonZeroU32::new(value)
.ok_or_else(|| invalid_input("approval inbox limit must be positive"))?;
if value.get() > 1_000 {
return Err(invalid_input("approval inbox limit cannot exceed 1,000"));
}
Ok(Self(value))
}
pub const fn get(self) -> u32 {
self.0.get()
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct WorkflowTaskTombstoneRejectionReason(String);
impl WorkflowTaskTombstoneRejectionReason {
pub fn parse(value: impl Into<String>) -> Result<Self, WorkflowStoreError> {
let value = value.into();
if value.trim().is_empty() || value.len() > 1_024 || value.chars().any(char::is_control) {
return Err(invalid_input(
"purge rejection reason must contain 1..=1,024 printable bytes",
));
}
Ok(Self(value))
}
pub fn as_str(&self) -> &str {
&self.0
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct WorkflowTaskLegalHoldReason(String);
impl WorkflowTaskLegalHoldReason {
pub fn parse(value: impl Into<String>) -> Result<Self, WorkflowStoreError> {
let value = value.into();
if value.trim().is_empty() || value.len() > 1_024 {
return Err(invalid_input(
"Task legal-hold reason must contain 1..=1,024 bytes",
));
}
Ok(Self(value))
}
pub fn as_str(&self) -> &str {
&self.0
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct WorkflowTaskTombstoneExportReceipt(String);
impl WorkflowTaskTombstoneExportReceipt {
pub fn parse(value: impl Into<String>) -> Result<Self, WorkflowStoreError> {
let value = value.into();
if value.trim().is_empty() || value.len() > 512 || value.chars().any(char::is_control) {
return Err(invalid_input(
"Task tombstone export receipt must contain 1..=512 printable bytes",
));
}
Ok(Self(value))
}
pub fn as_str(&self) -> &str {
&self.0
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct WorkflowTaskTombstonePurgeId(CheckpointId);
impl WorkflowTaskTombstonePurgeId {
pub fn new() -> Self {
Self(CheckpointId::new())
}
pub const fn from_checkpoint_id(value: CheckpointId) -> Self {
Self(value)
}
pub const fn as_checkpoint_id(self) -> CheckpointId {
self.0
}
}
impl Default for WorkflowTaskTombstonePurgeId {
fn default() -> Self {
Self::new()
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct WorkflowTaskLegalHold {
pub checkpoint_id: CheckpointId,
pub tenant_id: WorkflowTenantId,
pub placed_by: WorkerId,
pub reason: WorkflowTaskLegalHoldReason,
pub placed_at_ms: u64,
pub released_by: Option<WorkerId>,
pub released_at_ms: Option<u64>,
}
impl WorkflowTaskLegalHold {
pub const fn is_active(&self) -> bool {
self.released_at_ms.is_none()
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct WorkflowTaskTombstoneExport {
pub tenant_id: WorkflowTenantId,
pub through: WorkflowTaskTombstoneCursor,
pub receipt: WorkflowTaskTombstoneExportReceipt,
pub confirmed_by: WorkerId,
pub confirmed_at_ms: u64,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct WorkflowTaskTombstonePurgeIntent {
pub purge_id: WorkflowTaskTombstonePurgeId,
pub tenant_id: WorkflowTenantId,
pub prepared_by: WorkerId,
pub tombstone_count: u32,
pub first_cursor: Option<WorkflowTaskTombstoneCursor>,
pub last_cursor: Option<WorkflowTaskTombstoneCursor>,
pub export_through: WorkflowTaskTombstoneCursor,
pub fingerprint: String,
pub prepared_at_ms: u64,
pub expires_at_ms: u64,
pub approved_by: Option<WorkerId>,
pub approved_at_ms: Option<u64>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum WorkflowTaskTombstoneApprovalState {
Pending,
Claimed,
Approved,
Rejected,
Expired,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct WorkflowTaskTombstoneApprovalInboxItem {
pub intent: WorkflowTaskTombstonePurgeIntent,
pub state: WorkflowTaskTombstoneApprovalState,
pub claimed_by: Option<WorkerId>,
pub claim_expires_at_ms: Option<u64>,
pub rejected_by: Option<WorkerId>,
pub rejection_reason: Option<WorkflowTaskTombstoneRejectionReason>,
pub rejected_at_ms: Option<u64>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct WorkflowTaskTombstoneApprovalLease {
pub tenant_id: WorkflowTenantId,
pub purge_id: WorkflowTaskTombstonePurgeId,
pub reviewer: WorkerId,
pub fencing_token: u64,
pub expires_at_ms: u64,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct WorkflowTaskTombstonePurgeEvidence {
pub purge_id: WorkflowTaskTombstonePurgeId,
pub tenant_id: WorkflowTenantId,
pub prepared_by: WorkerId,
pub approved_by: WorkerId,
pub executed_by: WorkerId,
pub tombstone_count: u32,
pub first_cursor: WorkflowTaskTombstoneCursor,
pub last_cursor: WorkflowTaskTombstoneCursor,
pub export_through: WorkflowTaskTombstoneCursor,
pub fingerprint: String,
pub executed_at_ms: u64,
}
pub trait WorkflowTaskTombstoneGovernanceStore: WorkflowTaskRetentionStore {
fn place_task_tombstone_hold(
&self,
tenant_id: WorkflowTenantId,
checkpoint_id: CheckpointId,
actor: WorkerId,
reason: WorkflowTaskLegalHoldReason,
) -> WorkflowStoreFuture<'_, Result<WorkflowTaskLegalHold, WorkflowStoreError>>;
fn release_task_tombstone_hold(
&self,
tenant_id: WorkflowTenantId,
checkpoint_id: CheckpointId,
actor: WorkerId,
) -> WorkflowStoreFuture<'_, Result<WorkflowTaskLegalHold, WorkflowStoreError>>;
fn confirm_task_tombstone_export(
&self,
tenant_id: WorkflowTenantId,
through: WorkflowTaskTombstoneCursor,
receipt: WorkflowTaskTombstoneExportReceipt,
actor: WorkerId,
) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstoneExport, WorkflowStoreError>>;
fn prepare_task_tombstone_purge(
&self,
lease: WorkflowTaskCleanupLease,
retention: WorkflowTaskTombstoneRetention,
limit: WorkflowTaskTombstonePurgeLimit,
approval_window: WorkflowTaskTombstoneApprovalWindow,
) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeIntent, WorkflowStoreError>>;
fn approve_task_tombstone_purge(
&self,
tenant_id: WorkflowTenantId,
purge_id: WorkflowTaskTombstonePurgeId,
approver: WorkerId,
) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeIntent, WorkflowStoreError>>;
fn list_task_tombstone_purge_approvals(
&self,
tenant_id: WorkflowTenantId,
limit: WorkflowTaskTombstoneApprovalInboxLimit,
) -> WorkflowStoreFuture<
'_,
Result<Vec<WorkflowTaskTombstoneApprovalInboxItem>, WorkflowStoreError>,
>;
fn claim_task_tombstone_purge_approval(
&self,
tenant_id: WorkflowTenantId,
reviewer: WorkerId,
lease: crate::LeaseDuration,
) -> WorkflowStoreFuture<
'_,
Result<Option<WorkflowTaskTombstoneApprovalLease>, WorkflowStoreError>,
>;
fn approve_claimed_task_tombstone_purge(
&self,
lease: WorkflowTaskTombstoneApprovalLease,
) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeIntent, WorkflowStoreError>>;
fn reject_claimed_task_tombstone_purge(
&self,
lease: WorkflowTaskTombstoneApprovalLease,
reason: WorkflowTaskTombstoneRejectionReason,
) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstoneApprovalInboxItem, WorkflowStoreError>>;
fn execute_task_tombstone_purge(
&self,
lease: WorkflowTaskCleanupLease,
purge_id: WorkflowTaskTombstonePurgeId,
) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeEvidence, WorkflowStoreError>>;
fn get_task_tombstone_purge_evidence(
&self,
tenant_id: WorkflowTenantId,
purge_id: WorkflowTaskTombstonePurgeId,
) -> WorkflowStoreFuture<
'_,
Result<Option<WorkflowTaskTombstonePurgeEvidence>, WorkflowStoreError>,
>;
}
fn invalid_input(message: &'static str) -> WorkflowStoreError {
WorkflowStoreError::new(WorkflowStoreErrorKind::InvalidInput, message)
}