runledger_postgres/jobs/admin/
payload.rs1use 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
41pub 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 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}