runledger_postgres/jobs/workflows/
enqueue.rs1use runledger_core::jobs::{WorkflowRunEnqueue, validate_workflow_run_enqueue};
2
3use super::super::workflow_types::{EnqueueActiveWorkflowOutcome, WorkflowRunDbRecord};
4use super::active_claims::{
5 insert_workflow_active_claim_tx, load_existing_active_workflow_run_tx,
6 lock_workflow_active_key_tx,
7};
8use super::enqueue_persistence::{
9 insert_workflow_run_record_tx, try_load_existing_idempotent_workflow_run_tx,
10 validate_existing_idempotent_workflow_run,
11};
12use super::errors::{
13 workflow_active_key_api_required_error, workflow_active_key_required_error,
14 workflow_internal_state_error,
15};
16use super::read::load_workflow_run_by_id_tx;
17use super::release::enqueue_root_steps_tx;
18use super::runtime::recompute_workflow_run_status_tx;
19use super::snapshot::canonical_workflow_enqueue_request;
20use super::steps::{
21 WorkflowStepDependencyWriteContext, fetch_job_definition_defaults_tx,
22 insert_workflow_step_dependencies_tx, insert_workflow_steps_tx,
23};
24use super::validation::workflow_dag_validation_error;
25use crate::jobs::transaction_isolation::{ReadCommittedTx, ensure_read_committed_tx};
26use crate::{DbPool, DbTx, Error, Result};
27
28#[doc(alias = "dag")]
40#[doc(alias = "orchestration")]
41#[doc(alias = "dependencies")]
42pub async fn enqueue_workflow_run(
43 pool: &DbPool,
44 payload: &WorkflowRunEnqueue<'_>,
45) -> Result<WorkflowRunDbRecord> {
46 if payload.active_key().is_some() {
47 return Err(workflow_active_key_api_required_error());
48 }
49 let mut tx = pool
50 .begin()
51 .await
52 .map_err(|error| Error::ConnectionError(error.to_string()))?;
53 let outcome = enqueue_workflow_run_classified_tx(&mut tx, payload).await?;
54 let workflow_run = workflow_run_from_classified_outcome(outcome)?;
55 tx.commit()
56 .await
57 .map_err(|error| Error::ConnectionError(error.to_string()))?;
58 Ok(workflow_run)
59}
60
61#[doc(alias = "dag")]
77#[doc(alias = "orchestration")]
78#[doc(alias = "dependencies")]
79pub async fn enqueue_workflow_run_tx(
80 tx: &mut DbTx<'_>,
81 payload: &WorkflowRunEnqueue<'_>,
82) -> Result<WorkflowRunDbRecord> {
83 if payload.active_key().is_some() {
84 return Err(workflow_active_key_api_required_error());
85 }
86 let outcome = enqueue_workflow_run_classified_tx(tx, payload).await?;
87 workflow_run_from_classified_outcome(outcome)
88}
89
90fn workflow_run_from_classified_outcome(
91 outcome: EnqueueActiveWorkflowOutcome,
92) -> Result<WorkflowRunDbRecord> {
93 match outcome {
94 EnqueueActiveWorkflowOutcome::Inserted(run)
95 | EnqueueActiveWorkflowOutcome::ExistingIdempotent(run) => Ok(run),
96 EnqueueActiveWorkflowOutcome::ExistingActive(_) => Err(workflow_internal_state_error(
97 "active workflow collision was returned for a payload without active_key",
98 )),
99 }
100}
101
102pub async fn enqueue_or_get_active_workflow(
112 pool: &DbPool,
113 payload: &WorkflowRunEnqueue<'_>,
114) -> Result<EnqueueActiveWorkflowOutcome> {
115 if payload.active_key().is_none() {
116 return Err(workflow_active_key_required_error());
117 }
118 let mut tx = pool
119 .begin()
120 .await
121 .map_err(|error| Error::ConnectionError(error.to_string()))?;
122 let outcome = enqueue_workflow_run_classified_tx(&mut tx, payload).await?;
123 tx.commit()
124 .await
125 .map_err(|error| Error::ConnectionError(error.to_string()))?;
126 Ok(outcome)
127}
128
129pub async fn enqueue_or_get_active_workflow_tx(
134 tx: &mut DbTx<'_>,
135 payload: &WorkflowRunEnqueue<'_>,
136) -> Result<EnqueueActiveWorkflowOutcome> {
137 if payload.active_key().is_none() {
138 return Err(workflow_active_key_required_error());
139 }
140 enqueue_workflow_run_classified_tx(tx, payload).await
141}
142
143async fn enqueue_workflow_run_classified_tx(
144 tx: &mut DbTx<'_>,
145 payload: &WorkflowRunEnqueue<'_>,
146) -> Result<EnqueueActiveWorkflowOutcome> {
147 validate_workflow_run_enqueue(payload).map_err(workflow_dag_validation_error)?;
148 if payload.idempotency_key().is_some() || payload.active_key().is_some() {
149 let mut read_committed_tx = ensure_read_committed_tx(
150 tx,
151 "workflow coordinated enqueue",
152 "workflow.enqueue_idempotency_unsupported_isolation",
153 "Workflow idempotent or active-key enqueue requires READ COMMITTED transaction isolation.",
154 )
155 .await?;
156
157 return enqueue_coordinated_workflow_run_read_committed_tx(&mut read_committed_tx, payload)
158 .await;
159 }
160
161 enqueue_workflow_run_classified_tx_inner(tx, payload).await
162}
163
164async fn enqueue_coordinated_workflow_run_read_committed_tx(
165 tx: &mut ReadCommittedTx<'_, '_>,
166 payload: &WorkflowRunEnqueue<'_>,
167) -> Result<EnqueueActiveWorkflowOutcome> {
168 debug_assert!(payload.idempotency_key().is_some() || payload.active_key().is_some());
169 enqueue_workflow_run_classified_tx_inner(tx.as_tx(), payload).await
170}
171
172async fn enqueue_workflow_run_classified_tx_inner(
173 tx: &mut DbTx<'_>,
174 payload: &WorkflowRunEnqueue<'_>,
175) -> Result<EnqueueActiveWorkflowOutcome> {
176 let enqueue_request = canonical_workflow_enqueue_request(payload)?;
177
178 if let Some(active_key) = payload.active_key() {
179 lock_workflow_active_key_tx(tx, payload.organization_id(), active_key).await?;
180
181 if let Some(idempotency_key) = payload.idempotency_key() {
182 if let Some(existing) = try_load_existing_idempotent_workflow_run_tx(
183 tx,
184 payload,
185 idempotency_key,
186 &enqueue_request,
187 )
188 .await?
189 {
190 validate_existing_idempotent_workflow_run(&existing)?;
191 return Ok(EnqueueActiveWorkflowOutcome::ExistingIdempotent(
192 existing.into_record()?,
193 ));
194 }
195 }
196
197 if let Some(existing) =
198 load_existing_active_workflow_run_tx(tx, payload.organization_id(), active_key).await?
199 {
200 return Ok(EnqueueActiveWorkflowOutcome::ExistingActive(existing));
201 }
202 }
203
204 let workflow_run_insert = insert_workflow_run_record_tx(tx, payload, &enqueue_request).await?;
205 let workflow_run = workflow_run_insert.record;
206 if !workflow_run_insert.inserted {
207 return Ok(EnqueueActiveWorkflowOutcome::ExistingIdempotent(
210 workflow_run,
211 ));
212 }
213 if let Some(active_key) = payload.active_key() {
214 insert_workflow_active_claim_tx(tx, payload.organization_id(), active_key, workflow_run.id)
215 .await?;
216 }
217
218 let defaults_by_job_type = fetch_job_definition_defaults_tx(tx, payload.steps()).await?;
219 let step_id_by_key =
220 insert_workflow_steps_tx(tx, payload, workflow_run.id, &defaults_by_job_type).await?;
221 insert_workflow_step_dependencies_tx(
222 tx,
223 payload.steps(),
224 workflow_run.id,
225 &step_id_by_key,
226 WorkflowStepDependencyWriteContext::InitialEnqueue,
227 )
228 .await?;
229
230 enqueue_root_steps_tx(tx, workflow_run.id).await?;
231 recompute_workflow_run_status_tx(tx, workflow_run.id).await?;
232
233 let workflow_run = load_workflow_run_by_id_tx(
234 tx,
235 workflow_run.id,
236 "load workflow run after enqueue recompute",
237 )
238 .await?;
239 Ok(EnqueueActiveWorkflowOutcome::Inserted(workflow_run))
240}
241
242#[cfg(test)]
243mod tests {
244 use runledger_core::jobs::{
245 JobStage, JobType, StepKey, WorkflowRunEnqueueBuilder, WorkflowStepEnqueueBuilder,
246 WorkflowType,
247 };
248 use serde_json::json;
249 use sqlx::types::Uuid;
250
251 use super::super::snapshot::canonical_workflow_enqueue_request;
252
253 #[test]
254 fn canonical_workflow_enqueue_request_matches_golden_snapshot() {
255 let run_org = Uuid::now_v7();
256 let step_org = Uuid::now_v7();
257 let metadata = json!({"kind": "golden"});
258 let root_payload = json!({"step": "root"});
259 let child_payload = json!({"step": "child"});
260 let root = WorkflowStepEnqueueBuilder::new_external(StepKey::new("root"), &root_payload)
261 .try_build()
262 .expect("build root step");
263 let child = WorkflowStepEnqueueBuilder::new(
264 StepKey::new("child"),
265 JobType::new("jobs.test.child"),
266 &child_payload,
267 )
268 .organization_id(step_org)
269 .priority(7)
270 .max_attempts(2)
271 .timeout_seconds(45)
272 .stage(JobStage::Scheduled)
273 .depends_on_success(&[StepKey::new("root")])
274 .try_build()
275 .expect("build child step");
276 let workflow =
277 WorkflowRunEnqueueBuilder::new(WorkflowType::new("workflow.test.golden"), &metadata)
278 .organization_id(run_org)
279 .step(child)
280 .step(root)
281 .try_build()
282 .expect("build workflow");
283
284 let canonical =
285 canonical_workflow_enqueue_request(&workflow).expect("canonicalize workflow enqueue");
286
287 assert_eq!(
288 canonical,
289 json!({
290 "metadata": {"kind": "golden"},
291 "steps": [
292 {
293 "step_key": "child",
294 "execution_kind": "JOB",
295 "job_type": "jobs.test.child",
296 "organization_id": step_org,
297 "payload": {"step": "child"},
298 "priority": 7,
299 "max_attempts": 2,
300 "timeout_seconds": 45,
301 "stage": "scheduled",
302 "dependencies": [
303 {
304 "prerequisite_step_key": "root",
305 "release_mode": "ON_SUCCESS"
306 }
307 ]
308 },
309 {
310 "step_key": "root",
311 "execution_kind": "EXTERNAL",
312 "job_type": null,
313 "organization_id": run_org,
314 "payload": {"step": "root"},
315 "priority": null,
316 "max_attempts": null,
317 "timeout_seconds": null,
318 "stage": null,
319 "dependencies": []
320 }
321 ]
322 })
323 );
324 }
325
326 #[test]
327 fn canonical_workflow_enqueue_request_includes_result_step_key_when_present() {
328 let metadata = json!({});
329 let payload = json!({"step": "result"});
330 let result = WorkflowStepEnqueueBuilder::new(
331 StepKey::new("result"),
332 JobType::new("jobs.test.result"),
333 &payload,
334 )
335 .try_build()
336 .expect("build result step");
337 let workflow =
338 WorkflowRunEnqueueBuilder::new(WorkflowType::new("workflow.test.result"), &metadata)
339 .step(result)
340 .try_result_step_key("result")
341 .expect("set result step key")
342 .try_build()
343 .expect("build workflow");
344
345 let canonical =
346 canonical_workflow_enqueue_request(&workflow).expect("canonicalize workflow enqueue");
347
348 assert_eq!(canonical.get("result_step_key"), Some(&json!("result")));
349 }
350}