use std::{
collections::{BTreeMap, BTreeSet},
future::Future,
pin::Pin,
sync::Arc,
};
use runifold_core::CheckpointId;
use thiserror::Error;
use crate::{
LeaseDuration, WorkerId, WorkflowStoreError, WorkflowTaskCleanupLease, WorkflowTaskLegalHold,
WorkflowTaskLegalHoldReason, WorkflowTaskTombstone, WorkflowTaskTombstoneApprovalInboxItem,
WorkflowTaskTombstoneApprovalInboxLimit, WorkflowTaskTombstoneApprovalLease,
WorkflowTaskTombstoneApprovalWindow, WorkflowTaskTombstoneCursor,
WorkflowTaskTombstoneExportReceipt, WorkflowTaskTombstoneGovernanceStore,
WorkflowTaskTombstoneLimit, WorkflowTaskTombstonePurgeEvidence, WorkflowTaskTombstonePurgeId,
WorkflowTaskTombstonePurgeIntent, WorkflowTaskTombstonePurgeLimit,
WorkflowTaskTombstoneRejectionReason, WorkflowTaskTombstoneRetention, WorkflowTenantId,
};
#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
#[non_exhaustive]
pub enum WorkflowTaskGovernancePermission {
PlaceHold,
ReleaseHold,
Export,
PreparePurge,
ApprovePurge,
ReadApprovalInbox,
ClaimPurgeApproval,
RejectPurge,
ExecutePurge,
ReadEvidence,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum WorkflowTaskGovernanceOutcome {
Succeeded,
Denied,
AuthorizationError,
StoreError,
ArchiveError,
}
pub trait WorkflowTaskGovernanceObserver: Send + Sync {
fn observe(
&self,
permission: WorkflowTaskGovernancePermission,
outcome: WorkflowTaskGovernanceOutcome,
);
}
#[derive(Debug, Default)]
struct NoopWorkflowTaskGovernanceObserver;
impl WorkflowTaskGovernanceObserver for NoopWorkflowTaskGovernanceObserver {
fn observe(
&self,
_permission: WorkflowTaskGovernancePermission,
_outcome: WorkflowTaskGovernanceOutcome,
) {
}
}
#[derive(Clone, Debug, Error, Eq, PartialEq)]
#[error("Task tombstone governance authorization failed: {message}")]
pub struct WorkflowTaskGovernanceAuthorizationError {
message: String,
}
impl WorkflowTaskGovernanceAuthorizationError {
pub fn new(message: impl Into<String>) -> Self {
let mut message = message.into();
message.truncate(512);
Self { message }
}
}
pub type WorkflowTaskGovernanceAuthorizationFuture<'a> = Pin<
Box<dyn Future<Output = Result<bool, WorkflowTaskGovernanceAuthorizationError>> + Send + 'a>,
>;
pub trait WorkflowTaskGovernanceAuthorizer: Send + Sync {
fn authorize(
&self,
principal: &WorkerId,
tenant_id: &WorkflowTenantId,
permission: WorkflowTaskGovernancePermission,
) -> WorkflowTaskGovernanceAuthorizationFuture<'_>;
}
#[derive(Clone, Debug, Default)]
pub struct StaticWorkflowTaskGovernanceAuthorizer {
grants: BTreeMap<(WorkerId, WorkflowTenantId), BTreeSet<WorkflowTaskGovernancePermission>>,
}
impl StaticWorkflowTaskGovernanceAuthorizer {
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn with_grant(
mut self,
principal: WorkerId,
tenant_id: WorkflowTenantId,
permission: WorkflowTaskGovernancePermission,
) -> Self {
self.grants
.entry((principal, tenant_id))
.or_default()
.insert(permission);
self
}
}
impl WorkflowTaskGovernanceAuthorizer for StaticWorkflowTaskGovernanceAuthorizer {
fn authorize(
&self,
principal: &WorkerId,
tenant_id: &WorkflowTenantId,
permission: WorkflowTaskGovernancePermission,
) -> WorkflowTaskGovernanceAuthorizationFuture<'_> {
let allowed = self
.grants
.get(&(principal.clone(), tenant_id.clone()))
.is_some_and(|grants| grants.contains(&permission));
Box::pin(async move { Ok(allowed) })
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum WorkflowTaskTombstoneArchiveErrorKind {
Configuration,
Authorization,
Timeout,
Unavailable,
Integrity,
Ambiguous,
Other,
}
#[derive(Clone, Debug, Error, Eq, PartialEq)]
#[error("Task tombstone archive failed ({kind:?}): {message}")]
pub struct WorkflowTaskTombstoneArchiveError {
kind: WorkflowTaskTombstoneArchiveErrorKind,
message: String,
}
impl WorkflowTaskTombstoneArchiveError {
pub fn new(message: impl Into<String>) -> Self {
Self::with_kind(WorkflowTaskTombstoneArchiveErrorKind::Other, message)
}
pub fn with_kind(
kind: WorkflowTaskTombstoneArchiveErrorKind,
message: impl Into<String>,
) -> Self {
let mut message = message.into();
message.truncate(512);
Self { kind, message }
}
pub const fn kind(&self) -> WorkflowTaskTombstoneArchiveErrorKind {
self.kind
}
}
#[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)]
pub struct WorkflowTaskTombstoneArchiveBatchId(String);
impl WorkflowTaskTombstoneArchiveBatchId {
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(WorkflowStoreError::new(
crate::WorkflowStoreErrorKind::InvalidInput,
"Task tombstone archive batch ID must contain 1..=512 printable bytes",
));
}
Ok(Self(value))
}
fn from_batch(
tenant_id: &WorkflowTenantId,
first: WorkflowTaskTombstoneCursor,
last: WorkflowTaskTombstoneCursor,
) -> Self {
Self(format!(
"{}:{}:{}",
tenant_id.as_str(),
first.get(),
last.get()
))
}
pub fn as_str(&self) -> &str {
&self.0
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct WorkflowTaskTombstoneArchiveBatch {
pub batch_id: WorkflowTaskTombstoneArchiveBatchId,
pub tenant_id: WorkflowTenantId,
pub tombstones: Vec<WorkflowTaskTombstone>,
}
pub type WorkflowTaskTombstoneArchiveFuture<'a> = Pin<
Box<
dyn Future<
Output = Result<
WorkflowTaskTombstoneExportReceipt,
WorkflowTaskTombstoneArchiveError,
>,
> + Send
+ 'a,
>,
>;
pub trait WorkflowTaskTombstoneArchive: Send + Sync {
fn archive(
&self,
batch: WorkflowTaskTombstoneArchiveBatch,
) -> WorkflowTaskTombstoneArchiveFuture<'_>;
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct WorkflowTaskTombstoneArchiveReport {
pub batch_id: Option<WorkflowTaskTombstoneArchiveBatchId>,
pub tombstones_exported: u32,
pub through: Option<WorkflowTaskTombstoneCursor>,
pub receipt: Option<WorkflowTaskTombstoneExportReceipt>,
}
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum WorkflowTaskGovernanceError {
#[error("principal is not authorized for {permission:?}")]
Denied {
permission: WorkflowTaskGovernancePermission,
},
#[error(transparent)]
Authorization(#[from] WorkflowTaskGovernanceAuthorizationError),
#[error(transparent)]
Store(#[from] WorkflowStoreError),
#[error(transparent)]
Archive(#[from] WorkflowTaskTombstoneArchiveError),
#[error("cleanup lease owner does not match the authenticated principal")]
LeasePrincipalMismatch,
#[error("approval lease reviewer does not match the authenticated principal")]
ApprovalLeasePrincipalMismatch,
}
pub struct WorkflowTaskGovernanceControlPlane<S, A> {
store: Arc<S>,
authorizer: Arc<A>,
observer: Arc<dyn WorkflowTaskGovernanceObserver>,
}
impl<S, A> WorkflowTaskGovernanceControlPlane<S, A>
where
S: WorkflowTaskTombstoneGovernanceStore,
A: WorkflowTaskGovernanceAuthorizer,
{
pub fn new(store: Arc<S>, authorizer: Arc<A>) -> Self {
Self {
store,
authorizer,
observer: Arc::new(NoopWorkflowTaskGovernanceObserver),
}
}
#[must_use]
pub fn with_observer(mut self, observer: Arc<dyn WorkflowTaskGovernanceObserver>) -> Self {
self.observer = observer;
self
}
pub async fn place_hold(
&self,
principal: &WorkerId,
tenant_id: WorkflowTenantId,
checkpoint_id: CheckpointId,
reason: WorkflowTaskLegalHoldReason,
) -> Result<WorkflowTaskLegalHold, WorkflowTaskGovernanceError> {
let permission = WorkflowTaskGovernancePermission::PlaceHold;
self.authorize(principal, &tenant_id, permission).await?;
let result = self
.store
.place_task_tombstone_hold(tenant_id, checkpoint_id, principal.clone(), reason)
.await;
self.store_result(permission, result)
}
pub async fn release_hold(
&self,
principal: &WorkerId,
tenant_id: WorkflowTenantId,
checkpoint_id: CheckpointId,
) -> Result<WorkflowTaskLegalHold, WorkflowTaskGovernanceError> {
let permission = WorkflowTaskGovernancePermission::ReleaseHold;
self.authorize(principal, &tenant_id, permission).await?;
let result = self
.store
.release_task_tombstone_hold(tenant_id, checkpoint_id, principal.clone())
.await;
self.store_result(permission, result)
}
pub async fn export_next_page<R>(
&self,
principal: &WorkerId,
tenant_id: WorkflowTenantId,
after: Option<WorkflowTaskTombstoneCursor>,
limit: WorkflowTaskTombstoneLimit,
archive: &R,
) -> Result<WorkflowTaskTombstoneArchiveReport, WorkflowTaskGovernanceError>
where
R: WorkflowTaskTombstoneArchive,
{
let permission = WorkflowTaskGovernancePermission::Export;
self.authorize(principal, &tenant_id, permission).await?;
let tombstones = match self
.store
.list_task_tombstones(tenant_id.clone(), after, limit)
.await
{
Ok(tombstones) => tombstones,
Err(error) => return Err(self.store_error(permission, error)),
};
let (Some(first), Some(last)) = (tombstones.first(), tombstones.last()) else {
self.observer
.observe(permission, WorkflowTaskGovernanceOutcome::Succeeded);
return Ok(WorkflowTaskTombstoneArchiveReport {
batch_id: None,
tombstones_exported: 0,
through: None,
receipt: None,
});
};
let batch_id =
WorkflowTaskTombstoneArchiveBatchId::from_batch(&tenant_id, first.cursor, last.cursor);
let through = last.cursor;
let count = u32::try_from(tombstones.len()).unwrap_or(u32::MAX);
let receipt = match archive
.archive(WorkflowTaskTombstoneArchiveBatch {
batch_id: batch_id.clone(),
tenant_id: tenant_id.clone(),
tombstones,
})
.await
{
Ok(receipt) => receipt,
Err(error) => {
self.observer
.observe(permission, WorkflowTaskGovernanceOutcome::ArchiveError);
return Err(error.into());
}
};
if let Err(error) = self
.store
.confirm_task_tombstone_export(tenant_id, through, receipt.clone(), principal.clone())
.await
{
return Err(self.store_error(permission, error));
}
self.observer
.observe(permission, WorkflowTaskGovernanceOutcome::Succeeded);
Ok(WorkflowTaskTombstoneArchiveReport {
batch_id: Some(batch_id),
tombstones_exported: count,
through: Some(through),
receipt: Some(receipt),
})
}
pub async fn prepare_purge(
&self,
principal: &WorkerId,
lease: WorkflowTaskCleanupLease,
retention: WorkflowTaskTombstoneRetention,
limit: WorkflowTaskTombstonePurgeLimit,
approval_window: WorkflowTaskTombstoneApprovalWindow,
) -> Result<WorkflowTaskTombstonePurgeIntent, WorkflowTaskGovernanceError> {
let permission = WorkflowTaskGovernancePermission::PreparePurge;
self.require_lease_principal(principal, &lease, permission)?;
self.authorize(principal, &lease.tenant_id, permission)
.await?;
let result = self
.store
.prepare_task_tombstone_purge(lease, retention, limit, approval_window)
.await;
self.store_result(permission, result)
}
pub async fn approve_purge(
&self,
principal: &WorkerId,
tenant_id: WorkflowTenantId,
purge_id: WorkflowTaskTombstonePurgeId,
) -> Result<WorkflowTaskTombstonePurgeIntent, WorkflowTaskGovernanceError> {
let permission = WorkflowTaskGovernancePermission::ApprovePurge;
self.authorize(principal, &tenant_id, permission).await?;
let result = self
.store
.approve_task_tombstone_purge(tenant_id, purge_id, principal.clone())
.await;
self.store_result(permission, result)
}
pub async fn list_purge_approvals(
&self,
principal: &WorkerId,
tenant_id: WorkflowTenantId,
limit: WorkflowTaskTombstoneApprovalInboxLimit,
) -> Result<Vec<WorkflowTaskTombstoneApprovalInboxItem>, WorkflowTaskGovernanceError> {
let permission = WorkflowTaskGovernancePermission::ReadApprovalInbox;
self.authorize(principal, &tenant_id, permission).await?;
let result = self
.store
.list_task_tombstone_purge_approvals(tenant_id, limit)
.await;
self.store_result(permission, result)
}
pub async fn claim_purge_approval(
&self,
principal: &WorkerId,
tenant_id: WorkflowTenantId,
lease: LeaseDuration,
) -> Result<Option<WorkflowTaskTombstoneApprovalLease>, WorkflowTaskGovernanceError> {
let permission = WorkflowTaskGovernancePermission::ClaimPurgeApproval;
self.authorize(principal, &tenant_id, permission).await?;
let result = self
.store
.claim_task_tombstone_purge_approval(tenant_id, principal.clone(), lease)
.await;
self.store_result(permission, result)
}
pub async fn approve_claimed_purge(
&self,
principal: &WorkerId,
lease: WorkflowTaskTombstoneApprovalLease,
) -> Result<WorkflowTaskTombstonePurgeIntent, WorkflowTaskGovernanceError> {
let permission = WorkflowTaskGovernancePermission::ApprovePurge;
self.require_approval_principal(principal, &lease, permission)?;
self.authorize(principal, &lease.tenant_id, permission)
.await?;
let result = self.store.approve_claimed_task_tombstone_purge(lease).await;
self.store_result(permission, result)
}
pub async fn reject_claimed_purge(
&self,
principal: &WorkerId,
lease: WorkflowTaskTombstoneApprovalLease,
reason: WorkflowTaskTombstoneRejectionReason,
) -> Result<WorkflowTaskTombstoneApprovalInboxItem, WorkflowTaskGovernanceError> {
let permission = WorkflowTaskGovernancePermission::RejectPurge;
self.require_approval_principal(principal, &lease, permission)?;
self.authorize(principal, &lease.tenant_id, permission)
.await?;
let result = self
.store
.reject_claimed_task_tombstone_purge(lease, reason)
.await;
self.store_result(permission, result)
}
pub async fn execute_purge(
&self,
principal: &WorkerId,
lease: WorkflowTaskCleanupLease,
purge_id: WorkflowTaskTombstonePurgeId,
) -> Result<WorkflowTaskTombstonePurgeEvidence, WorkflowTaskGovernanceError> {
let permission = WorkflowTaskGovernancePermission::ExecutePurge;
self.require_lease_principal(principal, &lease, permission)?;
self.authorize(principal, &lease.tenant_id, permission)
.await?;
let result = self
.store
.execute_task_tombstone_purge(lease, purge_id)
.await;
self.store_result(permission, result)
}
pub async fn get_evidence(
&self,
principal: &WorkerId,
tenant_id: WorkflowTenantId,
purge_id: WorkflowTaskTombstonePurgeId,
) -> Result<Option<WorkflowTaskTombstonePurgeEvidence>, WorkflowTaskGovernanceError> {
let permission = WorkflowTaskGovernancePermission::ReadEvidence;
self.authorize(principal, &tenant_id, permission).await?;
let result = self
.store
.get_task_tombstone_purge_evidence(tenant_id, purge_id)
.await;
self.store_result(permission, result)
}
async fn authorize(
&self,
principal: &WorkerId,
tenant_id: &WorkflowTenantId,
permission: WorkflowTaskGovernancePermission,
) -> Result<(), WorkflowTaskGovernanceError> {
match self
.authorizer
.authorize(principal, tenant_id, permission)
.await
{
Ok(true) => Ok(()),
Ok(false) => {
self.observer
.observe(permission, WorkflowTaskGovernanceOutcome::Denied);
Err(WorkflowTaskGovernanceError::Denied { permission })
}
Err(error) => {
self.observer.observe(
permission,
WorkflowTaskGovernanceOutcome::AuthorizationError,
);
Err(error.into())
}
}
}
fn require_lease_principal(
&self,
principal: &WorkerId,
lease: &WorkflowTaskCleanupLease,
permission: WorkflowTaskGovernancePermission,
) -> Result<(), WorkflowTaskGovernanceError> {
if lease.owner == *principal {
Ok(())
} else {
self.observer
.observe(permission, WorkflowTaskGovernanceOutcome::Denied);
Err(WorkflowTaskGovernanceError::LeasePrincipalMismatch)
}
}
fn require_approval_principal(
&self,
principal: &WorkerId,
lease: &WorkflowTaskTombstoneApprovalLease,
permission: WorkflowTaskGovernancePermission,
) -> Result<(), WorkflowTaskGovernanceError> {
if lease.reviewer == *principal {
Ok(())
} else {
self.observer
.observe(permission, WorkflowTaskGovernanceOutcome::Denied);
Err(WorkflowTaskGovernanceError::ApprovalLeasePrincipalMismatch)
}
}
fn store_result<T>(
&self,
permission: WorkflowTaskGovernancePermission,
result: Result<T, WorkflowStoreError>,
) -> Result<T, WorkflowTaskGovernanceError> {
match result {
Ok(value) => {
self.observer
.observe(permission, WorkflowTaskGovernanceOutcome::Succeeded);
Ok(value)
}
Err(error) => Err(self.store_error(permission, error)),
}
}
fn store_error(
&self,
permission: WorkflowTaskGovernancePermission,
error: WorkflowStoreError,
) -> WorkflowTaskGovernanceError {
self.observer
.observe(permission, WorkflowTaskGovernanceOutcome::StoreError);
error.into()
}
}
impl<S, A> std::fmt::Debug for WorkflowTaskGovernanceControlPlane<S, A> {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("WorkflowTaskGovernanceControlPlane")
.field("store", &"<governance-store>")
.field("authorizer", &"<governance-authorizer>")
.field("observer", &"<governance-observer>")
.finish()
}
}
#[cfg(test)]
mod tests {
use super::{WorkflowTaskTombstoneArchiveError, WorkflowTaskTombstoneArchiveErrorKind};
#[test]
fn archive_error_preserves_kind_and_bounds_safe_message() {
let error = WorkflowTaskTombstoneArchiveError::with_kind(
WorkflowTaskTombstoneArchiveErrorKind::Ambiguous,
"x".repeat(1_024),
);
assert_eq!(
error.kind(),
WorkflowTaskTombstoneArchiveErrorKind::Ambiguous
);
assert!(error.to_string().len() < 600);
assert_eq!(
WorkflowTaskTombstoneArchiveError::new("fallback").kind(),
WorkflowTaskTombstoneArchiveErrorKind::Other
);
}
}