Skip to main content

runledger_postgres/jobs/admin/
payload.rs

1use runledger_core::jobs::{JobStatus, JobType};
2use serde_json::Value;
3use sqlx::types::Uuid;
4
5use crate::{DbPool, DbTx, Error, Result};
6
7use super::super::transaction_settings::{cap_local_lock_timeout_tx, set_local_lock_timeout_tx};
8
9const JOB_PAYLOAD_UUID_ARRAY_FIELD_UPDATE_LOCK_TIMEOUT: &str = "1s";
10const JOB_PAYLOAD_UUID_ARRAY_FIELD_UPDATE_LOCK_TIMEOUT_MS: i64 = 1_000;
11
12#[derive(Debug, Clone, Copy, PartialEq, Eq)]
13#[must_use = "callers must inspect Updated/NotFound/Rejected"]
14#[non_exhaustive]
15pub enum JobPayloadUuidArrayFieldUpdate {
16    Updated,
17    NotFound,
18    Rejected {
19        reason: JobPayloadUuidArrayFieldUpdateRejection,
20    },
21}
22
23#[derive(Debug, Clone, Copy, PartialEq, Eq)]
24#[non_exhaustive]
25pub enum JobPayloadUuidArrayFieldUpdateRejection {
26    WorkflowManaged,
27    IdempotentRequestSnapshot,
28    NotPendingOrClaimed,
29}
30
31#[derive(sqlx::FromRow)]
32struct JobPayloadUuidArrayFieldUpdateCandidate {
33    status: String,
34    worker_id: Option<String>,
35    lease_expires_at: Option<chrono::DateTime<chrono::Utc>>,
36    workflow_managed: bool,
37    idempotency_key: Option<String>,
38    enqueue_request: Option<Value>,
39}
40
41/// Updates one UUID-array payload field on a direct, unclaimed pending job.
42///
43/// Returns a classified rejection when the row is already claimed or terminal,
44/// belongs to a workflow step, or has an idempotency request snapshot that this
45/// API cannot keep consistent.
46pub async fn update_job_payload_uuid_array_field(
47    pool: &DbPool,
48    organization_id: Uuid,
49    job_id: Uuid,
50    job_type: JobType<'_>,
51    payload_field: &str,
52    values: &[Uuid],
53) -> Result<JobPayloadUuidArrayFieldUpdate> {
54    let mut tx = pool.begin().await.map_err(|error| {
55        Error::from_query_sqlx_with_context(
56            "begin job payload uuid array update transaction",
57            error,
58        )
59    })?;
60
61    let previous_lock_timeout =
62        cap_job_payload_uuid_array_field_update_lock_timeout_tx(&mut tx).await?;
63
64    let row_result = sqlx::query_as::<_, JobPayloadUuidArrayFieldUpdateCandidate>(
65        "SELECT
66             status::text AS status,
67             worker_id,
68             lease_expires_at,
69             EXISTS (
70                 SELECT 1
71                 FROM workflow_steps ws
72                 WHERE ws.job_id = job_queue.id
73             ) AS workflow_managed,
74             idempotency_key,
75             enqueue_request
76           FROM job_queue
77           WHERE id = $1
78             AND organization_id = $2
79             AND job_type = $3
80           FOR UPDATE",
81    )
82    .bind(job_id)
83    .bind(organization_id)
84    .bind(job_type)
85    .fetch_optional(&mut *tx)
86    .await;
87
88    let row = match row_result {
89        Ok(row) => {
90            set_local_lock_timeout_tx(
91                &mut tx,
92                &previous_lock_timeout,
93                "restore job payload uuid array update lock timeout",
94            )
95            .await?;
96            row
97        }
98        Err(error) => {
99            return Err(Error::from_query_sqlx_with_context(
100                "classify job payload uuid array update",
101                error,
102            ));
103        }
104    };
105
106    let Some(row) = row else {
107        tx.commit().await.map_err(|error| {
108            Error::from_query_sqlx_with_context(
109                "commit job payload uuid array update transaction",
110                error,
111            )
112        })?;
113        return Ok(JobPayloadUuidArrayFieldUpdate::NotFound);
114    };
115
116    // Order matters: workflow-managed jobs can also carry request snapshots, so
117    // return the ownership rejection before the snapshot-consistency rejection.
118    let rejection = if row.workflow_managed {
119        Some(JobPayloadUuidArrayFieldUpdateRejection::WorkflowManaged)
120    } else if row.idempotency_key.is_some() || row.enqueue_request.is_some() {
121        Some(JobPayloadUuidArrayFieldUpdateRejection::IdempotentRequestSnapshot)
122    } else if row.status != JobStatus::Pending.as_db_value()
123        || row.worker_id.is_some()
124        || row.lease_expires_at.is_some()
125    {
126        Some(JobPayloadUuidArrayFieldUpdateRejection::NotPendingOrClaimed)
127    } else {
128        None
129    };
130
131    if let Some(reason) = rejection {
132        tx.commit().await.map_err(|error| {
133            Error::from_query_sqlx_with_context(
134                "commit job payload uuid array update transaction",
135                error,
136            )
137        })?;
138        return Ok(JobPayloadUuidArrayFieldUpdate::Rejected { reason });
139    }
140
141    sqlx::query!(
142        "UPDATE job_queue
143         SET
144             payload = jsonb_set(
145                 payload,
146                 ARRAY[$4::text],
147                 to_jsonb($5::uuid[]),
148                 true
149             ),
150             updated_at = now()
151         WHERE id = $1
152           AND organization_id = $2
153           AND job_type = $3",
154        job_id,
155        organization_id,
156        job_type as _,
157        payload_field,
158        values,
159    )
160    .execute(&mut *tx)
161    .await
162    .map_err(|error| {
163        Error::from_query_sqlx_with_context("update job payload uuid array field", error)
164    })?;
165
166    tx.commit().await.map_err(|error| {
167        Error::from_query_sqlx_with_context(
168            "commit job payload uuid array update transaction",
169            error,
170        )
171    })?;
172    Ok(JobPayloadUuidArrayFieldUpdate::Updated)
173}
174
175async fn cap_job_payload_uuid_array_field_update_lock_timeout_tx(
176    tx: &mut DbTx<'_>,
177) -> Result<String> {
178    cap_local_lock_timeout_tx(
179        tx,
180        JOB_PAYLOAD_UUID_ARRAY_FIELD_UPDATE_LOCK_TIMEOUT,
181        JOB_PAYLOAD_UUID_ARRAY_FIELD_UPDATE_LOCK_TIMEOUT_MS,
182        "set job payload uuid array update lock timeout",
183    )
184    .await
185}