Skip to main content

runledger_postgres/jobs/types/
reaper.rs

1use chrono::{DateTime, Utc};
2use runledger_core::jobs::{JobFailure, JobTypeName};
3use serde_json::Value;
4use sqlx::types::Uuid;
5
6#[derive(Clone, Debug, Eq, PartialEq)]
7pub struct ReapedTerminalLeaseRecord {
8    pub job_id: Uuid,
9    pub job_type: JobTypeName,
10    pub organization_id: Option<Uuid>,
11    pub run_number: i32,
12    pub attempt: i32,
13    pub payload: Value,
14}
15
16#[derive(Clone, Debug, Eq, PartialEq)]
17#[non_exhaustive]
18pub enum ReapedLeaseDisposition {
19    ReleasedToPending,
20    RetryScheduled {
21        retry_delay_ms: i32,
22        next_run_at: DateTime<Utc>,
23    },
24    DeadLetteredTerminal {
25        payload: Value,
26    },
27}
28
29#[derive(Clone, Debug)]
30#[non_exhaustive]
31pub struct ReapedLeaseRecord {
32    pub job_id: Uuid,
33    pub job_type: JobTypeName,
34    pub organization_id: Option<Uuid>,
35    pub run_number: i32,
36    pub attempt: i32,
37    pub max_attempts: i32,
38    /// Checkpoint committed on the leased run before it was reaped.
39    pub checkpoint: Option<Value>,
40    pub worker_id: Option<String>,
41    pub started_without_renewal_heartbeat: bool,
42    pub failure: JobFailure,
43    pub disposition: ReapedLeaseDisposition,
44}
45
46#[derive(Clone, Debug)]
47#[non_exhaustive]
48pub struct ReapExpiredLeaseDeferredError {
49    pub job_id: Uuid,
50    pub run_number: i32,
51    pub attempt: i32,
52    pub error_code: String,
53    pub error_message: String,
54    pub sqlstate: Option<String>,
55}
56
57#[derive(Clone, Copy, Debug, Eq, PartialEq)]
58#[non_exhaustive]
59pub enum ReapExpiredLeaseCleanupOperation {
60    WorkflowActiveClaims,
61    ExecutionResourceClaims,
62}
63
64impl ReapExpiredLeaseCleanupOperation {
65    #[must_use]
66    pub const fn as_str(self) -> &'static str {
67        match self {
68            Self::WorkflowActiveClaims => "workflow_active_claims",
69            Self::ExecutionResourceClaims => "execution_resource_claims",
70        }
71    }
72}
73
74#[derive(Clone, Debug)]
75#[non_exhaustive]
76pub struct ReapExpiredLeaseCleanupError {
77    /// Bounded cleanup operation that failed after lease transitions committed.
78    pub operation: ReapExpiredLeaseCleanupOperation,
79    /// Persistence error text for trusted operator diagnostics.
80    pub error: String,
81}
82
83#[derive(Clone, Debug)]
84pub struct ReapExpiredLeasesResult {
85    pub processed: i64,
86    pub terminal_dead_lettered: Vec<ReapedTerminalLeaseRecord>,
87}
88
89/// Detailed lease-reaper outcome, including post-commit coordination cleanup.
90///
91/// Lease transitions commit before active/resource claim cleanup. A non-empty
92/// `cleanup_errors` therefore does not roll back the reaped jobs in `summary`.
93#[derive(Clone, Debug)]
94#[non_exhaustive]
95pub struct ReapExpiredLeasesDetailedResult {
96    pub summary: ReapExpiredLeasesResult,
97    pub reaped_leases: Vec<ReapedLeaseRecord>,
98    pub deferred_row_error_count: usize,
99    pub deferred_row_errors: Vec<ReapExpiredLeaseDeferredError>,
100    /// Quiesced reusable workflow claims removed in the bounded cleanup pass.
101    pub workflow_active_claims_released: u64,
102    /// Stale execution-resource claims removed in the bounded cleanup pass.
103    pub execution_resource_claims_released: u64,
104    pub cleanup_errors: Vec<ReapExpiredLeaseCleanupError>,
105}