1use std::{num::NonZeroU32, num::NonZeroU64, time::Duration};
4
5use runifold_core::CheckpointId;
6
7use crate::{
8 WorkerId, WorkflowStoreError, WorkflowStoreErrorKind, WorkflowStoreFuture,
9 WorkflowTaskCleanupLease, WorkflowTaskRetentionStore, WorkflowTaskTombstoneCursor,
10 WorkflowTenantId,
11};
12
13#[derive(Clone, Copy, Debug, Eq, PartialEq)]
15pub struct WorkflowTaskTombstoneRetention(NonZeroU64);
16
17impl WorkflowTaskTombstoneRetention {
18 pub fn new(duration: Duration) -> Result<Self, WorkflowStoreError> {
24 let millis = u64::try_from(duration.as_millis())
25 .ok()
26 .and_then(NonZeroU64::new)
27 .ok_or_else(|| {
28 invalid_input("Task tombstone retention must fit in positive whole milliseconds")
29 })?;
30 Ok(Self(millis))
31 }
32
33 pub const fn as_millis(self) -> u64 {
35 self.0.get()
36 }
37}
38
39#[derive(Clone, Copy, Debug, Eq, PartialEq)]
41pub struct WorkflowTaskTombstonePurgeLimit(NonZeroU32);
42
43impl WorkflowTaskTombstonePurgeLimit {
44 pub fn new(value: u32) -> Result<Self, WorkflowStoreError> {
50 let value =
51 NonZeroU32::new(value).ok_or_else(|| invalid_input("purge limit must be positive"))?;
52 if value.get() > 1_000 {
53 return Err(invalid_input("purge limit cannot exceed 1,000"));
54 }
55 Ok(Self(value))
56 }
57
58 pub const fn get(self) -> u32 {
60 self.0.get()
61 }
62}
63
64#[derive(Clone, Copy, Debug, Eq, PartialEq)]
66pub struct WorkflowTaskTombstoneApprovalWindow(NonZeroU64);
67
68impl WorkflowTaskTombstoneApprovalWindow {
69 pub fn new(duration: Duration) -> Result<Self, WorkflowStoreError> {
75 let millis = u64::try_from(duration.as_millis())
76 .ok()
77 .and_then(NonZeroU64::new)
78 .ok_or_else(|| {
79 invalid_input("purge approval window must fit in positive whole milliseconds")
80 })?;
81 Ok(Self(millis))
82 }
83
84 pub const fn as_millis(self) -> u64 {
86 self.0.get()
87 }
88}
89
90#[derive(Clone, Copy, Debug, Eq, PartialEq)]
92pub struct WorkflowTaskTombstoneApprovalInboxLimit(NonZeroU32);
93
94impl WorkflowTaskTombstoneApprovalInboxLimit {
95 pub fn new(value: u32) -> Result<Self, WorkflowStoreError> {
101 let value = NonZeroU32::new(value)
102 .ok_or_else(|| invalid_input("approval inbox limit must be positive"))?;
103 if value.get() > 1_000 {
104 return Err(invalid_input("approval inbox limit cannot exceed 1,000"));
105 }
106 Ok(Self(value))
107 }
108
109 pub const fn get(self) -> u32 {
111 self.0.get()
112 }
113}
114
115#[derive(Clone, Debug, Eq, PartialEq)]
117pub struct WorkflowTaskTombstoneRejectionReason(String);
118
119impl WorkflowTaskTombstoneRejectionReason {
120 pub fn parse(value: impl Into<String>) -> Result<Self, WorkflowStoreError> {
126 let value = value.into();
127 if value.trim().is_empty() || value.len() > 1_024 || value.chars().any(char::is_control) {
128 return Err(invalid_input(
129 "purge rejection reason must contain 1..=1,024 printable bytes",
130 ));
131 }
132 Ok(Self(value))
133 }
134
135 pub fn as_str(&self) -> &str {
137 &self.0
138 }
139}
140
141#[derive(Clone, Debug, Eq, PartialEq)]
143pub struct WorkflowTaskLegalHoldReason(String);
144
145impl WorkflowTaskLegalHoldReason {
146 pub fn parse(value: impl Into<String>) -> Result<Self, WorkflowStoreError> {
152 let value = value.into();
153 if value.trim().is_empty() || value.len() > 1_024 {
154 return Err(invalid_input(
155 "Task legal-hold reason must contain 1..=1,024 bytes",
156 ));
157 }
158 Ok(Self(value))
159 }
160
161 pub fn as_str(&self) -> &str {
163 &self.0
164 }
165}
166
167#[derive(Clone, Debug, Eq, PartialEq)]
169pub struct WorkflowTaskTombstoneExportReceipt(String);
170
171impl WorkflowTaskTombstoneExportReceipt {
172 pub fn parse(value: impl Into<String>) -> Result<Self, WorkflowStoreError> {
178 let value = value.into();
179 if value.trim().is_empty() || value.len() > 512 || value.chars().any(char::is_control) {
180 return Err(invalid_input(
181 "Task tombstone export receipt must contain 1..=512 printable bytes",
182 ));
183 }
184 Ok(Self(value))
185 }
186
187 pub fn as_str(&self) -> &str {
189 &self.0
190 }
191}
192
193#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
195pub struct WorkflowTaskTombstonePurgeId(CheckpointId);
196
197impl WorkflowTaskTombstonePurgeId {
198 pub fn new() -> Self {
200 Self(CheckpointId::new())
201 }
202
203 pub const fn from_checkpoint_id(value: CheckpointId) -> Self {
205 Self(value)
206 }
207
208 pub const fn as_checkpoint_id(self) -> CheckpointId {
210 self.0
211 }
212}
213
214impl Default for WorkflowTaskTombstonePurgeId {
215 fn default() -> Self {
216 Self::new()
217 }
218}
219
220#[derive(Clone, Debug, Eq, PartialEq)]
222pub struct WorkflowTaskLegalHold {
223 pub checkpoint_id: CheckpointId,
225 pub tenant_id: WorkflowTenantId,
227 pub placed_by: WorkerId,
229 pub reason: WorkflowTaskLegalHoldReason,
231 pub placed_at_ms: u64,
233 pub released_by: Option<WorkerId>,
235 pub released_at_ms: Option<u64>,
237}
238
239impl WorkflowTaskLegalHold {
240 pub const fn is_active(&self) -> bool {
242 self.released_at_ms.is_none()
243 }
244}
245
246#[derive(Clone, Debug, Eq, PartialEq)]
248pub struct WorkflowTaskTombstoneExport {
249 pub tenant_id: WorkflowTenantId,
251 pub through: WorkflowTaskTombstoneCursor,
253 pub receipt: WorkflowTaskTombstoneExportReceipt,
255 pub confirmed_by: WorkerId,
257 pub confirmed_at_ms: u64,
259}
260
261#[derive(Clone, Debug, Eq, PartialEq)]
263pub struct WorkflowTaskTombstonePurgeIntent {
264 pub purge_id: WorkflowTaskTombstonePurgeId,
266 pub tenant_id: WorkflowTenantId,
268 pub prepared_by: WorkerId,
270 pub tombstone_count: u32,
272 pub first_cursor: Option<WorkflowTaskTombstoneCursor>,
274 pub last_cursor: Option<WorkflowTaskTombstoneCursor>,
276 pub export_through: WorkflowTaskTombstoneCursor,
278 pub fingerprint: String,
280 pub prepared_at_ms: u64,
282 pub expires_at_ms: u64,
284 pub approved_by: Option<WorkerId>,
286 pub approved_at_ms: Option<u64>,
288}
289
290#[derive(Clone, Copy, Debug, Eq, PartialEq)]
292#[non_exhaustive]
293pub enum WorkflowTaskTombstoneApprovalState {
294 Pending,
296 Claimed,
298 Approved,
300 Rejected,
302 Expired,
304}
305
306#[derive(Clone, Debug, Eq, PartialEq)]
308pub struct WorkflowTaskTombstoneApprovalInboxItem {
309 pub intent: WorkflowTaskTombstonePurgeIntent,
311 pub state: WorkflowTaskTombstoneApprovalState,
313 pub claimed_by: Option<WorkerId>,
315 pub claim_expires_at_ms: Option<u64>,
317 pub rejected_by: Option<WorkerId>,
319 pub rejection_reason: Option<WorkflowTaskTombstoneRejectionReason>,
321 pub rejected_at_ms: Option<u64>,
323}
324
325#[derive(Clone, Debug, Eq, PartialEq)]
327pub struct WorkflowTaskTombstoneApprovalLease {
328 pub tenant_id: WorkflowTenantId,
330 pub purge_id: WorkflowTaskTombstonePurgeId,
332 pub reviewer: WorkerId,
334 pub fencing_token: u64,
336 pub expires_at_ms: u64,
338}
339
340#[derive(Clone, Debug, Eq, PartialEq)]
342pub struct WorkflowTaskTombstonePurgeEvidence {
343 pub purge_id: WorkflowTaskTombstonePurgeId,
345 pub tenant_id: WorkflowTenantId,
347 pub prepared_by: WorkerId,
349 pub approved_by: WorkerId,
351 pub executed_by: WorkerId,
353 pub tombstone_count: u32,
355 pub first_cursor: WorkflowTaskTombstoneCursor,
357 pub last_cursor: WorkflowTaskTombstoneCursor,
359 pub export_through: WorkflowTaskTombstoneCursor,
361 pub fingerprint: String,
363 pub executed_at_ms: u64,
365}
366
367pub trait WorkflowTaskTombstoneGovernanceStore: WorkflowTaskRetentionStore {
369 fn place_task_tombstone_hold(
371 &self,
372 tenant_id: WorkflowTenantId,
373 checkpoint_id: CheckpointId,
374 actor: WorkerId,
375 reason: WorkflowTaskLegalHoldReason,
376 ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskLegalHold, WorkflowStoreError>>;
377
378 fn release_task_tombstone_hold(
380 &self,
381 tenant_id: WorkflowTenantId,
382 checkpoint_id: CheckpointId,
383 actor: WorkerId,
384 ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskLegalHold, WorkflowStoreError>>;
385
386 fn confirm_task_tombstone_export(
388 &self,
389 tenant_id: WorkflowTenantId,
390 through: WorkflowTaskTombstoneCursor,
391 receipt: WorkflowTaskTombstoneExportReceipt,
392 actor: WorkerId,
393 ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstoneExport, WorkflowStoreError>>;
394
395 fn prepare_task_tombstone_purge(
397 &self,
398 lease: WorkflowTaskCleanupLease,
399 retention: WorkflowTaskTombstoneRetention,
400 limit: WorkflowTaskTombstonePurgeLimit,
401 approval_window: WorkflowTaskTombstoneApprovalWindow,
402 ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeIntent, WorkflowStoreError>>;
403
404 fn approve_task_tombstone_purge(
406 &self,
407 tenant_id: WorkflowTenantId,
408 purge_id: WorkflowTaskTombstonePurgeId,
409 approver: WorkerId,
410 ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeIntent, WorkflowStoreError>>;
411
412 fn list_task_tombstone_purge_approvals(
414 &self,
415 tenant_id: WorkflowTenantId,
416 limit: WorkflowTaskTombstoneApprovalInboxLimit,
417 ) -> WorkflowStoreFuture<
418 '_,
419 Result<Vec<WorkflowTaskTombstoneApprovalInboxItem>, WorkflowStoreError>,
420 >;
421
422 fn claim_task_tombstone_purge_approval(
424 &self,
425 tenant_id: WorkflowTenantId,
426 reviewer: WorkerId,
427 lease: crate::LeaseDuration,
428 ) -> WorkflowStoreFuture<
429 '_,
430 Result<Option<WorkflowTaskTombstoneApprovalLease>, WorkflowStoreError>,
431 >;
432
433 fn approve_claimed_task_tombstone_purge(
435 &self,
436 lease: WorkflowTaskTombstoneApprovalLease,
437 ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeIntent, WorkflowStoreError>>;
438
439 fn reject_claimed_task_tombstone_purge(
441 &self,
442 lease: WorkflowTaskTombstoneApprovalLease,
443 reason: WorkflowTaskTombstoneRejectionReason,
444 ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstoneApprovalInboxItem, WorkflowStoreError>>;
445
446 fn execute_task_tombstone_purge(
450 &self,
451 lease: WorkflowTaskCleanupLease,
452 purge_id: WorkflowTaskTombstonePurgeId,
453 ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeEvidence, WorkflowStoreError>>;
454
455 fn get_task_tombstone_purge_evidence(
457 &self,
458 tenant_id: WorkflowTenantId,
459 purge_id: WorkflowTaskTombstonePurgeId,
460 ) -> WorkflowStoreFuture<
461 '_,
462 Result<Option<WorkflowTaskTombstonePurgeEvidence>, WorkflowStoreError>,
463 >;
464}
465
466fn invalid_input(message: &'static str) -> WorkflowStoreError {
467 WorkflowStoreError::new(WorkflowStoreErrorKind::InvalidInput, message)
468}