Skip to main content

runledger_postgres/jobs/workflows/
enqueue.rs

1use 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/// Enqueues a workflow run in its own transaction.
29///
30/// Use this API for multi-step work with dependencies, fan-out/fan-in, external
31/// gates, cancellation as one logical run, or workflow-level idempotency. Build
32/// the payload with `WorkflowRunEnqueueBuilder` and
33/// `WorkflowStepEnqueueBuilder`.
34///
35/// Calls without an idempotency key always create a new workflow run. Calls
36/// with an idempotency key return the existing run only when the canonical
37/// enqueue request snapshot matches. Keyed rows without snapshots are rejected
38/// by the idempotency cutover.
39#[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/// Enqueues a workflow run and returns the existing run for an identical keyed retry.
62///
63/// Use this API when composing a workflow enqueue with other database writes in
64/// one transaction. For ordinary dependent work, prefer this workflow path over
65/// direct-job polling or handler-chained follow-up jobs.
66///
67/// Idempotency is strict for the submitted request snapshot. The snapshot is
68/// compared instead of live workflow step rows because steps and dependencies
69/// can be legitimately appended or mutated after initial enqueue. Strict
70/// workflow idempotency applies to keyed rows with an `enqueue_request`
71/// snapshot. New unkeyed rows also store snapshots for workflow recovery, while
72/// keyed legacy rows without snapshots are rejected by the idempotency cutover.
73/// Job-step stage is part of the canonical initial request after normalizing an
74/// omitted stage to `Queued`; changing the requested initial stage is treated as
75/// a different enqueue request.
76#[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
102/// Enqueues a workflow under a reusable active key and explicitly classifies
103/// insertion, active collision, and permanent idempotency collision.
104///
105/// Active-key scope is global when `organization_id` is absent and otherwise
106/// organization-local; workflow type is not part of the scope. Claims remain
107/// reserved until terminal work and canceled leases are quiescent, so
108/// [`EnqueueActiveWorkflowOutcome::ExistingActive`] may carry a terminal
109/// canceled run. Callers must make their decision from the outcome rather than
110/// status alone.
111pub 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
129/// Caller-transaction counterpart to [`enqueue_or_get_active_workflow`].
130///
131/// The transaction must use `READ COMMITTED`. This function neither commits nor
132/// rolls it back.
133pub 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        // Existing idempotent runs already have their steps and initial releases
208        // committed; never replay workflow initialization for a retry.
209        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}